@@ -444,38 +444,7 @@ class MessageContext
444444 // / Schedule a context object for sending.
445445 // / The object is considered complete at this point and is sent directly through the dispatcher callback
446446 // / of the context if initialized.
447- void schedule (Messages::value_type&& message)
448- {
449- auto const * header = message->header ();
450- if (header == nullptr ) {
451- throw std::logic_error (" No valid header message found" );
452- }
453- mScheduledMessages .emplace_back (std::move (message));
454- if (mDispatchControl .dispatch != nullptr ) {
455- // send all scheduled messages if there is no trigger callback or its result is true
456- if (mDispatchControl .trigger == nullptr || mDispatchControl .trigger (*header)) {
457- std::vector<fair::mq::Parts> outputsPerChannel;
458- outputsPerChannel.resize (mProxy .getNumOutputChannels ());
459- for (auto & message : mScheduledMessages ) {
460- fair::mq::Parts parts = message->finalize ();
461- assert (message->empty ());
462- assert (parts.Size () == 2 );
463- for (auto & part : parts) {
464- outputsPerChannel[mProxy .getOutputChannelIndex (message->route ()).value ].AddPart (std::move (part));
465- }
466- }
467- for (int ci = 0 ; ci < mProxy .getNumOutputChannels (); ++ci) {
468- auto & parts = outputsPerChannel[ci];
469- if (parts.Size () == 0 ) {
470- continue ;
471- }
472- mDispatchControl .dispatch (std::move (parts), ChannelIndex{ci}, DefaultChannelIndex);
473- }
474- mDidDispatch = mScheduledMessages .empty () == false ;
475- mScheduledMessages .clear ();
476- }
477- }
478- }
447+ void schedule (Messages::value_type&& message);
479448
480449 Messages getMessagesForSending ()
481450 {
0 commit comments