From bf93e5e23570f41edd8d4fe028226f600e4e8f71 Mon Sep 17 00:00:00 2001
From: Arturo Bernal
Date: Sat, 3 Oct 2026 17:29:17 +0200
Subject: [PATCH] Reclaim idle entries across routes at max capacity
---
.../pool/RouteSegmentedConnPoolBenchmark.java | 349 ++++++++++++++++++
.../hc/core5/pool/RouteSegmentedConnPool.java | 267 ++++++++++++--
.../pool/RouteSegmentedConnPoolTest.java | 137 +++++--
3 files changed, 703 insertions(+), 50 deletions(-)
create mode 100644 httpcore5-testing/src/test/java/org/apache/hc/core5/pool/RouteSegmentedConnPoolBenchmark.java
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);
+ }
+
}