Class SingleCoreIOReactor
- java.lang.Object
-
- org.apache.hc.core5.reactor.AbstractSingleCoreIOReactor
-
- org.apache.hc.core5.reactor.SingleCoreIOReactor
-
- All Implemented Interfaces:
java.io.Closeable,java.lang.AutoCloseable,ModalCloseable,ConnectionInitiator,IOReactor,IOWorkerStats
class SingleCoreIOReactor extends AbstractSingleCoreIOReactor implements ConnectionInitiator, IOWorkerStats
-
-
Field Summary
Fields Modifier and Type Field Description private java.util.Queue<ChannelEntry>channelQueueprivate java.util.Queue<IOSession>closedSessionsprivate IOEventHandlerFactoryeventHandlerFactoryprivate Decorator<IOSession>ioSessionDecoratorprivate longlastSelectMillisprivate longlastTimeoutCheckMillisprivate static intMAX_CHANNEL_REQUESTSprivate java.util.concurrent.atomic.AtomicIntegerprocessedRequestCountprivate IOReactorConfigreactorConfigprivate java.util.Queue<IOSessionRequest>requestQueueprivate longselectTimeoutMillisprivate IOSessionListenersessionListenerprivate Callback<IOSession>sessionShutdownCallbackprivate java.util.concurrent.atomic.AtomicBooleanshutdownInitiatedprivate IOReactorMetricsListenerthreadPoolListenerprivate java.util.concurrent.atomic.AtomicLongtotalWaitTime-
Fields inherited from class org.apache.hc.core5.reactor.AbstractSingleCoreIOReactor
selector
-
-
Constructor Summary
Constructors Constructor Description SingleCoreIOReactor(Callback<java.lang.Exception> exceptionCallback, IOEventHandlerFactory eventHandlerFactory, IOReactorConfig reactorConfig, Decorator<IOSession> ioSessionDecorator, IOSessionListener sessionListener, IOReactorMetricsListener threadPoolListener, Callback<IOSession> sessionShutdownCallback)
-
Method Summary
All Methods Static Methods Instance Methods Concrete Methods Modifier and Type Method Description private voidcheckTimeout(java.nio.channels.SelectionKey key, long nowMillis)private voidcloseOpenChannels()private voidclosePendingChannels()private voidclosePendingConnectionRequests()java.util.concurrent.Future<IOSession>connect(NamedEndpoint remoteEndpoint, java.net.SocketAddress remoteAddress, java.net.SocketAddress localAddress, Timeout timeout, java.lang.Object attachment, FutureCallback<IOSession> callback)Requests a connection to a remote host.(package private) voiddoExecute()(package private) voiddoTerminate()(package private) voidenqueueChannel(ChannelEntry entry)private voidinitiateSessionShutdown()longlastSelectMilli()private static java.nio.channels.SocketChannelopenSocketFor(java.net.SocketAddress remoteAddress)intpendingChannelCount()private voidprepareSocket(java.nio.channels.SocketChannel socketChannel)private voidprocessClosedSessions()private voidprocessConnectionRequest(java.nio.channels.SocketChannel socketChannel, IOSessionRequest sessionRequest)private voidprocessEvents(java.util.Set<java.nio.channels.SelectionKey> selectedKeys)private voidprocessPendingChannels()private voidprocessPendingConnectionRequests()private voidreportStatusToThreadPoolListener()Reports the current status of the I/O reactor's thread pool to the configured metrics listener.inttotalChannelCount()private voidvalidateActiveChannels()private voidvalidateAddress(java.net.SocketAddress address)-
Methods inherited from class org.apache.hc.core5.reactor.AbstractSingleCoreIOReactor
awaitShutdown, close, close, close, execute, getStatus, initiateShutdown, logException, toString
-
-
-
-
Field Detail
-
MAX_CHANNEL_REQUESTS
private static final int MAX_CHANNEL_REQUESTS
- See Also:
- Constant Field Values
-
eventHandlerFactory
private final IOEventHandlerFactory eventHandlerFactory
-
reactorConfig
private final IOReactorConfig reactorConfig
-
sessionListener
private final IOSessionListener sessionListener
-
closedSessions
private final java.util.Queue<IOSession> closedSessions
-
channelQueue
private final java.util.Queue<ChannelEntry> channelQueue
-
requestQueue
private final java.util.Queue<IOSessionRequest> requestQueue
-
shutdownInitiated
private final java.util.concurrent.atomic.AtomicBoolean shutdownInitiated
-
selectTimeoutMillis
private final long selectTimeoutMillis
-
lastTimeoutCheckMillis
private volatile long lastTimeoutCheckMillis
-
lastSelectMillis
private volatile long lastSelectMillis
-
threadPoolListener
private final IOReactorMetricsListener threadPoolListener
-
totalWaitTime
private final java.util.concurrent.atomic.AtomicLong totalWaitTime
-
processedRequestCount
private final java.util.concurrent.atomic.AtomicInteger processedRequestCount
-
-
Constructor Detail
-
SingleCoreIOReactor
SingleCoreIOReactor(Callback<java.lang.Exception> exceptionCallback, IOEventHandlerFactory eventHandlerFactory, IOReactorConfig reactorConfig, Decorator<IOSession> ioSessionDecorator, IOSessionListener sessionListener, IOReactorMetricsListener threadPoolListener, Callback<IOSession> sessionShutdownCallback)
-
-
Method Detail
-
enqueueChannel
void enqueueChannel(ChannelEntry entry) throws IOReactorShutdownException
- Throws:
IOReactorShutdownException
-
doTerminate
void doTerminate()
- Specified by:
doTerminatein classAbstractSingleCoreIOReactor
-
doExecute
void doExecute() throws java.io.IOException- Specified by:
doExecutein classAbstractSingleCoreIOReactor- Throws:
java.io.IOException
-
initiateSessionShutdown
private void initiateSessionShutdown()
-
validateActiveChannels
private void validateActiveChannels()
-
processEvents
private void processEvents(java.util.Set<java.nio.channels.SelectionKey> selectedKeys)
-
processPendingChannels
private void processPendingChannels() throws java.io.IOException- Throws:
java.io.IOException
-
processClosedSessions
private void processClosedSessions()
-
checkTimeout
private void checkTimeout(java.nio.channels.SelectionKey key, long nowMillis)
-
connect
public java.util.concurrent.Future<IOSession> connect(NamedEndpoint remoteEndpoint, java.net.SocketAddress remoteAddress, java.net.SocketAddress localAddress, Timeout timeout, java.lang.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
private void prepareSocket(java.nio.channels.SocketChannel socketChannel) throws java.io.IOException- Throws:
java.io.IOException
-
validateAddress
private void validateAddress(java.net.SocketAddress address) throws java.net.UnknownHostException- Throws:
java.net.UnknownHostException
-
processPendingConnectionRequests
private void processPendingConnectionRequests()
-
openSocketFor
private static java.nio.channels.SocketChannel openSocketFor(java.net.SocketAddress remoteAddress) throws java.io.IOException- Throws:
java.io.IOException
-
processConnectionRequest
private void processConnectionRequest(java.nio.channels.SocketChannel socketChannel, IOSessionRequest sessionRequest) throws java.io.IOException- Throws:
java.io.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
-
-