diff --git a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/AbstractH2StreamMultiplexer.java b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/AbstractH2StreamMultiplexer.java index c71026148..95df8395e 100644 --- a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/AbstractH2StreamMultiplexer.java +++ b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/AbstractH2StreamMultiplexer.java @@ -38,7 +38,6 @@ import java.util.Queue; import java.util.concurrent.ConcurrentLinkedDeque; import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Consumer; @@ -147,13 +146,13 @@ enum SettingsHandshake { READY, TRANSMITTED, ACKED } private volatile boolean peerNoRfc7540Priorities; - private static final long STREAM_TIMEOUT_GRANULARITY_NANOS = TimeUnit.SECONDS.toNanos(1); - private long lastStreamTimeoutCheckNanos; + private static final long STREAM_TIMEOUT_GRANULARITY_MILLIS = 1000; + private long lastStreamTimeoutCheckMillis; - private static final long VALIDATE_AFTER_INACTIVITY_GRANULARITY_NANOS = TimeUnit.SECONDS.toNanos(1); + private static final long VALIDATE_AFTER_INACTIVITY_GRANULARITY_MILLIS = 1000; private final Timeout validateAfterInactivity; private final Timeout pingAckTimeout; - private volatile long lastActivityNanos; + private volatile long lastActivityTime; AbstractH2StreamMultiplexer( final ProtocolIOSession ioSession, @@ -214,7 +213,7 @@ enum SettingsHandshake { READY, TRANSMITTED, ACKED } this.hPackDecoder.setMaxListSize(this.localConfig.getMaxHeaderListSize()); this.lowMark = H2Config.INIT.getInitialWindowSize() / 2; this.streamListener = streamListener; - this.lastActivityNanos = System.nanoTime(); + this.lastActivityTime = System.currentTimeMillis(); this.validateAfterInactivity = validateAfterInactivity; this.pingAckTimeout = Args.notNull(pingAckTimeout, "PING ACK timeout"); } @@ -560,8 +559,8 @@ public final void onOutput() throws HttpException, IOException { if (connState.compareTo(ConnectionHandshake.ACTIVE) <= 0 && remoteSettingState == SettingsHandshake.ACKED) { final long t = TimeValue.isPositive(validateAfterInactivity) ? - Math.max(validateAfterInactivity.toNanoseconds(), VALIDATE_AFTER_INACTIVITY_GRANULARITY_NANOS) : 0; - final boolean hasBeenIdleTooLong = t > 0 && System.nanoTime() - lastActivityNanos > t; + Math.max(validateAfterInactivity.toMilliseconds(), VALIDATE_AFTER_INACTIVITY_GRANULARITY_MILLIS) : 0; + final boolean hasBeenIdleTooLong = t > 0 && System.currentTimeMillis() - lastActivityTime > t; if (hasBeenIdleTooLong && ioSession.hasCommands() && pingHandlers.isEmpty()) { final Timeout socketTimeout = ioSession.getSocketTimeout(); ioSession.setSocketTimeout(pingAckTimeout); @@ -1545,7 +1544,7 @@ class H2StreamChannelImpl implements H2StreamChannel { private final AtomicInteger outputWindow; private volatile boolean localClosed; - private volatile long localResetNanos = Long.MIN_VALUE; + private volatile long localResetTime; H2StreamChannelImpl(final int id, final int initialInputWindowSize, final int initialOutputWindowSize) { this.id = id; @@ -1710,7 +1709,7 @@ public boolean localReset(final int code) throws IOException { return false; } localClosed = true; - localResetNanos = System.nanoTime(); + localResetTime = System.currentTimeMillis(); final RawFrame resetStream = frameFactory.createResetStream(id, code); commitFrameInternal(resetStream); @@ -1721,8 +1720,8 @@ public boolean localReset(final int code) throws IOException { } @Override - public long getLocalResetNanos() { - return localResetNanos; + public long getLocalResetTime() { + return localResetTime; } @Override @@ -1777,14 +1776,14 @@ private void checkStreamTimeouts(final long nowNanos) throws IOException { } private void validateStreamTimeouts() throws IOException { - final long nowNanos = System.nanoTime(); - if ((nowNanos - lastStreamTimeoutCheckNanos) >= STREAM_TIMEOUT_GRANULARITY_NANOS) { - lastStreamTimeoutCheckNanos = nowNanos; - checkStreamTimeouts(nowNanos); + final long nowMillis = System.currentTimeMillis(); + if ((nowMillis - lastStreamTimeoutCheckMillis) >= STREAM_TIMEOUT_GRANULARITY_MILLIS) { + lastStreamTimeoutCheckMillis = nowMillis; + checkStreamTimeouts(System.nanoTime()); } } private void updateLastActivity() { - this.lastActivityNanos = System.nanoTime(); + this.lastActivityTime = System.currentTimeMillis(); } } \ No newline at end of file diff --git a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2Stream.java b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2Stream.java index 88081f695..ad080cecc 100644 --- a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2Stream.java +++ b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2Stream.java @@ -32,7 +32,6 @@ import java.nio.channels.CancelledKeyException; import java.nio.charset.CharacterCodingException; import java.util.List; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; @@ -51,7 +50,7 @@ class H2Stream implements StreamControl { - private static final long LINGER_TIME_NANOS = TimeUnit.SECONDS.toNanos(1); + private static final long LINGER_TIME = 1000; // 1 second private final H2StreamChannel channel; private final H2StreamHandler handler; @@ -125,8 +124,8 @@ AtomicInteger getInputWindow() { } private boolean isPastLingerDeadline() { - final long localResetNanos = channel.getLocalResetNanos(); - return localResetNanos != Long.MIN_VALUE && System.nanoTime() - localResetNanos > LINGER_TIME_NANOS; + final long localResetTime = channel.getLocalResetTime(); + return localResetTime > 0 && localResetTime + LINGER_TIME < System.currentTimeMillis(); } boolean isClosedPastLingerDeadline() { diff --git a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2StreamChannel.java b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2StreamChannel.java index ed99bcf4f..d86159ef9 100644 --- a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2StreamChannel.java +++ b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2StreamChannel.java @@ -68,10 +68,10 @@ default void terminate() { } } - long getLocalResetNanos(); + long getLocalResetTime(); default boolean isLocalReset() { - return getLocalResetNanos() != Long.MIN_VALUE; + return getLocalResetTime() > 0; } } diff --git a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/nio/pool/H2ConnPool.java b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/nio/pool/H2ConnPool.java index fae155bd2..42d189650 100644 --- a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/nio/pool/H2ConnPool.java +++ b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/nio/pool/H2ConnPool.java @@ -143,8 +143,8 @@ protected void validateSession( final TimeValue timeValue = validateAfterInactivity; if (TimeValue.isNonNegative(timeValue)) { final long lastAccessTime = Math.min(ioSession.getLastReadTime(), ioSession.getLastWriteTime()); - final long deadline = lastAccessTime + timeValue.toNanoseconds(); - if (deadline <= System.nanoTime()) { + final long deadline = lastAccessTime + timeValue.toMilliseconds(); + if (deadline <= System.currentTimeMillis()) { ioSession.enqueue(new StaleCheckCommand(callback::execute), Command.Priority.NORMAL); return; } diff --git a/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/impl/nio/TestAbstractH2StreamMultiplexer.java b/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/impl/nio/TestAbstractH2StreamMultiplexer.java index 6b665c3f0..23a02b1c8 100644 --- a/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/impl/nio/TestAbstractH2StreamMultiplexer.java +++ b/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/impl/nio/TestAbstractH2StreamMultiplexer.java @@ -1893,23 +1893,23 @@ void testKeepAliveDisabledNeverEmitsPing() throws Exception { Assertions.assertTrue(frames.stream().noneMatch(FrameStub::isPing), "Disabled policy must never emit PING"); } - private static long getValidateAfterInactivityGranularityNanos() throws Exception { - final Field field = AbstractH2StreamMultiplexer.class.getDeclaredField("VALIDATE_AFTER_INACTIVITY_GRANULARITY_NANOS"); + private static long getValidateAfterInactivityGranularityMillis() throws Exception { + final Field field = AbstractH2StreamMultiplexer.class.getDeclaredField("VALIDATE_AFTER_INACTIVITY_GRANULARITY_MILLIS"); field.setAccessible(true); return field.getLong(null); } - private static void setLastActivityNanos(final AbstractH2StreamMultiplexer mux, final long nanos) throws Exception { - final Field field = AbstractH2StreamMultiplexer.class.getDeclaredField("lastActivityNanos"); + private static void setLastActivityTime(final AbstractH2StreamMultiplexer mux, final long millis) throws Exception { + final Field field = AbstractH2StreamMultiplexer.class.getDeclaredField("lastActivityTime"); field.setAccessible(true); - field.setLong(mux, nanos); + field.setLong(mux, millis); } private static void makeMuxIdle(final AbstractH2StreamMultiplexer mux, final Timeout validateAfterInactivity) throws Exception { - final long granularityNanos = getValidateAfterInactivityGranularityNanos(); - final long configuredNanos = validateAfterInactivity != null ? validateAfterInactivity.toNanoseconds() : 0; - final long effectiveNanos = configuredNanos > 0 ? Math.max(configuredNanos, granularityNanos) : 0; - setLastActivityNanos(mux, System.nanoTime() - effectiveNanos - TimeUnit.MILLISECONDS.toNanos(10)); + final long granularityMillis = getValidateAfterInactivityGranularityMillis(); + final long configuredMillis = validateAfterInactivity != null ? validateAfterInactivity.toMilliseconds() : 0; + final long effectiveMillis = configuredMillis > 0 ? Math.max(configuredMillis, granularityMillis) : 0; + setLastActivityTime(mux, System.currentTimeMillis() - effectiveMillis - 10); } @Test diff --git a/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/nio/pool/TestH2ConnPool.java b/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/nio/pool/TestH2ConnPool.java index 83651e603..b9bf2558c 100644 --- a/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/nio/pool/TestH2ConnPool.java +++ b/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/nio/pool/TestH2ConnPool.java @@ -28,7 +28,6 @@ import java.net.InetSocketAddress; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import org.apache.hc.core5.concurrent.FutureCallback; @@ -98,9 +97,8 @@ void testValidateSessionEnqueuesStaleCheck() { final IOSession session = Mockito.mock(IOSession.class); Mockito.when(session.isOpen()).thenReturn(true); - final long farPastNanos = System.nanoTime() - TimeUnit.DAYS.toNanos(1); - Mockito.when(session.getLastReadTime()).thenReturn(farPastNanos); - Mockito.when(session.getLastWriteTime()).thenReturn(farPastNanos); + Mockito.when(session.getLastReadTime()).thenReturn(0L); + Mockito.when(session.getLastWriteTime()).thenReturn(0L); @SuppressWarnings("unchecked") final Callback callback = (Callback) Mockito.mock(Callback.class); diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/AbstractSingleCoreIOReactor.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/AbstractSingleCoreIOReactor.java index 2d59f07e9..fc3de36f1 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/AbstractSingleCoreIOReactor.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/AbstractSingleCoreIOReactor.java @@ -109,14 +109,14 @@ public void execute() { @Override public final void awaitShutdown(final TimeValue waitTime) throws InterruptedException { Args.notNull(waitTime, "Wait time"); - final long deadlineNanos = System.nanoTime() + waitTime.toNanoseconds(); - long remainingNanos = waitTime.toNanoseconds(); + final long deadline = System.currentTimeMillis() + waitTime.toMilliseconds(); + long remaining = waitTime.toMilliseconds(); lock.lock(); try { while (this.status.get().compareTo(IOReactorStatus.SHUT_DOWN) < 0) { - condition.await(remainingNanos, TimeUnit.NANOSECONDS); - remainingNanos = deadlineNanos - System.nanoTime(); - if (remainingNanos <= 0) { + condition.await(remaining, TimeUnit.MILLISECONDS); + remaining = deadline - System.currentTimeMillis(); + if (remaining <= 0) { return; } } diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSession.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSession.java index 51c0aebeb..4c2421cc4 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSession.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSession.java @@ -208,26 +208,25 @@ enum Status { void setSocketTimeout(Timeout timeout); /** - * Returns monotonic nanosecond timestamp of the last read event. + * Returns the absolute millisecond timestamp of the last read event. * - * @return nanosecond timestamp obtained from {@link System#nanoTime()}. + * @return timestamp in milliseconds, compatible with {@link System#currentTimeMillis()}. */ long getLastReadTime(); /** - * Returns monotonic nanosecond timestamp of the last write event. + * Returns the absolute millisecond timestamp of the last write event. * - * @return nanosecond timestamp obtained from {@link System#nanoTime()}. + * @return timestamp in milliseconds, compatible with {@link System#currentTimeMillis()}. */ long getLastWriteTime(); /** - * Returns monotonic nanosecond timestamp of the last I/O event including - * socket timeout reset. + * Returns the absolute millisecond timestamp of the last I/O event including socket timeout reset. * * @see #getSocketTimeout() * - * @return nanosecond timestamp obtained from {@link System#nanoTime()}. + * @return timestamp in milliseconds, compatible with {@link System#currentTimeMillis()}. */ long getLastEventTime(); diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSessionImpl.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSessionImpl.java index 5ee78cd3d..d1421dc8c 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSessionImpl.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSessionImpl.java @@ -88,10 +88,10 @@ public IOSessionImpl(final String type, final SelectionKey key, final SocketChan this.id = String.format(type + "-%010d", COUNT.getAndIncrement()); this.handlerRef = new AtomicReference<>(); this.status = new AtomicReference<>(Status.ACTIVE); - final long nowNanos = System.nanoTime(); - this.lastReadTime = nowNanos; - this.lastWriteTime = nowNanos; - this.lastEventTime = nowNanos; + final long currentTimeMillis = System.currentTimeMillis(); + this.lastReadTime = currentTimeMillis; + this.lastWriteTime = currentTimeMillis; + this.lastEventTime = currentTimeMillis; } @Override @@ -221,7 +221,7 @@ public Timeout getSocketTimeout() { @Override public void setSocketTimeout(final Timeout timeout) { this.socketTimeout = Timeout.defaultsToInfinite(timeout); - this.lastEventTime = System.nanoTime(); + this.lastEventTime = System.currentTimeMillis(); } @Override @@ -255,13 +255,13 @@ public int write(final ByteBuffer src) throws IOException { @Override public void updateReadTime() { - lastReadTime = System.nanoTime(); + lastReadTime = System.currentTimeMillis(); lastEventTime = lastReadTime; } @Override public void updateWriteTime() { - lastWriteTime = System.nanoTime(); + lastWriteTime = System.currentTimeMillis(); lastEventTime = lastWriteTime; } diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSessionRequest.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSessionRequest.java index ae25e10a6..57f727194 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSessionRequest.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOSessionRequest.java @@ -73,7 +73,7 @@ public IOSessionRequest( this.closeableRef = new AtomicReference<>(); // Set the time when this request is created - this.enqueueTime = System.nanoTime(); + this.enqueueTime = System.currentTimeMillis(); } public void completed(final ProtocolIOSession ioSession) { @@ -136,7 +136,7 @@ public String toString() { // Getter for enqueueTime @Internal - public long getEnqueueNanos() { + public long getEnqueueTime() { return enqueueTime; } diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOWorkerStats.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOWorkerStats.java index f17e75d0a..75ec5fea7 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOWorkerStats.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/IOWorkerStats.java @@ -45,6 +45,6 @@ public interface IOWorkerStats { int pendingChannelCount(); // Cheap - long lastSelectNano(); + long lastSelectMilli(); } diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/InternalChannel.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/InternalChannel.java index ff18166a7..acf8c8afc 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/InternalChannel.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/InternalChannel.java @@ -57,12 +57,12 @@ final void handleIOEvent(final int ops) { } } - final boolean checkTimeout(final long nowNanos) { + final boolean checkTimeout(final long currentTimeMillis) { final Timeout timeout = getTimeout(); if (!timeout.isDisabled()) { - final long timeoutNanos = timeout.toNanoseconds(); - final long deadlineNanos = getLastEventTime() + timeoutNanos; - if (nowNanos > deadlineNanos) { + final long timeoutMillis = timeout.toMilliseconds(); + final long deadlineMillis = getLastEventTime() + timeoutMillis; + if (currentTimeMillis > deadlineMillis) { try { onTimeout(timeout); } catch (final CancelledKeyException ex) { diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/InternalConnectChannel.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/InternalConnectChannel.java index e90161611..a69ca198a 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/InternalConnectChannel.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/InternalConnectChannel.java @@ -44,7 +44,7 @@ final class InternalConnectChannel extends InternalChannel { private final InternalDataChannel dataChannel; private final IOEventHandlerFactory eventHandlerFactory; private final IOReactorConfig reactorConfig; - private final long creationNanos; + private final long creationTimeMillis; InternalConnectChannel( final SelectionKey key, @@ -60,7 +60,7 @@ final class InternalConnectChannel extends InternalChannel { this.dataChannel = dataChannel; this.eventHandlerFactory = eventHandlerFactory; this.reactorConfig = reactorConfig; - this.creationNanos = System.nanoTime(); + this.creationTimeMillis = System.currentTimeMillis(); } @Override @@ -70,8 +70,8 @@ void onIOEvent(final int readyOps) throws IOException { socketChannel.finishConnect(); } //check out connectTimeout - final long nowNanos = System.nanoTime(); - if (checkTimeout(nowNanos)) { + final long now = System.currentTimeMillis(); + if (checkTimeout(now)) { if (reactorConfig.getSocksProxyAddress() == null) { dataChannel.upgrade(eventHandlerFactory.createHandler(dataChannel, sessionRequest.attachment)); key.attach(dataChannel); @@ -95,7 +95,7 @@ Timeout getTimeout() { @Override long getLastEventTime() { - return creationNanos; + return creationTimeMillis; } @Override diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/MultiCoreIOReactor.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/MultiCoreIOReactor.java index 3799c9d7f..957a772eb 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/MultiCoreIOReactor.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/MultiCoreIOReactor.java @@ -88,25 +88,23 @@ public final void initiateShutdown() { @Override public final void awaitShutdown(final TimeValue waitTime) throws InterruptedException { Args.notNull(waitTime, "Wait time"); - final long deadlineNanos = System.nanoTime() + waitTime.toNanoseconds(); - long remainingNanos = waitTime.toNanoseconds(); + final long deadline = System.currentTimeMillis() + waitTime.toMilliseconds(); + long remaining = waitTime.toMilliseconds(); for (int i = 0; i < this.ioReactors.length; i++) { final IOReactor ioReactor = this.ioReactors[i]; if (ioReactor.getStatus().compareTo(IOReactorStatus.SHUT_DOWN) < 0) { - ioReactor.awaitShutdown(TimeValue.of(remainingNanos, TimeUnit.NANOSECONDS)); - remainingNanos = deadlineNanos - System.nanoTime(); - if (remainingNanos <= 0) { + ioReactor.awaitShutdown(TimeValue.of(remaining, TimeUnit.MILLISECONDS)); + remaining = deadline - System.currentTimeMillis(); + if (remaining <= 0) { return; } } } for (int i = 0; i < this.threads.length; i++) { final Thread thread = this.threads[i]; - final long millis = TimeUnit.NANOSECONDS.toMillis(remainingNanos); - final int nanos = (int) (remainingNanos - TimeUnit.MILLISECONDS.toNanos(millis)); - thread.join(millis, nanos); - remainingNanos = deadlineNanos - System.nanoTime(); - if (remainingNanos <= 0) { + thread.join(remaining); + remaining = deadline - System.currentTimeMillis(); + if (remaining <= 0) { return; } } diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/SingleCoreIOReactor.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/SingleCoreIOReactor.java index e5839b3e8..8e545eff9 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/SingleCoreIOReactor.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/SingleCoreIOReactor.java @@ -42,7 +42,6 @@ import java.util.Queue; import java.util.Set; import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.TimeUnit; import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -74,9 +73,8 @@ class SingleCoreIOReactor extends AbstractSingleCoreIOReactor implements Connect private final Queue requestQueue; private final AtomicBoolean shutdownInitiated; private final long selectTimeoutMillis; - private final long selectTimeoutNanos; - private volatile long lastTimeoutCheckNanos; - private volatile long lastSelectNanos; + private volatile long lastTimeoutCheckMillis; + private volatile long lastSelectMillis; private final IOReactorMetricsListener threadPoolListener; private final IOFunction socketChannelFactory; @@ -118,7 +116,6 @@ class SingleCoreIOReactor extends AbstractSingleCoreIOReactor implements Connect this.channelQueue = new ConcurrentLinkedQueue<>(); this.requestQueue = new ConcurrentLinkedQueue<>(); this.selectTimeoutMillis = this.reactorConfig.getSelectInterval().toMilliseconds(); - this.selectTimeoutNanos = TimeUnit.MILLISECONDS.toNanos(this.selectTimeoutMillis); } void enqueueChannel(final ChannelEntry entry) throws IOReactorShutdownException { @@ -154,7 +151,7 @@ void doExecute() throws IOException { } // Process selected I/O events - lastSelectNanos = System.nanoTime(); + lastSelectMillis = System.currentTimeMillis(); if (readyCount > 0) { processEvents(this.selector.selectedKeys()); } @@ -195,11 +192,11 @@ private void initiateSessionShutdown() { } private void validateActiveChannels() { - final long nowNanos = System.nanoTime(); - if ((nowNanos - this.lastTimeoutCheckNanos) >= this.selectTimeoutNanos) { - this.lastTimeoutCheckNanos = nowNanos; + final long currentTimeMillis = System.currentTimeMillis(); + if ((currentTimeMillis - this.lastTimeoutCheckMillis) >= this.selectTimeoutMillis) { + this.lastTimeoutCheckMillis = currentTimeMillis; for (final SelectionKey key : this.selector.keys()) { - checkTimeout(key, nowNanos); + checkTimeout(key, currentTimeMillis); } } } @@ -274,10 +271,10 @@ private void processClosedSessions() { } } - private void checkTimeout(final SelectionKey key, final long nowNanos) { + private void checkTimeout(final SelectionKey key, final long nowMillis) { final InternalChannel channel = (InternalChannel) key.attachment(); if (channel != null) { - channel.checkTimeout(nowNanos); + channel.checkTimeout(nowMillis); } } @@ -355,7 +352,7 @@ private void processPendingConnectionRequests() { for (int i = 0; i < MAX_CHANNEL_REQUESTS && (sessionRequest = this.requestQueue.poll()) != null; i++) { if (threadPoolListener != null) { // Calculate wait time safely without keeping long-lived state - final long waitTimeMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - sessionRequest.getEnqueueNanos()); + final long waitTimeMillis = System.currentTimeMillis() - sessionRequest.getEnqueueTime(); // Accumulate total wait time and increment count atomically totalWaitTime.addAndGet(waitTimeMillis); @@ -524,8 +521,8 @@ public int pendingChannelCount() { } @Override - public long lastSelectNano() { - return lastSelectNanos; + public long lastSelectMilli() { + return lastSelectMillis; } } diff --git a/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestInternalChannel.java b/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestInternalChannel.java index f747d46c2..cfd235609 100644 --- a/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestInternalChannel.java +++ b/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestInternalChannel.java @@ -28,7 +28,6 @@ import java.io.IOException; import java.nio.channels.CancelledKeyException; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; @@ -119,10 +118,10 @@ void handleIoEventClosesOnException() { @Test void checkTimeoutInvokesCallback() { try (TestChannel channel = new TestChannel()) { - channel.lastEventTime = System.nanoTime() - TimeUnit.SECONDS.toNanos(10); + channel.lastEventTime = 0; channel.timeout = Timeout.ofMilliseconds(1); - final boolean result = channel.checkTimeout(System.nanoTime()); + final boolean result = channel.checkTimeout(10); Assertions.assertFalse(result); Assertions.assertTrue(channel.timedOut.get()); @@ -134,7 +133,7 @@ void checkTimeoutSkipsWhenDisabled() { try (TestChannel channel = new TestChannel()) { channel.timeout = Timeout.DISABLED; - final boolean result = channel.checkTimeout(System.nanoTime()); + final boolean result = channel.checkTimeout(System.currentTimeMillis()); Assertions.assertTrue(result); Assertions.assertFalse(channel.timedOut.get()); diff --git a/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestSocksProxyProtocolHandler.java b/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestSocksProxyProtocolHandler.java index 64e06823d..9b6b1a182 100644 --- a/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestSocksProxyProtocolHandler.java +++ b/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestSocksProxyProtocolHandler.java @@ -136,7 +136,7 @@ private static final class TestIOSession implements IOSession { this.open = true; this.eventMask = 0; this.socketTimeout = Timeout.DISABLED; - this.lastReadTime = System.nanoTime(); + this.lastReadTime = System.currentTimeMillis(); this.lastWriteTime = this.lastReadTime; this.lastEventTime = this.lastReadTime; } @@ -153,7 +153,7 @@ public ByteChannel channel() { @Override public void setEventMask(final int ops) { this.eventMask = ops; - this.lastEventTime = System.nanoTime(); + this.lastEventTime = System.currentTimeMillis(); } @Override @@ -220,7 +220,7 @@ public Timeout getSocketTimeout() { @Override public void setSocketTimeout(final Timeout timeout) { this.socketTimeout = timeout; - this.lastEventTime = System.nanoTime(); + this.lastEventTime = System.currentTimeMillis(); } @Override @@ -240,13 +240,13 @@ public long getLastEventTime() { @Override public void updateReadTime() { - this.lastReadTime = System.nanoTime(); + this.lastReadTime = System.currentTimeMillis(); this.lastEventTime = this.lastReadTime; } @Override public void updateWriteTime() { - this.lastWriteTime = System.nanoTime(); + this.lastWriteTime = System.currentTimeMillis(); this.lastEventTime = this.lastWriteTime; }