UDP socket for use in a Reactor environment.
More...
#include <tsReactiveUDPSocket.h>
|
| | ReactiveUDPSocket (Reactor &reactor, UDPSocket &socket) |
| | Constructor.
|
| |
|
virtual | ~ReactiveUDPSocket () override |
| | Destructor.
|
| |
| void | cancelSendReceive (bool silent=false) |
| | Cancel any pending send or receive operation on this socket.
|
| |
| NonBlockingDevice & | device () |
| | Get a reference to the associated non-blocking device.
|
| |
| bool | isOpen () const |
| | Check if the reactive socket is open.
|
| |
| 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.
|
| |
| UDPSocket & | socket () |
| | Get a reference to the associated socket.
|
| |
| bool | startClose (ReactiveUDPHandlerInterface *handler, bool silent=false, const ObjectPtr &user_data=ObjectPtr()) |
| | Start closing the socket.
|
| |
| bool | startReceive (ReactiveUDPHandlerInterface *handler, size_t max_message_size=IP_MAX_PACKET_SIZE, const ObjectPtr &user_data=ObjectPtr()) |
| | Start the operation of receiving messages from the socket.
|
| |
| bool | startSend (ReactiveUDPHandlerInterface *handler, const void *data, size_t size, const IPSocketAddress &destination, const ObjectPtr &user_data=ObjectPtr()) |
| | Start the operation of sending a message to a destination address and port.
|
| |
| bool | startSend (ReactiveUDPHandlerInterface *handler, const void *data, size_t size, const ObjectPtr &user_data=ObjectPtr()) |
| | Start the operation of sending a message to the default destination address and port.
|
| |
|
| static int | SilentLevel (bool silent, int default_severity=Severity::Error) |
| | Compute a log severity level from a "silent" parameter.
|
| |
|
| 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.
|
| |
|
| 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.
|
| |
| 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 | 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.
|
| |
| std::shared_ptr< IOSB > | removeFromQueue (IOQueue &queue, IOSB *iosb) |
| | Search and remove a shared_ptr to IOSB, based on an IOSB address.
|
| |
| bool | uncheckedSignalQueuedOperations () |
| | Trigger the execution of processQueuedOperations() from another thread.
|
| |
UDP socket for use in a Reactor environment.
The class ReactiveUDPSocket is a wrapper around UDPSocket to handle reactive I/O.
The actual socket is a separate object. It is initialized and configured by the application. The application shall not directly call send(), receive(), or close() on this socket and delegate these operations to startSend(), startReceive(), and startClose() in class ReactiveUDPSocket.
◆ IOQueue
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
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.
◆ ReactiveUDPSocket()
| ts::ReactiveUDPSocket::ReactiveUDPSocket |
( |
Reactor & |
reactor, |
|
|
UDPSocket & |
socket |
|
) |
| |
Constructor.
- Parameters
-
| [in,out] | reactor | Associated reactor. The reactor object must remain valid as long as this object is valid. |
| [in,out] | socket | Associated socket. The socket object must remain valid as long as this object is valid. The ReactiveUDPSocket must be initialized before the socket is opened. |
◆ socket()
| UDPSocket & ts::ReactiveUDPSocket::socket |
( |
| ) |
|
|
inline |
Get a reference to the associated socket.
- Returns
- A reference to the associated socket.
◆ isOpen()
| bool ts::ReactiveUDPSocket::isOpen |
( |
| ) |
const |
|
inline |
Check if the reactive socket is open.
This is different from Socket::isOpen() during the closing phase, after startClose() has been called but before the underlying socket is fully closed.
- Returns
- True if the reactive socket is open, false if the underlying socket is closed or if startClose() has been called.
◆ startSend() [1/2]
Start the operation of sending a message to a destination address and port.
- Parameters
-
| [in] | handler | Handler class to call when the send operation completes. The method handleUDPSend() will be called. If nullptr, no handler is called. |
| [in] | data | Address of the message to send. The corresponding memory area must remain valid until the completion or cancelation of the send operation. |
| [in] | size | Size in bytes of the message to send. |
| [in] | destination | Socket address of the destination. Both address and port are mandatory in the socket address, they cannot be set to IPAddress::AnyAddress4 or IPSocketAddress::AnyPort. |
| [in] | user_data | A 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.
◆ startSend() [2/2]
Start the operation of sending a message to the default destination address and port.
- Parameters
-
| [in] | handler | Handler class to call when the send operation completes. The method handleUDPSend() will be called. If nullptr, no handler is called. |
| [in] | data | Address of the message to send. The corresponding memory area must remain valid until the completion or cancelation of the send operation. |
| [in] | size | Size in bytes of the message to send. |
| [in] | user_data | A 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.
◆ startReceive()
Start the operation of receiving messages from the socket.
- Parameters
-
| [in] | handler | Handler class to call each time a message is received. The method handleUDPReceive() will be called for each new datagram. Cannot be null. If a previous receive handler was registered, it is replaced. |
| [in] | max_message_size | Maximum incoming message size. Used as size of internal reception buffer. The default is the maximum IP packet size. |
| [in] | user_data | A 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.
◆ cancelSendReceive()
| void ts::ReactiveUDPSocket::cancelSendReceive |
( |
bool |
silent = false | ) |
|
Cancel any pending send or receive operation on this socket.
If a repeated reception operation is in progress, the repetition is canceled as well.
- Parameters
-
| [in] | silent | If true, do not report errors through the logger. |
◆ startClose()
Start closing the socket.
Pending asynchronous operations are canceled. The actual cancelation will take place later. In the meantime, the user's data buffers for these pending operations are busy and shall not be destroyed / deallocated by the application. The close operation terminates when the handler handleUDPClosed() is invoked. At this point, no more operation is pending and the application may get rid of data buffers.
- Parameters
-
| [in] | handler | Handler class to call when the close operation completes. The method handleUDPClosed() will be called. 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. |
- Returns
- True on success, false on error.
◆ device()
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] | queue | The queue from which to remove iosb. |
| [in] | iosb | Standard 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
-
| REQUEST | The subclass of Object which is set in react_data of all requests in inqueue. |
- Parameters
-
| [in,out] | inqueue | The queue from which all requests are removed. |
| [in,out] | outqueue | The 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] | silent | If 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] | silent | If 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] | silent | If 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] | silent | If true, do not report errors through the logger. |
◆ cancelAndWaitAsynchronousIO()
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] | iosb | The asynchronous I/O status block. |
| [in] | silent | If 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] | silent | If 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] | silent | If true, do not report errors through the logger. |
◆ createSignalQueuedOperations()
| bool ts::ReactiveBase::createSignalQueuedOperations |
( |
| ) |
|
|
protectedinherited |
◆ uncheckedSignalQueuedOperations()
| bool ts::ReactiveBase::uncheckedSignalQueuedOperations |
( |
| ) |
|
|
protectedinherited |
◆ handleUserEvent()
| virtual void ts::ReactiveBase::handleUserEvent |
( |
Reactor & |
reactor, |
|
|
EventId |
id |
|
) |
| |
|
overrideprotectedvirtualinherited |
◆ 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]
Associate this object with another Report to log errors.
- Parameters
-
| [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. |
- Returns
- The address of the previous Report object or a null pointer if there was none.
◆ setReport() [2/2]
Associate this object with another ReporterBase to log errors.
- Parameters
-
| [in] | delegate | Use 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] | mute | It 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] | silent | If true, do not report errors, report debug messages instead. |
| [in] | default_severity | Default 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] | reactor | Reactor into which the handler is invoked. |
| [in] | id | Id 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
-
◆ 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] | 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. |
The documentation for this class was generated from the following file: