| New file |
| | |
| | | /* |
| | | * The contents of this file are subject to the terms of the Common Development and |
| | | * Distribution License (the License). You may not use this file except in compliance with the |
| | | * License. |
| | | * |
| | | * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the |
| | | * specific language governing permission and limitations under the License. |
| | | * |
| | | * When distributing Covered Software, include this CDDL Header Notice in each file and include |
| | | * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL |
| | | * Header, with the fields enclosed by brackets [] replaced by your own identifying |
| | | * information: "Portions Copyright [year] [name of copyright owner]". |
| | | * |
| | | * Copyright 2026 3A Systems, LLC. |
| | | */ |
| | | package org.opends.server.backends.jdbc; |
| | | |
| | | import org.forgerock.opendj.ldap.ByteString; |
| | | import org.forgerock.opendj.server.config.server.JDBCBackendCfg; |
| | | import org.opends.server.DirectoryServerTestCase; |
| | | import org.opends.server.backends.pluggable.spi.AccessMode; |
| | | import org.opends.server.backends.pluggable.spi.Importer; |
| | | import org.opends.server.backends.pluggable.spi.StorageRuntimeException; |
| | | import org.opends.server.backends.pluggable.spi.TreeName; |
| | | import org.testng.annotations.AfterClass; |
| | | import org.testng.annotations.AfterMethod; |
| | | import org.testng.annotations.BeforeClass; |
| | | import org.testng.annotations.BeforeMethod; |
| | | import org.testng.annotations.Test; |
| | | |
| | | import java.sql.Connection; |
| | | import java.sql.Driver; |
| | | import java.sql.DriverManager; |
| | | import java.sql.DriverPropertyInfo; |
| | | import java.sql.PreparedStatement; |
| | | import java.sql.SQLException; |
| | | import java.util.ArrayList; |
| | | import java.util.Arrays; |
| | | import java.util.Collections; |
| | | import java.util.IdentityHashMap; |
| | | import java.util.LinkedHashSet; |
| | | import java.util.List; |
| | | import java.util.Map; |
| | | import java.util.Properties; |
| | | import java.util.Set; |
| | | import java.util.concurrent.BrokenBarrierException; |
| | | import java.util.concurrent.CountDownLatch; |
| | | import java.util.concurrent.CyclicBarrier; |
| | | import java.util.concurrent.ExecutorService; |
| | | import java.util.concurrent.Executors; |
| | | import java.util.concurrent.Future; |
| | | import java.util.concurrent.TimeUnit; |
| | | import java.util.concurrent.TimeoutException; |
| | | import java.util.concurrent.atomic.AtomicBoolean; |
| | | import java.util.concurrent.atomic.AtomicInteger; |
| | | import java.util.concurrent.atomic.AtomicReference; |
| | | import java.util.function.Supplier; |
| | | import java.util.logging.Logger; |
| | | |
| | | import static org.forgerock.opendj.config.ConfigurationMock.mockCfg; |
| | | import static org.mockito.Mockito.anyInt; |
| | | import static org.mockito.Mockito.anyString; |
| | | import static org.mockito.Mockito.doAnswer; |
| | | import static org.mockito.Mockito.doThrow; |
| | | import static org.mockito.Mockito.never; |
| | | import static org.mockito.Mockito.times; |
| | | import static org.mockito.Mockito.verify; |
| | | import static org.mockito.Mockito.when; |
| | | import static org.testng.Assert.assertEquals; |
| | | import static org.testng.Assert.assertFalse; |
| | | import static org.testng.Assert.assertTrue; |
| | | import static org.testng.Assert.fail; |
| | | |
| | | /** |
| | | * The {@link Importer} of the JDBC backend is used by several threads at once - phase two of an |
| | | * import runs one thread per tree - and {@code java.sql.Connection} is not thread-safe. Sharing |
| | | * one connection between those threads corrupts the driver rather than merely serializing the |
| | | * work: the sql server driver keeps the reconnect listeners of a connection in a plain |
| | | * {@code ArrayList} that every {@code prepareStatement()} and every statement {@code close()} |
| | | * mutates, and two import threads on one connection walked it past the end of its array |
| | | * (issue #891). |
| | | * <p> |
| | | * Needs no database: the connections are handed out by a driver of this test, which records every |
| | | * thread inside a statement of a connection - and which thread commits, rolls back or returns one |
| | | * while another is in there - so what it records is the defect itself, two threads on one |
| | | * connection at the same time, rather than an exception of one driver. It therefore holds for every |
| | | * dialect this backend takes. |
| | | * <p> |
| | | * Two ways of getting the threads to overlap, because trees that get connections of their own and |
| | | * trees that share one cannot be asked the same question. Where the connections differ, the driver |
| | | * holds the first statement of each thread at a rendezvous until its peers have one in flight of |
| | | * their own. Where they do not, a second thread never reaches such a rendezvous - the monitor of |
| | | * the connection is what this is about - so the threads line up before they write instead, and the |
| | | * first one in holds the connection for a moment to see whether the other turns up beside it. |
| | | */ |
| | | @SuppressWarnings("javadoc") |
| | | @Test(groups = { "precommit", "jdbc" }, sequential = true) |
| | | public class ImportConnectionsTestCase extends DirectoryServerTestCase { |
| | | |
| | | /** Long enough for a peer thread to reach a statement on a loaded machine, and no longer. */ |
| | | private static final long RENDEZVOUS_SECONDS = 30; |
| | | |
| | | /** |
| | | * How long the first thread inside a statement holds the connection it is on. Paid in full by |
| | | * every run where the importer serializes its threads as it should - which is the point: a peer |
| | | * that is let in arrives at once, having been waiting on the monitor since the starting line. |
| | | */ |
| | | private static final long HOLD_SECONDS = 2; |
| | | |
| | | /** How long a close of an import may take before a thread of the test is taken to be stuck. */ |
| | | private static final long CLOSE_SECONDS = 30; |
| | | |
| | | /** |
| | | * How long a thread parked inside a statement stays in there before it gives up on the test |
| | | * that parked it. Paid by no passing run - the test lets it go as soon as it has what it came |
| | | * for - and it is bounded so that a run which never gets there is a red test rather than a |
| | | * build that hangs on a monitor nobody is going to release. |
| | | */ |
| | | private static final long PARK_SECONDS = 30; |
| | | |
| | | private static final TreeName ID2ENTRY = new TreeName("dc=example,dc=com", "id2entry"); |
| | | private static final TreeName DN2ID = new TreeName("dc=example,dc=com", "dn2id"); |
| | | private static final TreeName STATE = new TreeName("dc=example,dc=com", "state"); |
| | | |
| | | private final StubDriver stub = new StubDriver(); |
| | | /** |
| | | * What the driver of this test answers a connect with, or null for a connection of its own. |
| | | * A failure per connect rather than one instance thrown over and over: a caller is free to |
| | | * give what it caught a cause and a suppressed exception of its own. |
| | | */ |
| | | private final AtomicReference<Supplier<SQLException>> refusal = new AtomicReference<>(); |
| | | |
| | | /** The connections an import issued a statement on, in no particular order. */ |
| | | private final Set<Connection> used = |
| | | Collections.synchronizedSet(Collections.newSetFromMap(new IdentityHashMap<Connection, Boolean>())); |
| | | /** The same, in the order the import first issued a statement on them. */ |
| | | private final List<Connection> usedInOrder = Collections.synchronizedList(new ArrayList<Connection>()); |
| | | /** Every commit of the import, in the order the connections took them. */ |
| | | private final List<Connection> commits = Collections.synchronizedList(new ArrayList<Connection>()); |
| | | /** |
| | | * The threads inside a statement of a connection right now - every one of them, rather than the |
| | | * first to arrive: a thread that finds a peer there records the pair and stays out of the map, |
| | | * and the peer's exit would then leave the map saying the connection is free while it is not. |
| | | * What comes after - another statement, a commit, the rollback of a return - would be recorded |
| | | * against nobody, which is the one thing this suite is here to notice. |
| | | */ |
| | | private final Map<Connection, Set<Thread>> inStatement = |
| | | Collections.synchronizedMap(new IdentityHashMap<Connection, Set<Thread>>()); |
| | | /** Every pair of threads that was inside a statement of one connection at the same time. */ |
| | | private final List<String> shared = Collections.synchronizedList(new ArrayList<String>()); |
| | | |
| | | /** Where the threads of a test meet, so that their statements are in flight at the same time. */ |
| | | private volatile CyclicBarrier rendezvous; |
| | | /** Whether this thread has already met its peers: a thread waits there once, whatever it issues after. */ |
| | | private final ThreadLocal<Boolean> met = ThreadLocal.withInitial(() -> Boolean.FALSE); |
| | | |
| | | /** |
| | | * Where the threads of a test line up before they write, which is in front of every monitor the |
| | | * importer takes - unlike {@link #rendezvous}, which they reach once they are already inside a |
| | | * statement, and which two threads sharing a connection therefore cannot both reach. |
| | | */ |
| | | private volatile CyclicBarrier startingLine; |
| | | /** |
| | | * Counted down by every thread that gets inside a statement, and waited on by the first of them: |
| | | * it holds the connection long enough for a peer to get in there beside it, if anything lets it. |
| | | */ |
| | | private volatile CountDownLatch insideAStatement; |
| | | /** Whether this thread has already held a connection: it holds the first statement it issues. */ |
| | | private final ThreadLocal<Boolean> held = ThreadLocal.withInitial(() -> Boolean.FALSE); |
| | | |
| | | /** |
| | | * Counted down by the first thread of a test to get inside a statement, and awaited by the test |
| | | * itself: what a write in flight looks like to the rest of the import - a connection whose |
| | | * statement has not come back, and which therefore has nothing recorded against it yet. |
| | | */ |
| | | private volatile CountDownLatch insideAStatementNow; |
| | | /** Counted down by the test to let the parked write finish, once it has seen what it came to see. */ |
| | | private volatile CountDownLatch letTheStatementFinish; |
| | | /** Whether this thread has already been parked: a thread is parked in the first statement it issues. */ |
| | | private final ThreadLocal<Boolean> parked = ThreadLocal.withInitial(() -> Boolean.FALSE); |
| | | |
| | | private ExecutorService threads; |
| | | /** The pool bound and its borrow deadline as this JVM had them, put back after every test. */ |
| | | private String poolMaxOfTheJvm; |
| | | private String poolTimeoutOfTheJvm; |
| | | |
| | | @BeforeClass |
| | | public void registerStubDriver() throws Exception { |
| | | DriverManager.registerDriver(stub); |
| | | } |
| | | |
| | | /** |
| | | * The bound of the pools of this suite, pinned: how many connections an import takes is clamped |
| | | * to it, so a value another suite of this package left set - CachedConnectionTestCase varies it |
| | | * - would decide the counts asserted here. A pool reads it once, when the first borrow of a |
| | | * connection string creates it, so it has to stand before the storage of a test is opened. |
| | | */ |
| | | @BeforeMethod |
| | | public void pinThePoolBound() { |
| | | poolMaxOfTheJvm = System.getProperty(CachedConnection.POOL_MAX_PROPERTY); |
| | | poolTimeoutOfTheJvm = System.getProperty(CachedConnection.POOL_TIMEOUT_PROPERTY); |
| | | System.setProperty(CachedConnection.POOL_MAX_PROPERTY, "16"); |
| | | } |
| | | |
| | | @AfterClass |
| | | public void deregisterStubDriver() throws Exception { |
| | | DriverManager.deregisterDriver(stub); |
| | | } |
| | | |
| | | @AfterMethod |
| | | public void forgetWhatTheTestRecorded() { |
| | | System.clearProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY); |
| | | // put back rather than cleared: a sibling suite of this package varies this property, and a |
| | | // pool is created with the value that stands when its connection string is first borrowed on |
| | | putBack(CachedConnection.POOL_MAX_PROPERTY, poolMaxOfTheJvm); |
| | | putBack(CachedConnection.POOL_TIMEOUT_PROPERTY, poolTimeoutOfTheJvm); |
| | | // the rendezvous of the next test is not the one this thread already went to |
| | | met.remove(); |
| | | held.remove(); |
| | | parked.remove(); |
| | | // let go of anything this test left parked before its threads are shut down: a write held |
| | | // inside a statement holds the monitor of a connection, and the next test of this suite |
| | | // starts by opening a storage on that pool |
| | | if (letTheStatementFinish != null) { |
| | | letTheStatementFinish.countDown(); |
| | | } |
| | | if (threads != null) { |
| | | threads.shutdownNow(); |
| | | threads = null; |
| | | } |
| | | refusal.set(null); |
| | | rendezvous = null; |
| | | startingLine = null; |
| | | insideAStatement = null; |
| | | insideAStatementNow = null; |
| | | letTheStatementFinish = null; |
| | | forgetTheConnections(); |
| | | inStatement.clear(); |
| | | shared.clear(); |
| | | } |
| | | |
| | | private static void putBack(String property, String value) { |
| | | if (value == null) { |
| | | System.clearProperty(property); |
| | | } else { |
| | | System.setProperty(property, value); |
| | | } |
| | | } |
| | | |
| | | private void forgetTheConnections() { |
| | | used.clear(); |
| | | usedInOrder.clear(); |
| | | commits.clear(); |
| | | } |
| | | |
| | | /** |
| | | * Two trees written at the same time are written through connections of their own. This is the |
| | | * contract {@code Importer} states and the defect of #891: every thread of phase two wrote |
| | | * through the one connection the importer borrowed in its constructor. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testTwoTreesAreWrittenThroughConnectionsOfTheirOwn() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-parallel"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importing(importer, () -> { |
| | | // both threads have to be inside a statement at once for the record below to say |
| | | // anything: a run where one finished before the other started shares no connection |
| | | // however the importer borrows them |
| | | rendezvous = new CyclicBarrier(2); |
| | | threads = daemonThreads(2, "import-writer"); |
| | | final Future<?> id2entry = put(importer, ID2ENTRY); |
| | | final Future<?> dn2id = put(importer, DN2ID); |
| | | id2entry.get(RENDEZVOUS_SECONDS * 2, TimeUnit.SECONDS); |
| | | dn2id.get(RENDEZVOUS_SECONDS * 2, TimeUnit.SECONDS); |
| | | |
| | | // belt and braces here rather than the cover of this suite: with a connection each |
| | | // the record below is empty however the import behaves, and a run where the two |
| | | // trees did share one fails at the rendezvous instead - the second thread would be |
| | | // on the monitor of the connection and never reach it. What pins the monitor is |
| | | // testTwoTreesOnOneConnectionAreNotWrittenAtTheSameTime |
| | | assertEquals(shared, Collections.emptyList(), |
| | | "two import threads issued a statement on one connection at the same time"); |
| | | assertEquals(used.size(), 2, "the two trees of the import did not get a connection each"); |
| | | }); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Two trees that share one connection are not written at the same time. The bound of an import, |
| | | * and a pool with nothing to spare, both put several trees on one connection - and what #891 is |
| | | * about is two threads issuing statements on one connection, not how many connections an import |
| | | * holds: the monitor of the connection is what has to keep the second thread out. |
| | | * <p> |
| | | * The two threads cannot be brought together at a barrier inside their statements the way the |
| | | * test above does it - the second one never gets in to reach it - so the first one in holds the |
| | | * connection for {@link #HOLD_SECONDS} instead, having lined up with its peer beforehand. A |
| | | * second thread let in beside it arrives at once, since it has been waiting on that monitor |
| | | * since the starting line, and both are then recorded in {@link #shared}. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testTwoTreesOnOneConnectionAreNotWrittenAtTheSameTime() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "1"); // both trees on one connection |
| | | final String url = StubDriver.PREFIX + "importer-serialized"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importing(importer, () -> { |
| | | startingLine = new CyclicBarrier(2); |
| | | insideAStatement = new CountDownLatch(2); |
| | | threads = daemonThreads(2, "import-writer"); |
| | | final Future<?> id2entry = put(importer, ID2ENTRY); |
| | | final Future<?> dn2id = put(importer, DN2ID); |
| | | id2entry.get(RENDEZVOUS_SECONDS * 2, TimeUnit.SECONDS); |
| | | dn2id.get(RENDEZVOUS_SECONDS * 2, TimeUnit.SECONDS); |
| | | |
| | | assertEquals(used.size(), 1, "the two trees of this import did not share one connection"); |
| | | assertEquals(shared, Collections.emptyList(), |
| | | "two import threads issued a statement on one connection at the same time"); |
| | | }); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * ... and no more connections than that: an import takes them for its whole duration, and the |
| | | * pool they come from is bounded and shared with the operations of every other backend on that |
| | | * database (#878). One thread per tree is what phase two runs, and a default backend has trees |
| | | * enough to empty a default pool. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testTreesShareTheConnectionsOfABoundAboveOne() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-bounded-two"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | // five trees over two connections: the bound is what the round robin wraps at |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(STATE, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(new TreeName("dc=example,dc=com", "id2childrencount"), |
| | | ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(new TreeName("dc=example,dc=com", "referral"), |
| | | ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | |
| | | assertEquals(used.size(), 2, "the import did not spread its trees over the connections of its bound"); |
| | | } finally { |
| | | importer.close(); |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * The bound an import takes where the operator set none, which is the branch every real install |
| | | * takes: half the bound of the pool, and no more than half of what the default bound would be. |
| | | * An import holds what it takes for its whole duration out of a pool that is shared with the |
| | | * LDAP traffic of every backend on that database (#878), so the deployment that raised the |
| | | * bound of its pool did not thereby ask for an import to take half of it. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testTheBoundOfAnImportDefaultsToHalfThePool() throws Exception { |
| | | // set by no test of this suite, and cleared after every one of them: this is the default |
| | | System.clearProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY); |
| | | // a bound of the pool that is not the default one: DEFAULT_POOL_MAX is |
| | | // max(16, processors * 2), so the 16 pinThePoolBound() puts on the pools of this suite is |
| | | // the default itself on a machine of eight cores or fewer - and half the pool and half the |
| | | // default are then the same number, which is no test of which of the two an import halves |
| | | System.setProperty(CachedConnection.POOL_MAX_PROPERTY, "6"); |
| | | final JDBCStorage half = importingStorage(StubDriver.PREFIX + "importer-default-bound"); |
| | | try { |
| | | assertEquals(half.importConnections(), 3, "an import did not default to half the bound of the pool"); |
| | | } finally { |
| | | half.close(); |
| | | } |
| | | System.setProperty(CachedConnection.POOL_MAX_PROPERTY, String.valueOf(CachedConnection.DEFAULT_POOL_MAX * 4)); |
| | | final JDBCStorage capped = importingStorage(StubDriver.PREFIX + "importer-default-bound-large-pool"); |
| | | try { |
| | | assertEquals(capped.importConnections(), CachedConnection.DEFAULT_POOL_MAX / 2, |
| | | "an import of a pool larger than the default took more than half of what the default bound would be"); |
| | | } finally { |
| | | capped.close(); |
| | | } |
| | | } |
| | | |
| | | /** ... and a bound the pool cannot honour is the pool's, not the property's. */ |
| | | @Test(timeOut = 120000) |
| | | public void testTheBoundOfAnImportIsClampedToTheBoundOfThePool() throws Exception { |
| | | System.setProperty(CachedConnection.POOL_MAX_PROPERTY, "2"); |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "8"); |
| | | final String url = StubDriver.PREFIX + "importer-clamped"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | try { |
| | | assertEquals(storage.importConnections(), 2, "the bound of an import was not clamped to the pool"); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | @Test(timeOut = 120000) |
| | | public void testTheConnectionsOfAnImportAreBounded() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "1"); |
| | | final String url = StubDriver.PREFIX + "importer-bounded"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(STATE, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | |
| | | assertEquals(used.size(), 1, "the import took more connections than the bound allows"); |
| | | } finally { |
| | | importer.close(); |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * The trees decide which connection a write goes through, not the threads that write them. The |
| | | * connections of an import are transactions of their own, so two of them writing one row would |
| | | * have the second wait for the first to commit - which an import does in {@code close()}, when |
| | | * the thread that would have to release the row is long done, and the bulk class carries no |
| | | * bound to break the wait. {@code setTrust()} writes the state tree of a container from an |
| | | * import thread in {@code beforePhaseOne} and from the closing thread in {@code afterPhaseTwo}, |
| | | * which is exactly that shape. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testOneTreeIsWrittenThroughOneConnectionWhoeverWritesIt() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "4"); |
| | | final String url = StubDriver.PREFIX + "importer-one-tree"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importing(importer, () -> { |
| | | // three threads, one tree: an import thread of phase one, another of phase two, and |
| | | // the thread that closes the importer, which is the one afterPhaseTwo runs on |
| | | putOnAThreadOfItsOwn(importer, STATE); |
| | | putOnAThreadOfItsOwn(importer, STATE); |
| | | importer.put(STATE, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | |
| | | assertEquals(used.size(), 1, "one tree was written through more than one connection"); |
| | | }); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Every connection an import took is committed and given back, not only the one its |
| | | * constructor borrowed: what a connection left behind holds is a transaction of the import and |
| | | * a permit of the pool, and a pool is never removed from the map - so a permit lost to an |
| | | * import is lost for the life of the server (#878). |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testEveryConnectionOfAnImportIsCommittedAndReturned() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "3"); |
| | | final String url = StubDriver.PREFIX + "importer-returned"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(STATE, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | |
| | | importer.close(); |
| | | |
| | | assertEquals(used.size(), 3, "the three trees of the import did not get a connection each"); |
| | | for (final Connection con : connectionsUsed()) { |
| | | verify(con, times(1)).commit(); |
| | | } |
| | | assertEquals(CachedConnection.poolOf(url).idleCount(), 3, "an import kept a connection of the pool"); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | assertEquals(CachedConnection.poolOf(url).meteredCount(), 0, "an import kept a permit of the pool"); |
| | | } |
| | | |
| | | /** |
| | | * A clear commits every connection of the import, not only the one whose tree it clears, and it |
| | | * does so before it empties anything. The one connection an import used to hold made that so of |
| | | * its own accord - {@code clearTree()} ends in a commit, and that commit made durable every |
| | | * write the import had made so far. What that is worth to a {@code rebuild-index} is the |
| | | * {@code setTrust(false)} of {@code RebuildIndexStrategy.beforePhaseOne}, which is written just |
| | | * in front of the clear of the tree it describes: left to {@code close()} that flag would still |
| | | * be uncommitted while the table was emptied and committed, and a server that stopped in |
| | | * between would come back to an index that is empty and marked trusted. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testAClearCommitsTheOtherConnectionsOfTheImport() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-clear-commits"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | // the state tree stands for the trust flag: written, and left uncommitted until |
| | | // something commits the connection it went to |
| | | importer.put(STATE, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.clearTree(ID2ENTRY); |
| | | |
| | | // asserted in order rather than counted: the destructive half of a clear is the one |
| | | // thing that must not be durable while a write it is meant to invalidate is not, so the |
| | | // commit of the peer belongs in front of the one the clear itself ends in - a count of |
| | | // two is the same whichever way round they happened |
| | | assertEquals(commits, Arrays.asList(usedInOrder.get(0), usedInOrder.get(1)), |
| | | "the clear did not commit the other connection of the import before it emptied its own tree"); |
| | | } finally { |
| | | importer.close(); |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * ... including a peer whose first write is still in flight, which is not a connection with |
| | | * nothing to commit but one whose write has not come back yet. |
| | | * <p> |
| | | * An import counts a write once its statement returns, so a connection written for the first |
| | | * time has nothing recorded against it for exactly as long as that statement takes. That window |
| | | * is reached by more than one thread at once: {@code beforePhaseOne} runs on the import threads |
| | | * - {@code OnDiskMergeImporter.processEntry} guards it with a latch per container, not one for |
| | | * all of them - so on an {@code import-ldif} of two base DNs the {@code setTrust(false)} of one |
| | | * container is written while the clears of the next are running. The one connection an import |
| | | * held before #891 carried such a write into the commit of the clear, being the same |
| | | * transaction and issued in front of it in the driver, so a clear that passed over it here |
| | | * would be a commit point the single connection did not have. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testAClearCommitsAPeerWhoseWriteIsStillInFlight() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-clear-commits-in-flight"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importing(importer, () -> { |
| | | insideAStatementNow = new CountDownLatch(1); |
| | | letTheStatementFinish = new CountDownLatch(1); |
| | | // the state tree stands for the trust flag of beforePhaseOne: written on a thread of |
| | | // the import, and held inside the statement it issues |
| | | final AtomicReference<Throwable> flagFailure = new AtomicReference<>(); |
| | | final Thread writing = daemon(() -> { |
| | | try { |
| | | importer.put(STATE, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | } catch (Throwable t) { |
| | | flagFailure.set(t); |
| | | } |
| | | }, "import-flag"); |
| | | // the clear of the next container, on a connection of its own |
| | | final AtomicReference<Throwable> clearFailure = new AtomicReference<>(); |
| | | final Thread clearing = daemon(() -> { |
| | | try { |
| | | importer.clearTree(ID2ENTRY); |
| | | } catch (Throwable t) { |
| | | clearFailure.set(t); |
| | | } |
| | | }, "import-clear"); |
| | | try { |
| | | writing.start(); |
| | | assertTrue(insideAStatementNow.await(RENDEZVOUS_SECONDS, TimeUnit.SECONDS), |
| | | "the write of the trust flag never got inside a statement"); |
| | | // where the clear gets to on its own: waiting on the monitor of the peer, which |
| | | // is the commit point being paid for - a clear that decided the peer had nothing |
| | | // to commit runs through to the end instead |
| | | clearing.start(); |
| | | untilBlockedOrDone(clearing); |
| | | } finally { |
| | | letTheStatementFinish.countDown(); |
| | | } |
| | | writing.join(TimeUnit.SECONDS.toMillis(RENDEZVOUS_SECONDS)); |
| | | clearing.join(TimeUnit.SECONDS.toMillis(RENDEZVOUS_SECONDS)); |
| | | reportFailureOf("the write of the trust flag", flagFailure); |
| | | reportFailureOf("the clear", clearFailure); |
| | | assertFalse(writing.isAlive(), "the write of the trust flag did not finish"); |
| | | assertFalse(clearing.isAlive(), "the clear did not finish"); |
| | | |
| | | assertEquals(usedInOrder.size(), 2, "the two trees of this import did not get a connection each"); |
| | | assertEquals(commits, Arrays.asList(usedInOrder.get(0), usedInOrder.get(1)), |
| | | "the clear passed over a peer whose write was still in flight"); |
| | | }); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * The connection of the last write is committed last. The connections of an import are |
| | | * transactions of their own, so {@code close()} cannot commit them as one - and the last thing |
| | | * an import writes is the flag that says the rest of it is good ({@code setTrust(true)} of |
| | | * {@code afterPhaseTwo}). Committed last, that flag is rolled back with the connection it is on |
| | | * whenever an earlier commit fails, so a failed import cannot leave an index marked trusted |
| | | * over data that never got there. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testTheConnectionsAreCommittedInTheOrderTheyWereLastWrittenTo() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "3"); |
| | | final String url = StubDriver.PREFIX + "importer-last-write"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(STATE, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | // the two writes that stand for the trust flags afterPhaseTwo writes, one per base DN: |
| | | // both have to be committed after the connection that holds nothing but data |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k2"), ByteString.valueOfUtf8("v")); |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k2"), ByteString.valueOfUtf8("v")); |
| | | |
| | | importer.close(); |
| | | |
| | | assertEquals(usedInOrder.size(), 3, "the three trees of the import did not get a connection each"); |
| | | assertEquals(commits, Arrays.asList(usedInOrder.get(2), usedInOrder.get(0), usedInOrder.get(1)), |
| | | "the connections were not committed in the order they were last written to"); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * {@code close()} waits for a statement of a straggler to finish rather than committing the |
| | | * connection beside it. |
| | | * <p> |
| | | * Phase two gives its threads five seconds to answer an interrupt and closes the importer |
| | | * whether they answered or not ({@code OnDiskMergeImporter.invokeParallel}), so a write of an |
| | | * import can still be in flight when its connections are committed and handed back. A commit |
| | | * issued beside that statement is two threads on one connection - the whole of #891, with the |
| | | * closing thread on one side of it. Every other test of this suite joins its writers before it |
| | | * closes, so this is the only one where anything of a close happens beside a statement, and it |
| | | * is what the other half of the detector is for: the one that watches a commit, a rollback or a |
| | | * return land while a thread is inside a statement. |
| | | * <p> |
| | | * The straggler is let go once the close has got as far as it can on its own, and its write is |
| | | * then committed rather than lost: a thread that was inside the import when {@code close()} |
| | | * began is part of it. That is the assertion which bites - a close that reads the state of a |
| | | * connection without waiting for the statement in flight on it finds nothing written and |
| | | * commits nothing, and the write of the straggler is rolled back by the return instead. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testACloseWaitsForAStatementOfAStragglerRatherThanCommittingBesideIt() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-close-waits"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importing(importer, () -> { |
| | | insideAStatementNow = new CountDownLatch(1); |
| | | letTheStatementFinish = new CountDownLatch(1); |
| | | final AtomicReference<Throwable> stragglerFailure = new AtomicReference<>(); |
| | | final Thread straggling = daemon(() -> { |
| | | try { |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | } catch (Throwable t) { |
| | | stragglerFailure.set(t); |
| | | } |
| | | }, "import-straggler"); |
| | | final AtomicReference<Throwable> closeFailure = new AtomicReference<>(); |
| | | final Thread closing = daemon(() -> { |
| | | try { |
| | | importer.close(); |
| | | } catch (Throwable t) { |
| | | closeFailure.set(t); |
| | | } |
| | | }, "import-close-beside-a-statement"); |
| | | try { |
| | | straggling.start(); |
| | | assertTrue(insideAStatementNow.await(RENDEZVOUS_SECONDS, TimeUnit.SECONDS), |
| | | "the write of the straggler never got inside a statement"); |
| | | // where the close gets to on its own: waiting on the monitor of the connection |
| | | // the straggler is writing through |
| | | closing.start(); |
| | | untilBlockedOrDone(closing); |
| | | } finally { |
| | | letTheStatementFinish.countDown(); |
| | | } |
| | | straggling.join(TimeUnit.SECONDS.toMillis(RENDEZVOUS_SECONDS)); |
| | | closing.join(TimeUnit.SECONDS.toMillis(CLOSE_SECONDS)); |
| | | reportFailureOf("the write of the straggler", stragglerFailure); |
| | | reportFailureOf("the close of the import", closeFailure); |
| | | assertFalse(straggling.isAlive(), "the write of the straggler did not finish"); |
| | | assertFalse(closing.isAlive(), "the import did not close"); |
| | | |
| | | assertEquals(shared, Collections.emptyList(), |
| | | "the close touched a connection while a thread of the import was inside a statement of it"); |
| | | assertEquals(usedInOrder.size(), 1, "the write of the straggler did not reach a connection"); |
| | | verify(usedInOrder.get(0), times(1)).commit(); |
| | | }); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | assertEquals(CachedConnection.poolOf(url).meteredCount(), 0, "the import kept a permit of the pool"); |
| | | } |
| | | |
| | | /** |
| | | * A write that arrives after {@code close()} takes no connection. Phase two gives its threads |
| | | * five seconds to answer an interrupt and then closes the importer whether they did or not, so |
| | | * a write can reach it once its connections are back in the pool - and one that borrowed there |
| | | * would take a permit nothing is left to give back, out of a pool that is never removed from |
| | | * the map (#878). |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testAWriteAfterCloseTakesNoConnection() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "4"); |
| | | final String url = StubDriver.PREFIX + "importer-after-close"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.close(); |
| | | |
| | | try { |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | fail("a write that arrives after the import is closed must be refused"); |
| | | } catch (StorageRuntimeException expected) { |
| | | // the import is over: there is no transaction left for this write to belong to |
| | | } |
| | | try { |
| | | // the same for a clear, which commits the other connections of the import before it |
| | | // empties a tree: those are back in the pool, serving whoever borrowed them next |
| | | importer.clearTree(ID2ENTRY); |
| | | fail("a clear that arrives after the import is closed must be refused"); |
| | | } catch (StorageRuntimeException expected) { |
| | | // as above |
| | | } |
| | | |
| | | assertEquals(used.size(), 1, "a write after close() took a connection of the pool"); |
| | | assertEquals(CachedConnection.poolOf(url).idleCount(), 1, "a write after close() kept a connection"); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | assertEquals(CachedConnection.poolOf(url).meteredCount(), 0, "a write after close() kept a permit"); |
| | | } |
| | | |
| | | /** |
| | | * A commit that fails still gives every connection of the import back, and takes the ones after |
| | | * it down with it. What a connection left behind holds is a permit of a pool that is never |
| | | * removed from the map, so a permit lost to a failed import is lost for the life of the server |
| | | * (#878) - and what a connection committed after the failure of an earlier one would hold is |
| | | * the trust flags of {@code afterPhaseTwo}, which {@code close()} commits last for exactly this |
| | | * reason: a failed import must not leave an index marked trusted over data that never got there. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testAFailedCommitStillReturnsEveryConnection() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "3"); |
| | | final String url = StubDriver.PREFIX + "importer-failed-commit"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(STATE, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | doThrow(new SQLException("the socket went away")).when(usedInOrder.get(1)).commit(); |
| | | |
| | | try { |
| | | importer.close(); |
| | | fail("the failure of the commit was not reported"); |
| | | } catch (StorageRuntimeException expected) { |
| | | assertEquals(expected.getCause().getMessage(), "the socket went away"); |
| | | } |
| | | |
| | | // the connections are committed in the order they were last written to, so this one |
| | | // stands where the trust flags of afterPhaseTwo stand: behind the connection that just |
| | | // failed, and rolled back by the return rather than committed on top of a failed import |
| | | verify(usedInOrder.get(2), never()).commit(); |
| | | assertEquals(CachedConnection.poolOf(url).idleCount(), 3, "a failed import kept a connection of the pool"); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | assertEquals(CachedConnection.poolOf(url).meteredCount(), 0, "a failed import kept a permit of the pool"); |
| | | } |
| | | |
| | | /** |
| | | * A tree that cannot be given a connection of its own is written through one the import already |
| | | * holds, rather than failing the import. The connections are taken as the trees are first |
| | | * touched, so a pool at its bound would otherwise stop an import halfway through - with the |
| | | * clears of {@code beforePhaseOne} already committed - where the one connection an import held |
| | | * before #891 would have carried it to the end. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testATreeSharesAConnectionWhenThePoolHasNoneToSpare() throws Exception { |
| | | System.setProperty(CachedConnection.POOL_MAX_PROPERTY, "2"); |
| | | System.setProperty(CachedConnection.POOL_TIMEOUT_PROPERTY, "1"); // nothing is going to be returned |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-pool-full"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | try { |
| | | final Importer importer = storage.startImport(); // one of the two connections of the pool |
| | | // and the other one to an operation of the server, so that the pool has none to spare |
| | | final Connection heldByAnOperation = CachedConnection.getConnection(url); |
| | | try { |
| | | importing(importer, () -> { |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | // on a thread of its own: a thread that already holds a connection of the pool is |
| | | // exempt from waiting at its bound, and every thread of phase two starts holding |
| | | // none |
| | | putOnAThreadOfItsOwn(importer, DN2ID); |
| | | |
| | | assertEquals(used.size(), 1, "the tree the pool had no connection for did not share one"); |
| | | }); |
| | | } finally { |
| | | heldByAnOperation.close(); |
| | | } |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * A second {@code close()} touches nothing: the connections of the import are back in the pool |
| | | * by then, and the pool hands them out again - so a close that walked its map a second time |
| | | * would roll back and re-pool a connection another borrower is holding. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testASecondCloseTouchesNothingTheImportGaveBack() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-second-close"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | try { |
| | | final Importer importer = storage.startImport(); |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | importer.close(); |
| | | |
| | | // the pool hands one of them straight back out, the way an LDAP operation would |
| | | final Connection borrowedAgain = CachedConnection.getConnection(url); |
| | | try { |
| | | importer.close(); |
| | | assertEquals(CachedConnection.poolOf(url).idleCount(), 1, |
| | | "a second close() returned a connection the pool had already handed to somebody else"); |
| | | } finally { |
| | | borrowedAgain.close(); |
| | | } |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | assertEquals(CachedConnection.poolOf(url).meteredCount(), 0, "a connection of the import was lost"); |
| | | } |
| | | |
| | | /** |
| | | * A tree the database refuses a connection for is written through one the import already holds, |
| | | * the way a tree the pool has nothing to spare for is (#1013). |
| | | * <p> |
| | | * The two ends of a borrow that can run out are not reported alike: the bound of the pool is |
| | | * answered with a {@code SQLTimeoutException}, while a database at a limit of its own is |
| | | * answered with whatever its driver says - and a code {@code CachedConnection.isWorthRetrying} |
| | | * does not recognize is raised at once rather than waited out. A mysql account with a |
| | | * {@code MAX_USER_CONNECTIONS} of its own answers 1226 on SQLState 42000, which is such a code |
| | | * (#1011). An import that failed on it would stop halfway through, with the clears of |
| | | * {@code beforePhaseOne} already committed, while every connection it holds still works. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testATreeSharesAConnectionWhenTheDatabaseRefusesOne() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-refused"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importing(importer, () -> { |
| | | // the first tree takes the connection the constructor borrowed, so that the second is |
| | | // one the import has to go to the database for |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | // armed only now: what this is about is the borrows an import makes as it goes, not |
| | | // the one it is started with - an import of a database that takes no connection at |
| | | // all is refused where it is started |
| | | refuseFurtherConnects("User 'opendj' has exceeded the 'max_user_connections' resource" |
| | | + " (current value: 1)", "42000", 1226); |
| | | // on a thread of its own, the way every thread of phase two starts: one that already |
| | | // holds a connection of the pool takes another path through the borrow |
| | | putOnAThreadOfItsOwn(importer, DN2ID); |
| | | |
| | | assertEquals(used.size(), 1, "the tree the database refused a connection for did not share one"); |
| | | }); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | assertEquals(CachedConnection.poolOf(url).meteredCount(), 0, "a refused borrow kept a permit of the pool"); |
| | | } |
| | | |
| | | /** |
| | | * ... and an import whose borrow is interrupted fails instead, keeping the interrupt. |
| | | * <p> |
| | | * The connection this import cannot have is the loss of a parallelism it asked for, and that is |
| | | * what sharing is the answer to. An interrupt is not that: phase two interrupts the threads of |
| | | * an import to stop it ({@code OnDiskMergeImporter.invokeParallel} gives them five seconds to |
| | | * answer), and a thread that answered by writing the tree through another connection would be |
| | | * carrying on with the work it was told to drop. The interrupt goes back on the thread as well: |
| | | * a borrow that throws {@code InterruptedException} has cleared it, and what reads it next is |
| | | * the executor of phase two. |
| | | */ |
| | | @Test(timeOut = 120000) |
| | | public void testAnInterruptedBorrowFailsTheImportRatherThanSharingAConnection() throws Exception { |
| | | System.setProperty(JDBCStorage.IMPORT_CONNECTIONS_PROPERTY, "2"); |
| | | final String url = StubDriver.PREFIX + "importer-interrupted"; |
| | | final JDBCStorage storage = importingStorage(url); |
| | | final Importer importer = storage.startImport(); |
| | | try { |
| | | importing(importer, () -> { |
| | | importer.put(ID2ENTRY, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | |
| | | final AtomicReference<Throwable> failure = new AtomicReference<>(); |
| | | final AtomicBoolean interruptKept = new AtomicBoolean(); |
| | | final Thread writer = daemon(() -> { |
| | | // set before the write rather than raced with it: the pool waits for an idle |
| | | // connection interruptibly, so a thread that carries an interrupt into the borrow |
| | | // gets the same InterruptedException as one interrupted while parked there |
| | | Thread.currentThread().interrupt(); |
| | | try { |
| | | importer.put(DN2ID, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | } catch (Throwable t) { |
| | | failure.set(t); |
| | | } finally { |
| | | interruptKept.set(Thread.currentThread().isInterrupted()); |
| | | } |
| | | }, "import-interrupted"); |
| | | writer.start(); |
| | | writer.join(TimeUnit.SECONDS.toMillis(RENDEZVOUS_SECONDS)); |
| | | assertFalse(writer.isAlive(), "the interrupted write did not finish"); |
| | | |
| | | final Throwable reported = failure.get(); |
| | | assertTrue(reported instanceof StorageRuntimeException, |
| | | "an interrupted borrow was not reported as a failure of the import: " + reported); |
| | | assertTrue(reported.getCause() instanceof InterruptedException, |
| | | "an interrupted borrow was reported as something else: " + reported.getCause()); |
| | | assertTrue(interruptKept.get(), "the borrow swallowed the interrupt of an import thread"); |
| | | }); |
| | | } finally { |
| | | storage.close(); |
| | | } |
| | | } |
| | | |
| | | /** A storage open for writing over the driver of this test, which is what an import needs. */ |
| | | private JDBCStorage importingStorage(String url) throws Exception { |
| | | // mockCfg rather than a bare mock: it answers every getter with the value declared in |
| | | // JDBCBackendConfiguration.xml, so a storage that comes to read a setting this test never |
| | | // thought of gets the default instead of a null - the way the sibling suites of this package |
| | | // build their configurations |
| | | final JDBCBackendCfg cfg = mockCfg(JDBCBackendCfg.class); |
| | | when(cfg.getDBDirectory()).thenReturn(url); |
| | | final JDBCStorage storage = new JDBCStorage(cfg, null); |
| | | storage.open(AccessMode.READ_WRITE); |
| | | // the connection the open borrowed and gave back is not one of the import's: only what a |
| | | // statement was issued on is recorded, and the open issues none |
| | | forgetTheConnections(); |
| | | return storage; |
| | | } |
| | | |
| | | /** |
| | | * Makes the driver of this test answer every further connect with a failure of the given shape, |
| | | * the way a database at a limit of its own does. |
| | | */ |
| | | private void refuseFurtherConnects(final String message, final String sqlState, final int vendorCode) { |
| | | refusal.set(() -> new SQLException(message, sqlState, vendorCode)); |
| | | } |
| | | |
| | | private Future<?> put(final Importer importer, final TreeName treeName) { |
| | | return threads.submit(() -> { |
| | | toTheStartingLine(); |
| | | importer.put(treeName, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | }); |
| | | } |
| | | |
| | | /** A write of a thread that has written nothing before, the way every thread of an import starts. */ |
| | | private void putOnAThreadOfItsOwn(final Importer importer, final TreeName treeName) throws Exception { |
| | | final AtomicReference<Throwable> failure = new AtomicReference<>(); |
| | | final Thread thread = daemon(() -> { |
| | | try { |
| | | importer.put(treeName, ByteString.valueOfUtf8("k"), ByteString.valueOfUtf8("v")); |
| | | } catch (Throwable t) { |
| | | failure.set(t); |
| | | } |
| | | }, "import-" + treeName.getIndexId()); |
| | | thread.start(); |
| | | thread.join(TimeUnit.SECONDS.toMillis(RENDEZVOUS_SECONDS)); |
| | | if (failure.get() != null) { |
| | | throw new IllegalStateException("the write of " + treeName + " failed", failure.get()); |
| | | } |
| | | // asserted rather than left to the join: a write still blocked - on the monitor of a |
| | | // connection, or on a borrow - would otherwise leave the assertions of the test standing on |
| | | // a run that never made the state they are about |
| | | assertFalse(thread.isAlive(), "the write of " + treeName + " did not finish"); |
| | | } |
| | | |
| | | /** Raises what a thread of this test caught, under a name that says which thread it was. */ |
| | | private static void reportFailureOf(String what, AtomicReference<Throwable> failure) { |
| | | final Throwable caught = failure.get(); |
| | | if (caught != null) { |
| | | throw new IllegalStateException(what + " failed", caught); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * A thread of this suite. Daemon, every one of them: a write left standing on the monitor of a |
| | | * connection by a failing run must not keep the JVM of the build alive once the test that |
| | | * started it has been reported. |
| | | */ |
| | | private static Thread daemon(Runnable body, String name) { |
| | | final Thread thread = new Thread(body, name); |
| | | thread.setDaemon(true); |
| | | return thread; |
| | | } |
| | | |
| | | /** The threads a test writes through, so that a stuck one is not a build that never ends. */ |
| | | private static ExecutorService daemonThreads(int count, final String name) { |
| | | final AtomicInteger numbered = new AtomicInteger(); |
| | | return Executors.newFixedThreadPool(count, |
| | | runnable -> daemon(runnable, name + "-" + numbered.incrementAndGet())); |
| | | } |
| | | |
| | | /** The body of a test that writes through an import - see {@link #importing(Importer, ImportBody)}. */ |
| | | private interface ImportBody { |
| | | void run() throws Exception; |
| | | } |
| | | |
| | | /** |
| | | * Runs the body of a test against an import and closes that import afterwards, whatever the |
| | | * body did - without the close having the last word on what went wrong. |
| | | * <p> |
| | | * The close used to stand in a {@code finally} of the test, and it is here instead because the |
| | | * two things that can go wrong there are not disjoint: {@code close()} commits and returns every |
| | | * connection under the monitor of that connection, so a write of this test left standing on one |
| | | * is exactly what makes the close wait - and a close that gave up waiting and threw out of that |
| | | * {@code finally} replaced the failure the test came for with one naming the monitor. Here the |
| | | * body has the first say and the trouble of the close rides along with it as suppressed. |
| | | */ |
| | | private void importing(Importer importer, ImportBody body) throws Exception { |
| | | Throwable found = null; |
| | | try { |
| | | body.run(); |
| | | } catch (Throwable t) { |
| | | found = t; |
| | | } |
| | | final Throwable reported = closeWithoutHanging(importer, found); |
| | | if (reported instanceof Error) { |
| | | throw (Error) reported; |
| | | } |
| | | if (reported instanceof Exception) { |
| | | throw (Exception) reported; |
| | | } |
| | | if (reported != null) { |
| | | throw new IllegalStateException("the import failed to close", reported); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Closes an import without waiting for a straggling thread of the test forever, and returns |
| | | * what is to be reported: what the body of the test found where it found something, with the |
| | | * trouble of the close - a close that never finished, or one that failed - suppressed into it. |
| | | */ |
| | | private Throwable closeWithoutHanging(Importer importer, Throwable found) { |
| | | final AtomicReference<Throwable> failure = new AtomicReference<>(); |
| | | final Thread closing = daemon(() -> { |
| | | try { |
| | | importer.close(); |
| | | } catch (Throwable t) { |
| | | failure.set(t); |
| | | } |
| | | }, "import-close"); |
| | | closing.start(); |
| | | Throwable trouble; |
| | | try { |
| | | closing.join(TimeUnit.SECONDS.toMillis(CLOSE_SECONDS)); |
| | | trouble = closing.isAlive() |
| | | ? new IllegalStateException("the import did not close within " + CLOSE_SECONDS |
| | | + "s: a thread of this test is holding the monitor of one of its connections") |
| | | : failure.get(); |
| | | } catch (InterruptedException e) { |
| | | Thread.currentThread().interrupt(); |
| | | trouble = new IllegalStateException("interrupted while closing the import", e); |
| | | } |
| | | if (found == null) { |
| | | return trouble; |
| | | } |
| | | if (trouble != null) { |
| | | found.addSuppressed(trouble); |
| | | } |
| | | return found; |
| | | } |
| | | |
| | | private List<Connection> connectionsUsed() { |
| | | synchronized (used) { |
| | | return new ArrayList<>(used); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * A connection of this test: it answers every statement with one mock of its own and records |
| | | * the thread that asked for it while that thread is inside the call. |
| | | */ |
| | | private Connection newConnection() throws SQLException { |
| | | final Connection con = org.mockito.Mockito.mock(Connection.class); |
| | | final PreparedStatement statement = org.mockito.Mockito.mock(PreparedStatement.class); |
| | | when(con.isValid(anyInt())).thenReturn(true); |
| | | // the statement names the connection it runs on, as CachedConnection.prepareStatement() has |
| | | // it: the bound of a statement is arbitrated on the connection, which is read off the statement |
| | | when(statement.getConnection()).thenReturn(con); |
| | | // a thread is inside a statement of a connection for longer than the prepareStatement that |
| | | // hands it one: what the sql server driver walks past the end of its array is the list of |
| | | // listeners that the execution and the close of a statement mutate as readily (#891), so |
| | | // those are recorded here too - the window this suite watches is the whole of a statement |
| | | doAnswer(invocation -> { |
| | | enterAStatement(con); |
| | | try { |
| | | return 0; // what a mock of this returns anyway: no row of this test is counted |
| | | } finally { |
| | | leaveAStatement(con); |
| | | } |
| | | }).when(statement).executeUpdate(); |
| | | doAnswer(invocation -> { |
| | | enterAStatement(con); |
| | | try { |
| | | return Boolean.FALSE; // as above: nothing of this test reads a result set |
| | | } finally { |
| | | leaveAStatement(con); |
| | | } |
| | | }).when(statement).execute(); |
| | | doAnswer(invocation -> { |
| | | enterAStatement(con); |
| | | try { |
| | | return null; |
| | | } finally { |
| | | leaveAStatement(con); |
| | | } |
| | | }).when(statement).close(); |
| | | // the commits of the import in the order they happen: which connection is committed when is |
| | | // what keeps a failed import from leaving an index marked trusted over data that is not there |
| | | doAnswer(invocation -> { |
| | | recordAThreadInsideAStatement(con, "commit"); |
| | | commits.add(con); |
| | | return null; |
| | | }).when(con).commit(); |
| | | // the return of a connection rolls it back and hands it to the next borrower, so neither may |
| | | // happen while a statement of an import thread is in flight on it |
| | | doAnswer(invocation -> { |
| | | recordAThreadInsideAStatement(con, "rollback"); |
| | | return null; |
| | | }).when(con).rollback(); |
| | | doAnswer(invocation -> { |
| | | recordAThreadInsideAStatement(con, "close"); |
| | | return null; |
| | | }).when(con).close(); |
| | | when(con.prepareStatement(anyString())).thenAnswer(invocation -> { |
| | | if (used.add(con)) { |
| | | usedInOrder.add(con); |
| | | } |
| | | enterAStatement(con); |
| | | try { |
| | | // the hold outside the monitor of the map above, and after this thread has entered |
| | | // it: what it is waiting for is a peer thread getting in here and recording itself |
| | | parkInsideTheStatement(); |
| | | holdTheConnection(); |
| | | meetPeers(); |
| | | return statement; |
| | | } finally { |
| | | leaveAStatement(con); |
| | | } |
| | | }); |
| | | return con; |
| | | } |
| | | |
| | | /** Records this thread as inside a statement of the given connection, with whoever is already there. */ |
| | | private void enterAStatement(Connection con) { |
| | | synchronized (inStatement) { |
| | | final Set<Thread> inside = inStatement.get(con); |
| | | final Set<Thread> threadsInside = inside != null ? inside : new LinkedHashSet<Thread>(); |
| | | for (final Thread peer : threadsInside) { |
| | | if (peer != Thread.currentThread()) { |
| | | shared.add(peer.getName() + " and " + Thread.currentThread().getName()); |
| | | } |
| | | } |
| | | threadsInside.add(Thread.currentThread()); |
| | | inStatement.put(con, threadsInside); |
| | | } |
| | | } |
| | | |
| | | private void leaveAStatement(Connection con) { |
| | | synchronized (inStatement) { |
| | | final Set<Thread> threadsInside = inStatement.get(con); |
| | | if (threadsInside != null && threadsInside.remove(Thread.currentThread()) && threadsInside.isEmpty()) { |
| | | inStatement.remove(con); |
| | | } |
| | | } |
| | | } |
| | | |
| | | /** Records a connection touched from one thread while another is inside a statement on it. */ |
| | | private void recordAThreadInsideAStatement(Connection con, String what) { |
| | | synchronized (inStatement) { |
| | | final Set<Thread> threadsInside = inStatement.get(con); |
| | | if (threadsInside == null) { |
| | | return; |
| | | } |
| | | for (final Thread inside : threadsInside) { |
| | | if (inside != Thread.currentThread()) { |
| | | shared.add(what + " of " + Thread.currentThread().getName() + " while " + inside.getName() |
| | | + " was inside a statement"); |
| | | } |
| | | } |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Parks a thread inside the first statement it issues and tells the test it is in there, until |
| | | * the test lets it go. |
| | | * <p> |
| | | * A write in flight, as the rest of the import sees one: the connection holds a statement that |
| | | * has not come back, so nothing is recorded against it yet - {@code put()} counts a write once |
| | | * the statement returns - and its monitor is held for as long as this lasts. What a test does |
| | | * from there is ask another thread of the import a question about that connection. |
| | | */ |
| | | private void parkInsideTheStatement() { |
| | | final CountDownLatch letGo = letTheStatementFinish; |
| | | if (letGo == null || parked.get()) { |
| | | return; |
| | | } |
| | | parked.set(Boolean.TRUE); |
| | | insideAStatementNow.countDown(); |
| | | try { |
| | | letGo.await(PARK_SECONDS, TimeUnit.SECONDS); |
| | | } catch (InterruptedException e) { |
| | | Thread.currentThread().interrupt(); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Waits until a thread of this test has got as far as it is going to get on its own: either it |
| | | * is blocked on a monitor - the thing this suite is about - or it is done. |
| | | * <p> |
| | | * Waited for rather than slept past: what the tests below ask is whether a thread stops at a |
| | | * monitor or walks through it, and a fixed pause would answer that with the load of the machine |
| | | * on a slow run. |
| | | */ |
| | | private static void untilBlockedOrDone(Thread thread) throws InterruptedException { |
| | | final long deadline = System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(RENDEZVOUS_SECONDS); |
| | | while (System.currentTimeMillis() < deadline) { |
| | | final Thread.State state = thread.getState(); |
| | | if (state == Thread.State.BLOCKED || state == Thread.State.TERMINATED) { |
| | | return; |
| | | } |
| | | Thread.sleep(10); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Holds the first statement of a thread on the connection it is on, until a peer thread has one |
| | | * of its own in flight there or the hold runs out. |
| | | * <p> |
| | | * How a test asks whether two threads can be inside a statement of one connection at the same |
| | | * time: they cannot be brought together at a barrier the way the threads of two connections are |
| | | * - a second thread held on the monitor of the connection never reaches one - so the first |
| | | * thread in waits there instead, and the peer either turns up beside it, which is #891, or does |
| | | * not, which is the hold being paid in full. |
| | | */ |
| | | private void holdTheConnection() { |
| | | final CountDownLatch inside = insideAStatement; |
| | | if (inside == null || held.get()) { |
| | | return; |
| | | } |
| | | held.set(Boolean.TRUE); |
| | | inside.countDown(); |
| | | try { |
| | | inside.await(HOLD_SECONDS, TimeUnit.SECONDS); |
| | | } catch (InterruptedException e) { |
| | | Thread.currentThread().interrupt(); |
| | | } |
| | | } |
| | | |
| | | /** Where the threads of a test start together, in front of the monitors of the importer. */ |
| | | private void toTheStartingLine() { |
| | | final CyclicBarrier line = startingLine; |
| | | if (line == null) { |
| | | return; |
| | | } |
| | | try { |
| | | line.await(RENDEZVOUS_SECONDS, TimeUnit.SECONDS); |
| | | } catch (InterruptedException | TimeoutException | BrokenBarrierException e) { |
| | | throw new IllegalStateException("the threads of this test did not reach the starting line", e); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Holds the first statement of a thread until every peer of the rendezvous has one in flight |
| | | * too. Reported as a failure of the statement rather than waited out: a rendezvous nobody else |
| | | * reaches is an import that serialized its threads, which is what the parallel test is about. |
| | | */ |
| | | private void meetPeers() throws SQLException { |
| | | final CyclicBarrier barrier = rendezvous; |
| | | if (barrier == null || met.get()) { |
| | | return; |
| | | } |
| | | met.set(Boolean.TRUE); |
| | | try { |
| | | barrier.await(RENDEZVOUS_SECONDS, TimeUnit.SECONDS); |
| | | } catch (TimeoutException | BrokenBarrierException e) { |
| | | throw new SQLException("no peer thread had a statement of its own in flight within " |
| | | + RENDEZVOUS_SECONDS + "s: the import did not let its threads write at the same time", e); |
| | | } catch (InterruptedException e) { |
| | | Thread.currentThread().interrupt(); |
| | | throw new SQLException("interrupted while waiting for the peer threads of the import", e); |
| | | } |
| | | } |
| | | |
| | | /** A driver of this test, so that the connections of an import need no database behind them. */ |
| | | private final class StubDriver implements Driver { |
| | | static final String PREFIX = "jdbc:opendj-import-stub:"; |
| | | |
| | | @Override |
| | | public Connection connect(String url, Properties info) throws SQLException { |
| | | if (!acceptsURL(url)) { |
| | | return null; |
| | | } |
| | | final Supplier<SQLException> refused = refusal.get(); |
| | | if (refused != null) { |
| | | throw refused.get(); |
| | | } |
| | | return newConnection(); |
| | | } |
| | | |
| | | @Override |
| | | public boolean acceptsURL(String url) { |
| | | return url != null && url.startsWith(PREFIX); |
| | | } |
| | | |
| | | @Override |
| | | public DriverPropertyInfo[] getPropertyInfo(String url, Properties info) { |
| | | return new DriverPropertyInfo[0]; |
| | | } |
| | | |
| | | @Override |
| | | public int getMajorVersion() { |
| | | return 1; |
| | | } |
| | | |
| | | @Override |
| | | public int getMinorVersion() { |
| | | return 0; |
| | | } |
| | | |
| | | @Override |
| | | public boolean jdbcCompliant() { |
| | | return false; |
| | | } |
| | | |
| | | @Override |
| | | public Logger getParentLogger() { |
| | | return Logger.getLogger(StubDriver.class.getName()); |
| | | } |
| | | } |
| | | } |