![]() |
TSDuck v3.45-4766
MPEG Transport Stream Toolkit
|
Using non-blocking devices with StreamInterface in a Reactor environment. More...
#include <tsReactiveStream.h>


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. | |
| NonBlockingDevice & | device () |
| Get a reference to the associated non-blocking device. | |
| bool | muteReport (bool mute) |
| Temporarily mute the associated report. | |
| 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. | |
| 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. | |
| StreamInterface & | stream () |
| 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< IOSB > | removeFromQueue (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. | |
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.
|
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.
|
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.
| ts::ReactiveStream::ReactiveStream | ( | Reactor & | reactor, |
| NonBlockingDevice & | device, | ||
| StreamInterface & | stream | ||
| ) |
Constructor.
| [in,out] | reactor | Associated reactor. The reactor object must remain valid as long as this object is valid. |
| [in,out] | device | Associated non-blocking device. The device object must remain valid as long as this object is valid. |
| [in,out] | stream | Associated 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. |
|
inline |
|
virtual |
Start the operation of reading data from the stream.
Reading operation is permanent and handler is called whenever incoming data are available.
| [in] | handler | Handler 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_size | Size of input buffers to receive data. |
| [in] | user_data | A shared pointer which will be passed unmodified to handler. |
Reimplemented in ts::ReactiveTLSConnection.
|
virtual |
Start the operation of writing data to the stream.
| [in] | handler | Handler class to call when the send operation completes. The method handleWriteStream() will be called. If nullptr, no handler is called. |
| [in] | data | Address of the data to write. The corresponding memory area must remain valid until the completion or cancelation of the send operation. |
| [in] | size | Size in bytes of the data to send. |
| [in] | user_data | A shared pointer which will be passed unmodified to handler. |
Reimplemented in ts::ReactiveTLSConnection.
|
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.
| [in] | silent | If true, do not report errors through the logger. |
Reimplemented in ts::ReactiveTLSConnection.
|
protectedvirtual |
Check if some I/O operations are in progress.
|
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.
|
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).
| [in] | handler | Handler 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] | silent | If true, do not report errors through the logger. |
| [in] | user_data | A shared pointer which will be passed unmodified to handler. |
|
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.
| [in] | silent | If true, do not report errors through the logger. |
Reimplemented in ts::ReactiveTCPConnection.
|
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.
| [in] | iosb | Shared pointer to an IOSB. Specific data can be saved in iosb->react_data. |
| [in] | in_reactor | Set to true if we are in a reactor context, meaning directly called from a reactor handler, not from the application. |
|
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.
| [in] | iosb | Shared pointer to an IOSB. Specific data can be saved in iosb->react_data. |
Reimplemented in ts::ReactiveTCPConnection.
|
protected |
Invoke the receive handler as many times as possible on a data buffer.
| [in,out] | data | Data buffer containing the received data. On output, data which are processed by the handler are removed. |
| [in,out] | control | Input control. On input, this is the previously returned value from the handler. Modified by the handler. |
| [in] | handler | Application handler. |
| [in] | error_code | Receive error code. If not success, the handler is called exactly once. |
| [in] | user_data | User data for the 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.
Reimplemented in ts::ReactiveTCPConnection, and ts::ReactiveTLSConnection.
|
overrideprotectedvirtual |
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 from ts::ReactorHandlerInterface.
|
overrideprotectedvirtual |
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 from ts::ReactorHandlerInterface.
Reimplemented in ts::ReactiveTCPConnection.
|
overrideprotectedvirtual |
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 from ts::ReactorHandlerInterface.
Reimplemented in ts::ReactiveTCPConnection.
|
inlineinherited |
Get a reference to the associated non-blocking device.
|
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.
| [in,out] | queue | The queue from which to remove iosb. |
| [in] | iosb | Standard pointer to an IOSB to search and remove. |
|
protectedinherited |
Transfer all requests from one queue to another and mark all I/O as canceled.
| REQUEST | The subclass of Object which is set in react_data of all requests in inqueue. |
| [in,out] | inqueue | The queue from which all requests are removed. |
| [in,out] | outqueue | The queue which receives all canceled requests. |
|
protectedinherited |
Activate read-ready notification for non-blocking I/O.
|
protectedinherited |
Deactivate read-ready notification for non-blocking I/O.
| [in] | silent | If true, do not report errors through the logger. |
|
protectedinherited |
Activate write-ready notification for non-blocking I/O.
|
protectedinherited |
Deactivate write-ready notification for non-blocking I/O.
| [in] | silent | If true, do not report errors through the logger. |
|
protectedinherited |
Activate notification for asynchronous I/O.
|
protectedinherited |
Deactivate notification for asynchronous I/O.
| [in] | silent | If true, do not report errors through the logger. |
|
protectedinherited |
Cancel all asynchronous I/O in progress.
The cancelation occurs in the background and end of canceled asynchronous I/O will be notified.
| [in] | silent | If true, do not report errors through the logger. |
|
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.
| [in,out] | iosb | The asynchronous I/O status block. |
| [in] | silent | If true, do not report errors through the logger. |
|
protectedinherited |
Deactivate all registrations for non-blocking and asynchronous I/O.
| [in] | silent | If true, do not report errors through the logger. |
|
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. |
|
staticconstexpr |
Default buffer size for receive operations.