| |
| <h1>broker.cpp</h1> |
| <div class="highlight"><pre><span></span><span class="cp">#include</span> <span class="cpf">"options.hpp"</span><span class="cp"></span> |
| |
| <span class="cp">#include</span> <span class="cpf"><proton/connection.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/connection_options.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/container.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/delivery.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/error_condition.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/listen_handler.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/listener.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/message.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/messaging_handler.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/receiver_options.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/sender_options.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/source_options.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/target.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/target_options.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/tracker.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/transport.hpp></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><proton/work_queue.hpp></span><span class="cp"></span> |
| |
| <span class="cp">#include</span> <span class="cpf"><deque></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><iostream></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><map></span><span class="cp"></span> |
| <span class="cp">#include</span> <span class="cpf"><string></span><span class="cp"></span> |
| |
| <span class="cp">#if PN_CPP_SUPPORTS_THREADS</span> |
| <span class="cp">#include</span> <span class="cpf"><thread></span><span class="cp"></span> |
| <span class="cp">#endif</span> |
| |
| <span class="cp">#include</span> <span class="cpf">"fake_cpp11.hpp"</span><span class="cp"></span> |
| |
| <span class="c1">// This is a simplified model for a message broker, that only allows for messages to go to a</span> |
| <span class="c1">// single receiver.</span> |
| <span class="c1">//</span> |
| <span class="c1">// This broker is multithread safe and if compiled with C++11 with a multithreaded Proton</span> |
| <span class="c1">// binding library will use as many threads as there are thread resources available (usually</span> |
| <span class="c1">// cores)</span> |
| <span class="c1">//</span> |
| <span class="c1">// Queues are only created and never destroyed</span> |
| <span class="c1">//</span> |
| <span class="c1">// Broker Entities (that need to be individually serialised)</span> |
| <span class="c1">// QueueManager - Creates new queues, finds queues</span> |
| <span class="c1">// Queue - Queues msgs, records subscribers, sends msgs to subscribers</span> |
| <span class="c1">// Connection - Receives Messages from network, sends messages to network.</span> |
| |
| <span class="c1">// Work</span> |
| <span class="c1">// FindQueue(queueName, connection) - From a Connection to the QueueManager</span> |
| <span class="c1">// This will create the queue if it doesn't already exist and send a BoundQueue</span> |
| <span class="c1">// message back to the connection.</span> |
| <span class="c1">// BoundQueue(queue) - From the QueueManager to a Connection</span> |
| <span class="c1">//</span> |
| <span class="c1">// QueueMsg(msg) - From a Connection (receiver) to a Queue</span> |
| <span class="c1">// Subscribe(sender) - From a Connection (sender) to a Queue</span> |
| <span class="c1">// Flow(sender, credit) - From a Connection (sender) to a Queue</span> |
| <span class="c1">// Unsubscribe(sender) - From a Connection (sender) to a Queue</span> |
| <span class="c1">//</span> |
| <span class="c1">// SendMsg(msg) - From a Queue to a Connection (sender)</span> |
| <span class="c1">// Unsubscribed() - From a Queue to a Connection (sender)</span> |
| |
| |
| <span class="c1">// Simple debug output</span> |
| <span class="kt">bool</span> <span class="n">verbose</span><span class="p">;</span> |
| <span class="cp">#define DOUT(x) do {if (verbose) {x};} while (false)</span> |
| |
| <span class="k">class</span> <span class="nc">Queue</span><span class="p">;</span> |
| <span class="k">class</span> <span class="nc">Sender</span><span class="p">;</span> |
| |
| <span class="k">typedef</span> <span class="n">std</span><span class="o">::</span><span class="n">map</span><span class="o"><</span><span class="n">proton</span><span class="o">::</span><span class="n">sender</span><span class="p">,</span> <span class="n">Sender</span><span class="o">*></span> <span class="n">senders</span><span class="p">;</span> |
| |
| <span class="k">class</span> <span class="nc">Sender</span> <span class="o">:</span> <span class="k">public</span> <span class="n">proton</span><span class="o">::</span><span class="n">messaging_handler</span> <span class="p">{</span> |
| <span class="k">friend</span> <span class="k">class</span> <span class="nc">connection_handler</span><span class="p">;</span> |
| |
| <span class="n">proton</span><span class="o">::</span><span class="n">sender</span> <span class="n">sender_</span><span class="p">;</span> |
| <span class="n">senders</span><span class="o">&</span> <span class="n">senders_</span><span class="p">;</span> |
| <span class="n">proton</span><span class="o">::</span><span class="n">work_queue</span><span class="o">&</span> <span class="n">work_queue_</span><span class="p">;</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">queue_name_</span><span class="p">;</span> |
| <span class="n">Queue</span><span class="o">*</span> <span class="n">queue_</span><span class="p">;</span> |
| <span class="kt">int</span> <span class="n">pending_credit_</span><span class="p">;</span> |
| |
| <span class="c1">// Messaging handlers</span> |
| <span class="kt">void</span> <span class="nf">on_sendable</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender</span> <span class="o">&</span><span class="n">sender</span><span class="p">)</span> <span class="n">OVERRIDE</span><span class="p">;</span> |
| <span class="kt">void</span> <span class="nf">on_sender_close</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender</span> <span class="o">&</span><span class="n">sender</span><span class="p">)</span> <span class="n">OVERRIDE</span><span class="p">;</span> |
| |
| <span class="k">public</span><span class="o">:</span> |
| <span class="n">Sender</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender</span> <span class="n">s</span><span class="p">,</span> <span class="n">senders</span><span class="o">&</span> <span class="n">ss</span><span class="p">)</span> <span class="o">:</span> |
| <span class="n">sender_</span><span class="p">(</span><span class="n">s</span><span class="p">),</span> <span class="n">senders_</span><span class="p">(</span><span class="n">ss</span><span class="p">),</span> <span class="n">work_queue_</span><span class="p">(</span><span class="n">s</span><span class="p">.</span><span class="n">work_queue</span><span class="p">()),</span> <span class="n">queue_</span><span class="p">(</span><span class="mi">0</span><span class="p">),</span> <span class="n">pending_credit_</span><span class="p">(</span><span class="mi">0</span><span class="p">)</span> |
| <span class="p">{}</span> |
| |
| <span class="kt">bool</span> <span class="n">add</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">work</span> <span class="n">f</span><span class="p">)</span> <span class="p">{</span> |
| <span class="k">return</span> <span class="n">work_queue_</span><span class="p">.</span><span class="n">add</span><span class="p">(</span><span class="n">f</span><span class="p">);</span> |
| <span class="p">}</span> |
| |
| |
| <span class="kt">void</span> <span class="n">boundQueue</span><span class="p">(</span><span class="n">Queue</span><span class="o">*</span> <span class="n">q</span><span class="p">,</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">qn</span><span class="p">);</span> |
| <span class="kt">void</span> <span class="nf">sendMsg</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">message</span> <span class="n">m</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Sender: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">" sending</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="n">sender_</span><span class="p">.</span><span class="n">send</span><span class="p">(</span><span class="n">m</span><span class="p">);</span> |
| <span class="p">}</span> |
| <span class="kt">void</span> <span class="nf">unsubscribed</span><span class="p">()</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Sender: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">" deleting</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="k">delete</span> <span class="k">this</span><span class="p">;</span> |
| <span class="p">}</span> |
| <span class="p">};</span> |
| |
| <span class="c1">// Queue - round robin subscriptions</span> |
| <span class="k">class</span> <span class="nc">Queue</span> <span class="p">{</span> |
| <span class="n">proton</span><span class="o">::</span><span class="n">work_queue</span> <span class="n">work_queue_</span><span class="p">;</span> |
| <span class="k">const</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">name_</span><span class="p">;</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">deque</span><span class="o"><</span><span class="n">proton</span><span class="o">::</span><span class="n">message</span><span class="o">></span> <span class="n">messages_</span><span class="p">;</span> |
| <span class="k">typedef</span> <span class="n">std</span><span class="o">::</span><span class="n">map</span><span class="o"><</span><span class="n">Sender</span><span class="o">*</span><span class="p">,</span> <span class="kt">int</span><span class="o">></span> <span class="n">subscriptions</span><span class="p">;</span> <span class="c1">// With credit</span> |
| <span class="n">subscriptions</span> <span class="n">subscriptions_</span><span class="p">;</span> |
| <span class="n">subscriptions</span><span class="o">::</span><span class="n">iterator</span> <span class="n">current_</span><span class="p">;</span> |
| |
| <span class="kt">void</span> <span class="nf">tryToSend</span><span class="p">()</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Queue: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">" tryToSend: "</span> <span class="o"><<</span> <span class="n">subscriptions_</span><span class="p">.</span><span class="n">size</span><span class="p">(););</span> |
| <span class="c1">// Starting at current_, send messages to subscriptions with credit:</span> |
| <span class="c1">// After each send try to find another subscription; Wrap around;</span> |
| <span class="c1">// Finish when we run out of messages or credit.</span> |
| <span class="kt">size_t</span> <span class="n">outOfCredit</span> <span class="o">=</span> <span class="mi">0</span><span class="p">;</span> |
| <span class="k">while</span> <span class="p">(</span><span class="o">!</span><span class="n">messages_</span><span class="p">.</span><span class="n">empty</span><span class="p">()</span> <span class="o">&&</span> <span class="n">outOfCredit</span><span class="o"><</span><span class="n">subscriptions_</span><span class="p">.</span><span class="n">size</span><span class="p">())</span> <span class="p">{</span> |
| <span class="c1">// If we got the end (or haven't started yet) start at the beginning</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">current_</span><span class="o">==</span><span class="n">subscriptions_</span><span class="p">.</span><span class="n">end</span><span class="p">())</span> <span class="p">{</span> |
| <span class="n">current_</span><span class="o">=</span><span class="n">subscriptions_</span><span class="p">.</span><span class="n">begin</span><span class="p">();</span> |
| <span class="p">}</span> |
| <span class="c1">// If we have credit send the message</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"("</span> <span class="o"><<</span> <span class="n">current_</span><span class="o">-></span><span class="n">second</span> <span class="o"><<</span> <span class="s">") "</span><span class="p">;);</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">current_</span><span class="o">-></span><span class="n">second</span><span class="o">></span><span class="mi">0</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="n">current_</span><span class="o">-></span><span class="n">first</span> <span class="o"><<</span> <span class="s">" "</span><span class="p">;);</span> |
| <span class="n">current_</span><span class="o">-></span><span class="n">first</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Sender</span><span class="o">::</span><span class="n">sendMsg</span><span class="p">,</span> <span class="n">current_</span><span class="o">-></span><span class="n">first</span><span class="p">,</span> <span class="n">messages_</span><span class="p">.</span><span class="n">front</span><span class="p">()));</span> |
| <span class="n">messages_</span><span class="p">.</span><span class="n">pop_front</span><span class="p">();</span> |
| <span class="o">--</span><span class="n">current_</span><span class="o">-></span><span class="n">second</span><span class="p">;</span> |
| <span class="o">++</span><span class="n">current_</span><span class="p">;</span> |
| <span class="p">}</span> <span class="k">else</span> <span class="p">{</span> |
| <span class="o">++</span><span class="n">outOfCredit</span><span class="p">;</span> |
| <span class="p">}</span> |
| <span class="p">}</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="p">}</span> |
| |
| <span class="k">public</span><span class="o">:</span> |
| <span class="n">Queue</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">container</span><span class="o">&</span> <span class="n">c</span><span class="p">,</span> <span class="k">const</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span><span class="o">&</span> <span class="n">n</span><span class="p">)</span> <span class="o">:</span> |
| <span class="n">work_queue_</span><span class="p">(</span><span class="n">c</span><span class="p">),</span> <span class="n">name_</span><span class="p">(</span><span class="n">n</span><span class="p">),</span> <span class="n">current_</span><span class="p">(</span><span class="n">subscriptions_</span><span class="p">.</span><span class="n">end</span><span class="p">())</span> |
| <span class="p">{}</span> |
| |
| <span class="kt">bool</span> <span class="n">add</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">work</span> <span class="n">f</span><span class="p">)</span> <span class="p">{</span> |
| <span class="k">return</span> <span class="n">work_queue_</span><span class="p">.</span><span class="n">add</span><span class="p">(</span><span class="n">f</span><span class="p">);</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">queueMsg</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">message</span> <span class="n">m</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Queue: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">"("</span> <span class="o"><<</span> <span class="n">name_</span> <span class="o"><<</span> <span class="s">") queueMsg</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="n">messages_</span><span class="p">.</span><span class="n">push_back</span><span class="p">(</span><span class="n">m</span><span class="p">);</span> |
| <span class="n">tryToSend</span><span class="p">();</span> |
| <span class="p">}</span> |
| <span class="kt">void</span> <span class="n">flow</span><span class="p">(</span><span class="n">Sender</span><span class="o">*</span> <span class="n">s</span><span class="p">,</span> <span class="kt">int</span> <span class="n">c</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Queue: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">"("</span> <span class="o"><<</span> <span class="n">name_</span> <span class="o"><<</span> <span class="s">") flow: "</span> <span class="o"><<</span> <span class="n">c</span> <span class="o"><<</span> <span class="s">" to "</span> <span class="o"><<</span> <span class="n">s</span> <span class="o"><<</span> <span class="s">"</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="n">subscriptions_</span><span class="p">[</span><span class="n">s</span><span class="p">]</span> <span class="o">=</span> <span class="n">c</span><span class="p">;</span> |
| <span class="n">tryToSend</span><span class="p">();</span> |
| <span class="p">}</span> |
| <span class="kt">void</span> <span class="n">subscribe</span><span class="p">(</span><span class="n">Sender</span><span class="o">*</span> <span class="n">s</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Queue: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">"("</span> <span class="o"><<</span> <span class="n">name_</span> <span class="o"><<</span> <span class="s">") subscribe Sender: "</span> <span class="o"><<</span> <span class="n">s</span> <span class="o"><<</span> <span class="s">"</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="n">subscriptions_</span><span class="p">[</span><span class="n">s</span><span class="p">]</span> <span class="o">=</span> <span class="mi">0</span><span class="p">;</span> |
| <span class="p">}</span> |
| <span class="kt">void</span> <span class="n">unsubscribe</span><span class="p">(</span><span class="n">Sender</span><span class="o">*</span> <span class="n">s</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Queue: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">"("</span> <span class="o"><<</span> <span class="n">name_</span> <span class="o"><<</span> <span class="s">") unsubscribe Sender: "</span> <span class="o"><<</span> <span class="n">s</span> <span class="o"><<</span> <span class="s">"</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="c1">// If we're about to erase the current subscription move on</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">current_</span> <span class="o">!=</span> <span class="n">subscriptions_</span><span class="p">.</span><span class="n">end</span><span class="p">()</span> <span class="o">&&</span> <span class="n">current_</span><span class="o">-></span><span class="n">first</span><span class="o">==</span><span class="n">s</span><span class="p">)</span> <span class="o">++</span><span class="n">current_</span><span class="p">;</span> |
| <span class="n">subscriptions_</span><span class="p">.</span><span class="n">erase</span><span class="p">(</span><span class="n">s</span><span class="p">);</span> |
| <span class="n">s</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Sender</span><span class="o">::</span><span class="n">unsubscribed</span><span class="p">,</span> <span class="n">s</span><span class="p">));</span> |
| <span class="p">}</span> |
| <span class="p">};</span> |
| |
| <span class="c1">// We have credit to send a message.</span> |
| <span class="kt">void</span> <span class="n">Sender</span><span class="o">::</span><span class="n">on_sendable</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender</span> <span class="o">&</span><span class="n">sender</span><span class="p">)</span> <span class="p">{</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">queue_</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">queue_</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Queue</span><span class="o">::</span><span class="n">flow</span><span class="p">,</span> <span class="n">queue_</span><span class="p">,</span> <span class="k">this</span><span class="p">,</span> <span class="n">sender</span><span class="p">.</span><span class="n">credit</span><span class="p">()));</span> |
| <span class="p">}</span> <span class="k">else</span> <span class="p">{</span> |
| <span class="n">pending_credit_</span> <span class="o">=</span> <span class="n">sender</span><span class="p">.</span><span class="n">credit</span><span class="p">();</span> |
| <span class="p">}</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">Sender</span><span class="o">::</span><span class="n">on_sender_close</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender</span> <span class="o">&</span><span class="n">sender</span><span class="p">)</span> <span class="p">{</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">queue_</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">queue_</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Queue</span><span class="o">::</span><span class="n">unsubscribe</span><span class="p">,</span> <span class="n">queue_</span><span class="p">,</span> <span class="k">this</span><span class="p">));</span> |
| <span class="p">}</span> <span class="k">else</span> <span class="p">{</span> |
| <span class="c1">// TODO: Is it possible to be closed before we get the queue allocated?</span> |
| <span class="c1">// If so, we should have a way to mark the sender deleted, so we can delete</span> |
| <span class="c1">// on queue binding</span> |
| <span class="p">}</span> |
| <span class="n">senders_</span><span class="p">.</span><span class="n">erase</span><span class="p">(</span><span class="n">sender</span><span class="p">);</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">Sender</span><span class="o">::</span><span class="n">boundQueue</span><span class="p">(</span><span class="n">Queue</span><span class="o">*</span> <span class="n">q</span><span class="p">,</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">qn</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Sender: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">" bound to Queue: "</span> <span class="o"><<</span> <span class="n">q</span> <span class="o"><<</span><span class="s">"("</span> <span class="o"><<</span> <span class="n">qn</span> <span class="o"><<</span> <span class="s">")</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="n">queue_</span> <span class="o">=</span> <span class="n">q</span><span class="p">;</span> |
| <span class="n">queue_name_</span> <span class="o">=</span> <span class="n">qn</span><span class="p">;</span> |
| |
| <span class="n">q</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Queue</span><span class="o">::</span><span class="n">subscribe</span><span class="p">,</span> <span class="n">q</span><span class="p">,</span> <span class="k">this</span><span class="p">));</span> |
| <span class="n">sender_</span><span class="p">.</span><span class="n">open</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender_options</span><span class="p">()</span> |
| <span class="p">.</span><span class="n">source</span><span class="p">((</span><span class="n">proton</span><span class="o">::</span><span class="n">source_options</span><span class="p">().</span><span class="n">address</span><span class="p">(</span><span class="n">queue_name_</span><span class="p">)))</span> |
| <span class="p">.</span><span class="n">handler</span><span class="p">(</span><span class="o">*</span><span class="k">this</span><span class="p">));</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">pending_credit_</span><span class="o">></span><span class="mi">0</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">queue_</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Queue</span><span class="o">::</span><span class="n">flow</span><span class="p">,</span> <span class="n">queue_</span><span class="p">,</span> <span class="k">this</span><span class="p">,</span> <span class="n">pending_credit_</span><span class="p">));</span> |
| <span class="p">}</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cout</span> <span class="o"><<</span> <span class="s">"sending from "</span> <span class="o"><<</span> <span class="n">queue_name_</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span><span class="p">;</span> |
| <span class="p">}</span> |
| |
| <span class="k">class</span> <span class="nc">Receiver</span> <span class="o">:</span> <span class="k">public</span> <span class="n">proton</span><span class="o">::</span><span class="n">messaging_handler</span> <span class="p">{</span> |
| <span class="k">friend</span> <span class="k">class</span> <span class="nc">connection_handler</span><span class="p">;</span> |
| |
| <span class="n">proton</span><span class="o">::</span><span class="n">receiver</span> <span class="n">receiver_</span><span class="p">;</span> |
| <span class="n">proton</span><span class="o">::</span><span class="n">work_queue</span><span class="o">&</span> <span class="n">work_queue_</span><span class="p">;</span> |
| <span class="n">Queue</span><span class="o">*</span> <span class="n">queue_</span><span class="p">;</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">deque</span><span class="o"><</span><span class="n">proton</span><span class="o">::</span><span class="n">message</span><span class="o">></span> <span class="n">messages_</span><span class="p">;</span> |
| |
| <span class="c1">// A message is received.</span> |
| <span class="kt">void</span> <span class="nf">on_message</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">delivery</span> <span class="o">&</span><span class="p">,</span> <span class="n">proton</span><span class="o">::</span><span class="n">message</span> <span class="o">&</span><span class="n">m</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="n">messages_</span><span class="p">.</span><span class="n">push_back</span><span class="p">(</span><span class="n">m</span><span class="p">);</span> |
| |
| <span class="k">if</span> <span class="p">(</span><span class="n">queue_</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">queueMsgs</span><span class="p">();</span> |
| <span class="p">}</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="nf">queueMsgs</span><span class="p">()</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Receiver: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">" queueing "</span> <span class="o"><<</span> <span class="n">messages_</span><span class="p">.</span><span class="n">size</span><span class="p">()</span> <span class="o"><<</span> <span class="s">" msgs to: "</span> <span class="o"><<</span> <span class="n">queue_</span> <span class="o"><<</span> <span class="s">"</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="k">while</span> <span class="p">(</span><span class="o">!</span><span class="n">messages_</span><span class="p">.</span><span class="n">empty</span><span class="p">())</span> <span class="p">{</span> |
| <span class="n">queue_</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Queue</span><span class="o">::</span><span class="n">queueMsg</span><span class="p">,</span> <span class="n">queue_</span><span class="p">,</span> <span class="n">messages_</span><span class="p">.</span><span class="n">front</span><span class="p">()));</span> |
| <span class="n">messages_</span><span class="p">.</span><span class="n">pop_front</span><span class="p">();</span> |
| <span class="p">}</span> |
| <span class="p">}</span> |
| |
| <span class="k">public</span><span class="o">:</span> |
| <span class="n">Receiver</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">receiver</span> <span class="n">r</span><span class="p">)</span> <span class="o">:</span> |
| <span class="n">receiver_</span><span class="p">(</span><span class="n">r</span><span class="p">),</span> <span class="n">work_queue_</span><span class="p">(</span><span class="n">r</span><span class="p">.</span><span class="n">work_queue</span><span class="p">()),</span> <span class="n">queue_</span><span class="p">(</span><span class="mi">0</span><span class="p">)</span> |
| <span class="p">{}</span> |
| |
| <span class="kt">bool</span> <span class="n">add</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">work</span> <span class="n">f</span><span class="p">)</span> <span class="p">{</span> |
| <span class="k">return</span> <span class="n">work_queue_</span><span class="p">.</span><span class="n">add</span><span class="p">(</span><span class="n">f</span><span class="p">);</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">boundQueue</span><span class="p">(</span><span class="n">Queue</span><span class="o">*</span> <span class="n">q</span><span class="p">,</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">qn</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"Receiver: "</span> <span class="o"><<</span> <span class="k">this</span> <span class="o"><<</span> <span class="s">" bound to Queue: "</span> <span class="o"><<</span> <span class="n">q</span> <span class="o"><<</span> <span class="s">"("</span> <span class="o"><<</span> <span class="n">qn</span> <span class="o"><<</span> <span class="s">")</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="n">queue_</span> <span class="o">=</span> <span class="n">q</span><span class="p">;</span> |
| <span class="n">receiver_</span><span class="p">.</span><span class="n">open</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">receiver_options</span><span class="p">()</span> |
| <span class="p">.</span><span class="n">source</span><span class="p">((</span><span class="n">proton</span><span class="o">::</span><span class="n">source_options</span><span class="p">().</span><span class="n">address</span><span class="p">(</span><span class="n">qn</span><span class="p">)))</span> |
| <span class="p">.</span><span class="n">handler</span><span class="p">(</span><span class="o">*</span><span class="k">this</span><span class="p">));</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cout</span> <span class="o"><<</span> <span class="s">"receiving to "</span> <span class="o"><<</span> <span class="n">qn</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span><span class="p">;</span> |
| |
| <span class="n">queueMsgs</span><span class="p">();</span> |
| <span class="p">}</span> |
| <span class="p">};</span> |
| |
| <span class="k">class</span> <span class="nc">QueueManager</span> <span class="p">{</span> |
| <span class="n">proton</span><span class="o">::</span><span class="n">container</span><span class="o">&</span> <span class="n">container_</span><span class="p">;</span> |
| <span class="n">proton</span><span class="o">::</span><span class="n">work_queue</span> <span class="n">work_queue_</span><span class="p">;</span> |
| <span class="k">typedef</span> <span class="n">std</span><span class="o">::</span><span class="n">map</span><span class="o"><</span><span class="n">std</span><span class="o">::</span><span class="n">string</span><span class="p">,</span> <span class="n">Queue</span><span class="o">*></span> <span class="n">queues</span><span class="p">;</span> |
| <span class="n">queues</span> <span class="n">queues_</span><span class="p">;</span> |
| <span class="kt">int</span> <span class="n">next_id_</span><span class="p">;</span> <span class="c1">// Use to generate unique queue IDs.</span> |
| |
| <span class="k">public</span><span class="o">:</span> |
| <span class="n">QueueManager</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">container</span><span class="o">&</span> <span class="n">c</span><span class="p">)</span> <span class="o">:</span> |
| <span class="n">container_</span><span class="p">(</span><span class="n">c</span><span class="p">),</span> <span class="n">work_queue_</span><span class="p">(</span><span class="n">c</span><span class="p">),</span> <span class="n">next_id_</span><span class="p">(</span><span class="mi">0</span><span class="p">)</span> |
| <span class="p">{}</span> |
| |
| <span class="kt">bool</span> <span class="n">add</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">work</span> <span class="n">f</span><span class="p">)</span> <span class="p">{</span> |
| <span class="k">return</span> <span class="n">work_queue_</span><span class="p">.</span><span class="n">add</span><span class="p">(</span><span class="n">f</span><span class="p">);</span> |
| <span class="p">}</span> |
| |
| <span class="k">template</span> <span class="o"><</span><span class="k">class</span> <span class="nc">T</span><span class="o">></span> |
| <span class="kt">void</span> <span class="n">findQueue</span><span class="p">(</span><span class="n">T</span><span class="o">&</span> <span class="n">connection</span><span class="p">,</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span><span class="o">&</span> <span class="n">qn</span><span class="p">)</span> <span class="p">{</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">qn</span><span class="p">.</span><span class="n">empty</span><span class="p">())</span> <span class="p">{</span> |
| <span class="c1">// Dynamic queue creation</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">ostringstream</span> <span class="n">os</span><span class="p">;</span> |
| <span class="n">os</span> <span class="o"><<</span> <span class="s">"_dynamic_"</span> <span class="o"><<</span> <span class="n">next_id_</span><span class="o">++</span><span class="p">;</span> |
| <span class="n">qn</span> <span class="o">=</span> <span class="n">os</span><span class="p">.</span><span class="n">str</span><span class="p">();</span> |
| <span class="p">}</span> |
| <span class="n">Queue</span><span class="o">*</span> <span class="n">q</span> <span class="o">=</span> <span class="mi">0</span><span class="p">;</span> |
| <span class="n">queues</span><span class="o">::</span><span class="n">iterator</span> <span class="n">i</span> <span class="o">=</span> <span class="n">queues_</span><span class="p">.</span><span class="n">find</span><span class="p">(</span><span class="n">qn</span><span class="p">);</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">i</span><span class="o">==</span><span class="n">queues_</span><span class="p">.</span><span class="n">end</span><span class="p">())</span> <span class="p">{</span> |
| <span class="n">q</span> <span class="o">=</span> <span class="k">new</span> <span class="n">Queue</span><span class="p">(</span><span class="n">container_</span><span class="p">,</span> <span class="n">qn</span><span class="p">);</span> |
| <span class="n">queues_</span><span class="p">[</span><span class="n">qn</span><span class="p">]</span> <span class="o">=</span> <span class="n">q</span><span class="p">;</span> |
| <span class="p">}</span> <span class="k">else</span> <span class="p">{</span> |
| <span class="n">q</span> <span class="o">=</span> <span class="n">i</span><span class="o">-></span><span class="n">second</span><span class="p">;</span> |
| <span class="p">}</span> |
| <span class="n">connection</span><span class="p">.</span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">T</span><span class="o">::</span><span class="n">boundQueue</span><span class="p">,</span> <span class="o">&</span><span class="n">connection</span><span class="p">,</span> <span class="n">q</span><span class="p">,</span> <span class="n">qn</span><span class="p">));</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">findQueueSender</span><span class="p">(</span><span class="n">Sender</span><span class="o">*</span> <span class="n">s</span><span class="p">,</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">qn</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">findQueue</span><span class="p">(</span><span class="o">*</span><span class="n">s</span><span class="p">,</span> <span class="n">qn</span><span class="p">);</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">findQueueReceiver</span><span class="p">(</span><span class="n">Receiver</span><span class="o">*</span> <span class="n">r</span><span class="p">,</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">qn</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">findQueue</span><span class="p">(</span><span class="o">*</span><span class="n">r</span><span class="p">,</span> <span class="n">qn</span><span class="p">);</span> |
| <span class="p">}</span> |
| <span class="p">};</span> |
| |
| <span class="k">class</span> <span class="nc">connection_handler</span> <span class="o">:</span> <span class="k">public</span> <span class="n">proton</span><span class="o">::</span><span class="n">messaging_handler</span> <span class="p">{</span> |
| <span class="n">QueueManager</span><span class="o">&</span> <span class="n">queue_manager_</span><span class="p">;</span> |
| <span class="n">senders</span> <span class="n">senders_</span><span class="p">;</span> |
| |
| <span class="k">public</span><span class="o">:</span> |
| <span class="n">connection_handler</span><span class="p">(</span><span class="n">QueueManager</span><span class="o">&</span> <span class="n">qm</span><span class="p">)</span> <span class="o">:</span> |
| <span class="n">queue_manager_</span><span class="p">(</span><span class="n">qm</span><span class="p">)</span> |
| <span class="p">{}</span> |
| |
| <span class="kt">void</span> <span class="n">on_connection_open</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">connection</span><span class="o">&</span> <span class="n">c</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="n">c</span><span class="p">.</span><span class="n">open</span><span class="p">();</span> <span class="c1">// Accept the connection</span> |
| <span class="p">}</span> |
| |
| <span class="c1">// A sender sends messages from a queue to a subscriber.</span> |
| <span class="kt">void</span> <span class="n">on_sender_open</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender</span> <span class="o">&</span><span class="n">sender</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">qn</span> <span class="o">=</span> <span class="n">sender</span><span class="p">.</span><span class="n">source</span><span class="p">().</span><span class="n">dynamic</span><span class="p">()</span> <span class="o">?</span> <span class="s">""</span> <span class="o">:</span> <span class="n">sender</span><span class="p">.</span><span class="n">source</span><span class="p">().</span><span class="n">address</span><span class="p">();</span> |
| <span class="n">Sender</span><span class="o">*</span> <span class="n">s</span> <span class="o">=</span> <span class="k">new</span> <span class="n">Sender</span><span class="p">(</span><span class="n">sender</span><span class="p">,</span> <span class="n">senders_</span><span class="p">);</span> |
| <span class="n">senders_</span><span class="p">[</span><span class="n">sender</span><span class="p">]</span> <span class="o">=</span> <span class="n">s</span><span class="p">;</span> |
| <span class="n">queue_manager_</span><span class="p">.</span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">QueueManager</span><span class="o">::</span><span class="n">findQueueSender</span><span class="p">,</span> <span class="o">&</span><span class="n">queue_manager_</span><span class="p">,</span> <span class="n">s</span><span class="p">,</span> <span class="n">qn</span><span class="p">));</span> |
| <span class="p">}</span> |
| |
| <span class="c1">// A receiver receives messages from a publisher to a queue.</span> |
| <span class="kt">void</span> <span class="n">on_receiver_open</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">receiver</span> <span class="o">&</span><span class="n">receiver</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">qname</span> <span class="o">=</span> <span class="n">receiver</span><span class="p">.</span><span class="n">target</span><span class="p">().</span><span class="n">address</span><span class="p">();</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">qname</span> <span class="o">==</span> <span class="s">"shutdown"</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cout</span> <span class="o"><<</span> <span class="s">"broker shutting down"</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span><span class="p">;</span> |
| <span class="c1">// Sending to the special "shutdown" queue stops the broker.</span> |
| <span class="n">receiver</span><span class="p">.</span><span class="n">connection</span><span class="p">().</span><span class="n">container</span><span class="p">().</span><span class="n">stop</span><span class="p">(</span> |
| <span class="n">proton</span><span class="o">::</span><span class="n">error_condition</span><span class="p">(</span><span class="s">"shutdown"</span><span class="p">,</span> <span class="s">"stop broker"</span><span class="p">));</span> |
| <span class="p">}</span> <span class="k">else</span> <span class="p">{</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">qname</span><span class="p">.</span><span class="n">empty</span><span class="p">())</span> <span class="p">{</span> |
| <span class="n">DOUT</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"ODD - trying to attach to a empty address</span><span class="se">\n</span><span class="s">"</span><span class="p">;);</span> |
| <span class="p">}</span> |
| <span class="n">Receiver</span><span class="o">*</span> <span class="n">r</span> <span class="o">=</span> <span class="k">new</span> <span class="n">Receiver</span><span class="p">(</span><span class="n">receiver</span><span class="p">);</span> |
| <span class="n">queue_manager_</span><span class="p">.</span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">QueueManager</span><span class="o">::</span><span class="n">findQueueReceiver</span><span class="p">,</span> <span class="o">&</span><span class="n">queue_manager_</span><span class="p">,</span> <span class="n">r</span><span class="p">,</span> <span class="n">qname</span><span class="p">));</span> |
| <span class="p">}</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">on_session_close</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">session</span> <span class="o">&</span><span class="n">session</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="c1">// Unsubscribe all senders that belong to session.</span> |
| <span class="k">for</span> <span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender_iterator</span> <span class="n">i</span> <span class="o">=</span> <span class="n">session</span><span class="p">.</span><span class="n">senders</span><span class="p">().</span><span class="n">begin</span><span class="p">();</span> <span class="n">i</span> <span class="o">!=</span> <span class="n">session</span><span class="p">.</span><span class="n">senders</span><span class="p">().</span><span class="n">end</span><span class="p">();</span> <span class="o">++</span><span class="n">i</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">senders</span><span class="o">::</span><span class="n">iterator</span> <span class="n">j</span> <span class="o">=</span> <span class="n">senders_</span><span class="p">.</span><span class="n">find</span><span class="p">(</span><span class="o">*</span><span class="n">i</span><span class="p">);</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">j</span> <span class="o">==</span> <span class="n">senders_</span><span class="p">.</span><span class="n">end</span><span class="p">())</span> <span class="k">continue</span><span class="p">;</span> |
| <span class="n">Sender</span><span class="o">*</span> <span class="n">s</span> <span class="o">=</span> <span class="n">j</span><span class="o">-></span><span class="n">second</span><span class="p">;</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">s</span><span class="o">-></span><span class="n">queue_</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">s</span><span class="o">-></span><span class="n">queue_</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Queue</span><span class="o">::</span><span class="n">unsubscribe</span><span class="p">,</span> <span class="n">s</span><span class="o">-></span><span class="n">queue_</span><span class="p">,</span> <span class="n">s</span><span class="p">));</span> |
| <span class="p">}</span> |
| <span class="n">senders_</span><span class="p">.</span><span class="n">erase</span><span class="p">(</span><span class="n">j</span><span class="p">);</span> |
| <span class="p">}</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">on_error</span><span class="p">(</span><span class="k">const</span> <span class="n">proton</span><span class="o">::</span><span class="n">error_condition</span><span class="o">&</span> <span class="n">e</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"error: "</span> <span class="o"><<</span> <span class="n">e</span><span class="p">.</span><span class="n">what</span><span class="p">()</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span><span class="p">;</span> |
| <span class="p">}</span> |
| |
| <span class="c1">// The container calls on_transport_close() last.</span> |
| <span class="kt">void</span> <span class="n">on_transport_close</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">transport</span><span class="o">&</span> <span class="n">t</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="c1">// Unsubscribe all senders.</span> |
| <span class="k">for</span> <span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">sender_iterator</span> <span class="n">i</span> <span class="o">=</span> <span class="n">t</span><span class="p">.</span><span class="n">connection</span><span class="p">().</span><span class="n">senders</span><span class="p">().</span><span class="n">begin</span><span class="p">();</span> <span class="n">i</span> <span class="o">!=</span> <span class="n">t</span><span class="p">.</span><span class="n">connection</span><span class="p">().</span><span class="n">senders</span><span class="p">().</span><span class="n">end</span><span class="p">();</span> <span class="o">++</span><span class="n">i</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">senders</span><span class="o">::</span><span class="n">iterator</span> <span class="n">j</span> <span class="o">=</span> <span class="n">senders_</span><span class="p">.</span><span class="n">find</span><span class="p">(</span><span class="o">*</span><span class="n">i</span><span class="p">);</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">j</span> <span class="o">==</span> <span class="n">senders_</span><span class="p">.</span><span class="n">end</span><span class="p">())</span> <span class="k">continue</span><span class="p">;</span> |
| <span class="n">Sender</span><span class="o">*</span> <span class="n">s</span> <span class="o">=</span> <span class="n">j</span><span class="o">-></span><span class="n">second</span><span class="p">;</span> |
| <span class="k">if</span> <span class="p">(</span><span class="n">s</span><span class="o">-></span><span class="n">queue_</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">s</span><span class="o">-></span><span class="n">queue_</span><span class="o">-></span><span class="n">add</span><span class="p">(</span><span class="n">make_work</span><span class="p">(</span><span class="o">&</span><span class="n">Queue</span><span class="o">::</span><span class="n">unsubscribe</span><span class="p">,</span> <span class="n">s</span><span class="o">-></span><span class="n">queue_</span><span class="p">,</span> <span class="n">s</span><span class="p">));</span> |
| <span class="p">}</span> |
| <span class="p">}</span> |
| <span class="k">delete</span> <span class="k">this</span><span class="p">;</span> <span class="c1">// All done.</span> |
| <span class="p">}</span> |
| <span class="p">};</span> |
| |
| <span class="k">class</span> <span class="nc">broker</span> <span class="p">{</span> |
| <span class="k">public</span><span class="o">:</span> |
| <span class="n">broker</span><span class="p">(</span><span class="k">const</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">addr</span><span class="p">)</span> <span class="o">:</span> |
| <span class="n">container_</span><span class="p">(</span><span class="s">"broker"</span><span class="p">),</span> <span class="n">queues_</span><span class="p">(</span><span class="n">container_</span><span class="p">),</span> <span class="n">listener_</span><span class="p">(</span><span class="n">queues_</span><span class="p">)</span> |
| <span class="p">{</span> |
| <span class="n">container_</span><span class="p">.</span><span class="n">listen</span><span class="p">(</span><span class="n">addr</span><span class="p">,</span> <span class="n">listener_</span><span class="p">);</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">run</span><span class="p">()</span> <span class="p">{</span> |
| <span class="cp">#if PN_CPP_SUPPORTS_THREADS</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cout</span> <span class="o"><<</span> <span class="s">"starting "</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="kr">thread</span><span class="o">::</span><span class="n">hardware_concurrency</span><span class="p">()</span> <span class="o"><<</span> <span class="s">" listening threads</span><span class="se">\n</span><span class="s">"</span><span class="p">;</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cout</span><span class="p">.</span><span class="n">flush</span><span class="p">();</span> |
| <span class="n">container_</span><span class="p">.</span><span class="n">run</span><span class="p">(</span><span class="n">std</span><span class="o">::</span><span class="kr">thread</span><span class="o">::</span><span class="n">hardware_concurrency</span><span class="p">());</span> |
| <span class="cp">#else</span> |
| <span class="n">container_</span><span class="p">.</span><span class="n">run</span><span class="p">();</span> |
| <span class="cp">#endif</span> |
| <span class="p">}</span> |
| |
| <span class="k">private</span><span class="o">:</span> |
| <span class="k">struct</span> <span class="nl">listener</span> <span class="p">:</span> <span class="k">public</span> <span class="n">proton</span><span class="o">::</span><span class="n">listen_handler</span> <span class="p">{</span> |
| <span class="n">listener</span><span class="p">(</span><span class="n">QueueManager</span><span class="o">&</span> <span class="n">c</span><span class="p">)</span> <span class="o">:</span> <span class="n">queues_</span><span class="p">(</span><span class="n">c</span><span class="p">)</span> <span class="p">{}</span> |
| |
| <span class="n">proton</span><span class="o">::</span><span class="n">connection_options</span> <span class="n">on_accept</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">listener</span><span class="o">&</span><span class="p">)</span> <span class="n">OVERRIDE</span><span class="p">{</span> |
| <span class="k">return</span> <span class="n">proton</span><span class="o">::</span><span class="n">connection_options</span><span class="p">().</span><span class="n">handler</span><span class="p">(</span><span class="o">*</span><span class="p">(</span><span class="k">new</span> <span class="n">connection_handler</span><span class="p">(</span><span class="n">queues_</span><span class="p">)));</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">on_open</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">listener</span><span class="o">&</span> <span class="n">l</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cout</span> <span class="o"><<</span> <span class="s">"broker listening on "</span> <span class="o"><<</span> <span class="n">l</span><span class="p">.</span><span class="n">port</span><span class="p">()</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span><span class="p">;</span> |
| <span class="p">}</span> |
| |
| <span class="kt">void</span> <span class="n">on_error</span><span class="p">(</span><span class="n">proton</span><span class="o">::</span><span class="n">listener</span><span class="o">&</span><span class="p">,</span> <span class="k">const</span> <span class="n">std</span><span class="o">::</span><span class="n">string</span><span class="o">&</span> <span class="n">s</span><span class="p">)</span> <span class="n">OVERRIDE</span> <span class="p">{</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"listen error: "</span> <span class="o"><<</span> <span class="n">s</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span><span class="p">;</span> |
| <span class="k">throw</span> <span class="n">std</span><span class="o">::</span><span class="n">runtime_error</span><span class="p">(</span><span class="n">s</span><span class="p">);</span> |
| <span class="p">}</span> |
| <span class="n">QueueManager</span><span class="o">&</span> <span class="n">queues_</span><span class="p">;</span> |
| <span class="p">};</span> |
| |
| <span class="n">proton</span><span class="o">::</span><span class="n">container</span> <span class="n">container_</span><span class="p">;</span> |
| <span class="n">QueueManager</span> <span class="n">queues_</span><span class="p">;</span> |
| <span class="n">listener</span> <span class="n">listener_</span><span class="p">;</span> |
| <span class="p">};</span> |
| |
| <span class="kt">int</span> <span class="nf">main</span><span class="p">(</span><span class="kt">int</span> <span class="n">argc</span><span class="p">,</span> <span class="kt">char</span> <span class="o">**</span><span class="n">argv</span><span class="p">)</span> <span class="p">{</span> |
| <span class="c1">// Command line options</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">string</span> <span class="n">address</span><span class="p">(</span><span class="s">"0.0.0.0"</span><span class="p">);</span> |
| <span class="n">example</span><span class="o">::</span><span class="n">options</span> <span class="n">opts</span><span class="p">(</span><span class="n">argc</span><span class="p">,</span> <span class="n">argv</span><span class="p">);</span> |
| |
| <span class="n">opts</span><span class="p">.</span><span class="n">add_flag</span><span class="p">(</span><span class="n">verbose</span><span class="p">,</span> <span class="sc">'v'</span><span class="p">,</span> <span class="s">"verbose"</span><span class="p">,</span> <span class="s">"verbose (debugging) output"</span><span class="p">);</span> |
| <span class="n">opts</span><span class="p">.</span><span class="n">add_value</span><span class="p">(</span><span class="n">address</span><span class="p">,</span> <span class="sc">'a'</span><span class="p">,</span> <span class="s">"address"</span><span class="p">,</span> <span class="s">"listen on URL"</span><span class="p">,</span> <span class="s">"URL"</span><span class="p">);</span> |
| |
| <span class="k">try</span> <span class="p">{</span> |
| <span class="n">verbose</span> <span class="o">=</span> <span class="nb">false</span><span class="p">;</span> |
| <span class="n">opts</span><span class="p">.</span><span class="n">parse</span><span class="p">();</span> |
| <span class="n">broker</span><span class="p">(</span><span class="n">address</span><span class="p">).</span><span class="n">run</span><span class="p">();</span> |
| <span class="k">return</span> <span class="mi">0</span><span class="p">;</span> |
| <span class="p">}</span> <span class="k">catch</span> <span class="p">(</span><span class="k">const</span> <span class="n">example</span><span class="o">::</span><span class="n">bad_option</span><span class="o">&</span> <span class="n">e</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cout</span> <span class="o"><<</span> <span class="n">opts</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span> <span class="o"><<</span> <span class="n">e</span><span class="p">.</span><span class="n">what</span><span class="p">()</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span><span class="p">;</span> |
| <span class="p">}</span> <span class="k">catch</span> <span class="p">(</span><span class="k">const</span> <span class="n">std</span><span class="o">::</span><span class="n">exception</span><span class="o">&</span> <span class="n">e</span><span class="p">)</span> <span class="p">{</span> |
| <span class="n">std</span><span class="o">::</span><span class="n">cerr</span> <span class="o"><<</span> <span class="s">"broker shutdown: "</span> <span class="o"><<</span> <span class="n">e</span><span class="p">.</span><span class="n">what</span><span class="p">()</span> <span class="o"><<</span> <span class="n">std</span><span class="o">::</span><span class="n">endl</span><span class="p">;</span> |
| <span class="p">}</span> |
| <span class="k">return</span> <span class="mi">1</span><span class="p">;</span> |
| <span class="p">}</span> |
| </pre></div> |
| |
| <p><a href="broker.cpp">Download this file</a></p> |