Class SingleCoreIOReactor

java.lang.Object
org.apache.hc.core5.reactor.AbstractSingleCoreIOReactor
org.apache.hc.core5.reactor.SingleCoreIOReactor
All Implemented Interfaces:
Closeable, AutoCloseable, ModalCloseable, ConnectionInitiator, IOReactor, IOWorkerStats

class SingleCoreIOReactor extends AbstractSingleCoreIOReactor implements ConnectionInitiator, IOWorkerStats
  • Field Details

    • MAX_CHANNEL_REQUESTS

      private static final int MAX_CHANNEL_REQUESTS
      See Also:
    • eventHandlerFactory

      private final IOEventHandlerFactory eventHandlerFactory
    • reactorConfig

      private final IOReactorConfig reactorConfig
    • ioSessionDecorator

      private final Decorator<IOSession> ioSessionDecorator
    • sessionListener

      private final IOSessionListener sessionListener
    • sessionShutdownCallback

      private final Callback<IOSession> sessionShutdownCallback
    • closedSessions

      private final Queue<IOSession> closedSessions
    • channelQueue

      private final Queue<ChannelEntry> channelQueue
    • requestQueue

      private final Queue<IOSessionRequest> requestQueue
    • shutdownInitiated

      private final 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 AtomicLong totalWaitTime
    • processedRequestCount

      private final AtomicInteger processedRequestCount
  • Constructor Details

  • Method Details

    • enqueueChannel

      void enqueueChannel(ChannelEntry entry) throws IOReactorShutdownException
      Throws:
      IOReactorShutdownException
    • doTerminate

      void doTerminate()
      Specified by:
      doTerminate in class AbstractSingleCoreIOReactor
    • doExecute

      void doExecute() throws IOException
      Specified by:
      doExecute in class AbstractSingleCoreIOReactor
      Throws:
      IOException
    • initiateSessionShutdown

      private void initiateSessionShutdown()
    • validateActiveChannels

      private void validateActiveChannels()
    • processEvents

      private void processEvents(Set<SelectionKey> selectedKeys)
    • processPendingChannels

      private void processPendingChannels() throws IOException
      Throws:
      IOException
    • processClosedSessions

      private void processClosedSessions()
    • checkTimeout

      private void checkTimeout(SelectionKey key, long nowMillis)
    • connect

      public Future<IOSession> connect(NamedEndpoint remoteEndpoint, SocketAddress remoteAddress, SocketAddress localAddress, Timeout timeout, Object attachment, FutureCallback<IOSession> callback) throws IOReactorShutdownException
      Description copied from interface: ConnectionInitiator
      Requests 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 Future interface.

      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 FutureCallback instance to get notified of events related to session requests, such as request completion, cancellation, failure or timeout.

      Specified by:
      connect in interface ConnectionInitiator
      Parameters:
      remoteEndpoint - name of the remote host.
      remoteAddress - remote socket address.
      localAddress - local socket address. Can be null, in which can the default local address and a random port will be used.
      timeout - connect timeout.
      attachment - the attachment object. Can be null.
      callback - interface. Can be null.
      Returns:
      session request object.
      Throws:
      IOReactorShutdownException
    • prepareSocket

      private void prepareSocket(SocketChannel socketChannel) throws IOException
      Throws:
      IOException
    • validateAddress

      private void validateAddress(SocketAddress address) throws UnknownHostException
      Throws:
      UnknownHostException
    • processPendingConnectionRequests

      private void processPendingConnectionRequests()
    • openSocketFor

      private static SocketChannel openSocketFor(SocketAddress remoteAddress) throws IOException
      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:
      totalChannelCount in interface IOWorkerStats
    • pendingChannelCount

      public int pendingChannelCount()
      Specified by:
      pendingChannelCount in interface IOWorkerStats
    • lastSelectMilli

      public long lastSelectMilli()
      Specified by:
      lastSelectMilli in interface IOWorkerStats