diff --git a/client/src/main/java/org/asynchttpclient/netty/channel/Http2ConnectionState.java b/client/src/main/java/org/asynchttpclient/netty/channel/Http2ConnectionState.java index 580ab0799..9fb93e625 100644 --- a/client/src/main/java/org/asynchttpclient/netty/channel/Http2ConnectionState.java +++ b/client/src/main/java/org/asynchttpclient/netty/channel/Http2ConnectionState.java @@ -129,18 +129,20 @@ public boolean offerPendingOpener(Runnable opener) { * Race-free against {@link #failPendingOpeners}: that method sets {@code closed} and drains the queue under * {@code pendingLock}. An opener enqueued before the drain runs is caught by the drain; an enqueue attempt * sequenced after it observes {@code closed} here (the lock provides the happens-before) and is rejected. - * Either way no opener is left stranded. + * Either way no opener is left stranded. Opener callbacks always run after releasing {@code pendingLock}, + * because opening a stream can invoke user code and must not serialize other request submissions. * * @return {@code true} if the opener was run inline or queued; {@code false} if rejected because the * connection is draining/closed or the pending queue is full (caller must fail the request) */ public boolean offerPendingOpener(NettyResponseFuture future, Runnable opener) { + boolean runOpener = false; synchronized (pendingLock) { if (draining.get() || closed.get()) { return false; } if (tryAcquireStream()) { - opener.run(); + runOpener = true; } else { if (pendingCount >= MAX_PENDING_OPENERS) { return false; @@ -148,23 +150,27 @@ public boolean offerPendingOpener(NettyResponseFuture future, Runnable opener pendingOpeners.add(new PendingOpener(future, opener)); pendingCount++; } - return true; } + if (runOpener) { + opener.run(); + } + return true; } private void drainPendingOpeners() { - synchronized (pendingLock) { - // Open as many queued requests as there are now-free stream slots. A single stream completion - // frees exactly one slot (so this usually runs one opener), but a SETTINGS frame that RAISES - // SETTINGS_MAX_CONCURRENT_STREAMS frees several at once — drain them all here rather than waking - // only one and stalling the rest until the next completion (a missed-wakeup; the Issue #2160 - // silent-timeout class). tryAcquireStream() enforces the cap and the draining/closed gate, so - // this never over-opens; every poll is under pendingLock, so a non-empty queue always yields a - // non-null opener. - while (!pendingOpeners.isEmpty() && tryAcquireStream()) { + while (true) { + PendingOpener pending; + synchronized (pendingLock) { + // A SETTINGS increase can free several slots at once. Reserve and dequeue one opener + // atomically, then repeat so an exception cannot strand an already-dequeued batch. + if (pendingOpeners.isEmpty() || !tryAcquireStream()) { + return; + } pendingCount--; - pendingOpeners.poll().opener.run(); + pending = pendingOpeners.poll(); } + // Stream opening can invoke user callbacks such as onRequestSend. + pending.opener.run(); } } diff --git a/client/src/test/java/org/asynchttpclient/netty/channel/Http2ConnectionStateTest.java b/client/src/test/java/org/asynchttpclient/netty/channel/Http2ConnectionStateTest.java index ed7a051c2..b6e4ad8d3 100644 --- a/client/src/test/java/org/asynchttpclient/netty/channel/Http2ConnectionStateTest.java +++ b/client/src/test/java/org/asynchttpclient/netty/channel/Http2ConnectionStateTest.java @@ -24,6 +24,7 @@ import java.util.concurrent.CyclicBarrier; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -270,6 +271,102 @@ public void pendingOpenerRunsOnRelease() { assertEquals(1, executionCount.get(), "Pending opener should have been executed on release"); } + @Test + public void immediateOpenerRunsOutsidePendingLock() throws Exception { + Http2ConnectionState state = new Http2ConnectionState(); + state.updateMaxConcurrentStreams(2); + CountDownLatch openerStarted = new CountDownLatch(1); + CountDownLatch releaseOpener = new CountDownLatch(1); + ExecutorService executor = Executors.newFixedThreadPool(2); + + try { + Future blockingOffer = executor.submit(() -> + state.offerPendingOpener(blockingOpener(openerStarted, releaseOpener))); + assertTrue(openerStarted.await(5, TimeUnit.SECONDS), "first opener should start"); + + Future competingOffer = executor.submit(() -> state.offerPendingOpener(() -> { })); + assertTrue(competingOffer.get(5, TimeUnit.SECONDS), + "a running opener must not hold pendingLock"); + + releaseOpener.countDown(); + assertTrue(blockingOffer.get(5, TimeUnit.SECONDS)); + state.releaseStream(); + state.releaseStream(); + } finally { + releaseOpener.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + @Test + public void queuedOpenerRunsOutsidePendingLock() throws Exception { + Http2ConnectionState state = new Http2ConnectionState(); + state.updateMaxConcurrentStreams(1); + assertTrue(state.tryAcquireStream()); + CountDownLatch openerStarted = new CountDownLatch(1); + CountDownLatch releaseOpener = new CountDownLatch(1); + state.addPendingOpener(blockingOpener(openerStarted, releaseOpener)); + ExecutorService executor = Executors.newFixedThreadPool(2); + + try { + Future drain = executor.submit(state::releaseStream); + assertTrue(openerStarted.await(5, TimeUnit.SECONDS), "queued opener should start"); + + Future competingOffer = executor.submit(() -> state.offerPendingOpener(() -> { })); + assertTrue(competingOffer.get(5, TimeUnit.SECONDS), + "a drained opener must not hold pendingLock"); + + releaseOpener.countDown(); + drain.get(5, TimeUnit.SECONDS); + state.releaseStream(); + state.releaseStream(); + } finally { + releaseOpener.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + @Test + public void raisedLimitDrainsMultiplePendingOpenersInOrder() { + Http2ConnectionState state = new Http2ConnectionState(); + state.updateMaxConcurrentStreams(1); + assertTrue(state.tryAcquireStream()); + List executionOrder = new ArrayList<>(); + state.addPendingOpener(() -> executionOrder.add(1)); + state.addPendingOpener(() -> executionOrder.add(2)); + state.addPendingOpener(() -> executionOrder.add(3)); + + state.updateMaxConcurrentStreams(4); + + assertEquals(List.of(1, 2, 3), executionOrder); + assertEquals(4, state.getActiveStreams()); + } + + @Test + public void throwingBatchOpenerLeavesRemainingQueueDrainable() { + Http2ConnectionState state = new Http2ConnectionState(); + state.updateMaxConcurrentStreams(1); + assertTrue(state.tryAcquireStream()); + List executionOrder = new ArrayList<>(); + state.addPendingOpener(() -> { + executionOrder.add(1); + throw new IllegalStateException("boom"); + }); + state.addPendingOpener(() -> executionOrder.add(2)); + state.addPendingOpener(() -> executionOrder.add(3)); + + assertThrows(IllegalStateException.class, () -> state.updateMaxConcurrentStreams(4)); + assertEquals(List.of(1), executionOrder); + assertEquals(2, state.getActiveStreams()); + + state.releaseStream(); + + assertEquals(List.of(1, 2, 3), executionOrder); + assertEquals(3, state.getActiveStreams()); + } + @Test public void multiplePendingOpenersExecuteInOrder() { Http2ConnectionState state = new Http2ConnectionState(); @@ -897,4 +994,18 @@ public void releasePermitOnceIsAtomicUnderConcurrency() throws InterruptedExcept } assertEquals(rounds, totalReleases.get(), "exactly one release per round"); } + + private static Runnable blockingOpener(CountDownLatch started, CountDownLatch release) { + return () -> { + started.countDown(); + try { + if (!release.await(10, TimeUnit.SECONDS)) { + throw new AssertionError("timed out waiting to release opener"); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError("interrupted while waiting to release opener", e); + } + }; + } }