Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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");
}
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand All @@ -1721,8 +1720,8 @@ public boolean localReset(final int code) throws IOException {
}

@Override
public long getLocalResetNanos() {
return localResetNanos;
public long getLocalResetTime() {
return localResetTime;
}

@Override
Expand Down Expand Up @@ -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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,10 +68,10 @@ default void terminate() {
}
}

long getLocalResetNanos();
long getLocalResetTime();

default boolean isLocalReset() {
return getLocalResetNanos() != Long.MIN_VALUE;
return getLocalResetTime() > 0;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Boolean> callback = (Callback<Boolean>) Mockito.mock(Callback.class);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -136,7 +136,7 @@ public String toString() {

// Getter for enqueueTime
@Internal
public long getEnqueueNanos() {
public long getEnqueueTime() {
return enqueueTime;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,6 @@ public interface IOWorkerStats {
int pendingChannelCount();

// Cheap
long lastSelectNano();
long lastSelectMilli();

}
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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);
Expand All @@ -95,7 +95,7 @@ Timeout getTimeout() {

@Override
long getLastEventTime() {
return creationNanos;
return creationTimeMillis;
}

@Override
Expand Down
Loading
Loading