From cf2068420f92f25985c22a6cdb16c17d9ceb1efb Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Sat, 05 Sep 2026 18:16:09 +0000
Subject: [PATCH] [#878] Bound the JDBC connection pool and expire its connections one by one (#884)

---
 opendj-server-legacy/src/main/java/org/opends/server/backends/jdbc/CachedConnection.java |  682 +++++++++++++++++++++++++++++++++++++++++++++++++-------
 1 files changed, 598 insertions(+), 84 deletions(-)

diff --git a/opendj-server-legacy/src/main/java/org/opends/server/backends/jdbc/CachedConnection.java b/opendj-server-legacy/src/main/java/org/opends/server/backends/jdbc/CachedConnection.java
index 22d63d4..d59e1f4 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/backends/jdbc/CachedConnection.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/backends/jdbc/CachedConnection.java
@@ -15,14 +15,12 @@
  */
 package org.opends.server.backends.jdbc;
 
-import com.github.benmanes.caffeine.cache.Caffeine;
-import com.github.benmanes.caffeine.cache.LoadingCache;
-import com.github.benmanes.caffeine.cache.RemovalCause;
 import org.forgerock.i18n.LocalizableMessage;
 import org.forgerock.i18n.slf4j.LocalizedLogger;
+import org.opends.server.api.WorkQueue;
+import org.opends.server.core.DirectoryServer;
 
 import java.sql.*;
-import java.time.Duration;
 import java.util.ArrayDeque;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -36,6 +34,8 @@
 import java.util.Properties;
 import java.util.Set;
 import java.util.concurrent.*;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.AtomicLong;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
@@ -121,6 +121,32 @@
     /** 57P03, cannot_connect_now: postgresql starting up, shutting down or in recovery. */
     private static final String NOT_ACCEPTING_YET_SQL_STATE = "57P03";
 
+    /**
+     * The greatest number of connections one pool holds to one database; 0 for no bound. Read once
+     * per pool, when the first borrow of a connection string creates it, unlike the bounds of a
+     * borrow above: a pool is never removed from the map, and the permits of one already created
+     * are not resized, so this one takes a restart of the server to change.
+     */
+    static final String POOL_MAX_PROPERTY = "org.openidentityplatform.opendj.jdbc.pool.max";
+    /**
+     * Sized like the worker thread pool the server sizes for itself
+     * ({@code Platform.computeNumberOfThreads(16, 2)}), since an operation borrows one connection for
+     * its duration: the bound is there to keep a burst from opening as many connections as the
+     * database will accept, not to throttle steady traffic.
+     * <p>
+     * That formula is only what {@code WorkQueue.computeNumWorkerThreads} falls back to. A configured
+     * {@code ds-cfg-num-worker-threads} replaces it outright, and there is no default that can follow
+     * it: this bound belongs to a database that two backends may share, while that count belongs to
+     * the server. So an installation that raised it is told at open where the two stand, by
+     * {@link #reportBoundBelowBorrowers}, rather than left to find the wait in a latency graph.
+     */
+    static final int DEFAULT_POOL_MAX = Math.max(16, Runtime.getRuntime().availableProcessors() * 2);
+
+    /** How long a borrow waits for a connection to be returned before looking at the pool again. */
+    private static final long POOL_FULL_POLL_MS = 250;
+    /** The sweep runs at half the TTL, and no more often than this. */
+    private static final long MIN_SWEEP_INTERVAL_MS = 1000;
+
     static final long MAX_BACKOFF_MS = 1000;
     static final long STALL_WARNING_AFTER_MS = 1000;
     static final long STALL_WARNING_INTERVAL_MS = 10000;
@@ -171,36 +197,49 @@
      */
     private static final Map<String, Long> poolDistrustedAt = new ConcurrentHashMap<>();
 
+    // Throttled like the stall warning above, and keyed the same way: the bound is one setting, so
+    // one line per interval says so - but it is a setting of one pool, and two backends standing
+    // full at once each have their own to report. A single timestamp would let the pool that
+    // reported first silence the other, whose operations are failing with nothing in the log
+    // naming the database behind them.
+    private static final Map<String, AtomicLong> lastPoolFullWarning = new ConcurrentHashMap<>();
+
     final Connection parent;
 
-    // A deque handed out from the end it is returned to: the connection borrowed next is the one
-    // returned last, so under any load the pool keeps reusing its hottest connections instead of
-    // walking round every one it ever opened. That is what gives the window above anything to
-    // bypass - a connection reached only after a whole cycle of the pool has been idle far longer
-    // than the window - and it leaves the connections nothing needs at the cold end of the deque,
-    // where the per-connection idle expiry of #878 can find them. Until that lands, the cold end
-    // is reached only when the whole pool expires, after DEFAULT_TTL_MS with the backend idle.
-    // A deque takes one lock for both of its ends where the queue it replaces took one for each,
-    // so a borrow and a return no longer proceed side by side - against the round trip the window
-    // above saves, and the connect the reuse saves, that lock is not worth a FIFO handoff.
-    static LoadingCache<String, BlockingDeque<CachedConnection>> cached = Caffeine.newBuilder()
-        .expireAfterAccess(Duration.ofMillis(getCacheTtlMillis()))
-        .removalListener((String key, BlockingDeque<CachedConnection> value, RemovalCause cause) -> {
-            for (CachedConnection con : value) {
-                try {
-                    if (!con.isClosed()) {
-                        con.parent.close();
-                    }
-                } catch (SQLException e) {
-                    // ignore
-                }
-            }
-        })
-        .build(conStr -> new LinkedBlockingDeque<>());
+    /** The pool this connection belongs to, held directly so that the return needs no lookup. */
+    private final Pool pool;
+    /** Whether this connection holds a permit of its pool: a reentrant borrow does not. */
+    private final boolean metered;
+    /**
+     * The depth counter of the thread that borrowed it, lowered by the return. Held rather than the
+     * thread itself: a return made on another thread has to lower the depth of the borrower all the
+     * same, and a check of the returning thread against the borrowing one left that depth standing -
+     * the borrower was then taken for a nested borrow for the life of the server, exempt from the
+     * wait at the bound and opening an unmetered connection, destroyed on return, per operation
+     * (issue #878).
+     */
+    private volatile AtomicInteger depth;
+    /** When it was last returned to the pool, which is what the TTL is measured from. */
+    volatile long returnedAtMillis;
+    private final AtomicBoolean permitReleased = new AtomicBoolean();
+    /** Whether it has been handed back already: JDBC makes close() on a closed connection a no-op. */
+    private final AtomicBoolean returned = new AtomicBoolean();
+
+    /** The pool of every connection string in use, kept until the last storage using it closes. */
+    static final ConcurrentMap<String, Pool> pools = new ConcurrentHashMap<>();
+
+    /** The sweep that closes connections nothing has borrowed for the TTL, started with the first pool. */
+    private static volatile ScheduledExecutorService sweeper;
+
+    /** Where the sweep closes what it reaped, so that a close which does not return keeps it: see {@link Pool#sweep}. */
+    private static volatile Executor closer = DIRECT_EXECUTOR;
 
     /**
      * Returns the time after which an idle pooled connection is closed, as configured by the
      * {@value #TTL_PROPERTY} system property. An invalid value is ignored in favor of the default.
+     * <p>
+     * Read on every borrow and every sweep rather than once, so that it can be changed on a running
+     * server the way the bounds of a borrow can.
      */
     private static long getCacheTtlMillis() {
         return getNonNegativeProperty(TTL_PROPERTY, DEFAULT_TTL_MS, "ms");
@@ -215,8 +254,10 @@
      * value the unit conversion saturates on from leaving every connection of the pool trusted for
      * the life of the server.
      * <p>
-     * Read at class initialization, like the ttl it is clamped to, so a value set after that
-     * changes neither.
+     * Read at class initialization, so a value of this property set after that does not change the
+     * window. The ttl is not: {@link #getCacheTtlMillis()} is read on every borrow and every sweep,
+     * and the clamp above is not applied again - the window keeps the value it was computed with,
+     * so a ttl lowered on a running server does not lower the window with it.
      */
     static long getAliveBypassMillis() {
         long configured = getNonNegativeProperty(ALIVE_BYPASS_PROPERTY, DEFAULT_ALIVE_BYPASS_MS, "ms");
@@ -240,6 +281,386 @@
         return configured;
     }
 
+    /** The pool of a connection string, created on first use. */
+    static Pool poolOf(String connectionString) {
+        final Pool pool = pools.computeIfAbsent(connectionString, Pool::new);
+        startSweeper();
+        return pool;
+    }
+
+    private static void startSweeper() {
+        if (sweeper != null) {
+            return;
+        }
+        synchronized (pools) {
+            if (sweeper == null) {
+                // A thread per close in flight, and none while nothing is being closed. One thread
+                // shared by all of them would only move the head of the line, which is the point
+                // of not closing on the sweeper in the first place.
+                closer = Executors.newCachedThreadPool(runnable -> {
+                    final Thread thread = new Thread(runnable, "JDBC backend connection pool closer");
+                    thread.setDaemon(true);
+                    return thread;
+                });
+                final ScheduledExecutorService service = Executors.newSingleThreadScheduledExecutor(runnable -> {
+                    final Thread thread = new Thread(runnable, "JDBC backend connection pool sweeper");
+                    thread.setDaemon(true);
+                    return thread;
+                });
+                sweeper = service;
+                // Rescheduled after each run rather than left at a fixed delay: the interval comes
+                // from the ttl, and the ttl is read on every borrow and every sweep so that it can
+                // be changed on a running server. A delay computed once would keep the sweeper of a
+                // lowered ttl waking as rarely as the old one, so connections would go on being
+                // reaped no sooner than the setting the operator replaced (issue #878).
+                scheduleNextSweep(service);
+            }
+        }
+    }
+
+    /** Half the ttl, and no more often than {@value #MIN_SWEEP_INTERVAL_MS} ms. */
+    private static long sweepIntervalMillis() {
+        return Math.max(MIN_SWEEP_INTERVAL_MS, getCacheTtlMillis() / 2);
+    }
+
+    /**
+     * Books the next sweep, and the one after it out of its own run. Every run books its successor
+     * in a finally: a sweep that ends in a Throwable the per-pool guard did not catch would
+     * otherwise stop the expiry of every pool in the JVM, the way a task thrown out of
+     * scheduleWithFixedDelay does.
+     */
+    private static void scheduleNextSweep(ScheduledExecutorService service) {
+        try {
+            service.schedule(() -> {
+                try {
+                    sweep();
+                } finally {
+                    scheduleNextSweep(service);
+                }
+            }, sweepIntervalMillis(), TimeUnit.MILLISECONDS);
+        } catch (RejectedExecutionException e) {
+            // the sweeper is shutting down: there is nothing left to book a run on
+            logger.traceException(e);
+        }
+    }
+
+    // Expiry has to happen without a borrow behind it. Caffeine was left without a scheduler, so an
+    // entry was only ever expired by a later cache operation - and a backend that has gone idle,
+    // the one case the TTL exists for, performs none (issue #878).
+    static void sweep() {
+        final long ttlMillis = getCacheTtlMillis();
+        final Executor closeOn = closer;
+        for (final Pool pool : pools.values()) {
+            try {
+                pool.sweep(ttlMillis, closeOn);
+            } catch (Throwable t) {
+                // Error included: scheduleWithFixedDelay cancels a task that throws, so anything
+                // escaping here would stop the expiry of every pool in the JVM for good - and
+                // silently, which is the failure mode the hand-off of the close exists to avoid.
+                logger.traceException(t);
+            }
+        }
+    }
+
+    /**
+     * Registers a storage as a user of the pool of a connection string. Reference counted because a
+     * pool belongs to a database rather than to a backend: two backends may address one database,
+     * and closing one of them must not take the connections of the other with it.
+     */
+    static void openPool(String connectionString) {
+        final Pool pool = poolOf(connectionString);
+        pool.addUser();
+        reportBoundBelowBorrowers(connectionString, pool);
+    }
+
+    /**
+     * Reports a bound smaller than the number of worker threads. An operation borrows one connection
+     * for its duration, so the worker threads are the borrowers this default is sized against - and
+     * it is sized against the count the server computes for itself, not against a
+     * {@code ds-cfg-num-worker-threads} the operator set, which replaces that count outright.
+     * <p>
+     * A lower bound than that is what is reported, not every way past it: the replay threads of
+     * replication default to the same count again and borrow on top of the workers, and an import or
+     * a rebuild borrows besides. So this names one difference the operator can act on rather than
+     * standing for the whole demand on the pool.
+     * <p>
+     * Nothing fails for the difference alone: the surplus waits for a connection to be returned,
+     * which is what the bound is there for. But every one of those waits is paid on an operation,
+     * and past {@value #POOL_TIMEOUT_PROPERTY} the operation fails - on a setting whose effect on
+     * this backend the operator had no reason to expect (issue #878).
+     */
+    private static void reportBoundBelowBorrowers(String connectionString, Pool pool) {
+        final WorkQueue<?> workQueue = DirectoryServer.getWorkQueue();
+        if (workQueue == null) {
+            // an offline tool, or the server before its work queue is up: no borrowers to count
+            return;
+        }
+        final int borrowers = workQueue.getNumWorkerThreads();
+        if (borrowers <= pool.max()) {
+            return;
+        }
+        final long poolTimeoutSeconds = getNonNegativeProperty(POOL_TIMEOUT_PROPERTY, DEFAULT_POOL_TIMEOUT_SECONDS, "s");
+        final String wait = poolTimeoutSeconds == 0
+            ? "waits for one to be returned for as long as that takes"
+            : "waits up to " + poolTimeoutSeconds + "s for one to be returned and fails if none is";
+        warnOnce(safeUrl(connectionString) + "|bound-below-borrowers",
+            "the connection pool of %s holds at most %d connections while %d worker threads may each borrow one:"
+                + " an operation finding it at its bound %s (raise %s to allow more connections, or lower"
+                + " ds-cfg-num-worker-threads)",
+            safeUrl(connectionString), pool.max(), borrowers, wait, POOL_MAX_PROPERTY);
+    }
+
+    /** Unregisters a storage; the connections are released once the last user is gone. */
+    static void closePool(String connectionString) {
+        final Pool pool = pools.get(connectionString);
+        if (pool != null) {
+            pool.removeUser();
+        }
+    }
+
+    /**
+     * The connections of one connection string.
+     * <p>
+     * This replaces the cache entry that used to hold them. That one carried the TTL on the pool
+     * rather than on a connection - {@code expireAfterAccess} keyed by the connection string, reset
+     * by every borrow and every return - so under continuous traffic nothing ever expired and the
+     * peak count of a burst stayed open for as long as the backend saw any traffic at all. It also
+     * had no bound, so the only ceiling on the connections of a backend was the {@code
+     * max_connections} of the database itself (issue #878).
+     */
+    static final class Pool {
+        final String connectionString;
+        /** Idle connections, most recently returned first: the ones a burst opened sink to the bottom, where the sweep finds them. */
+        private final LinkedBlockingDeque<CachedConnection> idle = new LinkedBlockingDeque<>();
+        /** One permit per live connection, borrowed or idle. Sized once: this is how large the pool may grow, not a rate. */
+        private final Semaphore permits;
+        private final int max;
+        /**
+         * How many connections of this pool the current thread holds. A borrow made while one is
+         * already held may exceed the bound, because the two are held at the same time and waiting
+         * for the first to be returned would wait for this very thread:
+         * {@code PersistentCompressedSchema.store()} opens a write of its own - the definition has
+         * to commit independently of the entry - and {@code EntryContainer.modifyDN} reaches it
+         * from inside a transaction, having encoded the entry there. The exemption is from the
+         * wait rather than from the pool: a nested borrow served out of the idle deque carries the
+         * permit that connection already holds and is pooled again on return like any other. Only
+         * one that had to establish a connection of its own, because the pool stood at its bound,
+         * holds no permit - and that one is closed rather than pooled when it comes back, so the
+         * pool does not grow past its bound.
+         * <p>
+         * Counted per pool rather than per thread, because that deadlock only exists within one
+         * pool: a count shared by all of them would judge a thread holding a connection to one
+         * database reentrant while it borrows from another, passing the bound of a pool it holds
+         * nothing of and destroying the connection instead of pooling it, on every operation.
+         */
+        private final ThreadLocal<AtomicInteger> held = ThreadLocal.withInitial(AtomicInteger::new);
+        /** Open storages using this pool, guarded by this. */
+        private int users;
+        /**
+         * Set when the last storage using this pool closed. A pool no storage ever registered with -
+         * a borrow made straight through {@link CachedConnection#getConnection}, as the tests do -
+         * is not closed and pools normally; only one that had a user and lost it stops keeping
+         * connections for a borrower that is not going to come.
+         */
+        private volatile boolean closed;
+
+        Pool(String connectionString) {
+            this.connectionString = connectionString;
+            final long configured = getNonNegativeProperty(POOL_MAX_PROPERTY, DEFAULT_POOL_MAX, "connections");
+            this.max = (configured == 0 || configured > Integer.MAX_VALUE) ? Integer.MAX_VALUE : (int) configured;
+            this.permits = new Semaphore(max);
+        }
+
+        int max() {
+            return max;
+        }
+
+        /** Whether the calling thread already holds a connection of this pool. */
+        boolean heldByCurrentThread() {
+            return held.get().get() > 0;
+        }
+
+        /**
+         * Raises the depth of the borrowing thread and hands back the counter it was raised on, for
+         * the connection to lower on its return. The counter rather than the thread, because the
+         * return need not happen on the thread that borrowed - and the depth that has to come down
+         * is the borrower's either way. Read by that thread alone but written by whichever returns
+         * the connection, which is why it is an AtomicInteger and not an int.
+         */
+        AtomicInteger enter() {
+            final AtomicInteger depth = held.get();
+            depth.incrementAndGet();
+            return depth;
+        }
+
+        /** Lowers a depth this pool handed out, never below zero. */
+        static void leave(AtomicInteger depth) {
+            depth.updateAndGet(held -> held > 0 ? held - 1 : 0);
+        }
+
+        int idleCount() {
+            return idle.size();
+        }
+
+        /**
+         * The connections of this pool holding a permit, borrowed and idle together. Not every
+         * connection of the pool: a borrow nested in one this thread already holds goes on
+         * unmetered when the pool stands at its bound, so the connections this count misses are
+         * exactly the ones over the bound. They take no place in it and are closed rather than
+         * pooled when they come back, which makes this the count the bound is about - how much of
+         * it is taken - rather than the number of sockets open to the database.
+         */
+        int meteredCount() {
+            return max - permits.availablePermits();
+        }
+
+        synchronized void addUser() {
+            users++;
+            closed = false;
+        }
+
+        void removeUser() {
+            final boolean wasLast;
+            synchronized (this) {
+                wasLast = users > 0 && --users == 0;
+                if (wasLast) {
+                    closed = true;
+                }
+            }
+            if (wasLast) {
+                // Outside the monitor: closing a connection is a round trip, and an open of the
+                // same database has no reason to wait behind it. The borrowed ones are not here to
+                // be closed - give() closes them when they come back, since a pool nobody uses must
+                // not keep them for a borrower that is not going to come.
+                logger.trace(LocalizableMessage.raw("releasing %d pooled connections of %s: its last user closed",
+                    idle.size(), safeUrl(connectionString)));
+                drainIdle();
+            }
+        }
+
+        void drainIdle() {
+            for (CachedConnection con = idle.pollFirst(); con != null; con = idle.pollFirst()) {
+                destroy(con);
+            }
+        }
+
+        /**
+         * Takes a connection out of the pool, waiting up to waitMs for one to be returned, and
+         * discarding the ones that are broken or have been idle for longer than the TTL.
+         * <p>
+         * Bounded by the deadline of the borrow, and not only by waitMs: a poll of no duration
+         * still hands out whatever the deque holds, and discarding a connection whose socket is
+         * half-open costs the validation timeout apiece. The pool holds as many of those as its
+         * bound allows, so draining the deque overran the bound the operator set - by minutes on a
+         * large pool, before the connect that follows it had even started (issue #878).
+         */
+        CachedConnection pollIdle(long waitMs, long ttlMillis, long deadline, boolean trusted)
+                throws InterruptedException {
+            long remainingWait = waitMs;
+            while (true) {
+                final long polledAt = System.currentTimeMillis();
+                final CachedConnection con = idle.pollFirst(remainingWait, TimeUnit.MILLISECONDS);
+                if (con == null) {
+                    return null;
+                }
+                if (System.currentTimeMillis() - con.returnedAtMillis <= ttlMillis && isUsable(con, trusted)) {
+                    return con;
+                }
+                destroy(con);
+                final long remaining = deadline - System.currentTimeMillis();
+                if (remaining <= 0) {
+                    return null;
+                }
+                // one more look, since a connection may have been returned in the meantime
+                remainingWait = Math.min(Math.max(0, remainingWait - (System.currentTimeMillis() - polledAt)), remaining);
+            }
+        }
+
+        /** Takes the right to hold one more connection, or reports that the pool is full. */
+        boolean tryReserve() {
+            return permits.tryAcquire();
+        }
+
+        void cancelReservation() {
+            permits.release();
+        }
+
+        /** Hands a connection back, closing it rather than pooling it when it may not be kept. */
+        void give(CachedConnection con) {
+            // An unmetered connection holds no permit, so pooling it would put the pool one over its
+            // bound for good; and a closed pool has nobody left to hand it to.
+            if (con.metered && !closed) {
+                addIdle(con);
+                if (closed) {
+                    // The last user left while this one was on its way back, so it missed the drain.
+                    drainIdle();
+                }
+            } else {
+                destroy(con);
+            }
+        }
+
+        /** Puts a connection into the pool. The caller must hold the right to keep it there. */
+        void addIdle(CachedConnection con) {
+            con.returnedAtMillis = System.currentTimeMillis();
+            idle.addFirst(con);
+        }
+
+        void destroy(CachedConnection con) {
+            try {
+                closeQuietly(con.parent);
+            } finally {
+                // However the close went, the pool holds one connection fewer. A permit not given
+                // back here is given back by nothing at all: only a live connection carries one,
+                // and this one is gone (issue #878).
+                con.releasePermit();
+            }
+        }
+
+        void sweep(long ttlMillis) {
+            sweep(ttlMillis, DIRECT_EXECUTOR);
+        }
+
+        /**
+         * Closes the connections nothing has borrowed for the TTL, handing each to the executor
+         * given rather than closing it here. The sweep of every pool shares one thread and
+         * {@code scheduleWithFixedDelay} never overlaps its runs, so one close that does not
+         * return would stop the expiry of every pool in the JVM - and silently, since only a
+         * thrown exception is logged. Oracle logs off over the network, and the read bound of the
+         * login has been lifted by then (issue #878).
+         */
+        void sweep(long ttlMillis, Executor closeOn) {
+            final long deadline = System.currentTimeMillis() - ttlMillis;
+            // From the tail: the least recently returned connection is the first to have expired,
+            // and once one has not, neither has anything in front of it.
+            for (CachedConnection con = idle.peekLast(); con != null; con = idle.peekLast()) {
+                if (con.returnedAtMillis > deadline) {
+                    return;
+                }
+                if (!idle.removeLastOccurrence(con)) {
+                    // A borrow took it between the two. What is behind it may still have expired,
+                    // and ending the cycle here would leave every one of those open until the
+                    // next sweep.
+                    continue;
+                }
+                if (con.returnedAtMillis > deadline) {
+                    // A borrow took it between the peek and the removal and gave it back, so the
+                    // reading the decision was made on is not the one it carries now: closing it
+                    // would cost the next borrow a connect over a connection a moment old. Back to
+                    // the end it is returned to, where its refreshed reading belongs.
+                    idle.addFirst(con);
+                    return;
+                }
+                final CachedConnection expired = con;
+                try {
+                    closeOn.execute(() -> destroy(expired));
+                } catch (RuntimeException e) { // no thread to close it on: here rather than nowhere
+                    destroy(expired);
+                }
+            }
+        }
+    }
+
     /**
      * Returns the value of a numeric system property, ignoring a value that is not a non-negative
      * number in favor of the default. The unit is the one the property is read in, so that the
@@ -596,17 +1017,42 @@
      */
     private volatile long lastKnownAliveNanos;
 
+    /**
+     * A connection outside the accounting of its pool: it holds no permit and is never pooled - the
+     * flag says so as well as the accounting does, since a connection holding no permit is closed
+     * by {@link Pool#give} rather than kept whatever the flag says.
+     * <p>
+     * It still names a pool, because that is what closes it and what the sweep runs over, so the
+     * pool of this connection string is created here if it does not exist yet and the sweeper is
+     * started with it.
+     */
     public CachedConnection(String connectionString, Connection parent) {
-        this(connectionString, parent, true);
+        this(connectionString, parent, poolOf(connectionString), false, false);
     }
 
-    CachedConnection(String connectionString, Connection parent, boolean poolable) {
+    CachedConnection(String connectionString, Connection parent, Pool pool, boolean metered, boolean poolable) {
         this.connectionString = connectionString;
         this.parent = parent;
+        this.pool = pool;
+        this.metered = metered;
         this.poolable = poolable;
         this.lastKnownAliveNanos = System.nanoTime();
     }
 
+    /** Gives back the right to hold this connection, once and only if it was taken. */
+    void releasePermit() {
+        if (metered && permitReleased.compareAndSet(false, true)) {
+            pool.cancelReservation();
+        }
+    }
+
+    /** Records that the borrowing thread holds this connection, so a borrow nested in it is recognized. */
+    private static CachedConnection borrowed(CachedConnection con) {
+        con.returned.set(false);
+        con.depth = con.pool.enter();
+        return con;
+    }
+
     /**
      * Borrows a connection: a usable one out of the pool, or a newly established one. Bounded in
      * both phases - every operation of this backend, the open of a backend and the import
@@ -630,26 +1076,65 @@
      * them is one borrow of a cold path, where the round trip the window saves is worth nothing.
      */
     static Connection getConnection(String connectionString, boolean trusted) throws Exception {
+        final Pool pool = poolOf(connectionString);
         final ConnectDialect dialect = ConnectDialect.of(connectionString);
         reportUnknownDialect(connectionString, dialect);
         final long connectTimeoutSeconds = Math.min(
             getNonNegativeProperty(CONNECT_TIMEOUT_PROPERTY, DEFAULT_CONNECT_TIMEOUT_SECONDS, "s"),
             Integer.MAX_VALUE / 1000);
         final long poolTimeoutSeconds = getNonNegativeProperty(POOL_TIMEOUT_PROPERTY, DEFAULT_POOL_TIMEOUT_SECONDS, "s");
+        final long ttlMillis = getCacheTtlMillis();
         final long startedAt = System.currentTimeMillis();
         final long deadline = (poolTimeoutSeconds == 0 || poolTimeoutSeconds >= Long.MAX_VALUE / 1000)
             ? Long.MAX_VALUE : startedAt + poolTimeoutSeconds * 1000;
+        // A thread already holding a connection is not made to wait for one: the two are held at
+        // the same time, so waiting for the first to come back would wait for itself.
+        final boolean reentrant = pool.heldByCurrentThread();
         long waitMs = 0;
         long backoffMs = 0;
         int attempts = 0;
         while (true) {
-            final CachedConnection pooled = poll(connectionString, waitMs, deadline, trusted);
+            final CachedConnection pooled = pool.pollIdle(waitMs, ttlMillis, deadline, trusted);
             if (pooled != null) {
-                return pooled;
+                return borrowed(pooled);
+            }
+            // Asked for whether this borrow is nested or not: the exemption a nested one carries is
+            // from the wait, not from the pool. A nested borrow made while the pool has room takes a
+            // permit like any other and is pooled again on return; only one that finds the pool at its
+            // bound goes on unmetered, and that one is closed rather than pooled when it comes back.
+            final boolean metered = pool.tryReserve();
+            if (!metered && !reentrant) {
+                // The pool holds as many connections as it may: only a returned one can serve this
+                // borrow now, and the deadline decides how long that is worth waiting for. This is
+                // the point of the bound - without it the borrow would open one more connection,
+                // and the only ceiling left would be the max_connections of the database itself.
+                final long remaining = deadline - System.currentTimeMillis();
+                if (remaining <= 0) {
+                    // The restart is part of the remedy, so the message says so: the bound is read
+                    // once, when the pool is created, and a pool is never removed from the map - so
+                    // the property set on a running server changes nothing until it is read again.
+                    final String message = "no connection to " + safeUrl(connectionString)
+                        + " could be borrowed within " + poolTimeoutSeconds + "s: all " + pool.max()
+                        + " connections of the pool are in use (raise " + POOL_MAX_PROPERTY
+                        + " and restart the server to allow more)";
+                    // The one failure the bound introduces has to reach the server log too: an
+                    // installation whose peak sits above the default would otherwise see its
+                    // operations fail with nothing in the log naming the pool behind it.
+                    warnPoolFull(connectionString, message);
+                    throw new SQLTimeoutException(message);
+                }
+                waitMs = Math.min(POOL_FULL_POLL_MS, remaining);
+                continue;
             }
             attempts++;
+            CachedConnection established = null;
+            boolean handedOff = false;
             try {
-                return connect(connectionString, dialect, attemptSeconds(connectTimeoutSeconds, deadline));
+                established = connect(connectionString, dialect,
+                    attemptSeconds(connectTimeoutSeconds, deadline), pool, metered);
+                final CachedConnection con = borrowed(established);
+                handedOff = true;
+                return con;
             } catch (SQLException e) {
                 // A database that takes no connection for the moment is the failure worth waiting
                 // out: it is at its connection limit, and one of ours is going to come back to the
@@ -680,6 +1165,20 @@
                 // a driver reporting a connect it will not make as an unchecked failure carries the
                 // connection string of the backend in its message as readily as a SQLException does
                 throw reportedUnchecked(e, connectionString);
+            } finally {
+                // What the attempt took is given back on every way out of it, not only on the
+                // SQLException a driver is supposed to throw. DriverManager catches SQLException
+                // alone, so an unchecked failure of a driver reaches here - Connector/J hands a url
+                // with a "%" in it to URLDecoder, and this backend keeps its credentials in the url
+                // - and a permit left behind is left behind for good: only a live connection
+                // carries one, and a failed attempt has none to give (issue #878).
+                if (!handedOff) {
+                    if (established != null) {
+                        pool.destroy(established); // the permit went with it, and comes back with it
+                    } else if (metered) {
+                        pool.cancelReservation();
+                    }
+                }
             }
         }
     }
@@ -728,33 +1227,6 @@
         return Math.max(1, Math.min(bound, Integer.MAX_VALUE / 1000));
     }
 
-    /**
-     * Takes a usable connection out of the pool, waiting up to waitMs for one to be returned to it.
-     * The validation of a connection costs a round trip, and the pool has no upper bound on the
-     * number of them it holds, so draining a pool the database no longer answers is given the
-     * deadline of the borrow as well: past it, establishing a connection is the faster answer.
-     * The connection in hand is always looked at first - trusted or validated, see
-     * {@link #isKnownAlive} - whatever the deadline says: a database at its connection limit has
-     * no other source of connections than the ones coming back, and one returned to the pool a
-     * moment before the deadline is the very connection this borrow waited for. Only a connection
-     * the database no longer answers is closed here.
-     */
-    private static CachedConnection poll(String connectionString, long waitMs, long deadline, boolean trusted)
-            throws InterruptedException {
-        CachedConnection con = cached.get(connectionString).pollFirst(waitMs, TimeUnit.MILLISECONDS);
-        while (con != null) {
-            if (isUsable(con, trusted)) {
-                return con;
-            }
-            closeQuietly(con.parent);
-            if (System.currentTimeMillis() >= deadline) {
-                return null;
-            }
-            con = cached.get(connectionString).pollFirst();
-        }
-        return null;
-    }
-
     private static boolean isUsable(CachedConnection con, boolean trusted) {
         if (trusted && isKnownAlive(con)) {
             return true;
@@ -833,10 +1305,14 @@
         if (distrusted != null && provenAt - distrusted <= 0) { // the overflow safe form of the comparison
             return false;
         }
-        // What the validation this replaces also answered: the removalListener above closes every
-        // connection it finds in the deque when the pool expires, and it iterates a weakly
-        // consistent view, so a connection taken out by a borrow running at the same time can be
-        // closed under it. Answered by the driver out of a flag of its own, not by a round trip.
+        // What the validation this replaces also answered, asked of the driver out of a flag of its
+        // own rather than by a round trip: a connection the driver has already given up on - the
+        // database dropped it and the driver noticed - is not one to hand out on the strength of a
+        // window. It no longer stands for a drain closing a connection under its borrower, the way
+        // it did while the pool was a cache entry whose removalListener iterated a weakly consistent
+        // view of the deque: every path that destroys an idle connection now takes it out of the
+        // deque first (pollIdle, drainIdle, and the removeLastOccurrence of the sweep), so what a
+        // borrow holds is not there to be found (issue #878).
         return !isClosed(con.parent);
     }
 
@@ -900,8 +1376,8 @@
         return previous < 0 ? 0 : previous;
     }
 
-    static CachedConnection connect(String connectionString, ConnectDialect dialect, long connectTimeoutSeconds)
-            throws SQLException {
+    static CachedConnection connect(String connectionString, ConnectDialect dialect, long connectTimeoutSeconds,
+            Pool pool, boolean metered) throws SQLException {
         // A driver is free to write into the map it is handed, so it gets one of its own.
         final Properties properties = new Properties();
         final boolean readBoundSet = dialect != null && connectTimeoutSeconds > 0
@@ -926,7 +1402,7 @@
             closeQuietly(conNew);
             throw e;
         }
-        final CachedConnection established = new CachedConnection(connectionString, conNew, poolable);
+        final CachedConnection established = new CachedConnection(connectionString, conNew, pool, metered, poolable);
         established.lastKnownAliveNanos = provenAt;
         return established;
     }
@@ -1000,6 +1476,19 @@
         return false;
     }
 
+    // The bound of the pool is a reason for an operation to fail that no version before it had,
+    // so it belongs in the server log as well as in the error the client is given. Throttled like
+    // the stall warning: every worker thread reaches it at once when the pool stands full.
+    private static void warnPoolFull(String connectionString, String message) {
+        final long now = System.currentTimeMillis();
+        final AtomicLong lastOfThisUrl =
+            lastPoolFullWarning.computeIfAbsent(safeUrl(connectionString), url -> new AtomicLong());
+        final long last = lastOfThisUrl.get();
+        if (now - last >= STALL_WARNING_INTERVAL_MS && lastOfThisUrl.compareAndSet(last, now)) {
+            logger.warn(LocalizableMessage.raw("%s", message));
+        }
+    }
+
     // A stall has to reach the server log: without it a database accepting no further connection
     // is indistinguishable from a hang. Throttled, since every operation of the backend borrows
     // through here and would otherwise log a copy of its own.
@@ -1386,8 +1875,8 @@
     private static void closeQuietly(Connection con) {
         try {
             con.close();
-        } catch (SQLException e) {
-            // ignore: it is on its way out anyway
+        } catch (SQLException | RuntimeException e) {
+            // ignore: it is on its way out anyway, and the caller has a permit to give back
         }
     }
 
@@ -1433,21 +1922,46 @@
 
     @Override
     public void close() throws SQLException {
-        try {
-            rollback();
-        } catch (SQLException e) {
-            // A connection that cannot be rolled back must not be handed to the next borrower -
-            // and must not be dropped on the floor either: nothing else holds it any more.
-            closeQuietly(parent);
-            throw e;
-        }
-        if (!poolable) {
-            closeQuietly(parent);
+        // JDBC makes close() on a closed connection a no-op, and this one has to be one: a second
+        // return would put the same connection into the pool twice, to be handed to two borrowers.
+        if (!returned.compareAndSet(false, true)) {
             return;
         }
-        // Returned to the end the next borrow takes it from, so that the pool keeps reusing its
-        // hottest connections rather than cycling through every one it ever opened.
-        cached.get(connectionString).addFirst(this);
+        final AtomicInteger borrowerDepth = depth;
+        depth = null;
+        if (borrowerDepth != null) {
+            Pool.leave(borrowerDepth);
+        }
+        // Set before the hand-off rather than after it: from the moment give() is called the pool
+        // owns this connection, and a second destroy() of one that reached the idle deque would
+        // close a connection still waiting there to be handed out.
+        boolean handedToPool = false;
+        try {
+            rollback();
+            if (poolable) {
+                // Straight to the pool it came from rather than through a lookup of its connection
+                // string: the entry the lookup returned could be evicted between the two, leaving
+                // the connection in a queue nothing referred to any more - never handed out, never
+                // closed (issue #878).
+                handedToPool = true;
+                pool.give(this);
+            }
+        } finally {
+            // Every way out that is not a give(): the SQLException a rollback is supposed to throw,
+            // a connection that may not be pooled, and the unchecked failure a driver throws
+            // instead of a SQLException. The CAS above has already made this the one close() of
+            // this connection, so what leaves here through neither give() nor destroy() is closed
+            // by nothing at all - and its permit is released by nothing either, since destroy() is
+            // the only caller of releasePermit(). A pool is never removed from the static map, so
+            // that place in the bound would be gone for the life of the server, and enough of them
+            // leave every borrow to fail with a SQLTimeoutException (issue #878).
+            if (!handedToPool) {
+                // destroy() rather than a bare close: the permit this connection holds has to go
+                // back to the pool with it, or the bound loses a place for every connection kept
+                // out of it.
+                pool.destroy(this);
+            }
+        }
     }
 
     @Override

--
Gitblit v1.10.0