Fresco/Prague/src/IPC Dispatcher.cc,1.15,1.16

Stefan Seefeld <[email protected]> Wed, 07 May 2003 22:21:13 -0500
Newsgroups gmane.comp.video.fresco.cvs
Message-ID <[email protected]>
Update of /cvs/fresco/Fresco/Prague/src/IPC
In directory purcel:/tmp/cvs-serv22129/src/IPC

Modified Files:
	Dispatcher.cc 
Log Message:
some API clarifications for Threads, and minor updates to enhance coding standard conformity

Index: Dispatcher.cc
===================================================================
RCS file: /cvs/fresco/Fresco/Prague/src/IPC/Dispatcher.cc,v
retrieving revision 1.15
retrieving revision 1.16
diff -u -d -r1.15 -r1.16
--- Dispatcher.cc	25 Mar 2001 08:25:16 -0000	1.15
+++ Dispatcher.cc	8 May 2003 03:21:11 -0000	1.16
@@ -1,8 +1,8 @@
 /*$Id$
  *
- * This source file is a part of the Berlin Project.
- * Copyright (C) 1999, 2000 Stefan Seefeld <[email protected]> 
- * http://www.berlin-consortium.org
+ * This source file is a part of the Fresco Project.
+ * Copyright (C) 1999, 2000 Stefan Seefeld <[email protected]> 
+ * http://www.fresco.org
  *
  * This library is free software; you can redistribute it and/or
  * modify it under the terms of the GNU Library General Public
@@ -35,13 +35,15 @@
 Mutex Dispatcher::singletonMutex;
 Dispatcher::Cleaner Dispatcher::cleaner;
 
-struct SignalNotifier : Signal::Notifier
+struct Dispatcher::task
 {
-  virtual void notify(int signum)
-  {
-    std::cerr << Signal::name(signum) << std::endl;
-    exit(1);
-  }
+  task() : fd(-1), agent(0), mask(Agent::none), released(false) {}
+  task(int ffd, Agent *a, Agent::iomask m) : fd(ffd), agent(a), mask(m) {}
+  bool operator < (const task &t) const { return fd < t.fd;}
+  int             fd;
+  Agent          *agent;
+  Agent::iomask   mask;
+  bool            released;
 };
 
 Dispatcher::Cleaner::~Cleaner()
@@ -57,25 +59,14 @@
   return dispatcher;
 }
 
+//. create a queue of up to 64 tasks 
+//. and a thread pool with 16 threads
 Dispatcher::Dispatcher()
-  //. create a queue of up to 64 tasks 
-  //. and a thread pool with 16 threads
-  : notifier(new SignalNotifier),
-    tasks(64),
-    workers(tasks, acceptor, 4),
-    server(&Dispatcher::run, this)
+  : my_tasks(64),
+    my_workers(my_tasks, my_acceptor, 4),
+    my_server(&Dispatcher::run, this)
 {
   Signal::mask(Signal::pipe);
-//   Signal::set(Signal::hangup, notifier);
-//   Signal::set(Signal::interrupt, notifier);
-//   Signal::set(Signal::quit, notifier);
-//   Signal::set(Signal::illegal, notifier);
-//   Signal::set(Signal::abort, notifier);
-//   Signal::set(Signal::fpe, notifier);
-//   Signal::set(Signal::bus, notifier);
-//   //  Signal::set(Signal::segv, notifier);
-//   Signal::set(Signal::iotrap, notifier);
-//   Signal::set(Signal::terminate, notifier);
 }
 
 Dispatcher::~Dispatcher()
@@ -83,62 +74,69 @@
 }
 
 void Dispatcher::bind(Agent *agent, int fd, Agent::iomask mask)
+  throw(std::invalid_argument)
 {
   Trace trace("Dispatcher::bind");
-  if (server.state() != Thread::running)
+  if (my_server.state() == Thread::READY)
+  {
+    pipe(my_wakeup);
+    my_rfds.set(my_wakeup[0]);
+    my_server.start();
+  }
+  Prague::Guard<Mutex> guard(my_mutex);
+  if (find(my_agents.begin(), my_agents.end(), agent) == my_agents.end())
+  {
+    my_agents.push_back(agent);
+    agent->add_ref();
+  }
+  if (mask & Agent::in)
+  {
+    if (mask & Agent::inready)
     {
-      pipe(wakeup);
-      rfds.set(wakeup[0]);
-      server.start();
+      my_wfds.set(fd);
+      if (my_wchannel.find(fd) == my_wchannel.end())
+	my_wchannel[fd] = new task(fd, agent, Agent::inready);
+      else throw std::invalid_argument("file descriptor already in use");
     }
-  Prague::Guard<Mutex> guard(mutex);
-  if (find(agents.begin(), agents.end(), agent) == agents.end())
+    if (mask & Agent::inexc)
     {
-      agents.push_back(agent);
-      agent->add_ref();
+      my_xfds.set(fd);
+      if (my_xchannel.find(fd) == my_xchannel.end())
+	my_xchannel[fd] = new task(fd, agent, Agent::inexc);
     }
-  if (mask & Agent::in)
+  }
+  if (mask & Agent::out)
+  {
+    if (mask & Agent::outready)
     {
-      if (mask & Agent::inready)
-	{
-	  wfds.set(fd);
-	  if (wchannel.find(fd) == wchannel.end()) wchannel[fd] = new task(fd, agent, Agent::inready);
-	  else std::cerr << "Dispatcher::bind() : Error : file descriptor already in use" << std::endl;
-	}
-      if (mask & Agent::inexc)
-	{
-	  xfds.set(fd);
-	  if (xchannel.find(fd) == xchannel.end()) xchannel[fd] = new task(fd, agent, Agent::inexc);
-	}
+      my_rfds.set(fd);
+      if (my_rchannel.find(fd) == my_rchannel.end())
+	my_rchannel[fd] = new task(fd, agent, Agent::outready);
+      else throw std::invalid_argument("file descriptor already in use");
     }
-  if (mask & Agent::out)
+    if (mask & Agent::outexc)
     {
-      if (mask & Agent::outready)
-	{
-	  rfds.set(fd);
-	  if (rchannel.find(fd) == rchannel.end()) rchannel[fd] = new task(fd, agent, Agent::outready);
-	  else std::cerr << "Dispatcher::bind() : Error : file descriptor already in use" << std::endl;
-	}
-      if (mask & Agent::outexc)
-	{
-	  xfds.set(fd);
-	  if (xchannel.find(fd) == xchannel.end()) xchannel[fd] = new task(fd, agent, Agent::outexc);
-	}
+      my_xfds.set(fd);
+      if (my_xchannel.find(fd) == my_xchannel.end())
+	my_xchannel[fd] = new task(fd, agent, Agent::outexc);
     }
+  }
   if (mask & Agent::err)
+  {
+    if (mask & Agent::errready)
     {
-      if (mask & Agent::errready)
-	{
-	  rfds.set(fd);
-	  if (rchannel.find(fd) == rchannel.end()) rchannel[fd] = new task(fd, agent, Agent::errready);
-	  else std::cerr << "Dispatcher::bind() : Error : file descriptor already in use" << std::endl;
-	}
-      if (mask & Agent::errexc)
-	{
-	  xfds.set(fd);
-	  if (xchannel.find(fd) == xchannel.end()) xchannel[fd] = new task(fd, agent, Agent::errexc);
-	}
+      my_rfds.set(fd);
+      if (my_rchannel.find(fd) == my_rchannel.end())
+	my_rchannel[fd] = new task(fd, agent, Agent::errready);
+      else throw std::invalid_argument("file descriptor already in use");
     }
+    if (mask & Agent::errexc)
+    {
+      my_xfds.set(fd);
+      if (my_xchannel.find(fd) == my_xchannel.end())
+	my_xchannel[fd] = new task(fd, agent, Agent::errexc);
+    }
+  }
   notify();
 }
 
@@ -148,50 +146,50 @@
   /*
    * release file descriptors
    */
