| | |
| | | public class CachedConnection implements Connection { |
| | | private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass(); |
| | | |
| | | // What has been reported once already. Every one of these reports a setting rather than an |
| | | // event - a property that is not a number, a url no bound of this class can reach, a driver |
| | | // whose property names are not known here - so it does not become truer by being repeated, |
| | | // and every operation of the backend comes through here. |
| | | // Declared above every field whose initializer can reach warnOnce(): class variable |
| | | // initializers run in textual order (JLS 12.4.2), so a set declared below aliveBypassNanos |
| | | // would still be null the moment a property this class reports on carries a value worth |
| | | // warning about - a window longer than the ttl, or one that is not a number - and the report |
| | | // would leave the class uninitializable rather than merely configured oddly. |
| | | static final Set<String> warnedOnce = ConcurrentHashMap.newKeySet(); |
| | | |
| | | static final String TTL_PROPERTY = "org.openidentityplatform.opendj.jdbc.ttl"; |
| | | static final long DEFAULT_TTL_MS = 15000; |
| | | |
| | | /** |
| | | * How long a pooled connection is handed out without being validated after it was last proven |
| | | * alive, in ms; 0 validates every borrow, the way this pool did before the window existed. |
| | | * <p> |
| | | * The validation of a connection is a round trip of its own - an empty query on postgresql, a |
| | | * ping on mysql, a round trip of its own on oracle and sql server - and every operation of |
| | | * this backend pays it next to the single statement the operation came for. It earns that on a |
| | | * connection that has been sitting in the pool, which the database or a firewall may have |
| | | * dropped in the meantime; it earns nothing on one that answered a moment ago, which is most |
| | | * of them under load. So a connection proven alive within this window is trusted rather than |
| | | * validated, the way the aliveBypassWindow of HikariCP does it. |
| | | */ |
| | | static final String ALIVE_BYPASS_PROPERTY = "org.openidentityplatform.opendj.jdbc.alive.bypass"; |
| | | static final long DEFAULT_ALIVE_BYPASS_MS = 500; |
| | | |
| | | /** |
| | | * The longest window this class uses, whatever {@value #ALIVE_BYPASS_PROPERTY} and the |
| | | * {@value #TTL_PROPERTY} it is clamped to say. The clamp to the ttl alone does not bound it: |
| | | * the ttl has no upper bound of its own, and with both set high enough the conversion to |
| | | * nanoseconds saturates - the window then outlasts every reading it is compared against, and |
| | | * no connection of the pool is ever validated again. An hour is already far past what this |
| | | * window is about, which is a connection that answered a moment ago. |
| | | * <p> |
| | | * A compile-time constant, so that it holds its value wherever it is read from: the initializer |
| | | * of {@link #aliveBypassNanos} reaches it, and a field initialized in declaration order would |
| | | * still be 0 there if it were ever moved below (JLS 12.4.2). |
| | | */ |
| | | static final long MAX_ALIVE_BYPASS_MS = 60 * 60 * 1000L; |
| | | |
| | | // Read once, at class initialization: every operation of this backend borrows a connection, |
| | | // and the borrow is not the place to parse a system property. Not final so that a test can |
| | | // vary the window without a class loader of its own, and volatile because a non-final static |
| | | // long is written neither atomically nor visibly to the threads reading it (JLS 17.7) - every |
| | | // worker of the backend and every replay thread reads this one. |
| | | static volatile long aliveBypassNanos = TimeUnit.MILLISECONDS.toNanos(getAliveBypassMillis()); |
| | | |
| | | /** |
| | | * Bounds the connect and the login of one attempt to establish a connection, in seconds; 0 for |
| | | * no bound of its own - the deadline of {@value #POOL_TIMEOUT_PROPERTY} still bounds the |
| | | * attempt, since it stands for the whole borrow. Setting both to 0 is what leaves a connect |
| | |
| | | private static final Map<String, AtomicLong> lastStallWarning = new ConcurrentHashMap<>(); |
| | | private static final AtomicLong lastReadBoundWarning = new AtomicLong(); |
| | | |
| | | // What has been reported once already. Every one of these reports a setting rather than an |
| | | // event - a property that is not a number, a url no bound of this class can reach, a driver |
| | | // whose property names are not known here - so it does not become truer by being repeated, |
| | | // and every operation of the backend comes through here. |
| | | static final Set<String> warnedOnce = ConcurrentHashMap.newKeySet(); |
| | | /** |
| | | * When an operation last reported that the database had dropped a connection of a pool, as a |
| | | * {@link System#nanoTime()} reading per connection string. A connection proven alive before |
| | | * that moment is validated on its next borrow whatever the window says: whatever dropped one |
| | | * connection - a restart, a failover, a network that went away - dropped every connection |
| | | * established before it, and the window would otherwise hand out the rest of that generation |
| | | * one by one until the pool runs out of them. It is set by the caller that saw the failure |
| | | * ({@code JDBCStorage}), never by a validation that failed here: an idle connection the server |
| | | * reaped is a routine event, and it says nothing about the connection in use that the pool is |
| | | * about to hand out. |
| | | */ |
| | | private static final Map<String, Long> poolDistrustedAt = new ConcurrentHashMap<>(); |
| | | |
| | | final Connection parent; |
| | | |
| | | static LoadingCache<String, BlockingQueue<CachedConnection>> cached = Caffeine.newBuilder() |
| | | // 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, BlockingQueue<CachedConnection> value, RemovalCause cause) -> { |
| | | .removalListener((String key, BlockingDeque<CachedConnection> value, RemovalCause cause) -> { |
| | | for (CachedConnection con : value) { |
| | | try { |
| | | if (!con.isClosed()) { |
| | |
| | | } |
| | | } |
| | | }) |
| | | .build(conStr -> new LinkedBlockingQueue<>()); |
| | | .build(conStr -> new LinkedBlockingDeque<>()); |
| | | |
| | | /** |
| | | * Returns the time after which an idle pooled connection is closed, as configured by the |
| | |
| | | } |
| | | |
| | | /** |
| | | * Returns the alive window, clamped to the {@value #TTL_PROPERTY} an idle pooled connection is |
| | | * kept for and to {@link #MAX_ALIVE_BYPASS_MS} behind it. A window longer than the ttl is one |
| | | * the pool cannot back: it goes on trusting the last answer of a connection past the point the |
| | | * pool would have closed and replaced it, which is a claim about a connection that is no longer |
| | | * there. The ttl has no upper bound of its own, though, so the second clamp is what keeps a |
| | | * 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. |
| | | */ |
| | | static long getAliveBypassMillis() { |
| | | long configured = getNonNegativeProperty(ALIVE_BYPASS_PROPERTY, DEFAULT_ALIVE_BYPASS_MS, "ms"); |
| | | final long ttl = getCacheTtlMillis(); |
| | | if (configured > ttl) { |
| | | warnOnce(ALIVE_BYPASS_PROPERTY + "=" + configured + ">" + ttl, |
| | | "The %s window of %d ms is longer than the %d ms of %s a pooled connection is kept for," |
| | | + " and is used as %d ms: a connection trusted for longer than the pool keeps it would" |
| | | + " be trusted past the point the pool closed it", |
| | | ALIVE_BYPASS_PROPERTY, configured, ttl, TTL_PROPERTY, ttl); |
| | | configured = ttl; |
| | | } |
| | | if (configured > MAX_ALIVE_BYPASS_MS) { // the ttl it was just clamped to has no upper bound of its own |
| | | warnOnce(ALIVE_BYPASS_PROPERTY + ">" + MAX_ALIVE_BYPASS_MS, |
| | | "The %s window of %d ms is longer than the %d ms this pool trusts a connection for at most," |
| | | + " and is used as %d ms: a longer one saturates the arithmetic it is compared in and" |
| | | + " leaves every connection of the pool trusted for the life of the server", |
| | | ALIVE_BYPASS_PROPERTY, configured, MAX_ALIVE_BYPASS_MS, MAX_ALIVE_BYPASS_MS); |
| | | return MAX_ALIVE_BYPASS_MS; |
| | | } |
| | | return configured; |
| | | } |
| | | |
| | | /** |
| | | * 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 |
| | | * value the message names is not mistaken for another. |
| | |
| | | */ |
| | | private final boolean poolable; |
| | | |
| | | /** |
| | | * When this connection last answered the database, as a {@link System#nanoTime()} reading: |
| | | * established - the login and the two round trips that set it up have just answered - or |
| | | * validated. It is never stamped on the way back into the pool, although that is where a |
| | | * connection has most recently been used: {@link #close()} ends the transaction, and pgjdbc |
| | | * short-circuits both {@code rollback()} and {@code commit()} when the transaction state is |
| | | * IDLE, so on a borrow that issued no statement - {@code JDBCStorage.open()}, a configuration |
| | | * change that leaves the base DNs alone, an import of nothing - not a byte reaches the server |
| | | * and the stamp would prove nothing, while marking a connection the database may have dropped |
| | | * as the freshest one in the pool. Stamping proof rather than use makes the window mean |
| | | * "validated at most once per window", which is a claim this class can always back. |
| | | * <p> |
| | | * It stands for the moment the connection was <em>asked</em>, not the moment its answer was |
| | | * filed: {@link #distrustPool} is compared against it as an ordering of two moments, and a |
| | | * proof that took a second to arrive would otherwise outlive a drop reported while it was |
| | | * still in flight. Reading it early only ever ages the proof, which costs a validation and |
| | | * never skips one. |
| | | */ |
| | | private volatile long lastKnownAliveNanos; |
| | | |
| | | public CachedConnection(String connectionString, Connection parent) { |
| | | this(connectionString, parent, true); |
| | | } |
| | |
| | | this.connectionString = connectionString; |
| | | this.parent = parent; |
| | | this.poolable = poolable; |
| | | this.lastKnownAliveNanos = System.nanoTime(); |
| | | } |
| | | |
| | | /** |
| | |
| | | * not answer into a hang rather than into an error the caller can report. |
| | | */ |
| | | static Connection getConnection(String connectionString) throws Exception { |
| | | return getConnection(connectionString, true); |
| | | } |
| | | |
| | | /** |
| | | * Borrows a connection, either trusting the alive window of {@value #ALIVE_BYPASS_PROPERTY} or |
| | | * validating whatever comes out of the pool. |
| | | * |
| | | * @param trusted false for a borrow nothing compensates a dropped connection on. What the |
| | | * window trades away is the connection that breaks inside it, and {@code JDBCStorage} takes |
| | | * that off the caller where it can - a write is replayed, a read tells the pool - but the |
| | | * borrows that open a backend, remove its files or start an import have neither: they issue |
| | | * their statements far from the borrow, and the one that opens a backend issues none at all, |
| | | * so a dropped connection would surface out of the {@code rollback()} of its release. Each of |
| | | * 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 ConnectDialect dialect = ConnectDialect.of(connectionString); |
| | | reportUnknownDialect(connectionString, dialect); |
| | | final long connectTimeoutSeconds = Math.min( |
| | |
| | | long backoffMs = 0; |
| | | int attempts = 0; |
| | | while (true) { |
| | | final CachedConnection pooled = poll(connectionString, waitMs, deadline); |
| | | final CachedConnection pooled = poll(connectionString, waitMs, deadline, trusted); |
| | | if (pooled != null) { |
| | | return pooled; |
| | | } |
| | |
| | | * 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 validated first, 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. |
| | | * 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) throws InterruptedException { |
| | | CachedConnection con = cached.get(connectionString).poll(waitMs, TimeUnit.MILLISECONDS); |
| | | 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)) { |
| | | if (isUsable(con, trusted)) { |
| | | return con; |
| | | } |
| | | closeQuietly(con.parent); |
| | | if (System.currentTimeMillis() >= deadline) { |
| | | return null; |
| | | } |
| | | con = cached.get(connectionString).poll(); |
| | | con = cached.get(connectionString).pollFirst(); |
| | | } |
| | | return null; |
| | | } |
| | | |
| | | private static boolean isUsable(CachedConnection con) { |
| | | private static boolean isUsable(CachedConnection con, boolean trusted) { |
| | | if (trusted && isKnownAlive(con)) { |
| | | return true; |
| | | } |
| | | // The validation needs a bound of its own: isValid(0) means "no timeout" in the JDBC |
| | | // contract, and a connection whose socket is half-open answers it no sooner than it |
| | | // answers anything else. isValid(n) is not that bound on every driver either - the SQL |
| | |
| | | // avoid, and pooling it would hand out a connection carrying a bound of ours |
| | | return false; |
| | | } |
| | | // Read before the round trip rather than after it: this stamp is what distrustPool() is |
| | | // compared against, as an ordering of two moments. A validation is allowed |
| | | // VALIDATION_TIMEOUT_SECONDS, so a stamp filed once the answer is in can be younger than a |
| | | // drop another operation reported while it was still in flight - and the connection would |
| | | // then be trusted for the rest of the window by the very check that exists to stop it. |
| | | final long provenAt = System.nanoTime(); |
| | | boolean usable; |
| | | try { |
| | | usable = con.isValid(VALIDATION_TIMEOUT_SECONDS); |
| | |
| | | "the connection is closed rather than pooled")) { |
| | | return false; // it would carry the bound of the validation into every statement |
| | | } |
| | | con.lastKnownAliveNanos = provenAt; |
| | | return true; |
| | | } |
| | | |
| | |
| | | private static final int VALIDATION_BOUND_FAILED = -2; |
| | | |
| | | /** |
| | | * Whether a connection can be handed out on the strength of the last answer it gave, without a |
| | | * round trip to ask for another. Three things have to hold: the window is on, the answer is |
| | | * younger than it, and nothing has reported since that the database dropped a connection of |
| | | * this pool. |
| | | * <p> |
| | | * What the window trades away is the connection that breaks inside it: it is handed out, and |
| | | * the failure surfaces on the statement of the caller rather than on the borrow. That is where |
| | | * a connection breaking mid-operation surfaces anyway - but not every caller of this backend |
| | | * reports such a failure to the client, so the trade is not the caller's alone to bear. |
| | | * {@code JDBCStorage} answers it on both sides: a write is replayed on a connection the next |
| | | * attempt borrows of its own, and a read as much as a write marks the pool distrusted, which |
| | | * closes the window for the rest of the generation the dropped connection belonged to. |
| | | */ |
| | | private static boolean isKnownAlive(CachedConnection con) { |
| | | final long window = aliveBypassNanos; |
| | | if (window <= 0) { |
| | | return false; |
| | | } |
| | | final long provenAt = con.lastKnownAliveNanos; |
| | | if (System.nanoTime() - provenAt >= window) { // the overflow safe form of the comparison |
| | | return false; |
| | | } |
| | | final Long distrusted = poolDistrustedAt.get(con.connectionString); |
| | | 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. |
| | | return !isClosed(con.parent); |
| | | } |
| | | |
| | | /** Whether the driver reports the connection as closed; one that cannot say is not one to trust. */ |
| | | private static boolean isClosed(Connection con) { |
| | | try { |
| | | return con.isClosed(); |
| | | } catch (SQLException e) { |
| | | return true; |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Reports that the database dropped a connection of this pool, so that no connection proven |
| | | * alive before now is handed out unvalidated again. It is called by the operation that saw the |
| | | * failure: this class only ever learns of one from the statement it broke, since a borrow |
| | | * inside the window asks the database nothing. |
| | | */ |
| | | static void distrustPool(String connectionString) { |
| | | // merge(later of the two) rather than computeIfAbsent().set(): two operations reporting a |
| | | // drop at once would otherwise move the distrust point backwards - the later reading is |
| | | // written first and the earlier one overwrites it - and the AtomicLong of computeIfAbsent |
| | | // is published holding its initial 0 before set() runs, which a borrow racing it reads as |
| | | // "never". Not Math.max: nanoTime() has no defined origin, so the readings are compared by |
| | | // their difference, the way every other comparison of one in this class is. |
| | | poolDistrustedAt.merge(connectionString, System.nanoTime(), |
| | | (reported, now) -> now - reported > 0 ? now : reported); |
| | | } |
| | | |
| | | /** |
| | | * Bounds the socket of a pooled connection for the length of its validation, returning the |
| | | * network timeout to put back afterwards - or {@link #VALIDATION_BOUND_LEFT_ALONE} for a |
| | | * connection left alone, either because the driver does not take a network timeout or because |
| | |
| | | final Properties properties = new Properties(); |
| | | final boolean readBoundSet = dialect != null && connectTimeoutSeconds > 0 |
| | | && dialect.bound(connectionString, properties, connectTimeoutSeconds); |
| | | // Read before the connect rather than after it, for the reason isUsable() reads it before |
| | | // the validation: the login answered somewhere inside this attempt, and a stamp taken once |
| | | // it returned could outlive a drop reported while it was still going on. |
| | | final long provenAt = System.nanoTime(); |
| | | final Connection conNew = DriverManager.getConnection(connectionString, properties); |
| | | boolean poolable = true; |
| | | try { |
| | |
| | | closeQuietly(conNew); |
| | | throw e; |
| | | } |
| | | return new CachedConnection(connectionString, conNew, poolable); |
| | | final CachedConnection established = new CachedConnection(connectionString, conNew, poolable); |
| | | established.lastKnownAliveNanos = provenAt; |
| | | return established; |
| | | } |
| | | |
| | | // The second bound of the login is a socket read timeout on mysql, oracle and sql server, in |
| | |
| | | closeQuietly(parent); |
| | | return; |
| | | } |
| | | cached.get(connectionString).add(this); |
| | | // 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); |
| | | } |
| | | |
| | | @Override |