![]() |
TSDuck v3.45-4766
MPEG Transport Stream Toolkit
|
Generic message queue for use in a Reactor environment. More...
#include <tsReactiveMessageQueue.h>


Public Types | |
| using | MessageQueueType = MessageQueue< MSG > |
| Common renaming of the queue type. | |
Public Member Functions | |
| ReactiveMessageQueue (Reactor &reactor, MessageQueue< MSG > &queue) | |
| Constructor. | |
| virtual | ~ReactiveMessageQueue () override |
| Destructor. | |
| bool | muteReport (bool mute) |
| Temporarily mute the associated report. | |
| MessageQueue< MSG > & | queue () |
| Get a reference to the associated message queue. | |
| Reactor & | reactor () |
| Get a reference to the associated reactor. | |
| virtual Report & | report () const override |
| Access the Report which is associated with this object. | |
| Report * | setReport (Report *report) |
| Associate this object with another Report to log errors. | |
| ReporterBase * | setReport (ReporterBase *delegate) |
| Associate this object with another ReporterBase to log errors. | |
| bool | signalQueuedOperations () |
| Trigger the execution of processQueuedOperations() in the context of a Reactor handler. | |
| void | startReceive (ReactiveMessageQueueHandlerInterface< MSG > *handler, const ObjectPtr &user_data=ObjectPtr()) |
| Declare the handler to receive dequeued messages. | |
Static Public Member Functions | |
| static int | SilentLevel (bool silent, int default_severity=Severity::Error) |
| Compute a log severity level from a "silent" parameter. | |
Protected Member Functions | |
| bool | createSignalQueuedOperations () |
| Create, if necessary, the dedicated user event for signalQueuedOperations(). | |
| void | deactivateQueuedOperations (bool silent) |
| Deactivate the execution of processQueuedOperations() in the context of a Reactor handler. | |
| virtual void | handleAsynchronousIO (Reactor &reactor, EventId id, NonBlockingDevice::IOSB &iosb, size_t io_size) |
| Handle an asynchronous I/O completion event in a Reactor. | |
| virtual void | handleBroadcastEvent (Reactor &reactor, int error_code, const ObjectPtr &user_data) |
| Handle a broadcast event in a Reactor. | |
| virtual void | handleProcessTermination (Reactor &reactor, EventId id, int pid) |
| Handle a process termination event in a Reactor. | |
| virtual void | handleReadReady (Reactor &reactor, EventId id, int error_code) |
| Handle a read-ready event in a Reactor. | |
| virtual void | handleTimer (Reactor &reactor, EventId id) |
| Handle a timer in a Reactor. | |
| virtual void | handleUserEvent (Reactor &, EventId) override |
| Handle a user-defined event in a Reactor. | |
| virtual void | handleWriteReady (Reactor &reactor, EventId id, int error_code) |
| Handle a write-ready event in a Reactor. | |
| virtual void | processQueuedOperations () override |
| This virtual method processes operations in the context of a Reactor handler. | |
| bool | uncheckedSignalQueuedOperations () |
| Trigger the execution of processQueuedOperations() from another thread. | |
Generic message queue for use in a Reactor environment.
The template class ReactiveMessageQueue is a wrapper around MessageQueue to handle reactive events.
The actual message queue is a separate object. It is initialized and configured by the application. The application shall not directly call dequeue() on this message queue and delegate this operation to startReceive() in ReactiveMessageQueue.
| MSG | The type of the messages to receive. |
| ts::ReactiveMessageQueue< MSG >::ReactiveMessageQueue | ( | Reactor & | reactor, |
| MessageQueue< MSG > & | queue | ||
| ) |
Constructor.
| [in,out] | reactor | Associated reactor. The reactor object must remain valid as long as this object is valid. |
| [in,out] | queue | Associated message queue. The message queue object must remain valid as long as this object is valid. |
|
inline |
Get a reference to the associated message queue.
| void ts::ReactiveMessageQueue< MSG >::startReceive | ( | ReactiveMessageQueueHandlerInterface< MSG > * | handler, |
| const ObjectPtr & | user_data = ObjectPtr() |
||
| ) |
Declare the handler to receive dequeued messages.
| [in] | handler | Handler class to call each time a message is received. The method handleDequeuedMessage() will be called for each new message. If a previous receive handler was registered, it is replaced. |
| [in] | user_data | A shared pointer which will be passed unmodified to handler. |
|
overrideprotectedvirtual |
This virtual method processes operations in the context of a Reactor handler.
This is dedicated to operations which must be serialized from an application perspective. These operations are typically queued when triggered from a method which is called by the application. When the reactor processes events, we are sure that the application is not executing a handler. The default implementation does nothing. A subclass should override it if it calls signalQueuedOperations().
Reimplemented from ts::ReactiveBase.
|
inlineinherited |
Get a reference to the associated reactor.
|
inherited |
Trigger the execution of processQueuedOperations() in the context of a Reactor handler.
Create if necessary and then signal a dedicated user event.
|
protectedinherited |
Deactivate the execution of processQueuedOperations() in the context of a Reactor handler.
Deactivate and delete the dedicated user event.
| [in] | silent | If true, do not report errors through the logger. |
|
protectedinherited |
Create, if necessary, the dedicated user event for signalQueuedOperations().
Useless if signalQueuedOperations() is used. Only required with use of uncheckedSignalQueuedOperations().
|
protectedinherited |
Trigger the execution of processQueuedOperations() from another thread.
The event must have been previously created, either using createSignalQueuedOperations() or signalQueuedOperations().
|
overrideprotectedvirtualinherited |
Handle a user-defined event in a Reactor.
| [in,out] | reactor | Reactor into which the handler is invoked. |
| [in] | id | Id of the event which was signaled. |
Reimplemented from ts::ReactorHandlerInterface.
|
overridevirtualinherited |
Access the Report which is associated with this object.
Can be called from another thread only if the Report object is thread-safe.
Implements ts::ReporterInterface.
Associate this object with another Report to log errors.
| [in] | report | Where to report errors. The report object must remain valid as long as this object exists or setReport() is used with another Report object. If report is null, log messages are discarded. |
|
inherited |
Associate this object with another ReporterBase to log errors.
| [in] | delegate | Use the report of another ReporterBase. If delegate is null, the previous explicit Report is used.. |
|
inherited |
Temporarily mute the associated report.
| [in] | mute | It true, report() will return a null report (log messages are discarded), until muteReport() is invoked again with mute set to false. |
|
inlinestaticinherited |
Compute a log severity level from a "silent" parameter.
Some subclass methods have a "silent" parameter to avoid reporting errors which may be insignificant, typically when closing a device after an error, in which case the close operation may produce other errors if the previous error left the device in an inconsistent state. While those errors should not be displayed as errors, we still display them at debug level.
| [in] | silent | If true, do not report errors, report debug messages instead. |
| [in] | default_severity | Default severity, in non-silent mode (error by default). |
|
virtualinherited |
Handle a broadcast event in a Reactor.
A broadcast event is sent to all currently registered events in the reactor.
| [in,out] | reactor | Reactor into which the handler is invoked. |
| [in] | error_code | Application-specific error code which was passed to Reactor::signalBroadcastEvent(). |
| [in] | user_data | The user-data shared pointer which was passed to Reactor::signalBroadcastEvent(). |
|
virtualinherited |
Handle a process termination event in a Reactor.
This handler is invoked when the process is no longer there. It is possible that the process was already terminated for a while. It is even possible that the process never really started. There is no portable way to get the process termination status.
| [in,out] | reactor | Reactor into which the handler is invoked. |
| [in] | id | Id of the event which was signaled. |
| [in] | pid | Process id of the terminated process. |
|
virtualinherited |
Handle a read-ready event in a Reactor.
This handler is only invoked in the non-blocking I/O model.
| [in,out] | reactor | Reactor into which the handler is invoked. |
| [in] | id | Id of the event which was signaled. |
| [in] | error_code | System-specific error code, zero on success, SYS_ERROR in case of unknown error. |
Reimplemented in ts::ReactiveStream.
|
virtualinherited |
Handle a write-ready event in a Reactor.
This handler is only invoked in the non-blocking I/O model.
| [in,out] | reactor | Reactor into which the handler is invoked. |
| [in] | id | Id of the event which was signaled. |
| [in] | error_code | System-specific error code, zero on success, SYS_ERROR in case of unknown error. |
Reimplemented in ts::ReactiveStream, and ts::ReactiveTCPConnection.
|
virtualinherited |
Handle an asynchronous I/O completion event in a Reactor.
This handler is only invoked in the asynchronous I/O model.
| [in,out] | reactor | Reactor into which the handler is invoked. |
| [in] | id | Id of the event which was signaled. |
| [in,out] | iosb | IOSB structure which was used when the asynchronous I/O was started. A system-specific error code is in iosb, SYS_CANCELED if the I/O was canceled before completion. |
| [in] | io_size | Size of the I/O in bytes. |
Reimplemented in ts::ReactiveStream, and ts::ReactiveTCPConnection.