-  Prague::Guard<Mutex> guard(mutex);
-  for (repository_t::iterator i = rchannel.begin(); i != rchannel.end(); i++)
+  Prague::Guard<Mutex> guard(my_mutex);
+  for (repository_t::iterator i = my_rchannel.begin(); i != my_rchannel.end(); i++)
     if ((*i).second->agent == agent && (fd == -1 || fd == (*i).second->fd))
-      {
-	deactivate((*i).second);
-	(*i).second->released = true;
-	rchannel.erase(i);
-      }
-  for (repository_t::iterator i = wchannel.begin(); i != wchannel.end(); i++)
+    {
+      deactivate((*i).second);
+      (*i).second->released = true;
+      my_rchannel.erase(i);
+    }
+  for (repository_t::iterator i = my_wchannel.begin(); i != my_wchannel.end(); i++)
     if ((*i).second->agent == agent && (fd == -1 || fd == (*i).second->fd))
-      {
-	deactivate((*i).second);
-	(*i).second->released = true;
-	wchannel.erase(i);
-      }
-  for (repository_t::iterator i = xchannel.begin(); i != xchannel.end(); i++)
+    {
+      deactivate((*i).second);
+      (*i).second->released = true;
+      my_wchannel.erase(i);
+    }
+  for (repository_t::iterator i = my_xchannel.begin(); i != my_xchannel.end(); i++)
     if ((*i).second->agent == agent && (fd == -1 || fd == (*i).second->fd))
-      {
-	deactivate((*i).second);
-	(*i).second->released = true;
-	xchannel.erase(i);
-      }
+    {
+      deactivate((*i).second);
+      (*i).second->released = true;
+      my_xchannel.erase(i);
+    }
   /*
    * release Agent if no more file descriptors left
    */
-  for (repository_t::iterator i = rchannel.begin(); i != rchannel.end(); i++)
+  for (repository_t::iterator i = my_rchannel.begin(); i != my_rchannel.end(); i++)
     if ((*i).second->agent == agent) return;
-  for (repository_t::iterator i = wchannel.begin(); i != wchannel.end(); i++)
+  for (repository_t::iterator i = my_wchannel.begin(); i != my_wchannel.end(); i++)
     if ((*i).second->agent == agent) return;
-  for (repository_t::iterator i = xchannel.begin(); i != xchannel.end(); i++)
+  for (repository_t::iterator i = my_xchannel.begin(); i != my_xchannel.end(); i++)
     if ((*i).second->agent == agent) return;
 
-  alist_t::iterator i = find(agents.begin(), agents.end(), agent);
-  if (i != agents.end())
-    {
-      agents.erase(i);
-      agent->remove_ref();
-    }
+  alist_t::iterator i = find(my_agents.begin(), my_agents.end(), agent);
+  if (i != my_agents.end())
+  {
+    my_agents.erase(i);
+    agent->remove_ref();
+  }
 }
 
 void *Dispatcher::run(void *X)
 {
   Dispatcher *dispatcher = reinterpret_cast<Dispatcher *>(X);
-  dispatcher->workers.start();
+  dispatcher->my_workers.start();
   do dispatcher->wait();
   while (true);
   return 0;
@@ -200,7 +198,7 @@
 void Dispatcher::dispatch(task *t)
 {
   deactivate(t);
-  tasks.push(t);
+  my_tasks.push(t);
 }
 
 void Dispatcher::process(task *t)
@@ -212,72 +210,73 @@
   bool done = !agent->process(t->fd, t->mask);
   agent->remove_ref();
   // now look whether the agent is released and the task should be deleted
-  Prague::Guard<Mutex> guard(mutex);
+  Prague::Guard<Mutex> guard(my_mutex);
   if (!done)
-    {
-      if (t->released) delete t;
-      else activate(t);
-    }
+  {
+    if (t->released) delete t;
+    else activate(t);
+  }
 }
 
 void Dispatcher::deactivate(task *t)
 {
   switch (t->mask)
-    {
-    case Agent::inready: wfds.clear(t->fd); break;
+  {
+    case Agent::inready: my_wfds.clear(t->fd); break;
     case Agent::outready:
-    case Agent::errready: rfds.clear(t->fd); break;
+    case Agent::errready: my_rfds.clear(t->fd); break;
     case Agent::inexc:
     case Agent::outexc:
-    case Agent::errexc: xfds.clear(t->fd); break;
+    case Agent::errexc: my_xfds.clear(t->fd); break;
     default: break;
-    }
+  }
 }
 
 void Dispatcher::activate(task *t)
 {
   switch (t->mask)
-    {
-    case Agent::inready: wfds.set(t->fd); break;
+  {
+    case Agent::inready: my_wfds.set(t->fd); break;
     case Agent::outready:
-    case Agent::errready: rfds.set(t->fd); break;
+    case Agent::errready: my_rfds.set(t->fd); break;
     case Agent::inexc:
     case Agent::outexc:
-    case Agent::errexc: xfds.set(t->fd); break;
+    case Agent::errexc: my_xfds.set(t->fd); break;
     default: break;
-    }
+  }
   notify();
 }
 
 void Dispatcher::wait()
 {
   Trace trace("Dispatcher::wait");
-  FdSet tmprfds = rfds;
-  FdSet tmpwfds = wfds;
-  FdSet tmpxfds = xfds;
-  unsigned int fdsize = std::max(std::max(tmprfds.max(), tmpwfds.max()), tmpxfds.max()) + 1;
+  FdSet tmprfds = my_rfds;
+  FdSet tmpwfds = my_wfds;
+  FdSet tmpxfds = my_xfds;
+  unsigned int fdsize = std::max(std::max(tmprfds.max(), tmpwfds.max()),
+				 tmpxfds.max()) + 1;
   int nsel = select(fdsize, tmprfds, tmpwfds, tmpxfds, 0);
-  pthread_testcancel();
+  Thread::testcancel();
   if (nsel == -1)
-    {
-      if (errno == EINTR || errno == EAGAIN) errno = 0;
-    }
+  {
+    if (errno == EINTR || errno == EAGAIN) errno = 0;
+  }
   else if (nsel > 0 && fdsize)
+  {
+    Prague::Guard<Mutex> guard(my_mutex);
+    for (repository_t::iterator i = my_rchannel.begin(); i != my_rchannel.end(); i++)
+      if (tmprfds.isset((*i).first))
+	dispatch((*i).second);
+    for (repository_t::iterator i = my_wchannel.begin(); i != my_wchannel.end(); i++)
+      if (tmpwfds.isset((*i).first))
+	dispatch((*i).second);
+    for (repository_t::iterator i = my_xchannel.begin(); i != my_xchannel.end(); i++)
+      if (tmpxfds.isset((*i).first))
+	dispatch((*i).second);
+    if (tmprfds.isset(my_wakeup[0]))
     {
-      Prague::Guard<Mutex> guard(mutex);
-      for (repository_t::iterator i = rchannel.begin(); i != rchannel.end(); i++)
-	if (tmprfds.isset((*i).first))
-	  dispatch((*i).second);
-      for (repository_t::iterator i = wchannel.begin(); i != wchannel.end(); i++)
-	if (tmpwfds.isset((*i).first))
-	  dispatch((*i).second);
-      for (repository_t::iterator i = xchannel.begin(); i != xchannel.end(); i++)
-	if (tmpxfds.isset((*i).first))
-	  dispatch((*i).second);
-      if (tmprfds.isset(wakeup[0]))
-	{
-	  char c[1];
-	  read(wakeup[0], c, 1);
-	}
+      char c[1];
+      read(my_wakeup[0], c, 1);
     }
+  }
 }