SObjectizer 5.8
Loading...
Searching...
No Matches
mchain_details.hpp
Go to the documentation of this file.
1/*
2 * SObjectizer-5
3 */
4
5/*!
6 * \since
7 * v.5.5.13
8 *
9 * \file
10 * \brief Implementation details for message chains.
11 */
12
13#pragma once
14
15#include <so_5/mchain.hpp>
16#include <so_5/mchain_select_ifaces.hpp>
17#include <so_5/environment.hpp>
18
19#include <so_5/ret_code.hpp>
20#include <so_5/exception.hpp>
21#include <so_5/error_logger.hpp>
22
23#include <so_5/details/abort_on_fatal_error.hpp>
24#include <so_5/details/at_scope_exit.hpp>
25#include <so_5/details/safe_cv_wait_for.hpp>
26
27#include <deque>
28#include <vector>
29#include <mutex>
30#include <condition_variable>
31
32namespace so_5 {
33
34namespace mchain_props {
35
36namespace details {
37
38//
39// ensure_queue_not_empty
40//
41/*!
42 * \since
43 * v.5.5.13
44 *
45 * \brief Helper function which throws an exception if queue is empty.
46 */
47template< typename Q >
48void
50 {
51 if( queue.is_empty() )
54 "an attempt to get message from empty demand queue" );
55 }
56
57//
58// ensure_queue_not_full
59//
60/*!
61 * \since
62 * v.5.5.13
63 *
64 * \brief Helper function which throws an exception if queue is full.
65 */
66template< typename Q >
67void
69 {
70 if( queue.is_full() )
73 "an attempt to push a message to full demand queue" );
74 }
75
76//
77// unlimited_demand_queue
78//
79/*!
80 * \since
81 * v.5.5.13
82 *
83 * \brief Implementation of demands queue for size-unlimited message chain.
84 */
86 {
87 public :
88 /*!
89 * \note This constructor is necessary just for a convinience.
90 */
92
93 //! Is queue full?
94 /*!
95 * \note Unlimited queue can't be null. Because of that this
96 * method always returns \a false.
97 */
98 [[nodiscard]]
99 bool
100 is_full() const noexcept { return false; }
101
102 //! Is queue empty?
103 [[nodiscard]]
104 bool
105 is_empty() const noexcept { return m_queue.empty(); }
106
107 //! Access to front item from queue.
108 [[nodiscard]]
109 demand_t &
111 {
113 return m_queue.front();
114 }
115
116 //! Remove the front item from queue.
117 void
119 {
121 m_queue.pop_front();
122 }
123
124 //! Add a new item to the end of the queue.
125 void
126 push_back( demand_t && demand )
127 {
128 m_queue.push_back( std::move(demand) );
129 }
130
131 //! Size of the queue.
132 [[nodiscard]]
133 std::size_t
134 size() const noexcept { return m_queue.size(); }
135
136 private :
137 //! Queue's storage.
138 std::deque< demand_t > m_queue;
139 };
140
141//
142// limited_dynamic_demand_queue
143//
144/*!
145 * \since
146 * v.5.5.13
147 *
148 * \brief Implementation of demands queue for size-limited message chain with
149 * dynamically allocated storage.
150 */
152 {
153 public :
154 //! Initializing constructor.
156 const capacity_t & capacity )
157 : m_max_size{ capacity.max_size() }
158 {}
159
160 //! Is queue full?
161 [[nodiscard]]
162 bool
163 is_full() const noexcept { return m_max_size == m_queue.size(); }
164
165 //! Is queue empty?
166 [[nodiscard]]
167 bool
168 is_empty() const noexcept { return m_queue.empty(); }
169
170 //! Access to front item from queue.
171 [[nodiscard]]
172 demand_t &
174 {
176 return m_queue.front();
177 }
178
179 //! Remove the front item from queue.
180 void
182 {
184 m_queue.pop_front();
185 }
186
187 //! Add a new item to the end of the queue.
188 void
189 push_back( demand_t && demand )
190 {
192 m_queue.push_back( std::move(demand) );
193 }
194
195 //! Size of the queue.
196 [[nodiscard]]
197 std::size_t
198 size() const noexcept { return m_queue.size(); }
199
200 private :
201 //! Queue's storage.
202 std::deque< demand_t > m_queue;
203 //! Maximum size of the queue.
204 const std::size_t m_max_size;
205 };
206
207//
208// limited_preallocated_demand_queue
209//
210/*!
211 * \since
212 * v.5.5.13
213 *
214 * \brief Implementation of demands queue for size-limited message chain with
215 * preallocated storage.
216 */
218 {
219 public :
220 //! Initializing constructor.
222 const capacity_t & capacity )
223 : m_storage( capacity.max_size(), demand_t{} )
224 , m_max_size{ capacity.max_size() }
225 , m_head{ 0 }
226 , m_size{ 0 }
227 {}
228
229 //! Is queue full?
230 [[nodiscard]]
231 bool
232 is_full() const noexcept { return m_max_size == m_size; }
233
234 //! Is queue empty?
235 [[nodiscard]]
236 bool
237 is_empty() const noexcept { return 0 == m_size; }
238
239 //! Access to front item from queue.
240 [[nodiscard]]
241 demand_t &
243 {
245 return m_storage[ m_head ];
246 }
247
248 //! Remove the front item from queue.
249 void
251 {
254 m_head = (m_head + 1) % m_max_size;
255 --m_size;
256 }
257
258 //! Add a new item to the end of the queue.
259 void
260 push_back( demand_t && demand )
261 {
263 auto index = (m_head + m_size) % m_max_size;
264 m_storage[ index ] = std::move(demand);
265 ++m_size;
266 }
267
268 //! Size of the queue.
269 [[nodiscard]]
270 std::size_t
271 size() const noexcept { return m_size; }
272
273 private :
274 //! Queue's storage.
275 std::vector< demand_t > m_storage;
276 //! Maximum size of the queue.
277 const std::size_t m_max_size;
278
279 //! Index of the queue head.
280 std::size_t m_head;
281 //! The current size of the queue.
282 std::size_t m_size;
283 };
284
285//
286// status
287//
288/*!
289 * \since
290 * v.5.5.13
291 *
292 * \brief Status of the message chain.
293 */
294enum class status
295 {
296 //! Bag is open and can be used for message sending.
297 open,
298 //! Bag is closed. New messages cannot be sent to it.
299 closed
300 };
301
302} /* namespace details */
303
304//
305// mchain_template
306//
307/*!
308 * \since
309 * v.5.5.13
310 *
311 * \brief Template-based implementation of message chain.
312 *
313 * \tparam Queue type of demand queue for message chain.
314 * \tparam Tracing_Base type with message tracing implementation details.
315 */
316template< typename Queue, typename Tracing_Base >
319 , private Tracing_Base
320 {
321 public :
322 //! Initializing constructor.
323 template< typename... Tracing_Args >
325 //! SObjectizer Environment for which message chain is created.
326 so_5::environment_t & env,
327 //! Mbox ID for this chain.
328 mbox_id_t id,
329 //! Chain parameters.
330 const mchain_params_t & params,
331 //! Arguments for Tracing_Base's constructor.
332 Tracing_Args &&... tracing_args )
333 : Tracing_Base( std::forward<Tracing_Args>(tracing_args)... )
334 , m_env( env )
335 , m_id( id )
336 , m_capacity( params.capacity() )
338 , m_queue( params.capacity() )
339 {}
340
341 mbox_id_t
342 id() const override
343 {
344 return m_id;
345 }
346
347 void
349 const std::type_index & /*msg_type*/,
350 abstract_message_sink_t & /*subscriber*/ ) override
351 {
354 "mchain doesn't support subscription" );
355 }
356
357 void
359 const std::type_index & /*msg_type*/,
360 abstract_message_sink_t & /*subscriber*/ ) noexcept override
361 {}
362
363 std::string
364 query_name() const override
365 {
366 std::ostringstream s;
367 s << "<mchain:id=" << m_id << ">";
368
369 return s.str();
370 }
371
373 type() const override
374 {
376 }
377
378 void
380 message_delivery_mode_t delivery_mode,
381 const std::type_index & msg_type,
382 const message_ref_t & message,
383 unsigned int /*redirection_deep*/ ) override
384 {
385 switch( delivery_mode )
386 {
388 this->try_to_store_message_to_queue_ordinary_mode(
389 msg_type,
390 message );
391 break;
392
394 this->try_to_store_message_to_queue_nonblocking_mode(
395 msg_type,
396 message );
397 break;
398 }
399 }
400
401 /*!
402 * \attention Will throw an exception because delivery
403 * filter is not applicable to MPSC-mboxes.
404 */
405 void
407 const std::type_index & /*msg_type*/,
408 const delivery_filter_t & /*filter*/,
409 abstract_message_sink_t & /*subscriber*/ ) override
410 {
413 "set_delivery_filter is called for mchain" );
414 }
415
416 void
418 const std::type_index & /*msg_type*/,
419 abstract_message_sink_t & /*subscriber*/ ) noexcept override
420 {}
421
422 [[nodiscard]]
425 demand_t & dest,
426 duration_t empty_queue_timeout ) override
427 {
428 std::unique_lock< std::mutex > lock{ m_lock };
429
430 // If queue is empty we must wait for some time.
431 bool queue_empty = m_queue.is_empty();
432 if( queue_empty )
433 {
435 // Waiting for new messages has no sence because
436 // chain is closed.
438
439 auto predicate = [this, &queue_empty]() -> bool {
440 queue_empty = m_queue.is_empty();
441 return !queue_empty ||
443 };
444
445 // Count of sleeping thread must be incremented before
446 // going to sleep and decremented right after.
448 auto decrement_threads = so_5::details::at_scope_exit(
449 [this] { --m_threads_to_wakeup; } );
450
451 // Wait until arrival of any message or closing of chain.
452 ::so_5::details::wait_for_big_interval(
453 lock,
455 empty_queue_timeout,
456 predicate );
457 }
458
459 // If queue is still empty nothing can be extracted and
460 // we must stop operation.
461 if( queue_empty )
463 // The chain is still open so there must be this result
465 // The chain is closed and there must be different result
467
469 }
470
471 bool
472 empty() const override
473 {
474 return m_queue.is_empty();
475 }
476
477 std::size_t
478 size() const override
479 {
480 return m_queue.size();
481 }
482
484 environment() const noexcept override
485 {
486 return m_env;
487 }
488
489 protected :
490 [[nodiscard]]
493 demand_t & dest,
494 select_case_t & select_case ) override
495 {
496 std::unique_lock< std::mutex > lock{ m_lock };
497
498 const bool queue_empty = m_queue.is_empty();
499 if( queue_empty )
500 {
502 // There is no need to wait for something.
504
505 // In other cases select_tail must be modified.
506 select_case.set_next( m_select_tail );
507 m_select_tail = &select_case;
508
510 }
511 else
513 }
514
515 [[nodiscard]]
518 const std::type_index & msg_type,
519 const message_ref_t & message,
520 mchain_props::select_case_t & select_case ) override
521 {
522 typename Tracing_Base::deliver_op_tracer tracer{
523 *this, // as tracing base.
524 *this, // as chain.
525 msg_type,
526 message };
527
528 std::unique_lock< std::mutex > lock{ m_lock };
529
530 // Message cannot be stored to closed chain.
533
534 if( m_queue.is_full() )
535 {
536 // The select_case should be stored until there will
537 // be a free space in the chain (or chain will be closed).
538 select_case.set_next( m_select_tail );
539 m_select_tail = &select_case;
540
542 }
543 else
544 {
545 // Just store a new message to the queue.
547 tracer,
548 msg_type,
549 message );
551 }
552 }
553
554 void
556 select_case_t & select_case ) noexcept override
557 {
558 std::lock_guard< std::mutex > lock{ m_lock };
559
561 select_case_t * prev = nullptr;
562 while( c )
563 {
564 select_case_t * const next = c->query_next();
565 if( c == &select_case )
566 {
567 if( prev )
568 prev->set_next( next );
569 else
570 m_select_tail = next;
571
572 return;
573 }
574
575 prev = c;
576 c = next;
577 }
578 }
579
580 void
581 actual_close( close_mode_t mode ) override
582 {
583 std::lock_guard< std::mutex > lock{ m_lock };
584
586 return;
587
589
590 const bool was_full = m_queue.is_full();
591
592 if( close_mode_t::drop_content == mode )
593 {
594 while( !m_queue.is_empty() )
595 {
596 this->trace_demand_drop_on_close(
597 *this, m_queue.front() );
598 m_queue.pop_front();
599 }
600 }
601
602 // Since v.5.7.0 select operations must be notified
603 // always, even if the mchain is not empty.
605
607 // Someone is waiting on empty chain for new messages.
608 // It must be informed that no new messages will be here.
609 m_underflow_cond.notify_all();
610
611 if( was_full )
612 // Someone can wait on full chain for free place for new message.
613 // It must be informed that the chain is closed.
614 m_overflow_cond.notify_all();
615 }
616
617 private :
618 //! SObjectizer Environment for which message chain is created.
620
621 //! Status of the chain.
623
624 //! Mbox ID for chain.
625 const mbox_id_t m_id;
626
627 //! Chain capacity.
629
630 //! Optional notificator for 'not_empty' condition.
631 const not_empty_notification_func_t m_not_empty_notificator;
632
633 //! Chain's demands queue.
634 Queue m_queue;
635
636 //! Chain's lock.
637 std::mutex m_lock;
638
639 //! Condition variable for waiting on empty queue.
640 std::condition_variable m_underflow_cond;
641 //! Condition variable for waiting on full queue.
642 std::condition_variable m_overflow_cond;
643
644 /*!
645 * \brief Count of threads sleeping on empty mchain.
646 *
647 * This value is incremented before sleeping on m_underflow_cond and
648 * decremented just after a return from this sleep.
649 *
650 * \since
651 * v.5.5.16
652 */
653 std::size_t m_threads_to_wakeup = { 0 };
654
655 /*!
656 * \brief A queue of multi-chain selects in which this chain is used.
657 *
658 * \since
659 * v.5.5.16
660 */
662
663 //! Actual implementation of pushing message to the queue.
664 /*!
665 * \note
666 * This implementation must be used for ordinary delivery operations.
667 * For delivery operations from timer thread another method must be
668 * called (see try_to_store_message_to_queue_nonblocking_mode()).
669 */
670 void
672 const std::type_index & msg_type,
673 const message_ref_t & message )
674 {
675 typename Tracing_Base::deliver_op_tracer tracer{
676 *this, // as tracing base.
677 *this, // as chain.
678 msg_type,
679 message };
680
681 std::unique_lock< std::mutex > lock{ m_lock };
682
683 // Message cannot be stored to closed chain.
685 return;
686
687 // If queue full and waiting on full queue is enabled we
688 // must wait for some time until there will be some space in
689 // the queue.
690 bool queue_full = m_queue.is_full();
692 {
693 ::so_5::details::wait_for_big_interval(
694 lock,
697 [this, &queue_full] {
698 queue_full = m_queue.is_full();
699 return !queue_full ||
701 } );
702
703 // Message cannot be stored to closed chain.
704 //
705 // NOTE: this additional check is necessary after
706 // wait for overflow_timeout because the chain can
707 // be closed during that wait.
709 return;
710 }
711
712 // If queue still full we must perform some reaction.
713 if( queue_full )
714 {
715 const auto reaction = m_capacity.overflow_reaction();
716 if( overflow_reaction_t::drop_newest == reaction )
717 {
718 // New message must be simply ignored.
719 tracer.overflow_drop_newest();
720 return;
721 }
722 else if( overflow_reaction_t::remove_oldest == reaction )
723 {
724 // The oldest message must be simply removed.
725 tracer.overflow_remove_oldest( m_queue.front() );
726 m_queue.pop_front();
727 }
728 else if( overflow_reaction_t::throw_exception == reaction )
729 {
730 tracer.overflow_throw_exception();
733 "an attempt to push message to full mchain "
734 "with overflow_reaction_t::throw_exception policy" );
735 }
736 else
737 {
738 so_5::details::abort_on_fatal_error( [&] {
739 tracer.overflow_throw_exception();
740 SO_5_LOG_ERROR( m_env, log_stream ) {
741 log_stream << "overflow_reaction_t::abort_app "
742 "will be performed for mchain (id="
743 << m_id << "), msg_type: "
744 << msg_type.name()
745 << ". Application will be aborted"
746 << std::endl;
747 }
748 } );
749 }
750 }
751
753 tracer,
754 msg_type,
755 message );
756 }
757
758 /*!
759 * \brief An implementation of storing another message to
760 * chain for the case of delated/periodic messages.
761 *
762 * This implementation handles overloaded chains differently:
763 * - there is no waiting on overloaded chain (even if such waiting
764 * is specified in mchain params);
765 * - overflow_reaction_t::throw_exception is replaced by
766 * overflow_reaction_t::drop_newest.
767 *
768 * These defferences are necessary because the context of timer
769 * thread is very special: there can't be any long-time operation
770 * (like waiting for free space on overloaded chain) and there can't
771 * be an exception about mchain's overflow.
772 *
773 * \since
774 * v.5.5.18
775 */
776 void
778 const std::type_index & msg_type,
779 const message_ref_t & message )
780 {
781 typename Tracing_Base::deliver_op_tracer tracer{
782 *this, // as tracing base.
783 *this, // as chain.
784 msg_type,
785 message };
786
787 std::unique_lock< std::mutex > lock{ m_lock };
788
789 // Message cannot be stored to closed chain.
791 return;
792
793 bool queue_full = m_queue.is_full();
794 // NOTE: there is no awaiting on full mchain.
795 // If queue full we must perform some reaction.
796 if( queue_full )
797 {
798 const auto reaction = m_capacity.overflow_reaction();
799 if( overflow_reaction_t::drop_newest == reaction ||
801 {
802 // New message must be simply ignored.
803 tracer.overflow_drop_newest();
804 return;
805 }
806 else if( overflow_reaction_t::remove_oldest == reaction )
807 {
808 // The oldest message must be simply removed.
809 tracer.overflow_remove_oldest( m_queue.front() );
810 m_queue.pop_front();
811 }
812 else
813 {
814 so_5::details::abort_on_fatal_error( [&] {
815 tracer.overflow_throw_exception();
816 SO_5_LOG_ERROR( m_env, log_stream ) {
817 log_stream << "overflow_reaction_t::abort_app "
818 "will be performed for mchain (id="
819 << m_id << "), msg_type: "
820 << msg_type.name()
821 << ". Application will be aborted"
822 << std::endl;
823 }
824 } );
825 }
826 }
827
829 tracer,
830 msg_type,
831 message );
832 }
833
834 /*!
835 * \brief Implementation of extract operation for the case when
836 * message queue is not empty.
837 *
838 * \attention This helper method must be called when chain object
839 * is locked in some hi-level method.
840 *
841 * \since
842 * v.5.5.16
843 */
846 demand_t & dest )
847 {
848 // If queue was full then someone can wait on it.
849 const bool queue_was_full = m_queue.is_full();
850 dest = std::move( m_queue.front() );
851 m_queue.pop_front();
852
853 this->trace_extracted_demand( *this, dest );
854
855 if( queue_was_full )
856 {
857 // Since v.5.7.0 waiting select_cases should be
858 // notified too because they are send_cases.
860
861 m_overflow_cond.notify_all();
862 }
863
865 }
866
867 /*!
868 * \since
869 * v.5.5.16
870 */
871 void
873 {
874 if( m_select_tail )
875 {
876 auto old = m_select_tail;
877 m_select_tail = nullptr;
878 old->notify();
879 }
880 }
881
882 /*!
883 * \brief A reusable method with implementation of
884 * last part of storing a message into chain.
885 *
886 * \note
887 * Intended to be called from try_to_store_message_to_queue_ordinary_mode()
888 * and try_to_store_message_to_queue_nonblocking_mode().
889 *
890 * \since
891 * v.5.5.18
892 */
893 void
895 typename Tracing_Base::deliver_op_tracer & tracer,
896 const std::type_index & msg_type,
897 const message_ref_t & message )
898 {
899 const bool was_empty = m_queue.is_empty();
900
901 m_queue.push_back( demand_t{ msg_type, message } );
902
903 tracer.stored( m_queue );
904
905 // If chain was empty then multi-chain cases must be notified.
906 // And if not_empty_notificator is defined then it must be used too.
907 if( was_empty )
908 {
910 so_5::details::invoke_noexcept_code(
911 [this] { m_not_empty_notificator(); } );
912
914 }
915
916 // Should be wake up some sleeping thread?
918 // Someone is waiting on empty queue.
919 m_underflow_cond.notify_one();
920 }
921 };
922
923} /* namespace mchain_props */
924
925} /* namespace so_5 */
An interace of message chain.
Definition mchain.hpp:447
Interface for message sink.
static bool special_sink_ptr_compare(const abstract_message_sink_t *a, const abstract_message_sink_t *b) noexcept
virtual void push_event(mbox_id_t mbox_id, message_delivery_mode_t delivery_mode, const std::type_index &msg_type, const message_ref_t &message, unsigned int redirection_deep, const message_limit::impl::action_msg_tracer_t *tracer)=0
Get a message and push it to the appropriate destination.
A base class for agents.
Definition agent.hpp:673
Interface for creator of new mbox in OOP style.
virtual mbox_t create(const mbox_creation_data_t &data)=0
Creation of custom mbox.
An interface of delivery filter object.
Definition mbox.hpp:62
SObjectizer Environment.
Mixin to be used in implementation of MPSC mbox with message limits.
Definition mpsc_mbox.hpp:42
agent_t & query_owner_reference() const noexcept
Definition mpsc_mbox.hpp:48
abstract_message_sink_t & message_sink_to_use(const local_mbox_details::subscription_info_with_sink_t &info) const noexcept
Definition mpsc_mbox.hpp:55
limitful_mpsc_mbox_mixin_t(outliving_reference_t< agent_t > owner)
Definition mpsc_mbox.hpp:62
Mixin to be used in implementation of MPSC mbox without message limits.
Definition mpsc_mbox.hpp:77
message_sink_without_message_limit_t m_actual_sink
Actual message sink to be used.
Definition mpsc_mbox.hpp:79
agent_t & query_owner_reference() const noexcept
Definition mpsc_mbox.hpp:84
abstract_message_sink_t & message_sink_to_use(const local_mbox_details::subscription_info_with_sink_t &) noexcept
Definition mpsc_mbox.hpp:91
limitless_mpsc_mbox_mixin_t(outliving_reference_t< agent_t > owner)
Definition mpsc_mbox.hpp:98
A special container for holding subscriber_info objects.
subscriber_adaptive_container_t(const subscriber_adaptive_container_t &o)
Copy constructor.
storage_type m_storage
The current storage type to be used by container.
subscriber_adaptive_container_t & operator=(subscriber_adaptive_container_t &&o) noexcept
Move operator.
void insert_to_vector(abstract_message_sink_t &sink_as_key, subscription_info_with_sink_t &&info)
Insertion of new item to vector.
void emplace(abstract_message_sink_t &sink_as_key, Args &&... args)
friend void swap(subscriber_adaptive_container_t &a, subscriber_adaptive_container_t &b) noexcept
vector_type m_vector
Container for small amount of subscriber_infos.
void insert(abstract_message_sink_t &sink_as_key, subscription_info_with_sink_t info)
void insert_to_map(abstract_message_sink_t &sink_as_key, subscription_info_with_sink_t &&info)
Insertion of new item to map.
subscriber_adaptive_container_t(subscriber_adaptive_container_t &&o) noexcept
Move constructor.
void switch_storage_to_map()
Switching storage from vector to map.
iterator find(abstract_message_sink_t &subscriber)
map_type m_map
Container for large amount of subscriber_infos.
iterator find_in_vector(abstract_message_sink_t &subscriber)
subscriber_adaptive_container_t & operator=(const subscriber_adaptive_container_t &o)
Copy operator.
iterator find_in_map(abstract_message_sink_t &subscriber)
void switch_storage_to_vector()
Switching storage from map to vector.
An information block about one subscription to one message type with presence of message_sink.
void set_filter(const delivery_filter_t &filter)
Set the delivery filter for the subscriber.
abstract_message_sink_t & sink_reference() const noexcept
Get a reference to the subscribed sink.
abstract_message_sink_t * sink_pointer() const noexcept
Get a pointer to the subscribed sink.
void set_sink(abstract_message_sink_t &sink)
Inform about addition of a subscription.
A template with implementation of local mbox.
mbox_type_t type() const override
Get the type of message box.
environment_t & environment() const noexcept override
SObjectizer Environment for which the mbox is created.
mbox_id_t id() const override
Unique ID of this mbox.
void do_deliver_message_impl(typename Tracing_Base::deliver_op_tracer const &tracer, message_delivery_mode_t delivery_mode, const std::type_index &msg_type, const message_ref_t &message, unsigned int redirection_deep)
void set_delivery_filter(const std::type_index &msg_type, const delivery_filter_t &filter, abstract_message_sink_t &subscriber) override
Set a delivery filter for message type and subscriber.
void do_deliver_message_to_subscriber(const local_mbox_details::subscription_info_with_sink_t &subscriber_info, typename Tracing_Base::deliver_op_tracer const &tracer, message_delivery_mode_t delivery_mode, const std::type_index &msg_type, const message_ref_t &message, unsigned int redirection_deep) const
void ensure_immutable_message(const std::type_index &msg_type, const message_ref_t &what) const
Ensures that message is an immutable message.
local_mbox_template(mbox_id_t id, environment_t &env, Tracing_Args &&... args)
void modify_and_remove_subscriber_if_needed(const std::type_index &type_wrapper, abstract_message_sink_t &subscriber, Info_Changer changer)
void do_deliver_message(message_delivery_mode_t delivery_mode, const std::type_index &msg_type, const message_ref_t &message, unsigned int redirection_deep) override
Deliver message for all subscribers with respect to message limits.
void subscribe_event_handler(const std::type_index &type_wrapper, abstract_message_sink_t &subscriber) override
Add the message handler.
void unsubscribe_event_handler(const std::type_index &type_wrapper, abstract_message_sink_t &subscriber) noexcept override
Remove all message handlers.
std::string query_name() const override
Get the mbox name.
void insert_or_modify_subscriber(const std::type_index &type_wrapper, abstract_message_sink_t &subscriber, Info_Maker maker, Info_Changer changer)
void drop_delivery_filter(const std::type_index &msg_type, abstract_message_sink_t &subscriber) noexcept override
Removes delivery filter for message type and subscriber.
mbox_t create_ordinary_mpsc_mbox(environment_t &env, agent_t &owner)
Create mpsc_mbox that handles message limits.
mbox_t create_limitless_mpsc_mbox(environment_t &env, agent_t &owner)
Create mpsc_mbox that ignores message limits.
named_mboxes_dictionary_t m_named_mboxes_dictionary
Named mboxes.
mbox_t introduce_named_mbox(mbox_namespace_name_t mbox_namespace, nonempty_name_t mbox_name, const std::function< mbox_t() > &mbox_factory)
Introduce named mbox with user-provided factory.
outliving_reference_t< so_5::msg_tracing::holder_t > m_msg_tracing_stuff
Data related to message delivery tracing.
mbox_core_t(outliving_reference_t< so_5::msg_tracing::holder_t > msg_tracing_stuff)
Definition mbox_core.cpp:27
std::mutex m_dictionary_lock
Named mbox map's lock.
mbox_t create_custom_mbox(environment_t &env, ::so_5::custom_mbox_details::creator_iface_t &creator)
Create a custom mbox.
std::atomic< mbox_id_t > m_mbox_id_counter
A counter for mbox ID generation.
void destroy_mbox(const full_named_mbox_id_t &name) noexcept
Remove a reference to the named mbox.
mchain_t create_mchain(environment_t &env, const mchain_params_t &params)
Create message chain.
mbox_t create_mbox(environment_t &env)
Create local anonymous mbox.
Definition mbox_core.cpp:35
mbox_t create_mbox(environment_t &env, nonempty_name_t mbox_name)
Create local named mbox.
Definition mbox_core.cpp:46
mbox_id_t allocate_mbox_id() noexcept
Allocate an ID for a new custom mbox or mchain.
mbox_core_stats_t query_stats()
Get statistics for run-time monitoring.
A base class for message sinks to be used by agents.
void push_event(mbox_id_t mbox_id, message_delivery_mode_t, const std::type_index &msg_type, const message_ref_t &message, unsigned int, const message_limit::impl::action_msg_tracer_t *tracer) override
Get a message and push it to the appropriate destination.
mpsc_mbox_template_t(mbox_id_t id, environment_t &env, outliving_reference_t< agent_t > owner, Tracing_Args &&... tracing_args)
mbox_type_t type() const override
Get the type of message box.
void modify_and_remove_subscription_if_needed(const std::type_index &msg_type, Info_Changer changer)
Helper for modification and deletion of subscription info.
const mbox_id_t m_id
ID of this mbox.
void unsubscribe_event_handler(const std::type_index &msg_type, abstract_message_sink_t &subscriber) noexcept override
Remove all message handlers.
void do_delivery(const std::type_index &msg_type, const message_ref_t &message, typename Tracing_Base::deliver_op_tracer const &tracer, L l)
Helper method to do delivery actions under locked object.
subscriptions_map_t m_subscriptions
Information about the current subscriptions.
default_rw_spinlock_t m_lock
Protection of object from modification.
void drop_delivery_filter(const std::type_index &msg_type, abstract_message_sink_t &) noexcept override
Removes delivery filter for message type and subscriber.
void set_delivery_filter(const std::type_index &msg_type, const delivery_filter_t &filter, abstract_message_sink_t &subscriber) override
Set a delivery filter for message type and subscriber.
std::string query_name() const override
Get the mbox name.
void do_deliver_message(message_delivery_mode_t delivery_mode, const std::type_index &msg_type, const message_ref_t &message, unsigned int redirection_deep) override
Deliver message for all subscribers with respect to message limits.
void subscribe_event_handler(const std::type_index &msg_type, abstract_message_sink_t &subscriber) override
Add the message handler.
mbox_id_t id() const override
Unique ID of this mbox.
environment_t & environment() const noexcept override
SObjectizer Environment for which the mbox is created.
void insert_or_modify_subscription(const std::type_index &msg_type, Info_Maker maker, Info_Changer changer)
Helper for performing insertion or modification of subscription info.
environment_t & m_env
Environment in that the mbox was created.
Base class for a mbox for the case when message delivery tracing is enabled.
void set_delivery_filter(const std::type_index &msg_type, const delivery_filter_t &filter, abstract_message_sink_t &subscriber) override
Set a delivery filter for message type and subscriber.
mbox_id_t id() const override
Unique ID of this mbox.
void subscribe_event_handler(const std::type_index &type_wrapper, abstract_message_sink_t &subscriber) override
Add the message handler.
named_local_mbox_t(full_named_mbox_id_t full_name, const mbox_t &mbox, impl::mbox_core_t &mbox_core)
mbox_type_t type() const override
Get the type of message box.
void do_deliver_message(message_delivery_mode_t delivery_mode, const std::type_index &msg_type, const message_ref_t &message, unsigned int redirection_deep) override
Deliver message for all subscribers with respect to message limits.
impl::mbox_core_ref_t m_mbox_core
An utility for this mbox.
environment_t & environment() const noexcept override
SObjectizer Environment for which the mbox is created.
std::string query_name() const override
Get the mbox name.
void unsubscribe_event_handler(const std::type_index &type_wrapper, abstract_message_sink_t &subscriber) noexcept override
Remove all message handlers.
void drop_delivery_filter(const std::type_index &msg_type, abstract_message_sink_t &subscriber) noexcept override
Removes delivery filter for message type and subscriber.
const full_named_mbox_id_t m_name
Mbox name.
intrusive_ptr_t(T *obj) noexcept
Constructor for a raw pointer.
intrusive_ptr_t & operator=(intrusive_ptr_t &&o) noexcept
Move operator.
T & operator*() const noexcept
A class for the name of mbox_namespace.
std::string_view query_name() const noexcept(noexcept(std::string_view{m_name}))
Get the value.
Parameters for message chain.
Definition mchain.hpp:729
const mchain_props::not_empty_notification_func_t & not_empty_notificator() const
Get chain's notificator for 'not_empty' condition.
Definition mchain.hpp:777
const mchain_props::capacity_t & capacity() const
Get chain's capacity and related params.
Definition mchain.hpp:757
Parameters for defining chain size.
Definition mchain.hpp:229
memory_usage_t memory_usage() const
Memory allocation type for size-limited chain.
Definition mchain.hpp:332
bool unlimited() const
Is message chain have no size limit?
Definition mchain.hpp:318
overflow_reaction_t overflow_reaction() const
Overflow reaction for size-limited chain.
Definition mchain.hpp:339
duration_t overflow_timeout() const
Get the value of waiting timeout for overflow case.
Definition mchain.hpp:356
bool is_overflow_timeout_defined() const
Is waiting timeout for overflow case defined?
Definition mchain.hpp:346
std::size_t max_size() const
Max size for size-limited chain.
Definition mchain.hpp:325
Implementation of demands queue for size-limited message chain with dynamically allocated storage.
limited_dynamic_demand_queue(const capacity_t &capacity)
Initializing constructor.
std::size_t size() const noexcept
Size of the queue.
void pop_front()
Remove the front item from queue.
void push_back(demand_t &&demand)
Add a new item to the end of the queue.
const std::size_t m_max_size
Maximum size of the queue.
demand_t & front()
Access to front item from queue.
Implementation of demands queue for size-limited message chain with preallocated storage.
std::size_t size() const noexcept
Size of the queue.
void push_back(demand_t &&demand)
Add a new item to the end of the queue.
const std::size_t m_max_size
Maximum size of the queue.
limited_preallocated_demand_queue(const capacity_t &capacity)
Initializing constructor.
demand_t & front()
Access to front item from queue.
Implementation of demands queue for size-unlimited message chain.
demand_t & front()
Access to front item from queue.
void push_back(demand_t &&demand)
Add a new item to the end of the queue.
std::deque< demand_t > m_queue
Queue's storage.
void pop_front()
Remove the front item from queue.
std::size_t size() const noexcept
Size of the queue.
bool is_empty() const noexcept
Is queue empty?
Template-based implementation of message chain.
void unsubscribe_event_handler(const std::type_index &, abstract_message_sink_t &) noexcept override
Remove all message handlers.
std::string query_name() const override
Get the mbox name.
mbox_type_t type() const override
Get the type of message box.
environment_t & environment() const noexcept override
SObjectizer Environment for which the mbox is created.
void set_delivery_filter(const std::type_index &, const delivery_filter_t &, abstract_message_sink_t &) override
const capacity_t m_capacity
Chain capacity.
std::size_t size() const override
Count of messages in the chain.
void try_to_store_message_to_queue_nonblocking_mode(const std::type_index &msg_type, const message_ref_t &message)
An implementation of storing another message to chain for the case of delated/periodic messages.
void try_to_store_message_to_queue_ordinary_mode(const std::type_index &msg_type, const message_ref_t &message)
Actual implementation of pushing message to the queue.
const not_empty_notification_func_t m_not_empty_notificator
Optional notificator for 'not_empty' condition.
void subscribe_event_handler(const std::type_index &, abstract_message_sink_t &) override
Add the message handler.
extraction_status_t extract_demand_from_not_empty_queue(demand_t &dest)
Implementation of extract operation for the case when message queue is not empty.
details::status m_status
Status of the chain.
select_case_t * m_select_tail
A queue of multi-chain selects in which this chain is used.
bool empty() const override
Is message chain empty?
std::condition_variable m_overflow_cond
Condition variable for waiting on full queue.
Queue m_queue
Chain's demands queue.
void do_deliver_message(message_delivery_mode_t delivery_mode, const std::type_index &msg_type, const message_ref_t &message, unsigned int) override
Deliver message for all subscribers with respect to message limits.
std::size_t m_threads_to_wakeup
Count of threads sleeping on empty mchain.
void remove_from_select(select_case_t &select_case) noexcept override
Removement of mchain from multi chain select.
mchain_props::push_status_t push(const std::type_index &msg_type, const message_ref_t &message, mchain_props::select_case_t &select_case) override
An attempt to push a new message into the mchain.
mbox_id_t id() const override
Unique ID of this mbox.
const mbox_id_t m_id
Mbox ID for chain.
environment_t & m_env
SObjectizer Environment for which message chain is created.
extraction_status_t extract(demand_t &dest, duration_t empty_queue_timeout) override
void actual_close(close_mode_t mode) override
Close the chain.
void complete_store_message_to_queue(typename Tracing_Base::deliver_op_tracer &tracer, const std::type_index &msg_type, const message_ref_t &message)
A reusable method with implementation of last part of storing a message into chain.
std::condition_variable m_underflow_cond
Condition variable for waiting on empty queue.
mchain_template(so_5::environment_t &env, mbox_id_t id, const mchain_params_t &params, Tracing_Args &&... tracing_args)
Initializing constructor.
extraction_status_t extract(demand_t &dest, select_case_t &select_case) override
An extraction attempt as a part of multi chain select.
void drop_delivery_filter(const std::type_index &, abstract_message_sink_t &) noexcept override
Removes delivery filter for message type and subscriber.
Base class for representation of one case in multi chain select.
select_case_t * query_next() const noexcept
void set_next(select_case_t *next) noexcept
Set the next item in the current queue to which select_case belongs.
void notify() noexcept
Notification for all waiting select_cases.
A base class for agent messages.
Definition message.hpp:47
friend message_mutability_t message_mutability(const intrusive_ptr_t< message_t > &what) noexcept
Helper method for safe get of message mutability flag.
Definition message.hpp:74
Interface of holder of message tracer and message trace filter objects.
virtual bool is_msg_tracing_enabled() const noexcept=0
Is message tracing enabled?
A class for the name which cannot be empty.
std::string giveout_value() noexcept(noexcept(std::string{ std::move(m_nonempty_name) }))
Get the value away from the object.
Helper class for indication of long-lived reference via its type.
Definition outliving.hpp:98
T & get() const noexcept
outliving_reference_t(outliving_reference_t const &o) noexcept
Scoped guard for shared locks.
#define SO_5_LOG_ERROR(logger, var_name)
A special macro for helping error logging.
#define SO_5_THROW_EXCEPTION(error_code, desc)
Definition exception.hpp:74
Some reusable and low-level classes/functions which can be used in public header files.
auto invoke_noexcept_code(L lambda) noexcept -> decltype(lambda())
std::unique_ptr< abstract_message_box_t > make_actual_mbox(outliving_reference_t< so_5::msg_tracing::holder_t > msg_tracing_stuff, A &&... args)
void ensure_sink_for_same_owner(agent_t &actual_owner, abstract_message_sink_t &sink)
Helper the ensures that sink can be used with agent.
Implementation details for MPMC mboxes.
Various helpers for message delivery tracing mechanism.
Details of SObjectizer run-time implementations.
Definition agent.cpp:780
mchain_t make_mchain(outliving_reference_t< so_5::msg_tracing::holder_t > tracer, const mchain_params_t &params, A &&... args)
Helper function for creation of a new mchain with respect to message tracing.
std::string default_global_mbox_namespace()
Helper function that returns name of the default global namespace for named mboxes.
Implementation details.
Definition mchain.hpp:37
void ensure_queue_not_empty(Q &&queue)
Helper function which throws an exception if queue is empty.
status
Status of the message chain.
@ closed
Bag is closed. New messages cannot be sent to it.
@ open
Bag is open and can be used for message sending.
void ensure_queue_not_full(Q &&queue)
Helper function which throws an exception if queue is full.
Various properties and parameters of message chains.
Definition mchain.hpp:28
close_mode_t
What to do with chain's content at close.
Definition mchain.hpp:410
@ drop_content
All messages must be removed from chain.
overflow_reaction_t
What reaction must be performed on attempt to push new message to the full message chain.
Definition mchain.hpp:199
@ remove_oldest
Oldest message in chain must be removed.
@ drop_newest
New message must be ignored and droped.
@ throw_exception
An exception must be thrown.
extraction_status_t
Result of extraction of message from a message chain.
Definition mchain.hpp:371
@ msg_extracted
Message extracted successfully.
@ chain_closed
Message cannot be extracted because chain is closed.
@ no_messages
No available messages in the chain.
push_status_t
Result of attempt of pushing a message into a message chain.
Definition mchain.hpp:389
@ stored
Message stored into a message chain.
@ chain_closed
Message wasn't stored because chain is closed.
memory_usage_t
Memory allocation for storage for size-limited chains.
Definition mchain.hpp:182
@ dynamic
Storage can be allocated and deallocated dynamically.
Public part of message delivery tracing mechanism.
Private part of message limit implementation.
Definition agent.cpp:33
const int rc_nullptr_as_result_of_user_mbox_factory
nullptr returned by user-provided mbox factory.
Definition ret_code.hpp:468
message_delivery_mode_t
Possible modes of message/signal delivery.
Definition types.hpp:172
mbox_type_t
Type of the message box.
Definition mbox.hpp:163
const int rc_msg_chain_overflow
Definition ret_code.hpp:203
const int rc_msg_chain_is_full
Attempt to push a message to full message queue.
Definition ret_code.hpp:193
message_mutability_t
A enum with variants of message mutability or immutability.
Definition types.hpp:94
const int rc_msg_chain_doesnt_support_subscriptions
Attempt to make subscription for message chain.
Definition ret_code.hpp:196
delivery_possibility_t
Result of checking delivery posibility.
Definition mbox.hpp:39
const int rc_illegal_subscriber_for_mpsc_mbox
An attempt to create illegal subscription to mpsc_mbox.
Definition ret_code.hpp:115
const int rc_msg_chain_is_empty
Attempt to get message from empty message queue.
Definition ret_code.hpp:190
const int rc_msg_chain_doesnt_support_delivery_filters
Attempt to set delivery_filter for message chain.
Definition ret_code.hpp:199
const int rc_mutable_msg_cannot_be_delivered_via_mpmc_mbox
An attempt to deliver mutable message via MPMC mbox.
Definition ret_code.hpp:234
outliving_reference_t< T > outliving_mutable(T &r)
Make outliving_reference wrapper for mutable reference.
Full name for a named mbox.
Definition mbox_core.hpp:69
full_named_mbox_id_t(std::string mbox_namespace, std::string mbox_name)
Initializing constructor.
Definition mbox_core.hpp:86
A coolection of data required for local mbox implementation.
messages_table_t m_subscribers
Map of subscribers to messages.
data_t(mbox_id_t id, environment_t &env)
environment_t & m_env
Environment for which the mbox is created.
default_rw_spinlock_t m_lock
Object lock.
const mbox_id_t m_id
ID of this mbox.
bool operator()(abstract_message_sink_t *a, abstract_message_sink_t *b) const noexcept
bool operator()(const subscribers_vector_item_t &a, const subscribers_vector_item_t &b) const noexcept
subscription_info_with_sink_t m_info
Information about the subscription.
abstract_message_sink_t * m_sink_as_key
Pointer to sink that has to be used as search key.
subscribers_vector_item_t(abstract_message_sink_t &sink_as_key, subscription_info_with_sink_t info)
The normal initializing constructor.
Statistics from mbox_core for run-time monitoring.
Definition mbox_core.hpp:49
mbox_t m_mbox
Real mbox for that name.
unsigned int m_external_ref_count
Reference count by external mbox_refs.
Base class for a mbox for the case when message delivery tracing is disabled.
An information which is necessary for creation of a new mbox.
mbox_creation_data_t(outliving_reference_t< environment_t > env, mbox_id_t id, outliving_reference_t< msg_tracing::holder_t > tracer)
Initializing constructor.
Description of one demand in message chain.
Definition mchain.hpp:144
demand_t()
Default constructor.
Definition mchain.hpp:151
demand_t(std::type_index msg_type, so_5::message_ref_t message_ref)
Initializing constructor.
Definition mchain.hpp:155