diff --git a/httpcore5-testing/src/test/java/org/apache/hc/core5/pool/RouteSegmentedConnPoolBenchmark.java b/httpcore5-testing/src/test/java/org/apache/hc/core5/pool/RouteSegmentedConnPoolBenchmark.java new file mode 100644 index 000000000..17fbd73f4 --- /dev/null +++ b/httpcore5-testing/src/test/java/org/apache/hc/core5/pool/RouteSegmentedConnPoolBenchmark.java @@ -0,0 +1,349 @@ +/* + * ==================================================================== + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * ==================================================================== + * + * This software consists of voluntary contributions made by many + * individuals on behalf of the Apache Software Foundation. For more + * information on the Apache Software Foundation, please see + * . + * + */ +package org.apache.hc.core5.pool; + +import java.io.IOException; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.hc.core5.io.CloseMode; +import org.apache.hc.core5.io.ModalCloseable; +import org.apache.hc.core5.util.TimeValue; +import org.apache.hc.core5.util.Timeout; +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Param; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.TearDown; +import org.openjdk.jmh.annotations.Threads; +import org.openjdk.jmh.annotations.Warmup; +import org.openjdk.jmh.infra.Blackhole; + +@BenchmarkMode(Mode.Throughput) +@OutputTimeUnit(TimeUnit.SECONDS) +@Warmup(iterations = 5, time = 1) +@Measurement(iterations = 8, time = 1) +@Fork(2) +@Threads(8) +public class RouteSegmentedConnPoolBenchmark { + + private static final int MAX_TOTAL = 64; + private static final int ROUTE_COUNT = 1024; + private static final Timeout LEASE_TIMEOUT = Timeout.ofSeconds(5); + private static final TimeValue TTL = TimeValue.NEG_ONE_MILLISECOND; + private static final TimeValue KEEP_ALIVE = TimeValue.ofMinutes(10); + private static final DummyConnection CONNECTION = new DummyConnection(); + + public enum PoolType { + STRICT, + OFFLOCK + } + + static ManagedConnPool createPool( + final PoolType type, + final int maxPerRoute, + final int maxTotal) { + + final DisposalCallback disposal = + (connection, closeMode) -> connection.close(closeMode); + + switch (type) { + case STRICT: + return new StrictConnPool( + maxPerRoute, + maxTotal, + TTL, + PoolReusePolicy.LIFO, + disposal, + null); + case OFFLOCK: + return new RouteSegmentedConnPool( + maxPerRoute, + maxTotal, + TTL, + PoolReusePolicy.LIFO, + disposal, + null); + default: + throw new IllegalStateException("Unexpected pool type: " + type); + } + } + + static PoolEntry lease( + final ManagedConnPool pool, + final Integer route) throws Exception { + return pool.lease(route, null, LEASE_TIMEOUT, null).get(); + } + + static void makeReusable(final PoolEntry entry) { + if (!entry.hasConnection()) { + entry.assignConnection(CONNECTION); + } + entry.updateExpiry(KEEP_ALIVE); + } + + @State(Scope.Benchmark) + public static class HotState { + + @Param({"STRICT", "OFFLOCK"}) + public PoolType poolType; + + ManagedConnPool pool; + Integer route; + + @Setup(Level.Trial) + public void setup() throws Exception { + pool = createPool(poolType, 1, MAX_TOTAL); + route = Integer.valueOf(0); + + final PoolEntry entry = lease(pool, route); + makeReusable(entry); + pool.release(entry, true); + } + + @TearDown(Level.Trial) + public void tearDown() throws IOException { + pool.close(CloseMode.IMMEDIATE); + } + } + + @State(Scope.Benchmark) + public static class FullHotState { + + @Param({"STRICT", "OFFLOCK"}) + public PoolType poolType; + + ManagedConnPool pool; + Integer route; + + @Setup(Level.Trial) + public void setup() throws Exception { + pool = createPool(poolType, 1, MAX_TOTAL); + + for (int i = 0; i < MAX_TOTAL; i++) { + final PoolEntry entry = + lease(pool, Integer.valueOf(i)); + makeReusable(entry); + pool.release(entry, true); + } + + route = Integer.valueOf(0); + } + + @TearDown(Level.Trial) + public void tearDown() throws IOException { + pool.close(CloseMode.IMMEDIATE); + } + } + + @State(Scope.Thread) + public static class BacklogFairnessState { + + private static final Integer SLOW_ROUTE = Integer.valueOf(0); + private static final Integer COLD_ROUTE = Integer.valueOf(1); + + @Param({"STRICT", "OFFLOCK"}) + public PoolType poolType; + + ManagedConnPool pool; + PoolEntry slow1; + PoolEntry slow2; + Future> coldWaiter; + Future> sameRouteWaiter; + PoolEntry coldEntry; + + @Setup(Level.Invocation) + public void setup() throws Exception { + pool = createPool(poolType, 2, 2); + + // Warm the disposal path outside the measured region. OFFLOCK starts + // its disposer lazily on the first graceful discard; without this, + // every sample would include worker-thread startup because this + // state creates a fresh pool for every invocation. + final CountDownLatch disposed = new CountDownLatch(1); + final PoolEntry warmup = + lease(pool, Integer.valueOf(-1)); + warmup.assignConnection(new DummyConnection(disposed)); + pool.release(warmup, false); + if (!disposed.await(1, TimeUnit.SECONDS)) { + throw new IllegalStateException("Disposer warmup timed out"); + } + + slow1 = lease(pool, SLOW_ROUTE); + slow2 = lease(pool, SLOW_ROUTE); + + // Queue another route first, then keep the route that owns all + // allocated slots permanently backlogged. + coldWaiter = pool.lease(COLD_ROUTE, null, LEASE_TIMEOUT, null); + sameRouteWaiter = pool.lease(SLOW_ROUTE, null, LEASE_TIMEOUT, null); + + makeReusable(slow1); + } + + @TearDown(Level.Invocation) + public void tearDown() throws IOException { + if (sameRouteWaiter != null) { + sameRouteWaiter.cancel(true); + } + if (coldWaiter != null && !coldWaiter.isDone()) { + coldWaiter.cancel(true); + } + if (slow2 != null) { + pool.release(slow2, false); + } + if (coldEntry != null) { + pool.release(coldEntry, false); + } + pool.close(CloseMode.IMMEDIATE); + } + } + + @State(Scope.Benchmark) + public static class ColdFullState { + + @Param({"STRICT", "OFFLOCK"}) + public PoolType poolType; + + ManagedConnPool pool; + Integer[] routes; + AtomicInteger sequence; + + @Setup(Level.Trial) + public void setup() throws Exception { + pool = createPool(poolType, 1, MAX_TOTAL); + + routes = new Integer[ROUTE_COUNT]; + for (int i = 0; i < ROUTE_COUNT; i++) { + routes[i] = Integer.valueOf(i); + } + + // Fill the global pool entirely with reusable idle entries. + for (int i = 0; i < MAX_TOTAL; i++) { + final PoolEntry entry = + lease(pool, routes[i]); + makeReusable(entry); + pool.release(entry, true); + } + + // Start outside the initially populated route window so the first + // lease necessarily exercises cross-route capacity reclamation. + sequence = new AtomicInteger(MAX_TOTAL); + } + + Integer nextColdRoute() { + return routes[sequence.getAndIncrement() & (ROUTE_COUNT - 1)]; + } + + @TearDown(Level.Trial) + public void tearDown() throws IOException { + pool.close(CloseMode.IMMEDIATE); + } + } + + @Benchmark + public void hotRoute( + final HotState state, + final Blackhole blackhole) throws Exception { + final PoolEntry entry = + lease(state.pool, state.route); + blackhole.consume(entry); + state.pool.release(entry, true); + } + + @Benchmark + public void hotRouteAtFullPool( + final FullHotState state, + final Blackhole blackhole) throws Exception { + final PoolEntry entry = + lease(state.pool, state.route); + blackhole.consume(entry); + state.pool.release(entry, true); + } + + @Benchmark + @BenchmarkMode(Mode.SampleTime) + @OutputTimeUnit(TimeUnit.MICROSECONDS) + @Threads(1) + public void coldRouteProgressUnderSameRouteBacklog( + final BacklogFairnessState state, + final Blackhole blackhole) throws Exception { + state.pool.release(state.slow1, true); + state.slow1 = null; + + state.coldEntry = state.coldWaiter.get(1, TimeUnit.SECONDS); + blackhole.consume(state.coldEntry); + } + + @Benchmark + public void coldRouteAtFullPoolWithIdle( + final ColdFullState state, + final Blackhole blackhole) throws Exception { + final PoolEntry entry = + lease(state.pool, state.nextColdRoute()); + blackhole.consume(entry); + makeReusable(entry); + state.pool.release(entry, true); + } + + static final class DummyConnection implements ModalCloseable { + + private final CountDownLatch closed; + + DummyConnection() { + this(null); + } + + DummyConnection(final CountDownLatch closed) { + this.closed = closed; + } + + @Override + public void close(final CloseMode closeMode) { + signalClosed(); + } + + @Override + public void close() { + signalClosed(); + } + + private void signalClosed() { + if (closed != null) { + closed.countDown(); + } + } + } +} diff --git a/httpcore5/src/main/java/org/apache/hc/core5/pool/RouteSegmentedConnPool.java b/httpcore5/src/main/java/org/apache/hc/core5/pool/RouteSegmentedConnPool.java index 88f516053..765f2b368 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/pool/RouteSegmentedConnPool.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/pool/RouteSegmentedConnPool.java @@ -67,8 +67,8 @@ * *

Per-route state is kept in independent segments. Disposal of connections is offloaded * to a bounded executor so slow graceful closes do not block threads leasing on other routes. - * A minimal round-robin helper is engaged only when there are many pending routes and - * there is global headroom; it never scans all routes.

+ * A minimal round-robin helper is engaged for pending routes. When the global capacity is + * exhausted, reusable idle entries can be reclaimed across segments without a global lock.

* * @param route key type * @param connection type (must be {@link ModalCloseable}) @@ -118,6 +118,10 @@ public final class RouteSegmentedConnPool implement private final AtomicBoolean draining = new AtomicBoolean(false); private final AtomicInteger pendingRouteCount = new AtomicInteger(0); + // Segments that currently have at least one reusable idle entry. This is a + // best-effort hint queue used only when the global capacity is exhausted. + private final ConcurrentLinkedQueue reclaimQueue = new ConcurrentLinkedQueue<>(); + public RouteSegmentedConnPool( final int defaultMaxPerRoute, final int maxTotal, @@ -193,10 +197,16 @@ public RouteSegmentedConnPool( } final class Segment { + final R route; final ConcurrentLinkedDeque> available = new ConcurrentLinkedDeque<>(); final ConcurrentLinkedDeque waiters = new ConcurrentLinkedDeque<>(); final AtomicInteger allocated = new AtomicInteger(0); final AtomicBoolean enqueued = new AtomicBoolean(false); + final AtomicBoolean reclaimEnqueued = new AtomicBoolean(false); + + Segment(final R route) { + this.route = route; + } int limitPerRoute(final R route) { final Integer v = maxPerRoute.get(route); @@ -253,7 +263,7 @@ public Future> lease( final FutureCallback> callback) { ensureOpen(); - final Segment seg = segments.computeIfAbsent(route, r -> new Segment()); + final Segment seg = segments.computeIfAbsent(route, Segment::new); // 1) Try available PoolEntry hit; @@ -295,7 +305,7 @@ public Future> lease( // Late hit after enqueuing final PoolEntry late = pollAvailable(seg, state); if (late != null) { - if (w.complete(late)) { + if (seg.waiters.remove(w) && w.complete(late)) { cancelTimeout(w); fireOnLease(route); if (callback != null) { @@ -320,10 +330,33 @@ public Future> lease( } if (!handedOff) { offerAvailable(seg, late); + serveOnePendingForOtherRoutes(seg); + triggerDrainIfMany(); } } } + // Capacity may have appeared after the first allocation attempt but + // before this waiter was enqueued. Retry synchronously for this exact + // waiter so single-route progress does not depend on the asynchronous + // round-robin helper. + if (!w.isDone() && tryAllocateOne(route, seg)) { + final PoolEntry entry = new PoolEntry<>(route, timeToLive, disposal, clock); + if (seg.waiters.remove(w) && w.complete(entry)) { + fireOnLease(route); + dequeueIfDrained(seg); + if (callback != null) { + callback.completed(entry); + } + return w; + } + + // The waiter was completed or cancelled concurrently after the + // allocation was reserved. Return the reservation to the pool. + seg.allocated.decrementAndGet(); + totalAllocated.decrementAndGet(); + } + scheduleTimeout(w, seg); if (callback != null) { @@ -342,6 +375,8 @@ public Future> lease( }); } + // The background helper is a throughput optimization for genuinely + // multi-route contention, not a liveness mechanism for a single waiter. triggerDrainIfMany(); return w; } @@ -363,15 +398,38 @@ public void release(final PoolEntry entry, final boolean reusable) { final boolean stillValid = reusable && !isPastTtl(entry, now) && !entry.getExpiryDeadline().isBefore(now); if (stillValid) { - if (!handOffToCompatibleWaiter(entry, seg)) { + if (hasPendingForOtherRoutes(seg)) { + // Give a pending route outside this segment one opportunity to + // claim the returned global slot before handing the connection + // straight back to a waiter of the same route. Otherwise a + // permanently backlogged route can recycle all of its slots + // forever and starve every other route. offerAvailable(seg, entry); if (!seg.waiters.isEmpty()) { enqueueIfNeeded(route, seg); - triggerDrainIfMany(); } + serveOnePendingForOtherRoutes(seg); + + // If the cross-route attempt did not reclaim this particular + // entry, preserve route-local reuse and hand it directly to a + // compatible waiter. + if (seg.available.remove(entry)) { + if (!handOffToCompatibleWaiter(entry, seg)) { + offerAvailable(seg, entry); + } + } + triggerDrainIfMany(); + } else if (!handOffToCompatibleWaiter(entry, seg)) { + offerAvailable(seg, entry); + if (!seg.waiters.isEmpty()) { + enqueueIfNeeded(route, seg); + } + triggerDrainIfMany(); } } else { discardAndDecr(entry, CloseMode.GRACEFUL); + serveOnePendingIfPossible(); + triggerDrainIfMany(); } maybeCleanupSegment(route, seg); @@ -403,6 +461,7 @@ public void close(final CloseMode closeMode) { if (seg.enqueued.getAndSet(false)) { pendingRouteCount.decrementAndGet(); } + seg.reclaimEnqueued.set(false); for (final PoolEntry p : seg.available) { if (seg.available.remove(p)) { @@ -418,6 +477,7 @@ public void close(final CloseMode closeMode) { } segments.clear(); pendingQueue.clear(); + reclaimQueue.clear(); pendingRouteCount.set(0); // Let in-flight graceful closes progress; no blocking here. @@ -447,6 +507,10 @@ public void closeIdle(final TimeValue idleTime) { } } maybeCleanupSegment(route, seg); + if (processed > 0) { + serveOnePendingIfPossible(); + triggerDrainIfMany(); + } } } @@ -472,6 +536,10 @@ public void closeExpired() { } } maybeCleanupSegment(route, seg); + if (processed > 0) { + serveOnePendingIfPossible(); + triggerDrainIfMany(); + } } } @@ -494,7 +562,12 @@ public int getMaxTotal() { @Override public void setMaxTotal(final int max) { - maxTotal.set(Math.max(1, max)); + final int updated = Math.max(1, max); + final int previous = maxTotal.getAndSet(updated); + if (updated > previous) { + serveOnePendingIfPossible(); + triggerDrainIfMany(); + } } @Override @@ -579,6 +652,8 @@ private void scheduleTimeout(final Waiter w, final Segment seg) { if (p != null) { if (!handOffToCompatibleWaiter(p, seg)) { offerAvailable(seg, p); + serveOnePendingForOtherRoutes(seg); + triggerDrainIfMany(); } } }, w.requestTimeout.toMilliseconds(), TimeUnit.MILLISECONDS); @@ -597,6 +672,7 @@ private void offerAvailable(final Segment seg, final PoolEntry p) { } else { seg.available.addLast(p); } + markReclaimable(seg); } private PoolEntry pollAvailable(final Segment seg, final Object neededState) { @@ -664,22 +740,119 @@ private void discardAndDecr(final PoolEntry p, final CloseMode mode) { private void maybeCleanupSegment(final R route, final Segment seg) { if (seg.allocated.get() == 0 && seg.available.isEmpty() && seg.waiters.isEmpty()) { - segments.remove(route, seg); - if (seg.enqueued.getAndSet(false)) { - pendingRouteCount.decrementAndGet(); + if (segments.remove(route, seg)) { + if (seg.reclaimEnqueued.getAndSet(false)) { + reclaimQueue.remove(seg); + } + if (seg.enqueued.getAndSet(false)) { + pendingRouteCount.decrementAndGet(); + } + } + } + } + + private void markReclaimable(final Segment seg) { + if (!seg.reclaimEnqueued.get() + && !seg.available.isEmpty() + && seg.reclaimEnqueued.compareAndSet(false, true)) { + reclaimQueue.offer(seg); + } + } + + /** + * Replaces one reusable idle entry with an allocation for {@code route} + * without releasing the global slot in between. This keeps the full-pool + * cold-route path free of a global lock and prevents another allocator from + * stealing the reclaimed slot. + */ + private boolean tryReclaimAndAllocate(final R route, final Segment target) { + if (target.allocated.get() >= target.limitPerRoute(route)) { + return false; + } + + for (; ; ) { + if (totalAllocated.get() != maxTotal.get()) { + return false; + } + + final Segment victimSegment = reclaimQueue.poll(); + if (victimSegment == null) { + return false; + } + victimSegment.reclaimEnqueued.set(false); + + if (segments.get(victimSegment.route) != victimSegment) { + continue; + } + + // removeLast mirrors StrictConnPool's eviction side: for LIFO this + // is the least recently released entry; for FIFO it is the newest. + final PoolEntry victim = victimSegment.available.pollLast(); + if (victim == null) { + markReclaimable(victimSegment); + continue; + } + + markReclaimable(victimSegment); + + if (victimSegment == target) { + // Replacing an idle entry in the same segment leaves both the + // per-route and global allocation counts unchanged. + discardEntry(victim, CloseMode.GRACEFUL); + return true; + } + + victimSegment.allocated.decrementAndGet(); + + for (; ; ) { + final int per = target.allocated.get(); + if (per >= target.limitPerRoute(route)) { + // The victim has already gone away and the transfer can no + // longer be completed. Convert the reserved global slot + // into real headroom. + totalAllocated.decrementAndGet(); + discardEntry(victim, CloseMode.GRACEFUL); + maybeCleanupSegment(victimSegment.route, victimSegment); + triggerDrainIfMany(); + return false; + } + if (target.allocated.compareAndSet(per, per + 1)) { + // Global total intentionally stays unchanged: the victim's + // slot has been transferred to the target segment. + discardEntry(victim, CloseMode.GRACEFUL); + maybeCleanupSegment(victimSegment.route, victimSegment); + return true; + } } } } private boolean tryAllocateOne(final R route, final Segment seg) { for (; ; ) { + if (seg.allocated.get() >= seg.limitPerRoute(route)) { + return false; + } + + final int max = maxTotal.get(); final int tot = totalAllocated.get(); - if (tot >= maxTotal.get()) { + + if (tot > max) { return false; } + if (tot == max) { + if (tryReclaimAndAllocate(route, seg)) { + return true; + } + if (totalAllocated.get() < maxTotal.get()) { + continue; + } + return false; + } + if (!totalAllocated.compareAndSet(tot, tot + 1)) { continue; } + for (; ; ) { final int per = seg.allocated.get(); if (per >= seg.limitPerRoute(route)) { @@ -706,25 +879,54 @@ private void dequeueIfDrained(final Segment seg) { } } - private void triggerDrainIfMany() { - if (pendingRouteCount.get() < RR_MIN_PENDING_ROUTES) { + private boolean hasDrainOpportunity() { + return totalAllocated.get() < maxTotal.get() || !reclaimQueue.isEmpty(); + } + + private void serveOnePendingIfPossible() { + if (pendingRouteCount.get() == 0 || !hasDrainOpportunity()) { return; } - if (totalAllocated.get() >= maxTotal.get()) { + // Scan a small bounded number of route tokens, but complete at most one + // lease inline. This is enough for liveness without turning release or + // maintenance into a global drain. + serveRoundRobin(RR_INLINE_FALLBACK_BUDGET, 1); + } + + private boolean hasPendingForOtherRoutes(final Segment seg) { + final int ownPending = seg.enqueued.get() ? 1 : 0; + return pendingRouteCount.get() > ownPending; + } + + private boolean serveOnePendingForOtherRoutes(final Segment seg) { + if (!hasPendingForOtherRoutes(seg) || !hasDrainOpportunity()) { + return false; + } + return serveRoundRobin(RR_INLINE_FALLBACK_BUDGET, 1, seg) > 0; + } + + private void triggerDrainIfMany() { + if (pendingRouteCount.get() < RR_MIN_PENDING_ROUTES || !hasDrainOpportunity()) { return; } + triggerDrain(); + } + + private void triggerDrain() { if (!draining.compareAndSet(false, true)) { return; } final Runnable task = () -> { + int created = 0; try { - serveRoundRobin(RR_BUDGET); + created = serveRoundRobin(RR_BUDGET); } finally { draining.set(false); - if (pendingRouteCount.get() >= RR_MIN_PENDING_ROUTES - && totalAllocated.get() < maxTotal.get() - && !pendingQueue.isEmpty()) { + if (created > 0 + && pendingRouteCount.get() >= RR_MIN_PENDING_ROUTES + && !pendingQueue.isEmpty() + && hasDrainOpportunity()) { triggerDrainIfMany(); } } @@ -742,10 +944,19 @@ private void triggerDrainIfMany() { } } - private void serveRoundRobin(final int budget) { + private int serveRoundRobin(final int budget) { + return serveRoundRobin(budget, budget, null); + } + + private int serveRoundRobin(final int budget, final int maxCreated) { + return serveRoundRobin(budget, maxCreated, null); + } + + private int serveRoundRobin(final int budget, final int maxCreated, final Segment excludedSegment) { int created = 0; + final int attempts = Math.min(budget, pendingRouteCount.get()); - for (; created < budget; ) { + for (int i = 0; i < attempts && created < maxCreated; i++) { final R route = pendingQueue.poll(); if (route == null) { break; @@ -754,6 +965,10 @@ private void serveRoundRobin(final int budget) { if (seg == null) { continue; } + if (seg == excludedSegment) { + pendingQueue.offer(route); + continue; + } if (seg.waiters.isEmpty()) { if (seg.enqueued.getAndSet(false)) { pendingRouteCount.decrementAndGet(); @@ -773,9 +988,13 @@ private void serveRoundRobin(final int budget) { } else { final PoolEntry entry = new PoolEntry<>(route, timeToLive, disposal, clock); cancelTimeout(w); - w.complete(entry); - fireOnLease(w.route); - created++; + if (w.complete(entry)) { + fireOnLease(w.route); + created++; + } else { + seg.allocated.decrementAndGet(); + totalAllocated.decrementAndGet(); + } } if (!seg.waiters.isEmpty()) { @@ -785,7 +1004,9 @@ private void serveRoundRobin(final int budget) { pendingRouteCount.decrementAndGet(); } } + maybeCleanupSegment(route, seg); } + return created; } /** diff --git a/httpcore5/src/test/java/org/apache/hc/core5/pool/RouteSegmentedConnPoolTest.java b/httpcore5/src/test/java/org/apache/hc/core5/pool/RouteSegmentedConnPoolTest.java index ce4b9192b..ea645cbdb 100644 --- a/httpcore5/src/test/java/org/apache/hc/core5/pool/RouteSegmentedConnPoolTest.java +++ b/httpcore5/src/test/java/org/apache/hc/core5/pool/RouteSegmentedConnPoolTest.java @@ -319,43 +319,126 @@ void getRoutesCoversAllocatedAvailableAndWaiters() throws Exception { assertTrue(pool.getRoutes().isEmpty(), "Initially there should be no routes"); - // Allocate on rA + // Allocate on rA. final PoolEntry a = pool.lease("rA", null, Timeout.ofSeconds(1), null).get(1, TimeUnit.SECONDS); - assertEquals(new HashSet(Collections.singletonList("rA")), pool.getRoutes(), - "rA must be listed because it is leased (allocated > 0)"); - // Make rA available + assertEquals( + new HashSet<>(Collections.singletonList("rA")), + pool.getRoutes(), + "rA must be listed because it is leased"); + + // Make rA idle and reclaimable. a.assignConnection(new FakeConnection()); a.updateExpiry(TimeValue.ofSeconds(30)); pool.release(a, true); - assertEquals(new HashSet<>(Collections.singletonList("rA")), pool.getRoutes(), - "rA must be listed because it has AVAILABLE entries"); - // Enqueue waiter on rB (will time out) - final Future> waiterB = - pool.lease("rB", null, Timeout.ofMilliseconds(300), null); - final Set routesNow = pool.getRoutes(); - assertTrue(routesNow.contains("rA") && routesNow.contains("rB"), - "Both rA (available) and rB (waiter) must be listed"); + assertEquals( + new HashSet<>(Collections.singletonList("rA")), + pool.getRoutes(), + "rA must be listed because it has an available entry"); - // Let rB time out (do NOT free capacity before the timeout fires) - final ExecutionException ex = assertThrows( - ExecutionException.class, - () -> waiterB.get(600, TimeUnit.MILLISECONDS)); - assertInstanceOf(TimeoutException.class, ex.getCause()); - assertEquals("Lease timed out", ex.getCause().getMessage()); + // The pool is globally full, but rA has an idle entry. Leasing rB must + // reclaim rA's global slot instead of leaving rB pending. + final PoolEntry b = + pool.lease("rB", null, Timeout.ofSeconds(1), null).get(1, TimeUnit.SECONDS); - // Now drain rA by leasing and discarding to trigger segment cleanup - final PoolEntry a2 = - pool.lease("rA", null, Timeout.ofSeconds(1), null).get(1, TimeUnit.SECONDS); - pool.release(a2, false); // discard - final Set afterDropA = pool.getRoutes(); - assertFalse(afterDropA.contains("rA"), "rA segment should be cleaned up"); - assertFalse(afterDropA.contains("rB"), "rB waiter timed out; should not remain listed"); + assertEquals("rB", b.getRoute()); + + final Set routesAfterReclaim = pool.getRoutes(); + assertFalse(routesAfterReclaim.contains("rA"), + "rA should be removed after its idle entry is reclaimed"); + assertTrue(routesAfterReclaim.contains("rB"), + "rB must be listed because it owns the reclaimed allocation"); + + final PoolStats totalStats = pool.getTotalStats(); + assertEquals(1, totalStats.getLeased()); + assertEquals(0, totalStats.getAvailable()); + assertEquals(0, totalStats.getPending()); + + pool.release(b, false); + + assertTrue(pool.getRoutes().isEmpty(), + "All routes should be gone after the final allocation is discarded"); - // Final cleanup pool.close(CloseMode.IMMEDIATE); - assertTrue(pool.getRoutes().isEmpty(), "All routes must be gone after close()"); + assertTrue(pool.getRoutes().isEmpty(), + "All routes must be gone after close()"); } + + @Test + void hotRouteProgressesAfterItsIdleEntryWasReclaimed() throws Exception { + final RouteSegmentedConnPool pool = + newPool(2, 2, TimeValue.NEG_ONE_MILLISECOND, PoolReusePolicy.LIFO, FakeConnection::close); + + final PoolEntry slow1 = + pool.lease("slow", null, Timeout.ofSeconds(1), null).get(1, TimeUnit.SECONDS); + + final PoolEntry hot = + pool.lease("hot", null, Timeout.ofSeconds(1), null).get(1, TimeUnit.SECONDS); + hot.assignConnection(new FakeConnection()); + hot.updateExpiry(TimeValue.ofSeconds(30)); + pool.release(hot, true); + + // Fill the second slot of the slow route. At maxTotal this reclaims + // the hot route's idle entry, so "hot" no longer has a local idle. + final PoolEntry slow2 = + pool.lease("slow", null, Timeout.ofSeconds(1), null).get(1, TimeUnit.SECONDS); + assertEquals(0, pool.getStats("hot").getAvailable()); + + // Queue the other route first, then keep the slow route backlogged. + final Future> hotWaiter = + pool.lease("hot", null, Timeout.ofSeconds(1), null); + final Future> slowWaiter = + pool.lease("slow", null, Timeout.ofSeconds(1), null); + + slow1.assignConnection(new FakeConnection()); + slow1.updateExpiry(TimeValue.ofSeconds(30)); + pool.release(slow1, true); + + // A reusable release from the permanently backlogged slow route must + // give another route an opportunity before recycling the slot locally. + final PoolEntry hotAgain = + hotWaiter.get(500, TimeUnit.MILLISECONDS); + assertEquals("hot", hotAgain.getRoute()); + assertFalse(slowWaiter.isDone(), + "Same-route backlog must not starve an already pending other route"); + + slowWaiter.cancel(true); + pool.release(slow2, false); + pool.release(hotAgain, false); + pool.close(CloseMode.IMMEDIATE); + } + + @Test + void coldRouteProgressesWhenBackloggedRouteReleasesReusableEntry() throws Exception { + final RouteSegmentedConnPool pool = + newPool(2, 2, TimeValue.NEG_ONE_MILLISECOND, PoolReusePolicy.LIFO, FakeConnection::close); + + final PoolEntry slow1 = + pool.lease("slow", null, Timeout.ofSeconds(1), null).get(1, TimeUnit.SECONDS); + final PoolEntry slow2 = + pool.lease("slow", null, Timeout.ofSeconds(1), null).get(1, TimeUnit.SECONDS); + + final Future> coldWaiter = + pool.lease("cold", null, Timeout.ofSeconds(1), null); + final Future> slowWaiter = + pool.lease("slow", null, Timeout.ofSeconds(1), null); + + slow1.assignConnection(new FakeConnection()); + slow1.updateExpiry(TimeValue.ofSeconds(30)); + pool.release(slow1, true); + + final PoolEntry cold = + coldWaiter.get(500, TimeUnit.MILLISECONDS); + assertEquals("cold", cold.getRoute()); + assertFalse(slowWaiter.isDone(), + "Same-route backlog must not consume every reusable release"); + + slowWaiter.cancel(true); + pool.release(slow2, false); + pool.release(cold, false); + pool.close(CloseMode.IMMEDIATE); + } + }