![]() |
TSDuck v3.45-4776
MPEG Transport Stream Toolkit
|
Pool of worker threads in a Reactor environment. More...
#include <tsReactiveWorkerPool.h>


Public Member Functions | |
| ReactiveWorkerPool (Reactor &reactor) | |
| Constructor. | |
| virtual | ~ReactiveWorkerPool () override |
| Destructor. | |
| size_t | currentBusyThreads () const |
| Get the current number of worker threads currently busy executing worker tasks. | |
| size_t | currentThreads () const |
| Get the current number of worker threads in the pool. | |
| size_t | maxThreads () const |
| Get the maximum number of worker threads in the pool. | |
| 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. | |
| void | setMaxThreads (size_t count) |
| Set the maximum number of worker threads in the pool. | |
| 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. | |
| void | setStackSize (size_t size) |
| Set the stack size in bytes of worker threads in the pool. | |
| bool | signalQueuedOperations () |
| Trigger the execution of processQueuedOperations() in the context of a Reactor handler. | |
| bool | startWork (ReactiveWorkerInterface *work, ReactiveWorkerHandlerInterface *handler, const ObjectPtr &user_data=ObjectPtr()) |
| Start a lengthy task in a delegated worker thread and return immediately. | |
| void | waitForTermination (bool cancel_pending=false) |
| Wait for all worker threads to become idle and terminate all worker threads. | |
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_MAX_THREADS = 8 |
| Default maximum number of worker threads in the pool for newly created ReactiveWorkerPool instances. | |
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. | |
Pool of worker threads in a Reactor environment.
A Reactor event loop is a mono-thread synchronous environment. Each callback shall run without blocking I/O and without lengthy operations. When an application needs to run lengthy tasks in reactor callbacks, it must delegate these tasks in a worker thread. An instance of ReactiveWorkerPool can be used in association with a reactor to execute these lengthy tasks.
The instance of ReactiveWorkerPool is not thread-safe and should be used in the reactor thread only.
|
explicit |
Constructor.
| [in,out] | reactor | Associated reactor. The reactor object must remain valid as long as this object is valid. |
|
inline |
Set the maximum number of worker threads in the pool.
This value only applies to future threads which will be created if the load increases. If the pool already started more threads, the number of active threads is not reduced. The default maximum number of worker threads in the pool for newly created ReactiveWorkerPool instances is DEFAULT_MAX_THREADS.
| [in] | count | Maximum number of worker threads in the pool. |
|
inline |
Get the maximum number of worker threads in the pool.
|
inline |
Get the current number of worker threads in the pool.
|
inline |
Get the current number of worker threads currently busy executing worker tasks.
|
inline |
Set the stack size in bytes of worker threads in the pool.
This value only applies to future threads which will be created if the load increases. By default, the stack size of new threads depends on the operating system.
| [in] | size | Stack size in bytes of new worker threads in the pool. When zero, use the system default stack size. |
| bool ts::ReactiveWorkerPool::startWork | ( | ReactiveWorkerInterface * | work, |
| ReactiveWorkerHandlerInterface * | handler, | ||
| const ObjectPtr & | user_data = ObjectPtr() |
||
| ) |
Start a lengthy task in a delegated worker thread and return immediately.
If a worker thread is idle, the work is immediately passed to it. Otherwise, if the maximum number of worker threads is not reached yet, a new worker thread is created to execute the work. Otherwise, the work is delayed until a worker thread becomes idle.
| [in] | work | An object instance which executes the work. Cannot be null. |
| [in] | handler | Handler class to call when the work completes. If nullptr, no handler is called. |
| [in] | user_data | A shared pointer which will be passed unmodified to work and handler. |
| void ts::ReactiveWorkerPool::waitForTermination | ( | bool | cancel_pending = false | ) |
Wait for all worker threads to become idle and terminate all worker threads.
| [in] | cancel_pending | If true, cancel works which are not started yet. |
|
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.