TSDuck v3.45-4766
MPEG Transport Stream Toolkit
Loading...
Searching...
No Matches
ts::ReactiveStream Class Reference

Using non-blocking devices with StreamInterface in a Reactor environment. More...

#include <tsReactiveStream.h>

Inheritance diagram for ts::ReactiveStream:
Collaboration diagram for ts::ReactiveStream:

Public Member Functions

 ReactiveStream (Reactor &reactor, NonBlockingDevice &device, StreamInterface &stream)
 Constructor.
 
virtual void cancelReadWriteStream (bool silent=false)
 Cancel any pending read or write operation on this stream.
 
NonBlockingDevicedevice ()
 Get a reference to the associated non-blocking device.
 
bool muteReport (bool mute)
 Temporarily mute the associated report.
 
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.
 
virtual bool startReadStream (ReactiveStreamHandlerInterface *handler, size_t buffer_size=DEFAULT_RECEIVE_BUFFER_SIZE, const ObjectPtr &user_data=ObjectPtr())
 Start the operation of reading data from the stream.
 
virtual bool startWriteStream (ReactiveStreamHandlerInterface *handler, const void *data, size_t size, const ObjectPtr &user_data=ObjectPtr())
 Start the operation of writing data to the stream.
 
StreamInterfacestream ()
 Get a reference to the associated stream.
 

Static Public Member Functions

static int SilentLevel (bool silent, int default_severity=Severity::Error)
 Compute a log severity level from a "silent" parameter.
 

Static Public Attributes

static constexpr size_t DEFAULT_RECEIVE_BUFFER_SIZE = 4096
 Default buffer size for receive operations.
 

Protected Types

using IOQueue = std::list< std::shared_ptr< IOSB > >
 Queues of I/O requests are queues of shared_ptr to IOSB.
 
using IOSB = NonBlockingDevice::IOSB
 IOSB shortcut fpr subclasses.
 
using IOSet = std::set< std::shared_ptr< IOSB > >
 Unordered set of I/O requests, set of shared_ptr to IOSB.
 

Protected Member Functions

bool activateAsynchronousIO ()
 Activate notification for asynchronous I/O.
 
bool activateReadReady ()
 Activate read-ready notification for non-blocking I/O.
 
bool activateWriteReady ()
 Activate write-ready notification for non-blocking I/O.
 
bool cancelAndWaitAsynchronousIO (NonBlockingDevice::IOSB &iosb, bool silent)
 Cancel one specific pending asynchronous I/O and wait for its completion.
 
void cancelAsynchronousIO (bool silent)
 Cancel all asynchronous I/O in progress.
 
template<class REQUEST >
requires std::derived_from<REQUEST, ts::Object>
void cancelQueue (IOQueue &inqueue, IOQueue &outqueue)
 Transfer all requests from one queue to another and mark all I/O as canceled.
 
bool createSignalQueuedOperations ()
 Create, if necessary, the dedicated user event for signalQueuedOperations().
 
void deactivateAll (bool silent)
 Deactivate all registrations for non-blocking and asynchronous I/O.
 
void deactivateAsynchronousIO (bool silent)
 Deactivate notification for asynchronous I/O.
 
void deactivateQueuedOperations (bool silent)
 Deactivate the execution of processQueuedOperations() in the context of a Reactor handler.
 
void deactivateReadReady (bool silent)
 Deactivate read-ready notification for non-blocking I/O.
 
void deactivateWriteReady (bool silent)
 Deactivate write-ready notification for non-blocking I/O.
 
bool enqueueCompletedIO (const std::shared_ptr< IOSB > &iosb, bool in_reactor)
 Enqueue an IOSB in the completed I/O queue.
 
virtual void handleAsynchronousIO (Reactor &, EventId, IOSB &, size_t) override
 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 &, EventId, int) override
 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 &, EventId, int) override
 Handle a write-ready event in a Reactor.
 
virtual bool hasPendingIO ()
 Check if some I/O operations are in progress.
 
virtual bool needsWriteReady ()
 With non-blocking I/O, check if we need to be notified of write-ready conditions.
 
virtual int processCloseWriteStream (bool silent)
 Implement the close of the write direction of the stream.
 
virtual void processQueuedOperations () override
 This virtual method processes operations in the context of a Reactor handler.
 
