Subversion Repositories HomeAutomation

Rev

Rev 1309 | Go to most recent revision | Blame | Compare with Previous | Last modification | View Log | SVN | RSS feed

  1. /*
  2.  * Subscriber.cpp
  3.  *
  4.  *  Created on: Jul 21, 2009
  5.  *      Author: Mattias Runge
  6.  */
  7.  
  8. #include "Subscriber.h"
  9.  
  10. namespace atom {
  11. namespace broker {
  12.  
  13. Subscriber::Subscriber(boost::shared_ptr<Broker> broker, bool receiveFromMyself)
  14. {
  15.     LOG.setName("Subscriber");
  16.     this->myBroker = broker;
  17.     this->myReceiveFromMyself = receiveFromMyself;
  18.     this->myBroker->connect(boost::bind(&Subscriber::newMessageHandler, this, _1));
  19. }
  20.  
  21. Subscriber::~Subscriber()
  22. {
  23.     this->stop();
  24. }
  25.  
  26. void Subscriber::put(Message::pointer message)
  27. {
  28.     message->setOrigin(this);
  29.     this->myBroker->put(message);
  30. }
  31.  
  32. void Subscriber::onNewMessage(Message::pointer message)
  33. {
  34. }
  35.  
  36. void Subscriber::newMessageHandler(Message::pointer message)
  37. {
  38.     if (this->myReceiveFromMyself || !message->isOrigin(this))
  39.     {
  40.         this->myQueue.push(message);
  41.  
  42.         this->myCondition.notify_all();
  43.     }
  44. }
  45.  
  46. void Subscriber::run()
  47. {
  48.     while (true)
  49.     {
  50.         lock guard(this->myMutex);
  51.  
  52.         try
  53.         {
  54.             while (this->myQueue.size() > 0)
  55.             {
  56.                 this->onNewMessage(this->myQueue.pop());
  57.             }
  58.  
  59.             this->myCondition.wait(guard);
  60.         }
  61.         catch (boost::thread_interrupted e)
  62.         {
  63.             guard.unlock();
  64.             LOG.warn("Thread interrupted");
  65.             break;
  66.         }
  67.     }
  68. }
  69.  
  70. }
  71. }
  72.