Rev 1317 | Details | Compare with Previous | Last modification | View Log | SVN | RSS feed
| Rev | Author | Line No. | Line |
|---|---|---|---|
| 1309 | runge | 1 | /* |
| 2 | * Subscriber.cpp |
||
| 3 | * |
||
| 4 | * Created on: Jul 21, 2009 |
||
| 5 | * Author: Mattias Runge |
||
| 6 | */ |
||
| 7 | |||
| 8 | #include "Subscriber.h" |
||
| 9 | |||
| 1403 | runge | 10 | #include "broker/Broker.h" |
| 11 | |||
| 1309 | runge | 12 | namespace atom { |
| 13 | |||
| 1403 | runge | 14 | Subscriber::Subscriber(bool receive_from_self) |
| 1309 | runge | 15 | { |
| 16 | LOG.setName("Subscriber"); |
||
| 1403 | runge | 17 | |
| 18 | this->receive_from_self_ = receive_from_self; |
||
| 19 | |||
| 20 | Broker::GetInstance()->Connect(boost::bind(&Subscriber::Slot_OnMessage, this, _1)); |
||
| 1309 | runge | 21 | } |
| 22 | |||
| 23 | Subscriber::~Subscriber() |
||
| 24 | { |
||
| 1403 | runge | 25 | this->Stop(); |
| 1309 | runge | 26 | } |
| 27 | |||
| 1403 | runge | 28 | void Subscriber::Put(Message::Pointer message) |
| 1309 | runge | 29 | { |
| 1403 | runge | 30 | Broker::GetInstance()->Put(message); |
| 1309 | runge | 31 | } |
| 32 | |||
| 1403 | runge | 33 | Message::Pointer CreateMessage() |
| 1309 | runge | 34 | { |
| 1403 | runge | 35 | return Message::Pointer(new Message(static_cast<void *>this)); |
| 1309 | runge | 36 | } |
| 37 | |||
| 1403 | runge | 38 | void Subscriber::Slot_OnMessage(Message::Ă…ointer message) |
| 1309 | runge | 39 | { |
| 1403 | runge | 40 | if (this->receive_from_self_ || message->GetOrigin() != static_cast<void *>this) |
| 1309 | runge | 41 | { |
| 1403 | runge | 42 | this->queue_.push(message); |
| 43 | this->on_message_condition_.notify_all(); |
||
| 1309 | runge | 44 | } |
| 45 | } |
||
| 46 | |||
| 1403 | runge | 47 | void Subscriber::Run() |
| 1309 | runge | 48 | { |
| 49 | while (true) |
||
| 50 | { |
||
| 1403 | runge | 51 | boost::mutex::scoped_lock guard(this->guard_mutex_); |
| 1309 | runge | 52 | |
| 53 | try |
||
| 54 | { |
||
| 1403 | runge | 55 | while (this->queue_.size() > 0) |
| 1309 | runge | 56 | { |
| 1403 | runge | 57 | this->OnMessage(this->queue_.pop()); |
| 1309 | runge | 58 | } |
| 59 | |||
| 1403 | runge | 60 | this->on_message_condition_.wait(guard); |
| 1309 | runge | 61 | } |
| 62 | catch (boost::thread_interrupted e) |
||
| 63 | { |
||
| 64 | guard.unlock(); |
||
| 65 | LOG.warn("Thread interrupted"); |
||
| 66 | break; |
||
| 67 | } |
||
| 68 | } |
||
| 69 | } |
||
| 70 | |||
| 1403 | runge | 71 | } // namespace atom |