Class SingleCoreIOReactor
- All Implemented Interfaces:
Closeable, AutoCloseable, ModalCloseable, ConnectionInitiator, IOReactor, IOWorkerStats
-
Field Summary
FieldsModifier and TypeFieldDescriptionprivate final Queue<ChannelEntry> private final IOEventHandlerFactoryprivate longprivate longprivate static final intprivate final AtomicIntegerprivate final IOReactorConfigprivate final Queue<IOSessionRequest> private final longprivate final IOSessionListenerprivate final AtomicBooleanprivate final IOReactorMetricsListenerprivate final AtomicLongFields inherited from class AbstractSingleCoreIOReactor
selector -
Constructor Summary
ConstructorsConstructorDescriptionSingleCoreIOReactor(Callback<Exception> exceptionCallback, IOEventHandlerFactory eventHandlerFactory, IOReactorConfig reactorConfig, Decorator<IOSession> ioSessionDecorator, IOSessionListener sessionListener, IOReactorMetricsListener threadPoolListener, Callback<IOSession> sessionShutdownCallback) -
Method Summary
Modifier and TypeMethodDescriptionprivate voidcheckTimeout(SelectionKey key, long nowMillis) private voidprivate voidprivate voidconnect(NamedEndpoint remoteEndpoint, SocketAddress remoteAddress, SocketAddress localAddress, Timeout timeout, Object attachment, FutureCallback<IOSession> callback) Requests a connection to a remote host.(package private) void(package private) void(package private) voidenqueueChannel(ChannelEntry entry) private voidlongprivate static SocketChannelopenSocketFor(SocketAddress remoteAddress) intprivate voidprepareSocket(SocketChannel socketChannel) private voidprivate voidprocessConnectionRequest(SocketChannel socketChannel, IOSessionRequest sessionRequest) private voidprocessEvents(Set<SelectionKey> selectedKeys) private voidprivate voidprivate voidReports the current status of the I/O reactor's thread pool to the configured metrics listener.intprivate voidprivate voidvalidateAddress(SocketAddress address) Methods inherited from class AbstractSingleCoreIOReactor
awaitShutdown, close, close, close, execute, getStatus, initiateShutdown, logException, toString
-
Field Details
-
MAX_CHANNEL_REQUESTS
private static final int MAX_CHANNEL_REQUESTS- See Also:
-
eventHandlerFactory
-
reactorConfig
-
ioSessionDecorator
-
sessionListener
-
sessionShutdownCallback
-
closedSessions
-
channelQueue
-
requestQueue
-
shutdownInitiated
-
selectTimeoutMillis
private final long selectTimeoutMillis -
lastTimeoutCheckMillis
private volatile long lastTimeoutCheckMillis -
lastSelectMillis
private volatile long lastSelectMillis -
threadPoolListener
-
totalWaitTime
-
processedRequestCount
-
-
Constructor Details
-
SingleCoreIOReactor
SingleCoreIOReactor(Callback<Exception> exceptionCallback, IOEventHandlerFactory eventHandlerFactory, IOReactorConfig reactorConfig, Decorator<IOSession> ioSessionDecorator, IOSessionListener sessionListener, IOReactorMetricsListener threadPoolListener, Callback<IOSession> sessionShutdownCallback)
-
-
Method Details
-
enqueueChannel
- Throws:
IOReactorShutdownException
-
doTerminate
void doTerminate()- Specified by:
doTerminatein classAbstractSingleCoreIOReactor
-
doExecute
- Specified by:
doExecutein classAbstractSingleCoreIOReactor- Throws:
IOException
-
initiateSessionShutdown
private void initiateSessionShutdown() -
validateActiveChannels
private void validateActiveChannels() -
processEvents
-
processPendingChannels
- Throws:
IOException
-
processClosedSessions
private void processClosedSessions() -
checkTimeout
-
connect
public Future<IOSession> connect(NamedEndpoint remoteEndpoint, SocketAddress remoteAddress, SocketAddress localAddress, Timeout timeout, Object attachment, FutureCallback<IOSession> callback) throws IOReactorShutdownException Description copied from interface:ConnectionInitiatorRequests a connection to a remote host.Opening a connection to a remote host usually tends to be a time consuming process and may take a while to complete. One can monitor and control the process of session initialization by means of the
Futureinterface.There are several parameters one can use to exert a greater control over the process of session initialization:
A non-null local socket address parameter can be used to bind the socket to a specific local address.
An attachment object can added to the new session's context upon initialization. This object can be used to pass an initial processing state to the protocol handler.
It is often desirable to be able to react to the completion of a session request asynchronously without having to wait for it, blocking the current thread of execution. One can optionally provide an implementation
FutureCallbackinstance to get notified of events related to session requests, such as request completion, cancellation, failure or timeout.- Specified by:
connectin interfaceConnectionInitiator- Parameters:
remoteEndpoint- name of the remote host.remoteAddress- remote socket address.localAddress- local socket address. Can benull, in which can the default local address and a random port will be used.timeout- connect timeout.attachment- the attachment object. Can benull.callback- interface. Can benull.- Returns:
- session request object.
- Throws:
IOReactorShutdownException
-
prepareSocket
- Throws:
IOException
-
validateAddress
- Throws:
UnknownHostException
-
processPendingConnectionRequests
private void processPendingConnectionRequests() -
openSocketFor
- Throws:
IOException
-
processConnectionRequest
private void processConnectionRequest(SocketChannel socketChannel, IOSessionRequest sessionRequest) throws IOException - Throws:
IOException
-
closeOpenChannels
private void closeOpenChannels() -
closePendingChannels
private void closePendingChannels() -
closePendingConnectionRequests
private void closePendingConnectionRequests() -
reportStatusToThreadPoolListener
private void reportStatusToThreadPoolListener()Reports the current status of the I/O reactor's thread pool to the configured metrics listener.This method gathers three key metrics:
- Active Threads: The number of currently active threads handling I/O sessions.
- Pending Connections: The number of connection requests waiting to be processed.
- Saturation Percentage: The ratio of active threads to the
maximum allowed connections (defined by
MAX_CHANNEL_REQUESTS), expressed as a percentage. It provides insight into how saturated the thread pool is relative to its maximum capacity. The formula for calculating saturation is:saturationPercentage = (activeThreads / MAX_CHANNEL_REQUESTS) * 100.0
If the number of pending connections exceeds
MAX_CHANNEL_REQUESTS, resource starvation is detected, and an appropriate event is reported. -
totalChannelCount
public int totalChannelCount()- Specified by:
totalChannelCountin interfaceIOWorkerStats
-
pendingChannelCount
public int pendingChannelCount()- Specified by:
pendingChannelCountin interfaceIOWorkerStats
-
lastSelectMilli
public long lastSelectMilli()- Specified by:
lastSelectMilliin interfaceIOWorkerStats
-