void processReceiveBuffer (ByteBlock &data, ReactiveInputControl &control, ReactiveStreamHandlerInterface *handler, int error_code, const ObjectPtr &user_data)
 Invoke the receive handler as many times as possible on a data buffer.
 
std::shared_ptr< IOSBremoveFromQueue (IOQueue &queue, IOSB *iosb)
 Search and remove a shared_ptr to IOSB, based on an IOSB address.
 
virtual bool startCloseWriteStream (ReactiveStreamHandlerInterface *handler, bool silent, const ObjectPtr &user_data)
 Start closing the write direction of the stream.
 
virtual bool tryCompletedIO (const std::shared_ptr< IOSB > &iosb)
 Try to interpret an IOSB as a valid completed I/O request and process the completion.
 
bool uncheckedSignalQueuedOperations ()
 Trigger the execution of processQueuedOperations() from another thread.
 

Detailed Description

Using non-blocking devices with StreamInterface in a Reactor environment.

The class ReactiveStream is a wrapper around a stream device class to handle reactive I/O. Typically, the device class is a subclass of NonBlockingDevice which implements StreamInterface.

The stream device is a separate object. It is initialized and configured by the application. The application shall not directly call writeStream() or readStream() on this device and delegate these operations to startWriteStream() and startReadStream() in class ReactiveStream.

Member Typedef Documentation

◆ IOQueue

using ts::ReactiveDevice::IOQueue = std::list<std::shared_ptr<IOSB> >
protectedinherited

Queues of I/O requests are queues of shared_ptr to IOSB.

This is typically used with non-blocking I/O where we must process requests in order. Send and receive requests are structures which are stored in the react_data of the IOSB.

◆ IOSet

using ts::ReactiveDevice::IOSet = std::set<std::shared_ptr<IOSB> >
protectedinherited

Unordered set of I/O requests, set of shared_ptr to IOSB.

This is typically used with asynchronous I/O. The ordering is enforced because I/O are started in order of calls from applications. The completion processing is likely the same, but driven by the system I/O Completion Ports and we must not assume any order. Send and receive requests are structures which are stored in the react_data of the IOSB.

Constructor & Destructor Documentation

◆ ReactiveStream()

ts::ReactiveStream::ReactiveStream ( Reactor reactor,
NonBlockingDevice device,
StreamInterface stream 
)

Constructor.

Parameters
[in,out]reactorAssociated reactor. The reactor object must remain valid as long as this object is valid.
[in,out]deviceAssociated non-blocking device. The device object must remain valid as long as this object is valid.
[in,out]streamAssociated stream interface. Typically, the device class is a subclass of NonBlockingDevice which implements StreamInterface. Therefore, device and stream are two views of the same object. That device object must remain valid as long as this instance of ReactiveStream exists.

Member Function Documentation

◆ stream()

StreamInterface & ts::ReactiveStream::stream ( )
inline

Get a reference to the associated stream.

Typically, device() and stream() return two views of the same object.

Returns
A reference to the associated stream.

◆ startReadStream()

virtual bool ts::ReactiveStream::startReadStream ( ReactiveStreamHandlerInterface handler,
size_t  buffer_size = DEFAULT_RECEIVE_BUFFER_SIZE,
const ObjectPtr user_data = ObjectPtr() 
)
virtual

Start the operation of reading data from the stream.

Reading operation is permanent and handler is called whenever incoming data are available.

Parameters
[in]handlerHandler class to call each time data are received. The method handleReadStream() will be called for each new chunk of data. Cannot be null. If a previous receive handler was registered, it is replaced.
[in]buffer_sizeSize of input buffers to receive data.
[in]user_dataA shared pointer which will be passed unmodified to handler.
Returns
True on success, false on error. Success means that the I/O was successfully started. The final status of the I/O will be transmitted in the handler.

Reimplemented in ts::ReactiveTLSConnection.

◆ startWriteStream()

virtual bool ts::ReactiveStream::startWriteStream ( ReactiveStreamHandlerInterface handler,
const void *  data,
size_t  size,
const ObjectPtr user_data = ObjectPtr() 
)
virtual

Start the operation of writing data to the stream.

Parameters
[in]handlerHandler class to call when the send operation completes. The method handleWriteStream() will be called. If nullptr, no handler is called.
[in]dataAddress of the data to write. The corresponding memory area must remain valid until the completion or cancelation of the send operation.
[in]sizeSize in bytes of the data to send.
[in]user_dataA shared pointer which will be passed unmodified to handler.
Returns
True on success, false on error. Success means that the I/O was successfully started. The final status of the I/O will be transmitted in the handler.

Reimplemented in ts::ReactiveTLSConnection.

◆ cancelReadWriteStream()

virtual void ts::ReactiveStream::cancelReadWriteStream ( bool  silent = false)
virtual

Cancel any pending read or write operation on this stream.

If a repeated read operation is in progress, the repetition is canceled as well.

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

Reimplemented in ts::ReactiveTLSConnection.

◆ hasPendingIO()

virtual bool ts::ReactiveStream::hasPendingIO ( )
protectedvirtual

Check if some I/O operations are in progress.

Returns
True if there is any pending operation or permanent read operation in progress. False otherwise.

◆ needsWriteReady()

virtual bool ts::ReactiveStream::needsWriteReady ( )
protectedvirtual

With non-blocking I/O, check if we need to be notified of write-ready conditions.

If overriden by a subclass, the method must also call its superclass counterpart.

Returns
True if non-blocking write-ready notifications are required. False otherwise.

◆ startCloseWriteStream()

virtual bool ts::ReactiveStream::startCloseWriteStream ( ReactiveStreamHandlerInterface handler,
bool  silent,
const ObjectPtr user_data 
)
protectedvirtual

Start closing the write direction of the stream.

The request is enqueued and will be processed by calling processCloseWriteStream() after all previous send requests are completed. This feature can be used by subclasses for which "closing the write direction of the stream" means something (e.g. true for TCP socket, meaningless for files).

Parameters
[in]handlerHandler class to call when the close-writer operation completes. The method handleWriteStream() will be called with its parameter error_code containing SYS_EOF. If nullptr, no handler is called.
[in]silentIf true, do not report errors through the logger.
[in]user_dataA shared pointer which will be passed unmodified to handler.
Returns
True on success, false on error.

◆ processCloseWriteStream()

virtual int ts::ReactiveStream::processCloseWriteStream ( bool  silent)
protectedvirtual

Implement the close of the write direction of the stream.

Must be implemented by subclasses which use startCloseWriteStream(). The default implementation does nothing. This method is always called in the context of a Reactor handler, never as an indirect call from the application.

Parameters
[in]silentIf true, do not report errors through the logger.
Returns
A system-specific error code. SYS_SUCCESS in case of success.

Reimplemented in ts::ReactiveTCPConnection.

◆ enqueueCompletedIO()

bool ts::ReactiveStream::enqueueCompletedIO ( const std::shared_ptr< IOSB > &  iosb,
bool  in_reactor 
)
protected

Enqueue an IOSB in the completed I/O queue.

This method can be called by subclasses which need to process specific I/O completion, other than read or write on the stream.

Parameters
[in]iosbShared pointer to an IOSB. Specific data can be saved in iosb->react_data.
[in]in_reactorSet to true if we are in a reactor context, meaning directly called from a reactor handler, not from the application.
Returns
True on success, false on error.

◆ tryCompletedIO()

virtual bool ts::ReactiveStream::tryCompletedIO ( const std::shared_ptr< IOSB > &  iosb)
protectedvirtual

Try to interpret an IOSB as a valid completed I/O request and process the completion.

If overriden by a subclass, the method must also call its superclass counterpart.

Parameters
[in]iosbShared pointer to an IOSB. Specific data can be saved in iosb->react_data.
Returns
True if the iosb contained a recognized I/O completion request. False if the request is unknown.

Reimplemented in ts::ReactiveTCPConnection.

◆ processReceiveBuffer()

void ts::ReactiveStream::processReceiveBuffer ( ByteBlock data,
ReactiveInputControl control,
ReactiveStreamHandlerInterface handler,
int  error_code,
const ObjectPtr user_data 
)
protected

Invoke the receive handler as many times as possible on a data buffer.

Parameters
[in,out]dataData buffer containing the received data. On output, data which are processed by the handler are removed.
[in,out]controlInput control. On input, this is the previously returned value from the handler. Modified by the handler.
[in]handlerApplication handler.
[in]error_codeReceive error code. If not success, the handler is called exactly once.
[in]user_dataUser data for the handler.

◆ processQueuedOperations()

virtual void ts::ReactiveStream::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.

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

◆ handleReadReady()

virtual void ts::ReactiveStream::handleReadReady ( Reactor reactor,
EventId  id,
int  error_code 
)
overrideprotectedvirtual

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 from ts::ReactorHandlerInterface.

