TSDuck v3.45-4766
MPEG Transport Stream Toolkit
Loading...
Searching...
No Matches
ts::ReactiveMessageQueue< MSG > Class Template Reference

Generic message queue for use in a Reactor environment. More...

#include <tsReactiveMessageQueue.h>

Inheritance diagram for ts::ReactiveMessageQueue< MSG >:
Collaboration diagram for ts::ReactiveMessageQueue< MSG >:

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.
 
Reactorreactor ()
 Get a reference to the associated reactor.
 
virtual Reportreport () const override
 Access the Report which is associated with this object.
 
ReportsetReport (Report *report)
 Associate this object with another Report to log errors.
 
ReporterBasesetReport (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.
 

Detailed Description

template<typename MSG>
class ts::ReactiveMessageQueue< MSG >

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.

Template Parameters
MSGThe type of the messages to receive.

Constructor & Destructor Documentation

◆ ReactiveMessageQueue()

template<typename MSG >
ts::ReactiveMessageQueue< MSG >::ReactiveMessageQueue ( Reactor reactor,
MessageQueue< MSG > &  queue 
)

Constructor.

Parameters
[in,out]reactorAssociated reactor. The reactor object must remain valid as long as this object is valid.
[in,out]queueAssociated message queue. The message queue object must remain valid as long as this object is valid.

Member Function Documentation

◆ queue()

template<typename MSG >
MessageQueue< MSG > & ts::ReactiveMessageQueue< MSG >::queue ( )
inline

Get a reference to the associated message queue.

Returns
A reference to the associated socket.

◆ startReceive()

template<typename MSG >
void ts::ReactiveMessageQueue< MSG >::startReceive ( ReactiveMessageQueueHandlerInterface< MSG > *  handler,
const ObjectPtr user_data = ObjectPtr() 
)

Declare the handler to receive dequeued messages.

Parameters
[in]handlerHandler 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_dataA shared pointer which will be passed unmodified to handler.

◆ processQueuedOperations()

template<typename MSG >
void ts::ReactiveMessageQueue< MSG >::processQueuedOperations ( )
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.

◆ reactor()

Reactor & ts::ReactiveBase::reactor ( )
inlineinherited

Get a reference to the associated reactor.

Returns
A reference to the associated reactor.

◆ signalQueuedOperations()

bool ts::ReactiveBase::signalQueuedOperations ( )
inherited

Trigger the execution of processQueuedOperations() in the context of a Reactor handler.

Create if necessary and then signal a dedicated user event.

Returns
True on success, false on error.

◆ deactivateQueuedOperations()

void ts::ReactiveBase::deactivateQueuedOperations ( bool  silent)
protectedinherited

Deactivate the execution of processQueuedOperations() in the context of a Reactor handler.

Deactivate and delete the dedicated user event.

Parameters
[in]silentIf true, do not report errors through the logger.

◆ createSignalQueuedOperations()

bool ts::ReactiveBase::createSignalQueuedOperations ( )
protectedinherited

Create, if necessary, the dedicated user event for signalQueuedOperations().

Useless if signalQueuedOperations() is used. Only required with use of uncheckedSignalQueuedOperations().

Returns
True on success, false on error.

◆ uncheckedSignalQueuedOperations()

bool ts::ReactiveBase::uncheckedSignalQueuedOperations ( )
protectedinherited

Trigger the execution of processQueuedOperations() from another thread.

The event must have been previously created, either using createSignalQueuedOperations() or signalQueuedOperations().

Returns
True on success, false on error.

◆ handleUserEvent()

virtual void ts::ReactiveBase::handleUserEvent ( Reactor reactor,
EventId  id 
)
overrideprotectedvirtualinherited

Handle a user-defined event in a Reactor.

Parameters
[in,out]reactorReactor into which the handler is invoked.
[in]idId of the event which was signaled.

Reimplemented from ts::ReactorHandlerInterface.

◆ report()

virtual Report & ts::ReporterBase::report ( ) const
overridevirtualinherited

Access the Report which is associated with this object.

Can be called from another thread only if the Report object is thread-safe.

Returns
A reference to the associated report.

Implements ts::ReporterInterface.

◆ setReport() [1/2]

Report * ts::ReporterBase::setReport ( Report report)
inherited

Associate this object with another Report to log errors.

Parameters
[in]reportWhere 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.
Returns
The address of the previous Report object or a null pointer if there was none.

◆ setReport() [2/2]

ReporterBase * ts::ReporterBase::setReport ( ReporterBase delegate)
inherited

Associate this object with another ReporterBase to log errors.

Parameters
[in]delegateUse the report of another ReporterBase. If delegate is null, the previous explicit Report is used..
Returns
The address of the previous ReporterBase object or a null pointer if there was none.

◆ muteReport()

bool ts::ReporterBase::muteReport ( bool  mute)
inherited

Temporarily mute the associated report.

Parameters
[in]muteIt true, report() will return a null report (log messages are discarded), until muteReport() is invoked again with mute set to false.
Returns
Previous state of the mute field.

◆ SilentLevel()

static int ts::ReporterBase::SilentLevel ( bool  silent,
int  default_severity = Severity::Error 
)
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.

Parameters
[in]silentIf true, do not report errors, report debug messages instead.
[in]default_severityDefault severity, in non-silent mode (error by default).
Returns
Error when silent is false, Debug otherwise.

◆ handleTimer()

virtual void ts::ReactorHandlerInterface::handleTimer ( Reactor reactor,
EventId  id 
)
virtualinherited

Handle a timer in a Reactor.

Parameters
[in,out]reactorReactor into which the handler is invoked.
[in]idId of the timer which expires.

◆ handleBroadcastEvent()

virtual void ts::ReactorHandlerInterface::handleBroadcastEvent ( Reactor reactor,
int  error_code,
const ObjectPtr user_data 
)
virtualinherited

Handle a broadcast event in a Reactor.

A broadcast event is sent to all currently registered events in the reactor.

Parameters
[in,out]reactorReactor into which the handler is invoked.
[in]error_codeApplication-specific error code which was passed to Reactor::signalBroadcastEvent().
[in]user_dataThe user-data shared pointer which was passed to Reactor::signalBroadcastEvent().

◆ handleProcessTermination()

virtual void ts::ReactorHandlerInterface::handleProcessTermination ( Reactor reactor,
EventId  id,
int  pid 
)
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.

Parameters
[in,out]reactorReactor into which the handler is invoked.
[in]idId of the event which was signaled.
[in]pidProcess id of the terminated process.

◆ handleReadReady()

virtual void ts::ReactorHandlerInterface::handleReadReady ( Reactor reactor,
EventId  id,
int  error_code 
)
virtualinherited

Handle a read-ready event in a Reactor.

This handler is only invoked in the non-blocking I/O model.

Parameters
[in,out]reactorReactor into which the handler is invoked.
[in]idId of the event which was signaled.
[in]error_codeSystem-specific error code, zero on success, SYS_ERROR in case of unknown error.

Reimplemented in ts::ReactiveStream.

◆ handleWriteReady()

virtual void ts::ReactorHandlerInterface::handleWriteReady ( Reactor reactor,
EventId  id,
int  error_code 
)
virtualinherited

Handle a write-ready event in a Reactor.

This handler is only invoked in the non-blocking I/O model.

Parameters
[in,out]reactorReactor into which the handler is invoked.
[in]idId of the event which was signaled.
[in]error_codeSystem-specific error code, zero on success, SYS_ERROR in case of unknown error.

Reimplemented in ts::ReactiveStream, and ts::ReactiveTCPConnection.

◆ handleAsynchronousIO()

virtual void ts::ReactorHandlerInterface::handleAsynchronousIO ( Reactor reactor,
EventId  id,
NonBlockingDevice::IOSB iosb,
size_t  io_size 
)
virtualinherited

Handle an asynchronous I/O completion event in a Reactor.

This handler is only invoked in the asynchronous I/O model.

Parameters
[in,out]reactorReactor into which the handler is invoked.
[in]idId of the event which was signaled.
[in,out]iosbIOSB 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_sizeSize of the I/O in bytes.

Reimplemented in ts::ReactiveStream, and ts::ReactiveTCPConnection.


The documentation for this class was generated from the following file: