diff --git a/src/java.net.http/share/classes/jdk/internal/net/http/BufferingSubscriber.java b/src/java.net.http/share/classes/jdk/internal/net/http/BufferingSubscriber.java index 273f23c1c1c..c52b96ac828 100644 --- a/src/java.net.http/share/classes/jdk/internal/net/http/BufferingSubscriber.java +++ b/src/java.net.http/share/classes/jdk/internal/net/http/BufferingSubscriber.java @@ -1,5 +1,5 @@ /* - * Copyright (c) 2017, 2023, Oracle and/or its affiliates. All rights reserved. + * Copyright (c) 2017, 2026, Oracle and/or its affiliates. All rights reserved. * DO NOT ALTER OR REMOVE COPYRIGHT NOTICES OR THIS FILE HEADER. * * This code is free software; you can redistribute it and/or modify it @@ -183,7 +183,7 @@ public class BufferingSubscriber implements TrustedSubscriber } private final SequentialScheduler pushDemandedScheduler = - new SequentialScheduler(new PushDemandedTask()); + SequentialScheduler.lockingScheduler(new PushDemandedTask()); void pushDemanded() { if (cancelled.get()) @@ -191,7 +191,7 @@ public class BufferingSubscriber implements TrustedSubscriber pushDemandedScheduler.runOrSchedule(); } - class PushDemandedTask extends SequentialScheduler.CompleteRestartableTask { + class PushDemandedTask implements Runnable { @Override public void run() { try { diff --git a/src/java.net.http/share/classes/jdk/internal/net/http/HttpConnection.java b/src/java.net.http/share/classes/jdk/internal/net/http/HttpConnection.java index 0c1388cec8a..611c54a768b 100644 --- a/src/java.net.http/share/classes/jdk/internal/net/http/HttpConnection.java +++ b/src/java.net.http/share/classes/jdk/internal/net/http/HttpConnection.java @@ -1,5 +1,5 @@ /* - * Copyright (c) 2015, 2025, Oracle and/or its affiliates. All rights reserved. + * Copyright (c) 2015, 2026, Oracle and/or its affiliates. All rights reserved. * DO NOT ALTER OR REMOVE COPYRIGHT NOTICES OR THIS FILE HEADER. * * This code is free software; you can redistribute it and/or modify it @@ -52,7 +52,6 @@ import jdk.internal.net.http.common.Demand; import jdk.internal.net.http.common.FlowTube; import jdk.internal.net.http.common.Logger; import jdk.internal.net.http.common.SequentialScheduler; -import jdk.internal.net.http.common.SequentialScheduler.DeferredCompleter; import jdk.internal.net.http.common.Log; import jdk.internal.net.http.common.Utils; @@ -554,7 +553,7 @@ abstract class HttpConnection implements Closeable { volatile Flow.Subscriber> subscriber; volatile HttpWriteSubscription subscription; final SequentialScheduler writeScheduler = - new SequentialScheduler(this::flushTask); + SequentialScheduler.lockingScheduler(this::flushTask); @Override public void subscribe(Flow.Subscriber> subscriber) { synchronized (reading) { @@ -570,13 +569,9 @@ abstract class HttpConnection implements Closeable { signal(); } - void flushTask(DeferredCompleter completer) { - try { - HttpWriteSubscription sub = subscription; - if (sub != null) sub.flush(); - } finally { - completer.complete(); - } + void flushTask() { + HttpWriteSubscription sub = subscription; + if (sub != null) sub.flush(); } void signal() { diff --git a/src/java.net.http/share/classes/jdk/internal/net/http/PullPublisher.java b/src/java.net.http/share/classes/jdk/internal/net/http/PullPublisher.java index d1019c05629..3b0bcc006c9 100644 --- a/src/java.net.http/share/classes/jdk/internal/net/http/PullPublisher.java +++ b/src/java.net.http/share/classes/jdk/internal/net/http/PullPublisher.java @@ -1,5 +1,5 @@ /* - * Copyright (c) 2016, 2025, Oracle and/or its affiliates. All rights reserved. + * Copyright (c) 2016, 2026, Oracle and/or its affiliates. All rights reserved. * DO NOT ALTER OR REMOVE COPYRIGHT NOTICES OR THIS FILE HEADER. * * This code is free software; you can redistribute it and/or modify it @@ -85,7 +85,7 @@ class PullPublisher implements Flow.Publisher { private volatile boolean completed; private volatile boolean cancelled; private volatile Throwable error; - final SequentialScheduler pullScheduler = new SequentialScheduler(new PullTask()); + final SequentialScheduler pullScheduler = SequentialScheduler.lockingScheduler(new PullTask()); private final Demand demand = new Demand(); Subscription(Flow.Subscriber subscriber, @@ -96,9 +96,9 @@ class PullPublisher implements Flow.Publisher { this.error = throwable; } - final class PullTask extends SequentialScheduler.CompleteRestartableTask { + final class PullTask implements Runnable { @Override - protected void run() { + public void run() { if (completed || cancelled) { return; } diff --git a/src/java.net.http/share/classes/jdk/internal/net/http/SocketTube.java b/src/java.net.http/share/classes/jdk/internal/net/http/SocketTube.java index ef935b008d3..27b4ef37089 100644 --- a/src/java.net.http/share/classes/jdk/internal/net/http/SocketTube.java +++ b/src/java.net.http/share/classes/jdk/internal/net/http/SocketTube.java @@ -1,5 +1,5 @@ /* - * Copyright (c) 2017, 2025, Oracle and/or its affiliates. All rights reserved. + * Copyright (c) 2017, 2026, Oracle and/or its affiliates. All rights reserved. * DO NOT ALTER OR REMOVE COPYRIGHT NOTICES OR THIS FILE HEADER. * * This code is free software; you can redistribute it and/or modify it @@ -29,15 +29,12 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.util.List; import java.util.Objects; -import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.Flow; import java.util.concurrent.atomic.AtomicReference; import java.nio.channels.SelectableChannel; import java.nio.channels.SelectionKey; import java.nio.channels.SocketChannel; import java.util.ArrayList; -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; import java.util.function.Consumer; import java.util.function.Supplier; import jdk.internal.net.http.common.BufferSupplier; @@ -46,8 +43,6 @@ import jdk.internal.net.http.common.FlowTube; import jdk.internal.net.http.common.Log; import jdk.internal.net.http.common.Logger; import jdk.internal.net.http.common.SequentialScheduler; -import jdk.internal.net.http.common.SequentialScheduler.DeferredCompleter; -import jdk.internal.net.http.common.SequentialScheduler.RestartableTask; import jdk.internal.net.http.common.Utils; /** @@ -160,34 +155,6 @@ final class SocketTube implements FlowTube { new IOException("connection closed locally", cause)); } - /** - * A restartable task used to process tasks in sequence. - */ - private static class SocketFlowTask implements RestartableTask { - final Runnable task; - private final Lock lock = new ReentrantLock(); - SocketFlowTask(Runnable task) { - this.task = task; - } - @Override - public final void run(DeferredCompleter taskCompleter) { - try { - // The logics of the sequential scheduler should ensure that - // the restartable task is running in only one thread at - // a given time: there should never be contention. - boolean locked = lock.tryLock(); - assert locked : "contention detected in SequentialScheduler"; - try { - task.run(); - } finally { - if (locked) lock.unlock(); - } - } finally { - taskCompleter.complete(); - } - } - } - // This is best effort - there's no guarantee that the printed set of values // is consistent. It should only be considered as weakly accurate - in // particular in what concerns the events states, especially when displaying @@ -682,7 +649,7 @@ final class SocketTube implements FlowTube { private final AsyncEvent subscribeEvent; InternalReadSubscription() { - readScheduler = new SequentialScheduler(new SocketFlowTask(this::read)); + readScheduler = SequentialScheduler.lockingScheduler(this::read); subscribeEvent = new AsyncTriggerEvent(this::signalError, this::handleSubscribeEvent); readEvent = new ReadEvent(channel, this); diff --git a/src/java.net.http/share/classes/jdk/internal/net/http/common/SSLFlowDelegate.java b/src/java.net.http/share/classes/jdk/internal/net/http/common/SSLFlowDelegate.java index 0f97e191b37..2fa166028be 100644 --- a/src/java.net.http/share/classes/jdk/internal/net/http/common/SSLFlowDelegate.java +++ b/src/java.net.http/share/classes/jdk/internal/net/http/common/SSLFlowDelegate.java @@ -244,17 +244,10 @@ public class SSLFlowDelegate { final ReentrantLock readBufferLock = new ReentrantLock(); final Logger debugr = Utils.getDebugLogger(this::dbgString, Utils.DEBUG); - private final class ReaderDownstreamPusher implements Runnable { - @Override - public void run() { - processData(); - } - } - Reader() { super(); scheduler = SequentialScheduler.lockingScheduler( - new ReaderDownstreamPusher()); + this::processData); this.readBuf = ByteBuffer.allocate(1024); readBuf.limit(0); // keep in read mode } @@ -588,14 +581,10 @@ public class SSLFlowDelegate { volatile boolean completing; boolean completed; // only accessed in processData - class WriterDownstreamPusher extends SequentialScheduler.CompleteRestartableTask { - @Override public void run() { processData(); } - } - Writer() { super(); writeList = Collections.synchronizedList(new LinkedList<>()); - scheduler = new SequentialScheduler(new WriterDownstreamPusher()); + scheduler = SequentialScheduler.lockingScheduler(this::processData); } @Override diff --git a/src/java.net.http/share/classes/jdk/internal/net/http/common/SequentialScheduler.java b/src/java.net.http/share/classes/jdk/internal/net/http/common/SequentialScheduler.java index adc77f4d408..f5e896fc901 100644 --- a/src/java.net.http/share/classes/jdk/internal/net/http/common/SequentialScheduler.java +++ b/src/java.net.http/share/classes/jdk/internal/net/http/common/SequentialScheduler.java @@ -1,5 +1,5 @@ /* - * Copyright (c) 2016, 2023, Oracle and/or its affiliates. All rights reserved. + * Copyright (c) 2016, 2026, Oracle and/or its affiliates. All rights reserved. * DO NOT ALTER OR REMOVE COPYRIGHT NOTICES OR THIS FILE HEADER. * * This code is free software; you can redistribute it and/or modify it @@ -111,7 +111,7 @@ public final class SequentialScheduler { * later time, and maybe in different thread. This type exists for * readability purposes at use-sites only. */ - public abstract static class DeferredCompleter { + private abstract static class DeferredCompleter { /** Extensible from this (outer) class ONLY. */ private DeferredCompleter() { } @@ -124,7 +124,7 @@ public final class SequentialScheduler { * A restartable task. */ @FunctionalInterface - public interface RestartableTask { + private interface RestartableTask { /** * The body of the task. @@ -140,7 +140,7 @@ public final class SequentialScheduler { * A simple and self-contained task that completes once its {@code run} * method returns. */ - public abstract static class CompleteRestartableTask + private abstract static class CompleteRestartableTask implements RestartableTask { @Override @@ -161,7 +161,7 @@ public final class SequentialScheduler { * memory visibility between runs. Since the main loop can't run concurrently, * the lock shouldn't be contended and no deadlock should ever be possible. */ - public static final class LockingRestartableTask + private static final class LockingRestartableTask extends CompleteRestartableTask { private final Runnable mainLoop; @@ -208,7 +208,7 @@ public final class SequentialScheduler { } } - public SequentialScheduler(RestartableTask restartableTask) { + private SequentialScheduler(RestartableTask restartableTask) { this.restartableTask = requireNonNull(restartableTask); this.completer = new TryEndDeferredCompleter(); this.schedulableTask = new SchedulableTask(); diff --git a/src/java.net.http/share/classes/jdk/internal/net/http/websocket/TransportImpl.java b/src/java.net.http/share/classes/jdk/internal/net/http/websocket/TransportImpl.java index c78555c8f6e..fa3529a0729 100644 --- a/src/java.net.http/share/classes/jdk/internal/net/http/websocket/TransportImpl.java +++ b/src/java.net.http/share/classes/jdk/internal/net/http/websocket/TransportImpl.java @@ -1,5 +1,5 @@ /* - * Copyright (c) 2017, 2023, Oracle and/or its affiliates. All rights reserved. + * Copyright (c) 2017, 2026, Oracle and/or its affiliates. All rights reserved. * DO NOT ALTER OR REMOVE COPYRIGHT NOTICES OR THIS FILE HEADER. * * This code is free software; you can redistribute it and/or modify it @@ -29,7 +29,6 @@ import jdk.internal.net.http.common.Demand; import jdk.internal.net.http.common.Logger; import jdk.internal.net.http.common.MinimalFuture; import jdk.internal.net.http.common.SequentialScheduler; -import jdk.internal.net.http.common.SequentialScheduler.CompleteRestartableTask; import jdk.internal.net.http.common.Utils; import java.io.IOException; @@ -58,7 +57,7 @@ public class TransportImpl implements Transport { /* Used for correlating enters to and exists from a method */ private final AtomicLong counter = new AtomicLong(); - private final SequentialScheduler sendScheduler = new SequentialScheduler(new SendTask()); + private final SequentialScheduler sendScheduler = SequentialScheduler.lockingScheduler(new SendTask()); private final MessageQueue queue; private final MessageEncoder encoder = new MessageEncoder(); @@ -93,7 +92,7 @@ public class TransportImpl implements Transport { // To ensure the initial non-final `data` will be visible // (happens-before) when `readEvent.handle()` invokes `receiveScheduler` // the following assignment is done last: - receiveScheduler = new SequentialScheduler(new ReceiveTask()); + receiveScheduler = SequentialScheduler.lockingScheduler(new ReceiveTask()); } private ByteBuffer createWriteBuffer() { @@ -361,7 +360,7 @@ public class TransportImpl implements Transport { } @SuppressWarnings({"rawtypes"}) - private class SendTask extends CompleteRestartableTask { + private class SendTask implements Runnable { private final MessageQueue.QueueCallback encodingCallback = new MessageQueue.QueueCallback<>() { @@ -654,7 +653,7 @@ public class TransportImpl implements Transport { } } - private class ReceiveTask extends CompleteRestartableTask { + private class ReceiveTask implements Runnable { @Override public void run() { diff --git a/src/java.net.http/share/classes/jdk/internal/net/http/websocket/WebSocketImpl.java b/src/java.net.http/share/classes/jdk/internal/net/http/websocket/WebSocketImpl.java index aa9b027e200..98b6a47be58 100644 --- a/src/java.net.http/share/classes/jdk/internal/net/http/websocket/WebSocketImpl.java +++ b/src/java.net.http/share/classes/jdk/internal/net/http/websocket/WebSocketImpl.java @@ -1,5 +1,5 @@ /* - * Copyright (c) 2015, 2018, Oracle and/or its affiliates. All rights reserved. + * Copyright (c) 2015, 2026, Oracle and/or its affiliates. All rights reserved. * DO NOT ALTER OR REMOVE COPYRIGHT NOTICES OR THIS FILE HEADER. * * This code is free software; you can redistribute it and/or modify it @@ -116,7 +116,7 @@ public final class WebSocketImpl implements WebSocket { private final AtomicBoolean pendingPingOrPong = new AtomicBoolean(); private final Transport transport; private final SequentialScheduler receiveScheduler - = new SequentialScheduler(new ReceiveTask()); + = SequentialScheduler.lockingScheduler(new ReceiveTask()); private final Demand demand = new Demand(); private final Executor clientExecutor; @@ -416,7 +416,7 @@ public final class WebSocketImpl implements WebSocket { * - after the state has been observed as CLOSE/ERROR, the scheduler * is stopped */ - private class ReceiveTask extends SequentialScheduler.CompleteRestartableTask { + private class ReceiveTask implements Runnable { // Transport only asked here and nowhere else because we must make sure // onOpen is invoked first and no messages become pending before onOpen diff --git a/test/jdk/java/net/httpclient/whitebox/java.net.http/jdk/internal/net/http/SSLEchoTubeTest.java b/test/jdk/java/net/httpclient/whitebox/java.net.http/jdk/internal/net/http/SSLEchoTubeTest.java index 4316a46a25c..3041992ac2a 100644 --- a/test/jdk/java/net/httpclient/whitebox/java.net.http/jdk/internal/net/http/SSLEchoTubeTest.java +++ b/test/jdk/java/net/httpclient/whitebox/java.net.http/jdk/internal/net/http/SSLEchoTubeTest.java @@ -265,7 +265,7 @@ public class SSLEchoTubeTest extends AbstractSSLTubeTest { private final Queue queue = new ConcurrentLinkedQueue<>(); private final int maxQueueSize; private final SequentialScheduler processingScheduler = - new SequentialScheduler(createProcessingTask()); + SequentialScheduler.lockingScheduler(createProcessingTask()); /* Writing into this tube */ private volatile long requested; @@ -360,11 +360,11 @@ public class SSLEchoTubeTest extends AbstractSSLTubeTest { } int transmitted = 0; - private SequentialScheduler.RestartableTask createProcessingTask() { - return new SequentialScheduler.CompleteRestartableTask() { + private Runnable createProcessingTask() { + return new Runnable() { @Override - protected void run() { + public void run() { try { while (!cancelled.get()) { Object item = queue.peek(); @@ -374,39 +374,36 @@ public class SSLEchoTubeTest extends AbstractSSLTubeTest { requestMore(); return; } - try { - System.out.printf("EchoTube processing item, requested=%s, demand=%s, transmitted=%s%n", - requested, demand.get(), transmitted); - if (item instanceof List) { - if (!demand.tryDecrement()) { - System.out.println("EchoTube no demand"); - return; - } - @SuppressWarnings("unchecked") - List bytes = (List) item; - Object removed = queue.remove(); - assert removed == item; - System.out.println("EchoTube processing " - + Utils.remaining(bytes)); - transmitted++; - subscriber.onNext(bytes); - requestMore(); - } else if (item instanceof Throwable) { - cancelled.set(true); - Object removed = queue.remove(); - assert removed == item; - System.out.println("EchoTube processing " + item); - subscriber.onError((Throwable) item); - } else if (item == EOF) { - cancelled.set(true); - Object removed = queue.remove(); - assert removed == item; - System.out.println("EchoTube processing EOF"); - subscriber.onComplete(); - } else { - throw new InternalError(String.valueOf(item)); + System.out.printf("EchoTube processing item, requested=%s, demand=%s, transmitted=%s%n", + requested, demand.get(), transmitted); + if (item instanceof List) { + if (!demand.tryDecrement()) { + System.out.println("EchoTube no demand"); + return; } - } finally { + @SuppressWarnings("unchecked") + List bytes = (List) item; + Object removed = queue.remove(); + assert removed == item; + System.out.println("EchoTube processing " + + Utils.remaining(bytes)); + transmitted++; + subscriber.onNext(bytes); + requestMore(); + } else if (item instanceof Throwable) { + cancelled.set(true); + Object removed = queue.remove(); + assert removed == item; + System.out.println("EchoTube processing " + item); + subscriber.onError((Throwable) item); + } else if (item == EOF) { + cancelled.set(true); + Object removed = queue.remove(); + assert removed == item; + System.out.println("EchoTube processing EOF"); + subscriber.onComplete(); + } else { + throw new InternalError(String.valueOf(item)); } } } catch(Throwable t) {