◆ handleWriteReady()

virtual void ts::ReactiveStream::handleWriteReady ( Reactor reactor,
EventId  id,
int  error_code 
)
overrideprotectedvirtual

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 from ts::ReactorHandlerInterface.

Reimplemented in ts::ReactiveTCPConnection.

◆ handleAsynchronousIO()

virtual void ts::ReactiveStream::handleAsynchronousIO ( Reactor reactor,
EventId  id,
IOSB iosb,
size_t  io_size 
)
overrideprotectedvirtual

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 from ts::ReactorHandlerInterface.

Reimplemented in ts::ReactiveTCPConnection.

◆ device()

NonBlockingDevice & ts::ReactiveDevice::device ( )
inlineinherited

Get a reference to the associated non-blocking device.

Returns
A reference to the associated non-blocking device.

◆ removeFromQueue()

std::shared_ptr< IOSB > ts::ReactiveDevice::removeFromQueue ( IOQueue queue,
IOSB iosb 
)
protectedinherited

Search and remove a shared_ptr to IOSB, based on an IOSB address.

Search from the front (end) of the queue since a completed I/O is likely on the front.

Parameters
[in,out]queueThe queue from which to remove iosb.
[in]iosbStandard pointer to an IOSB to search and remove.
Returns
The removed shared_ptr to IOSB, or a null pointer if iosb is not found.

◆ cancelQueue()

template<class REQUEST >
requires std::derived_from<REQUEST, ts::Object>
void ts::ReactiveDevice::cancelQueue ( IOQueue inqueue,
IOQueue outqueue 
)
protectedinherited

Transfer all requests from one queue to another and mark all I/O as canceled.

Template Parameters
REQUESTThe subclass of Object which is set in react_data of all requests in inqueue.
Parameters
[in,out]inqueueThe queue from which all requests are removed.
[in,out]outqueueThe queue which receives all canceled requests.

◆ activateReadReady()

bool ts::ReactiveDevice::activateReadReady ( )
protectedinherited

Activate read-ready notification for non-blocking I/O.

Returns
True on success, false on error.

◆ deactivateReadReady()

void ts::ReactiveDevice::deactivateReadReady ( bool  silent)
protectedinherited

Deactivate read-ready notification for non-blocking I/O.

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

◆ activateWriteReady()

bool ts::ReactiveDevice::activateWriteReady ( )
protectedinherited

Activate write-ready notification for non-blocking I/O.

Returns
True on success, false on error.

◆ deactivateWriteReady()

void ts::ReactiveDevice::deactivateWriteReady ( bool  silent)
protectedinherited

Deactivate write-ready notification for non-blocking I/O.

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

◆ activateAsynchronousIO()

bool ts::ReactiveDevice::activateAsynchronousIO ( )
protectedinherited

Activate notification for asynchronous I/O.

Returns
True on success, false on error.

◆ deactivateAsynchronousIO()

void ts::ReactiveDevice::deactivateAsynchronousIO ( bool  silent)
protectedinherited

Deactivate notification for asynchronous I/O.

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

◆ cancelAsynchronousIO()

void ts::ReactiveDevice::cancelAsynchronousIO ( bool  silent)
protectedinherited

Cancel all asynchronous I/O in progress.

The cancelation occurs in the background and end of canceled asynchronous I/O will be notified.

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

◆ cancelAndWaitAsynchronousIO()

bool ts::ReactiveDevice::cancelAndWaitAsynchronousIO ( NonBlockingDevice::IOSB iosb,
bool  silent 
)
protectedinherited

Cancel one specific pending asynchronous I/O and wait for its completion.

Warning: This is a blocking call. It shall be used in case of trouble only.

Parameters
[in,out]iosbThe asynchronous I/O status block.
[in]silentIf true, do not report errors through the logger.
Returns
True on success, false on error.

◆ deactivateAll()

void ts::ReactiveDevice::deactivateAll ( bool  silent)
protectedinherited

Deactivate all registrations for non-blocking and asynchronous I/O.

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

◆ 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.

Member Data Documentation

◆ DEFAULT_RECEIVE_BUFFER_SIZE

constexpr size_t ts::ReactiveStream::DEFAULT_RECEIVE_BUFFER_SIZE = 4096
staticconstexpr

Default buffer size for receive operations.

See also
startReadStream()

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