| | |
| | | import java.sql.*; |
| | | import java.util.*; |
| | | import java.util.concurrent.ConcurrentHashMap; |
| | | import java.util.concurrent.Executor; |
| | | import java.util.concurrent.atomic.AtomicBoolean; |
| | | import java.util.function.Predicate; |
| | | |
| | | import static org.opends.server.backends.pluggable.spi.StorageUtils.addErrorMessage; |
| | |
| | | return ccr; |
| | | } |
| | | |
| | | ResultSet executeResultSet(PreparedStatement statement) throws SQLException { |
| | | /** |
| | | * What a statement of this backend may legitimately take, and the property bounding it. One |
| | | * value cannot serve both: an entry read is a single row of an index, while the count of a |
| | | * tree and the delete that empties one before an import are a scan and a rewrite of a whole |
| | | * table, which take minutes on a populated backend and are not a symptom of anything. |
| | | */ |
| | | enum StatementBound { |
| | | /** one row by primary key, or one batch of a cursor along its index */ |
| | | OPERATION("org.openidentityplatform.opendj.jdbc.query.timeout", 120), |
| | | /** |
| | | * a whole table at once: count(*), the delete of clearTree, the scan behind the highest |
| | | * entry id, create index, drop table, and every batch of a cursor walking a tree whole. |
| | | * This class ships <em>unbounded</em>: what such a statement legitimately takes follows the |
| | | * size of the backend and the speed of its database, neither of which can be guessed here, |
| | | * so the deployment that knows both sets the property - until it does, a create index |
| | | * waiting for a metadata lock still waits for as long as the engine lets it, and so do the |
| | | * walks a backend makes while it opens (the load of the compressed schema, the read that |
| | | * checks id2entry is there) and the export behind the generation ID of a replicated domain. |
| | | * That is what this backend did before any of these bounds existed; bounding them as the |
| | | * work of a client operation, which is the only other value there was to give them, stopped |
| | | * a large backend from opening at all. |
| | | */ |
| | | BULK("org.openidentityplatform.opendj.jdbc.bulk.timeout", 0); |
| | | |
| | | final String property; |
| | | final int defaultSeconds; |
| | | |
| | | StatementBound(String property, int defaultSeconds) { |
| | | this.property = property; |
| | | this.defaultSeconds = defaultSeconds; |
| | | } |
| | | |
| | | /** |
| | | * The bound in seconds, as configured by {@link #property}: 0, or a negative value, leaves |
| | | * the statement unbounded, as it was before this bound existed, while a value that is not a |
| | | * number is ignored in favour of {@link #defaultSeconds} - {@code Integer.getInteger()} |
| | | * falls back to its default rather than reading such a value as a zero. A value above |
| | | * {@link JDBCStorage#MAX_BOUND_SECONDS} is taken down to it, for the reason recorded there. |
| | | */ |
| | | int seconds() { |
| | | return clampSeconds(Integer.getInteger(property, defaultSeconds)); |
| | | } |
| | | } |
| | | |
| | | /** What a caller of {@link #executeResultSet} makes of the rows, while the bound is still armed. */ |
| | | interface RowsHandler<T> { |
| | | T handle(ResultSet rows) throws SQLException; |
| | | } |
| | | |
| | | /** |
| | | * The value of a row that is there. A row whose {@code v} is null is one this backend never |
| | | * wrote - the column is nullable, however it is written - and it must not be answered with the |
| | | * {@code null} a single-row read uses, which is already taken and means "no such key": read that |
| | | * way, a key that exists is reported as absent. Named rather than left to the bare |
| | | * {@code NullPointerException} of {@code ByteString.wrap}, which names neither the fault nor the |
| | | * table it is in, and a {@code RuntimeException} rather than an {@code SQLException}, so that a |
| | | * corrupt row is never weighed against the bound of the statement that read it and reported as a |
| | | * timeout of a property that would have changed nothing. The key is left out of the message for |
| | | * the reason {@link #timedOut} leaves the statement out of its own: it is entry data. |
| | | */ |
| | | static ByteString valueOfRow(ResultSet rows, String tableName) throws SQLException { |
| | | return ByteString.wrap(valueOfRow(rows.getBytes("v"), tableName)); |
| | | } |
| | | |
| | | /** |
| | | * The same check where a batch of a cursor reads the value beside its key, by position. Checked |
| | | * as the rows are taken off the statement rather than as they are handed out one by one: there |
| | | * the failure is inside the bound and inside the {@code catch} of the batch, while a batch |
| | | * buffered whole and unwrapped later fails from {@code advanceFromBuffer()} - outside both, and |
| | | * as the bare {@code NullPointerException} this exists to replace. |
| | | */ |
| | | static byte[] valueOfRow(byte[] value, String tableName) { |
| | | if (value == null) { |
| | | throw new StorageRuntimeException("jdbc: a row of "+tableName+" is present with no value"); |
| | | } |
| | | return value; |
| | | } |
| | | |
| | | <T> T executeResultSet(PreparedStatement statement, RowsHandler<T> rows) throws SQLException { |
| | | return executeResultSet(statement, StatementBound.OPERATION, rows); |
| | | } |
| | | |
| | | /** |
| | | * Runs a query under the bound of its class and hands the rows to {@code rows} while that bound |
| | | * is still armed. They are read there rather than after this method returns because a driver |
| | | * transfers them as they are asked for: read outside, the transfer - up to a whole batch of a |
| | | * cursor - would run with neither layer of the bound covering it, which is exactly where a |
| | | * database that stops answering mid-drain parks the worker thread. {@code setQueryTimeout} |
| | | * covering {@code ResultSet.next()} is optional in the JDBC contract ("drivers <em>may</em> |
| | | * also apply this limit"), and the two drivers of this backend that do not buffer a result |
| | | * whole - oracle prefetches ten rows at a time, mssql buffers adaptively - are the ones that |
| | | * do not. |
| | | */ |
| | | <T> T executeResultSet(PreparedStatement statement, StatementBound bound, RowsHandler<T> rows) throws SQLException { |
| | | if (logger.isTraceEnabled()) { |
| | | logger.trace(LocalizableMessage.raw("jdbc: %s",statement)); |
| | | } |
| | | return statement.executeQuery(); |
| | | return bounded(statement, bound, () -> { |
| | | try (final ResultSet rs=statement.executeQuery()) { |
| | | return rows.handle(rs); |
| | | } |
| | | }); |
| | | } |
| | | |
| | | int execute(PreparedStatement statement) throws SQLException { |
| | | return execute(statement, StatementBound.OPERATION); |
| | | } |
| | | |
| | | int execute(PreparedStatement statement, StatementBound bound) throws SQLException { |
| | | if (logger.isTraceEnabled()) { |
| | | logger.trace(LocalizableMessage.raw("jdbc: %s",statement)); |
| | | } |
| | | return statement.executeUpdate(); |
| | | return bounded(statement, bound, statement::executeUpdate); |
| | | } |
| | | |
| | | // unlike execute(), tolerates statements that return a result set ("analyze table" on mysql) |
| | | interface Execution<T> { |
| | | T run() throws SQLException; |
| | | } |
| | | |
| | | /** |
| | | * Runs a statement under the bound of its class. A statement of a class that carries one has to |
| | | * end: a row locked by an unrelated session, a table waiting for a metadata lock or a database |
| | | * that stops answering mid-query would otherwise park the worker thread that issued it for |
| | | * good. A class configured with no bound - which {@link StatementBound#BULK} ships as - takes |
| | | * neither of the two layers below and waits as this backend waited before they existed. |
| | | * <p> |
| | | * The bound is asked of the driver rather than of the session, because a pooled connection |
| | | * cannot carry a session setting - {@code CachedConnection.close()} only rolls back, so a |
| | | * {@code statement_timeout} of one operation would apply to whoever borrows the connection |
| | | * next - and it is applied in two layers, since the first one is not answered everywhere: |
| | | * {@code setQueryTimeout} cancels the statement and keeps the connection, while the socket read |
| | | * timeout behind it ends the wait even when the cancel is not acted upon. Oracle needs that |
| | | * second layer: a session blocked in a row-lock enqueue does not process the break its driver |
| | | * sends, so the timeout is armed and never arrives (the container suites cover it). That second |
| | | * layer belongs to the connection rather than to the statement, so it is arbitrated between the |
| | | * statements running on one - see {@link Backstop}. |
| | | */ |
| | | private <T> T bounded(PreparedStatement statement, StatementBound bound, Execution<T> execution) throws SQLException { |
| | | final int seconds=bound.seconds(); |
| | | // whether the cancel is in force: a driver is free to refuse the query timeout, and then the |
| | | // socket read timeout behind it is the only layer this statement has - one that arrives later |
| | | final boolean cancelArmed=seconds > 0 && setQueryTimeout(statement, seconds); |
| | | // an unbounded class is announced to the connection all the same: a statement told it may |
| | | // take as long as it needs must not be cut by the socket read timeout of a concurrent one |
| | | return bounded(connectionOf(statement), bound.property, seconds, cancelArmed, execution); |
| | | } |
| | | |
| | | /** |
| | | * Runs the catalog lookups of {@code openTree()} under the bound of their class. They ask |
| | | * {@code DatabaseMetaData}, which takes no query timeout, so the socket read timeout behind the |
| | | * cancel is the only layer they can be given - and they do need one: they run once per tree on |
| | | * every open of a backend, and the catalog is answered by the same engine, behind the same |
| | | * locks, as the {@code create table} they guard. |
| | | * <p> |
| | | * That layer is only as good as what it actually arms, which is not always something: a driver |
| | | * with no network timeout, a connection that failed the call, one already carrying a tighter |
| | | * timeout of a deployment's own, and a statement of an unbounded class running beside this one |
| | | * each leave such a lookup with no bound at all. It is then reported as what it is - see |
| | | * {@link #timedOut} - rather than as a property that bounded nothing. |
| | | */ |
| | | <T> T bounded(Connection con, StatementBound bound, Execution<T> execution) throws SQLException { |
| | | // no cancel to arm: DatabaseMetaData takes no query timeout, so the socket read timeout behind |
| | | // it is the only layer these have, and nothing ends their wait before the margin of that layer |
| | | return bounded(con, bound.property, bound.seconds(), false, execution); |
| | | } |
| | | |
| | | /** |
| | | * Runs a statement under a bound of its own rather than under the bound of a class, for the one |
| | | * statement that has a property of its own: the statistics refresh after an import, which |
| | | * legitimately takes as long as a scan of the table it describes. |
| | | */ |
| | | private <T> T bounded(Connection con, String property, int seconds, boolean cancelArmed, Execution<T> execution) |
| | | throws SQLException { |
| | | final long startedAt=nanoTime(); |
| | | final Backstop backstop=holdBackstop(con, seconds); |
| | | try { |
| | | return execution.run(); |
| | | }catch (SQLException e) { |
| | | // what the second layer carries is read here rather than at the top: it is arbitrated |
| | | // between the statements in flight, so it is the value at the moment of the failure that |
| | | // bounded this statement - and it is read before the release below takes it back off |
| | | throw timedOut(e, property, seconds, cancelArmed, armedMillis(backstop), startedAt); |
| | | }finally { |
| | | releaseBackstop(backstop, con, seconds); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * What the socket read timeout of a connection carries for the statements on it right now, or 0 |
| | | * where this layer is not in force for them at all. It is not enough that a bound was asked for: |
| | | * {@link #applyBackstop} arms nothing on a connection whose driver refused the call or has no |
| | | * network timeout to give, nothing on one already carrying a timeout of a deployment's own that |
| | | * is tighter than ours, and nothing while a statement of an unbounded class runs beside this one. |
| | | */ |
| | | private static int armedMillis(Backstop state) { |
| | | if (state == null) { |
| | | return 0; // no connection to arm it on: the cancel is the whole bound of such a statement |
| | | } |
| | | synchronized (state) { |
| | | return state.armed; |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Asks the driver to cancel the statement at the bound. Not every driver has one: the JDBC |
| | | * contract allows {@code SQLFeatureNotSupportedException} and this backend takes whatever URL a |
| | | * deployment configures, so a driver without it degrades to the socket read timeout behind it |
| | | * rather than failing every statement it is given. |
| | | */ |
| | | private boolean setQueryTimeout(PreparedStatement statement, int seconds) { |
| | | try { |
| | | statement.setQueryTimeout(seconds); |
| | | return true; |
| | | }catch (SQLException | RuntimeException e) { |
| | | if (queryTimeoutWarned.compareAndSet(false, true)) { |
| | | logger.warn(LocalizableMessage.raw("jdbc: the driver would not take a query timeout (%s): a statement of this" |
| | | + " backend is left to the socket read timeout behind it", e.getMessage())); |
| | | } |
| | | return false; |
| | | } |
| | | } |
| | | |
| | | private Connection connectionOf(PreparedStatement statement) { |
| | | try { |
| | | return statement.getConnection(); |
| | | }catch (SQLException | RuntimeException e) { |
| | | return null; // nothing to arm the backstop on; the cancel above is the whole bound |
| | | } |
| | | } |
| | | |
| | | /** How long the socket read timeout outlasts the cancel it backs up, giving it room to arrive. */ |
| | | static final int BACKSTOP_MARGIN_SECONDS = 30; |
| | | |
| | | /** How far under its bound a driver may report the cancel, its timer being kept in whole seconds. */ |
| | | static final long CLOCK_SLACK_MILLIS = 250; |
| | | |
| | | /** |
| | | * Ceiling of every bound this backend arms, in seconds - 24.9 days, which is what a socket read |
| | | * timeout can hold at all: {@code setNetworkTimeout} takes milliseconds of an {@code int}, and a |
| | | * bound past this one has no value of that layer to be given. It is <em>not</em> what keeps the |
| | | * arithmetic of {@link #backstopMillis} in range - the {@code long} multiply under the |
| | | * {@code Math.min} there does that on its own, up to the point where adding the margin overflows |
| | | * an {@code int} before the multiply ever runs - so a reader who later takes that {@code Math.min} |
| | | * away must not read this clamp as covering them. |
| | | * <p> |
| | | * Clamped rather than refused, and clamped rather than read as "no bound": a bound this large |
| | | * cancels nothing a database will not have ended first, so taking a nonsensical value down to it |
| | | * costs a deployment nothing, while reading it as an unbound would take a bound away from a |
| | | * deployment that asked for one. A property set to {@code Integer.MAX_VALUE} therefore bounds a |
| | | * statement at 24.9 days rather than leaving it unbounded; {@code 0} is what leaves it unbounded. |
| | | */ |
| | | static final int MAX_BOUND_SECONDS = Integer.MAX_VALUE/1000 - BACKSTOP_MARGIN_SECONDS; |
| | | |
| | | static int clampSeconds(int seconds) { |
| | | return Math.max(0, Math.min(MAX_BOUND_SECONDS, seconds)); |
| | | } |
| | | |
| | | /** |
| | | * What {@link #timedOut} calls the second layer when that layer is the only one a statement ran |
| | | * under, so that a test can tell the two apart in a message: a run where the first layer stopped |
| | | * working degrades to this one by design, silently, and a suite that only measures how long a |
| | | * statement waited would go green with the cancel gone entirely. |
| | | */ |
| | | static final String BACKSTOP_ALONE = "the socket read timeout behind "; |
| | | |
| | | // The clock a bound is measured on, in one place so that a test can drive it: the classification |
| | | // below turns on a few milliseconds either side of the bound, and a mock statement cannot be made |
| | | // to take a real second without the suite taking one too. Monotonic, so that a step of the wall |
| | | // clock can neither lengthen nor shorten what a statement is measured to have taken. |
| | | long nanoTime() { |
| | | return System.nanoTime(); |
| | | } |
| | | |
| | | // setNetworkTimeout() takes the executor its timeout handling runs on; the drivers of this |
| | | // backend only set a socket option in it, so it costs a call rather than a thread. |
| | | private static final Executor DIRECT_EXECUTOR = Runnable::run; |
| | | |
| | | // Set when the driver of this storage has no network timeout to give at all, which is a property |
| | | // of the driver rather than of a connection: asking it again would cost a throw per statement, |
| | | // and the entry a connection's Backstop lives in is gone as soon as nothing runs on it. Held per |
| | | // storage rather than per JVM, like the warnings below: a driver that will not take one of these |
| | | // says so once for every backend running on it, instead of one backend silencing it for all. |
| | | private final AtomicBoolean backstopUnsupported = new AtomicBoolean(); |
| | | private final AtomicBoolean backstopUnsupportedWarned = new AtomicBoolean(); |
| | | private final AtomicBoolean backstopFailedWarned = new AtomicBoolean(); |
| | | private final AtomicBoolean queryTimeoutWarned = new AtomicBoolean(); |
| | | |
| | | /** |
| | | * The socket read timeout of one connection, and the statements running on it. This second |
| | | * layer of the bound is a property of the socket rather than of a statement, so it cannot be |
| | | * armed and put back per statement wherever a connection carries more than one at a time: an |
| | | * {@code ImporterImpl} holds a single connection for the whole of an import and writes to it |
| | | * from every phase-one worker and every phase-two task, and there the first statement to finish |
| | | * would take the backstop away from every statement still in flight - while a statement whose |
| | | * class carries no bound at all would run under whatever value a concurrent one happened to |
| | | * arm, dying at it with nothing to say which property cut it, since such a statement never |
| | | * reaches {@link #timedOut}. |
| | | * <p> |
| | | * So the value armed is the loosest of the bounds of the statements in flight, and a statement |
| | | * with no bound of its own takes it off for as long as it runs: this backstop exists to end a |
| | | * wait nothing else would end, never to cut a statement that was told it may take as long as it |
| | | * needs. What the connection carried before is put back when the last of them is through. |
| | | */ |
| | | private static final class Backstop { |
| | | /** Bounds of the statements in flight, in milliseconds and by count, the loosest last. */ |
| | | final TreeMap<Integer,Integer> bounds=new TreeMap<>(); |
| | | /** Statements in flight with no bound of their own, which no backstop may cut short. */ |
| | | int unbounded; |
| | | /** Statements holding this entry, bounded or not: at zero it leaves {@link #backstops}. */ |
| | | int holders; |
| | | /** What the connection carried before the backstop armed it, and is given back afterwards. */ |
| | | int previous; |
| | | /** What the backstop has armed, or 0 when the connection carries {@link #previous}. */ |
| | | int armed; |
| | | /** |
| | | * Set when the driver would not take a network timeout on this connection: it is not asked |
| | | * again while the statements holding this entry run. A connection is the right scope for |
| | | * that: the common cause is a connection on its way out, and a driver that has no network |
| | | * timeout at all is remembered for the whole storage instead - see {@link #backstopUnsupported}. |
| | | */ |
| | | boolean failed; |
| | | } |
| | | |
| | | // Keyed by identity on the connection of the driver: CachedConnection.prepareStatement() hands |
| | | // the statement to the connection it wraps, so that is the one a statement reports, while the |
| | | // catalog lookups above hold the wrapper of that same connection - both have to find the same |
| | | // entry, so a wrapper is unwrapped on the way in. Static because the pool these connections |
| | | // come from is static; an entry lives only while statements are running on its connection. |
| | | private static final Map<Connection,Backstop> backstops = new IdentityHashMap<>(); |
| | | |
| | | private static Connection physical(Connection con) { |
| | | return con instanceof CachedConnection ? ((CachedConnection)con).parent : con; |
| | | } |
| | | |
| | | /** |
| | | * Puts the bound of a statement about to run on the connection that will run it, and makes the |
| | | * socket read timeout of that connection fit every statement in flight on it. Reaching this |
| | | * bound, unlike reaching the cancel it backs up, costs the connection: the driver closes it, |
| | | * which is the price of a wait the database was never going to end on its own. |
| | | */ |
| | | private Backstop holdBackstop(Connection con, int seconds) { |
| | | final Connection physical=physical(con); |
| | | if (physical == null) { |
| | | return null; |
| | | } |
| | | final Backstop state; |
| | | synchronized (backstops) { |
| | | state=backstops.computeIfAbsent(physical, c -> new Backstop()); |
| | | state.holders++; // held from here, so that the entry outlives a concurrent release |
| | | } |
| | | synchronized (state) { |
| | | if (seconds > 0) { |
| | | state.bounds.merge(backstopMillis(seconds), 1, Integer::sum); |
| | | }else { |
| | | state.unbounded++; |
| | | } |
| | | applyBackstop(physical, state); |
| | | } |
| | | return state; |
| | | } |
| | | |
| | | private void releaseBackstop(Backstop state, Connection con, int seconds) { |
| | | if (state == null) { |
| | | return; |
| | | } |
| | | final Connection physical=physical(con); |
| | | try { |
| | | synchronized (state) { |
| | | if (seconds > 0) { |
| | | final int millis=backstopMillis(seconds); |
| | | final Integer inFlight=state.bounds.get(millis); |
| | | if (inFlight == null || inFlight <= 1) { |
| | | state.bounds.remove(millis); |
| | | }else { |
| | | state.bounds.put(millis, inFlight-1); |
| | | } |
| | | }else { |
| | | state.unbounded--; |
| | | } |
| | | applyBackstop(physical, state); |
| | | } |
| | | }finally { // the entry is let go whatever the driver did, so that it cannot outlive its connection |
| | | synchronized (backstops) { |
| | | if (--state.holders <= 0) { // nothing is running on it: the connection is on its own again |
| | | backstops.remove(physical); |
| | | } |
| | | } |
| | | } |
| | | } |
| | | |
| | | private static int backstopMillis(int seconds) { |
| | | return (int) Math.min(Integer.MAX_VALUE, (seconds+BACKSTOP_MARGIN_SECONDS)*1000L); |
| | | } |
| | | |
| | | /** |
| | | * Makes the socket read timeout of the connection what the statements in flight on it need: the |
| | | * loosest of their bounds, or nothing of ours at all while one of them carries no bound. Called |
| | | * with the monitor of {@code state} held, since it both reads those counts and acts on the |
| | | * driver. |
| | | */ |
| | | private void applyBackstop(Connection con, Backstop state) { |
| | | if (state.failed || backstopUnsupported.get()) { |
| | | // but a connection this backstop has already armed does not keep carrying it: the entry |
| | | // remembering what it carried before is dropped when its last statement is through, and the |
| | | // value would go back to the pool as the connection's own read timeout. Reachable through |
| | | // the second guard, which is a latch of the whole storage: a connection armed before it was |
| | | // set would otherwise never be disarmed. Where nothing was armed this costs no call. |
| | | restorePrevious(con, state); |
| | | return; |
| | | } |
| | | final int wanted=state.unbounded > 0 || state.bounds.isEmpty() ? 0 : state.bounds.lastKey(); |
| | | try { |
| | | if (wanted == 0) { |
| | | if (state.armed != 0) { |
| | | con.setNetworkTimeout(DIRECT_EXECUTOR, state.previous); |
| | | state.armed=0; |
| | | } |
| | | return; |
| | | } |
| | | if (state.armed == 0) { |
| | | state.previous=con.getNetworkTimeout(); |
| | | } |
| | | // only ever tighten: a connection that already carries a read timeout carries one a |
| | | // deployment asked for, and this backstop exists to cap a cancel that is not acted |
| | | // upon, not to relax anything. 0 is "no timeout" in the JDBC contract, so it is the |
| | | // one value there is always something to gain by replacing. |
| | | if (state.previous > 0 && state.previous <= wanted) { |
| | | if (state.armed != 0) { |
| | | con.setNetworkTimeout(DIRECT_EXECUTOR, state.previous); |
| | | state.armed=0; |
| | | } |
| | | return; |
| | | } |
| | | if (state.armed != wanted) { |
| | | con.setNetworkTimeout(DIRECT_EXECUTOR, wanted); |
| | | state.armed=wanted; |
| | | } |
| | | }catch (SQLException | RuntimeException e) { |
| | | state.failed=true; // whatever the cause, this connection is not asked again while it runs |
| | | // and what it carried before goes back, while there is still an entry saying what that was: |
| | | // this one is dropped as soon as the last statement on the connection is through, and a |
| | | // backstop left armed would go back to the pool as the connection's own read timeout - which |
| | | // is exactly how the next borrower reads it, tightening to it and never replacing it. |
| | | restorePrevious(con, state); |
| | | // The two causes are told apart, because they deserve opposite treatment and one of them |
| | | // would otherwise spend the single warning the other needs: a driver with no network |
| | | // timeout at all says so through SQLFeatureNotSupportedException, and there is nothing to |
| | | // gain by asking it once per statement for the life of the storage, while a connection on |
| | | // its way out - it may be the one that reached this very timeout - says nothing about the |
| | | // driver and must not disable the backstop for the connections that are still healthy. |
| | | if (e instanceof SQLFeatureNotSupportedException) { |
| | | backstopUnsupported.set(true); |
| | | if (backstopUnsupportedWarned.compareAndSet(false, true)) { |
| | | logger.warn(LocalizableMessage.raw("jdbc: the driver takes no socket read timeout (%s): a statement the" |
| | | + " database does not cancel will wait for it indefinitely, unless the connect properties of the URL" |
| | | + " configured for this backend carry one", e.getMessage())); |
| | | } |
| | | }else if (backstopFailedWarned.compareAndSet(false, true)) { |
| | | logger.warn(LocalizableMessage.raw("jdbc: the socket read timeout backing up a cancelled statement could not" |
| | | + " be set on a connection (%s): a statement the database does not cancel will wait for it" |
| | | + " indefinitely there", e.getMessage())); |
| | | } |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Gives the connection back the read timeout it carried before this backstop armed one, and |
| | | * forgets having armed it. Best effort by construction: the caller reaches this from a driver |
| | | * call that has just failed, so the connection may well be gone - and where it is, it is the |
| | | * driver that closes it rather than this backend. |
| | | */ |
| | | private static void restorePrevious(Connection con, Backstop state) { |
| | | if (state.armed == 0) { |
| | | return; // the connection carries its own value already |
| | | } |
| | | try { |
| | | con.setNetworkTimeout(DIRECT_EXECUTOR, state.previous); |
| | | }catch (SQLException | RuntimeException ignored) { |
| | | // nothing further can be done for this connection here, and the failure to report is the |
| | | // one that brought us into the catch above |
| | | }finally { |
| | | state.armed=0; |
| | | } |
| | | } |
| | | |
| | | // Every driver reports a cancelled statement differently - postgresql as 57014, oracle as |
| | | // ORA-01013, and neither of them as a SQLTimeoutException - so the bound is recognized by the |
| | | // time the statement took rather than by the class or the state of its failure, and that time |
| | | // is taken from the monotonic clock, which a step of the wall clock can neither lengthen nor |
| | | // shorten. What it cannot tell apart is a failure of another kind arriving after the bound, |
| | | // which is why the failure it replaces is chained rather than swallowed. The SQL state and the |
| | | // error number are carried over as well, since a failure that arrives at the bound may still be |
| | | // one a caller classifies: a mysql lock wait, reported in class 40, ends inside a longer bound |
| | | // and stays the replayable conflict it is. The statement itself is left out of the message: a |
| | | // driver renders it with its parameters bound, and those are entry data. |
| | | private SQLException timedOut(SQLException e, String property, int seconds, boolean cancelArmed, |
| | | int backstopArmedMillis, long startedAt) { |
| | | if (seconds <= 0) { |
| | | return e; |
| | | } |
| | | // Which layer was really in force, and until when. Where the cancel is armed, the property |
| | | // ends the wait at its own value. Where it is not - a statement of DatabaseMetaData takes no |
| | | // query timeout, and a driver is free to refuse one - the socket read timeout behind it is the |
| | | // only layer there is, and that one arrives a margin later: measuring such a statement against |
| | | // the property alone reported a connection reset at 121 s as a query timeout of 120 s and sent |
| | | // the operator to a property that bounded nothing. Asking for that layer is not having it, |
| | | // which is why the value armed is passed in rather than derived from the property here: a |
| | | // driver with no network timeout, a connection that failed the call, one already carrying a |
| | | // tighter timeout of its own, and a statement of an unbounded class running beside this one |
| | | // each leave it unarmed. A statement neither layer bounded reached no bound of ours at all, so |
| | | // its failure is the driver's own and is left exactly as it is: naming a property that armed |
| | | // nothing sends an operator to raise a value that changes nothing about the wait they saw. |
| | | final long endsAfterMillis=cancelArmed ? seconds*1000L : backstopArmedMillis; |
| | | if (endsAfterMillis <= 0) { |
| | | return e; |
| | | } |
| | | final long elapsedMillis=(nanoTime()-startedAt)/1000000L; |
| | | // The bound is allowed a little slack under it: a driver keeps its timer in whole seconds and |
| | | // reports the cancel a few milliseconds before the bound is arithmetically due, and measured |
| | | // to the millisecond such a statement would arrive as a bare 57014 or ORA-01013, naming |
| | | // neither the property that cancelled it nor the fact that it was cancelled at all. |
| | | if (elapsedMillis < endsAfterMillis-CLOCK_SLACK_MILLIS) { |
| | | return e; |
| | | } |
| | | // The time is reported as measured rather than as the bound. Where the database does not act |
| | | // on the cancel - a session blocked in a row-lock enqueue on oracle - the wait ends at the |
| | | // socket read timeout, a margin past the property that armed it, and "did not finish within |
| | | // the 120s" of a statement that waited 150 s is a message an operator cannot put next to a |
| | | // clock. The property named still governs both layers, since backstopMillis() derives the |
| | | // second one from it, so raising it stays the remedy either way. |
| | | return new SQLTimeoutException("jdbc: the statement took "+elapsedMillis+" ms, reaching the " |
| | | +(endsAfterMillis/1000L)+"s of " |
| | | +(cancelArmed ? property : BACKSTOP_ALONE+property+" (the only layer bounding a statement that takes no" |
| | | +" query timeout; it is armed at the loosest bound of the statements sharing this connection, this" |
| | | +" one's being "+seconds+"s plus the margin of that layer)") |
| | | +": raise that property, or set it to 0 for no bound", e.getSQLState(), e.getErrorCode(), e); |
| | | } |
| | | |
| | | // Unlike execute(), tolerates a statement that returns a result set - the comment statement of |
| | | // mssql is a batch that ends in an exec - and, unlike it, carries no bound of its own: what is |
| | | // left of this method runs on a stamp connection, which is given a lock timeout of its own |
| | | // (Dialect.lockTimeoutSql) and a socket read timeout in its connect properties. |
| | | void executeAny(PreparedStatement statement) throws SQLException { |
| | | if (logger.isTraceEnabled()) { |
| | | logger.trace(LocalizableMessage.raw("jdbc: %s",statement)); |
| | |
| | | } |
| | | |
| | | Connection getConnection() throws Exception { |
| | | return CachedConnection.getConnection(config.getDBDirectory()); |
| | | return getConnection(true); |
| | | } |
| | | |
| | | /** |
| | |
| | | * per import or per removal buys back exactly what master did on every borrow. |
| | | */ |
| | | Connection getValidatedConnection() throws Exception { |
| | | return CachedConnection.getConnection(config.getDBDirectory(), false); |
| | | return getConnection(false); |
| | | } |
| | | |
| | | // The one borrow of this storage: both methods above go through it, so that whatever stands in |
| | | // for the pool stands in for every path that takes a connection. A stand-in of the trusted |
| | | // borrow alone let the open, the import and the removal - the three that ask for a validated |
| | | // one - reach a real database instead. |
| | | Connection getConnection(boolean trusted) throws Exception { |
| | | return CachedConnection.getConnection(config.getDBDirectory(), trusted); |
| | | } |
| | | |
| | | |
| | |
| | | // Asked of the very session that parses the literal: sql_mode is a session setting, and a |
| | | // session opened at another moment can have been given another value of it. |
| | | boolean isMysqlBackslashEscape(Connection con) throws SQLException { |
| | | try (final PreparedStatement statement=con.prepareStatement("select @@sql_mode"); |
| | | final ResultSet rs=executeResultSet(statement)) { |
| | | final String sqlMode=rs.next() ? rs.getString(1) : null; |
| | | try (final PreparedStatement statement=con.prepareStatement("select @@sql_mode")) { |
| | | final String sqlMode=executeResultSet(statement, rs -> rs.next() ? rs.getString(1) : null); |
| | | return sqlMode==null || !sqlMode.toUpperCase().contains("NO_BACKSLASH_ESCAPES"); |
| | | } |
| | | } |
| | |
| | | // A session setting must reach the server as a plain batch: the sql server driver runs a |
| | | // prepared statement through sp_executesql, and a setting made there is reverted when that |
| | | // call returns - before the statement it is meant to protect ever runs. |
| | | // |
| | | // Outside both layers of the bound, like the comment statement executeAny() runs, and for the |
| | | // same reason: this is issued from newStampConnection() on a stamp connection, whose connect |
| | | // properties carry a socket read timeout of their own (Dialect.connectProperties). |
| | | private void executeSessionStatement(Connection con, String sql) throws SQLException { |
| | | try (final Statement statement=con.createStatement()) { |
| | | if (logger.isTraceEnabled()) { |
| | |
| | | } |
| | | try (final PreparedStatement statement=con.prepareStatement(sql)) { |
| | | statement.setString(1,arg); |
| | | try (final ResultSet rs=executeResultSet(statement)) { |
| | | return rs.next() ? rs.getString(1) : null; |
| | | } |
| | | return executeResultSet(statement, rs -> rs.next() ? rs.getString(1) : null); |
| | | } |
| | | } |
| | | |
| | |
| | | if (dialect==null) { // no portable statistics refresh for other engines |
| | | return false; // nothing was refreshed: reporting success here would make the assertion of the tests vacuous |
| | | } |
| | | final int timeoutSeconds=Math.max(0,Integer.getInteger(STATISTICS_TIMEOUT_PROPERTY,STATISTICS_TIMEOUT_SECONDS_DEFAULT)); |
| | | final int timeoutSeconds=clampSeconds(Integer.getInteger(STATISTICS_TIMEOUT_PROPERTY,STATISTICS_TIMEOUT_SECONDS_DEFAULT)); |
| | | boolean allRefreshed=true; |
| | | for (final TreeName treeName : trees) { |
| | | final String tableName=getTableName(treeName); |
| | |
| | | throw new IllegalStateException("no statistics refresh for dialect "+dialect); |
| | | } |
| | | try (final PreparedStatement statement=con.prepareStatement(sql)) { |
| | | statement.setQueryTimeout(timeoutSeconds); // 0: wait without limit |
| | | // 0: wait without limit - and false where the driver would not take the cancel, which |
| | | // leaves the socket read timeout behind it as the only layer this refresh runs under |
| | | final boolean cancelArmed=timeoutSeconds>0 && setQueryTimeout(statement, timeoutSeconds); |
| | | for (int i=0;i<args.length;i++) { |
| | | statement.setString(i+1,args[i]); |
| | | } |
| | | if (dialect==Dialect.MYSQL) { // mysql reports analyze problems as a result row, not an SQLException |
| | | try (final ResultSet rs=executeResultSet(statement)) { |
| | | while (rs.next()) { |
| | | if ("error".equalsIgnoreCase(rs.getString("Msg_type"))) { |
| | | throw new SQLException(rs.getString("Msg_text")); |
| | | // Under the bound of the statistics refresh rather than under a class of |
| | | // StatementBound, which would put its own value over one this statement has a |
| | | // property for - but under both layers of it all the same: on oracle this is |
| | | // dbms_stats.gather_table_stats, the engine whose session does not act on the |
| | | // break its driver sends, and it runs at the very end of a successful import, |
| | | // where a cancel that never arrives would park it with the data already |
| | | // committed and nothing left to report. |
| | | bounded(con, STATISTICS_TIMEOUT_PROPERTY, timeoutSeconds, cancelArmed, () -> { |
| | | if (logger.isTraceEnabled()) { |
| | | logger.trace(LocalizableMessage.raw("jdbc: %s",statement)); |
| | | } |
| | | if (dialect==Dialect.MYSQL) { // mysql reports analyze problems as a result row, not an SQLException |
| | | try (final ResultSet rs=statement.executeQuery()) { |
| | | while (rs.next()) { |
| | | if ("error".equalsIgnoreCase(rs.getString("Msg_type"))) { |
| | | throw new SQLException(rs.getString("Msg_text")); |
| | | } |
| | | } |
| | | } |
| | | }else { // tolerates a statement that returns a result set, which execute() does not |
| | | statement.execute(); |
| | | } |
| | | }else { |
| | | executeAny(statement); |
| | | } |
| | | return null; |
| | | }); |
| | | con.commit(); |
| | | } |
| | | }catch (Exception e) { |
| | |
| | | try { |
| | | for (final TreeName treeName : trees) { |
| | | try (final PreparedStatement statement = con.prepareStatement("drop table " + getTableName(treeName))) { |
| | | execute(statement); |
| | | execute(statement, StatementBound.BULK); |
| | | } |
| | | } |
| | | con.commit(); |
| | |
| | | return driverNameOf(con).contains("microsoft") ? "cast(? as char(128))" : "?"; |
| | | } |
| | | |
| | | private class ReadableTransactionImpl implements ReadableTransaction { |
| | | class ReadableTransactionImpl implements ReadableTransaction { |
| | | final Connection con; |
| | | /** |
| | | * The class the statements of this transaction take. It follows who runs them rather than |
| | | * what they look like: an import issues the same select and the same upsert a client |
| | | * operation does, but nobody is waiting on it - and on mssql it works the table unindexed, |
| | | * {@code k} being a {@code varbinary(max)} that cannot be an index key - so bounding an |
| | | * import as an entry read fails an import that ran to the end before this bound existed. |
| | | * The catalog lookups of {@code openTree()} keep the operation class whoever runs them: they |
| | | * read a data dictionary rather than the data, so a wait there is another session's metadata |
| | | * lock, which is one of the waits this bound exists to end. |
| | | */ |
| | | final StatementBound bound; |
| | | boolean isReadOnly=true; |
| | | |
| | | public ReadableTransactionImpl(Connection con) { |
| | | this(con, StatementBound.OPERATION); |
| | | } |
| | | |
| | | ReadableTransactionImpl(Connection con, StatementBound bound) { |
| | | this.con=con; |
| | | this.bound=bound; |
| | | } |
| | | |
| | | @Override |
| | | public ByteString read(TreeName treeName, ByteSequence key) { |
| | | try (final PreparedStatement statement=con.prepareStatement("select v from "+readTableName(treeName)+" where h="+hashParam(con)+" and k=?")){ |
| | | // the non-enrolling name: a read must not put a tree this backend does not own - the |
| | | // shared compressed schema tree of #873 - up for removal |
| | | final String tableName=readTableName(treeName); |
| | | try (final PreparedStatement statement=con.prepareStatement("select v from "+tableName+" where h="+hashParam(con)+" and k=?")){ |
| | | statement.setString(1,key2hash.get(ByteBuffer.wrap(key.toByteArray()))); |
| | | statement.setBytes(2,real2db(key.toByteArray())); |
| | | try(ResultSet rc=executeResultSet(statement)) { |
| | | return rc.next() ? ByteString.wrap(rc.getBytes("v")) : null; |
| | | } |
| | | return executeResultSet(statement, bound, rc -> rc.next() ? valueOfRow(rc, tableName) : null); |
| | | }catch (SQLException e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | |
| | | |
| | | @Override |
| | | public Cursor<ByteString, ByteString> openCursor(TreeName treeName) { |
| | | return new CursorImpl(isReadOnly,con,treeName); |
| | | return new CursorImpl(isReadOnly,con,treeName,bound); |
| | | } |
| | | |
| | | /** |
| | | * {@inheritDoc} |
| | | * <p> |
| | | * The batches of such a cursor are bulk statements however ordinary they look: nobody is |
| | | * waiting on the walk, and on mssql it is not even a walk along an index - {@code k} is a |
| | | * {@code varbinary(max)} there, which cannot be an index key, so every batch is a scan and |
| | | * a sort of the table. Bounding those as entry reads aborted an export or a rebuild that |
| | | * ran to the end before this bound existed. |
| | | */ |
| | | @Override |
| | | public Cursor<ByteString, ByteString> openBulkCursor(TreeName treeName) { |
| | | return new CursorImpl(isReadOnly,con,treeName,StatementBound.BULK); |
| | | } |
| | | |
| | | /** |
| | | * {@inheritDoc} |
| | | * <p> |
| | | * Bulk whoever asks: {@code select count(*)} is a scan of the whole table on every engine |
| | | * here, so what it takes follows the size of the backend rather than the work of the caller |
| | | * that happens to ask. Its callers are administrative either way - {@code dbtest} through |
| | | * {@code BackendStat}, and the counts {@code verify-index} reports - so the override costs a |
| | | * client operation nothing. It is not the count behind {@code NOTE_BACKEND_STARTED}: that one |
| | | * is {@code BackendImpl.getEntryCount()} through {@code RootContainer.getEntryCount()}, which |
| | | * sums {@code id2childrenCount} and never reaches this method. |
| | | * <p> |
| | | * One of the places the class of the transaction is overridden downwards, the others being |
| | | * {@link #openBulkCursor(TreeName)}, {@link CursorImpl#positionToLastKey()} and the DDL a |
| | | * write transaction issues - the {@code create table} and the three {@code create index} of |
| | | * {@code openTree()}, the {@code delete from} of {@code clearTree()} and the {@code drop |
| | | * table} of {@code deleteTree()} - which is where an operation-class transaction, the one |
| | | * {@code write()} runs with, can take the shared backstop of its connection off. The count |
| | | * is deliberately not given here: whoever audits that list has to read it off the class |
| | | * rather than trust a number that a later hard-coded {@code BULK} would leave stale. |
| | | */ |
| | | @Override |
| | | public long getRecordCount(TreeName treeName) { |
| | | try (final PreparedStatement statement=con.prepareStatement("select count(*) from "+readTableName(treeName)); |
| | | final ResultSet rc=executeResultSet(statement)){ |
| | | return rc.next() ? rc.getLong(1) : 0; |
| | | try (final PreparedStatement statement=con.prepareStatement("select count(*) from "+readTableName(treeName))){ |
| | | return executeResultSet(statement, StatementBound.BULK, rc -> rc.next() ? rc.getLong(1) : 0); |
| | | }catch (SQLException e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | |
| | | // the writeable one inherits it. |
| | | boolean isExistsTable(TreeName treeName) { |
| | | final String tableName = readTableName(treeName); |
| | | // the catalog lookup guarding a create table is bounded as the operation it is, not as |
| | | // the bulk statement it guards, and not as the class of the transaction that happens to |
| | | // ask: it reads a data dictionary rather than the data, so a wait here is the metadata |
| | | // lock of another session |
| | | try { |
| | | final DatabaseMetaData metaData = con.getMetaData(); |
| | | // asked of the catalog by name: openTree(createOnDemand) calls this for every tree |
| | | // of the backend - about 25 of them for a stock suffix, on every open - and listing |
| | | // every table of the database each time costs the whole catalog once per tree, on a |
| | | // database this backend may well be sharing with something else |
| | | try (final ResultSet rs = metaData.getTables(null, null, |
| | | storedIdentifier(metaData, tableName), new String[]{"TABLE"})) { |
| | | while (rs.next()) { |
| | | // the name still has to be compared: "_" is a single-character wildcard in a |
| | | // metadata pattern, so "opendj_<hash>" also matches a table named "opendjX<hash>" |
| | | if (tableName.equalsIgnoreCase(rs.getString("TABLE_NAME"))) { |
| | | return true; |
| | | return bounded(con, StatementBound.OPERATION, () -> { |
| | | final DatabaseMetaData metaData = con.getMetaData(); |
| | | // asked of the catalog by name: openTree(createOnDemand) calls this for every tree |
| | | // of the backend - about 25 of them for a stock suffix, on every open - and listing |
| | | // every table of the database each time costs the whole catalog once per tree, on a |
| | | // database this backend may well be sharing with something else |
| | | try (final ResultSet rs = metaData.getTables(null, null, |
| | | storedIdentifier(metaData, tableName), new String[]{"TABLE"})) { |
| | | while (rs.next()) { |
| | | // the name still has to be compared: "_" is a single-character wildcard in a |
| | | // metadata pattern, so "opendj_<hash>" also matches a table named "opendjX<hash>" |
| | | if (tableName.equalsIgnoreCase(rs.getString("TABLE_NAME"))) { |
| | | return true; |
| | | } |
| | | } |
| | | } |
| | | } |
| | | return false; |
| | | }); |
| | | } catch (Exception e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | | return false; |
| | | } |
| | | } |
| | | /** |
| | |
| | | boolean partlyCommitted; |
| | | |
| | | public WriteableTransactionTransactionImpl(Connection con) { |
| | | super(con); |
| | | this(con, StatementBound.OPERATION); |
| | | } |
| | | |
| | | WriteableTransactionTransactionImpl(Connection con, StatementBound bound) { |
| | | super(con, bound); |
| | | //captured once rather than read per operation: the access mode of the storage is mutable state - |
| | | //ImporterImpl reopens the storage READ_WRITE under its caller - and a transaction has to keep the mode |
| | | //it was created with. It also drives isReadOnly, so that a cursor this transaction opens refuses |
| | |
| | | */ |
| | | private void commitStatement(String sql, boolean ddl) throws SQLException { |
| | | partlyCommitted|=ddl && commitsBeforeDdl(); |
| | | // Bulk, whatever class the transaction itself carries: every statement issued through here is |
| | | // one of the ones #877 names as overridden downwards - the create table and the three create |
| | | // index of openTree(), the delete from of clearTree() and the drop table of deleteTree() - |
| | | // and nobody is waiting on any of them. |
| | | try (final PreparedStatement statement=con.prepareStatement(sql)) { |
| | | execute(statement); |
| | | execute(statement, StatementBound.BULK); |
| | | partlyCommitted=true; // a commit that fails leaves the outcome unknown, which is no more replayable |
| | | con.commit(); |
| | | } |
| | |
| | | } |
| | | |
| | | boolean isExistsIndex(String tableName, String indexName) throws SQLException { |
| | | // approximate=true: with false the oracle driver runs ANALYZE on every call |
| | | try (final ResultSet rs = con.getMetaData().getIndexInfo(null, null, tableName, false, true)) { |
| | | while (rs.next()) { |
| | | if (indexName.equalsIgnoreCase(rs.getString("INDEX_NAME"))) { |
| | | return true; |
| | | return bounded(con, StatementBound.OPERATION, () -> { |
| | | // approximate=true: with false the oracle driver runs ANALYZE on every call |
| | | try (final ResultSet rs = con.getMetaData().getIndexInfo(null, null, tableName, false, true)) { |
| | | while (rs.next()) { |
| | | if (indexName.equalsIgnoreCase(rs.getString("INDEX_NAME"))) { |
| | | return true; |
| | | } |
| | | } |
| | | } |
| | | } |
| | | return false; |
| | | return false; |
| | | }); |
| | | } |
| | | |
| | | public void clearTree(TreeName treeName) { |
| | |
| | | statement.setString(1, key2hash.get(ByteBuffer.wrap(key.toByteArray()))); |
| | | statement.setBytes(2, real2db(key.toByteArray())); |
| | | statement.setBytes(3, value.toByteArray()); |
| | | return (execute(statement) == 1 && statement.getUpdateCount() > 0); |
| | | return (execute(statement, bound) == 1 && statement.getUpdateCount() > 0); |
| | | } |
| | | }else if (driverName.contains("mysql")) { //mysql upsert |
| | | try (final PreparedStatement statement = con.prepareStatement("insert into " + getTableName(treeName) + " (h,k,v) values (?,?,?) as new ON DUPLICATE KEY UPDATE v=new.v")) { |
| | | statement.setString(1, key2hash.get(ByteBuffer.wrap(key.toByteArray()))); |
| | | statement.setBytes(2, real2db(key.toByteArray())); |
| | | statement.setBytes(3, value.toByteArray()); |
| | | return (execute(statement) == 1 && statement.getUpdateCount() > 0); |
| | | return (execute(statement, bound) == 1 && statement.getUpdateCount() > 0); |
| | | } |
| | | }else if (driverName.contains("oracle")) { //ANSI MERGE without ; |
| | | try (final PreparedStatement statement = con.prepareStatement("merge into " + getTableName(treeName) + " old using (select ? h,? k,? v from dual) new on (old.h=new.h and old.k=new.k) WHEN MATCHED THEN UPDATE SET old.v=new.v WHEN NOT MATCHED THEN INSERT (h,k,v) VALUES (new.h,new.k,new.v)")) { |
| | | statement.setString(1, key2hash.get(ByteBuffer.wrap(key.toByteArray()))); |
| | | statement.setBytes(2, real2db(key.toByteArray())); |
| | | statement.setBytes(3, value.toByteArray()); |
| | | return (execute(statement) == 1 && statement.getUpdateCount() > 0); |
| | | return (execute(statement, bound) == 1 && statement.getUpdateCount() > 0); |
| | | } |
| | | }else if (driverName.contains("microsoft")) { //ANSI MERGE with ; WITH (HOLDLOCK) makes the upsert atomic: without it SQL Server MERGE can race two concurrent NOT MATCHED inserts of the same key into a PRIMARY KEY violation. UPDLOCK is required on top of it: with HOLDLOCK alone the search phase takes a shared lock that the WHEN MATCHED update then has to convert to an exclusive one, so two concurrent upserts of the same key deadlock on the conversion; an update lock is taken right away and makes the second transaction wait instead. h is cast back to char so that the join can seek the primary key instead of scanning the whole table under those locks, see hashParam() |
| | | try (final PreparedStatement statement = con.prepareStatement("merge into " + getTableName(treeName) + " WITH (HOLDLOCK, UPDLOCK) old using (select cast(? as char(128)) h,? k,? v) new on (old.h=new.h and old.k=new.k) WHEN MATCHED THEN UPDATE SET old.v=new.v WHEN NOT MATCHED THEN INSERT (h,k,v) VALUES (new.h,new.k,new.v);")) { |
| | | statement.setString(1, key2hash.get(ByteBuffer.wrap(key.toByteArray()))); |
| | | statement.setBytes(2, real2db(key.toByteArray())); |
| | | statement.setBytes(3, value.toByteArray()); |
| | | return (execute(statement) == 1 && statement.getUpdateCount() > 0); |
| | | return (execute(statement, bound) == 1 && statement.getUpdateCount() > 0); |
| | | } |
| | | }else { //ANSI SQL: try update before insert with not exists |
| | | return update(treeName,key,value) || insert(treeName,key,value); |
| | |
| | | statement.setBytes(3, value.toByteArray()); |
| | | statement.setString(4, key2hash.get(ByteBuffer.wrap(key.toByteArray()))); |
| | | statement.setBytes(5, real2db(key.toByteArray())); |
| | | return (execute(statement)==1 && statement.getUpdateCount()>0); |
| | | return (execute(statement, bound)==1 && statement.getUpdateCount()>0); |
| | | } |
| | | } |
| | | |
| | |
| | | statement.setBytes(1,value.toByteArray()); |
| | | statement.setString(2,key2hash.get(ByteBuffer.wrap(key.toByteArray()))); |
| | | statement.setBytes(3,real2db(key.toByteArray())); |
| | | return (execute(statement)==1 && statement.getUpdateCount()>0); |
| | | return (execute(statement, bound)==1 && statement.getUpdateCount()>0); |
| | | } |
| | | } |
| | | |
| | |
| | | try (final PreparedStatement statement=con.prepareStatement("delete from "+getTableName(treeName)+" where h="+hashParam(con)+" and k=?")){ |
| | | statement.setString(1,key2hash.get(ByteBuffer.wrap(key.toByteArray()))); |
| | | statement.setBytes(2,real2db(key.toByteArray())); |
| | | return (execute(statement)==1 && statement.getUpdateCount()>0); |
| | | return (execute(statement, bound)==1 && statement.getUpdateCount()>0); |
| | | }catch (SQLException e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | |
| | | ByteString currentValue; |
| | | boolean defined; |
| | | |
| | | public CursorImpl(boolean isReadOnly, Connection con, TreeName treeName) { |
| | | // The class of the statements this cursor issues, from whoever opened it: a search walks its |
| | | // index and has a client waiting, while an import, an export or a rebuild walks a whole tree |
| | | // with nobody waiting - and on mssql it walks it unindexed either way. It is not read off the |
| | | // shape of the statement, because the opening batch of every cursor is the same |
| | | // unconditioned "order by k" that positionToLastKey() issues, search or not. |
| | | final StatementBound batchBound; |
| | | |
| | | public CursorImpl(boolean isReadOnly, Connection con, TreeName treeName, StatementBound batchBound) { |
| | | this.isReadOnly=isReadOnly; |
| | | this.con=con; |
| | | this.treeName=treeName; |
| | | // the read statements below take the non-enrolling name: a cursor is how the migration |
| | | // of #873 reads the shared tree, and reading a tree must not put it up for removal |
| | | this.tableName=readTableName(treeName); |
| | | this.batchBound=batchBound; |
| | | this.limitClause=((CachedConnection)con).parent.getClass().getName().contains("mysql") |
| | | ? " limit ?,?" : " offset ? rows fetch next ? rows only"; |
| | | } |
| | |
| | | return size; |
| | | } |
| | | |
| | | boolean fetchBatch(String condition, byte[] dbKey, long offset, boolean descending, int limit) { |
| | | /** |
| | | * Reads one batch of the cursor. The class of the bound is the caller's: a batch taken |
| | | * along the index of the tree for a client is an operation, while a batch that has to look |
| | | * at the whole table to answer - the one behind {@link #positionToLastKey()} - and every |
| | | * batch of a cursor an import or a rebuild walks ({@link #batchBound}) is bulk work, and |
| | | * the two cannot share a value. |
| | | */ |
| | | boolean fetchBatch(String condition, byte[] dbKey, long offset, boolean descending, int limit, StatementBound bound) { |
| | | fetchCount++; |
| | | buffer.clear(); |
| | | try (final PreparedStatement statement=con.prepareStatement("select k,v from "+tableName |
| | |
| | | } |
| | | statement.setLong(i++,offset); |
| | | statement.setLong(i,limit); |
| | | try(final ResultSet rc=executeResultSet(statement)) { |
| | | return executeResultSet(statement, bound, rc -> { |
| | | while (rc.next()) { |
| | | buffer.add(new byte[][]{rc.getBytes(1),rc.getBytes(2)}); |
| | | buffer.add(new byte[][]{rc.getBytes(1),valueOfRow(rc.getBytes(2),tableName)}); |
| | | } |
| | | } |
| | | return !buffer.isEmpty(); |
| | | }); |
| | | }catch (SQLException e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | | return !buffer.isEmpty(); |
| | | } |
| | | |
| | | void advanceFromBuffer() { |
| | |
| | | |
| | | @Override |
| | | public boolean next() { |
| | | if (buffer.isEmpty() && !fetchBatch(currentKeyDb==null?null:">",currentKeyDb,0,false,adaptiveBatchSize())) { |
| | | if (buffer.isEmpty() && !fetchBatch(currentKeyDb==null?null:">",currentKeyDb,0,false,adaptiveBatchSize(),batchBound)) { |
| | | defined=false; |
| | | return false; |
| | | } |
| | |
| | | try (final PreparedStatement statement=con.prepareStatement("delete from "+writeTableName+" where h="+hashParam(con)+" and k=?")){ |
| | | statement.setString(1,key2hash.get(ByteBuffer.wrap(db2real(currentKeyDb)))); |
| | | statement.setBytes(2,currentKeyDb); |
| | | execute(statement); |
| | | execute(statement, batchBound); |
| | | }catch (SQLException e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | |
| | | if (!buffer.isEmpty()) { // jumped outside the buffered range: random access, back to small batches |
| | | nextBatchSize=initialBatchSize; |
| | | } |
| | | if (fetchBatch(">=",target,0,false,adaptiveBatchSize())) { |
| | | if (fetchBatch(">=",target,0,false,adaptiveBatchSize(),batchBound)) { |
| | | advanceFromBuffer(); |
| | | return true; |
| | | } |
| | |
| | | @Override |
| | | public boolean positionToKey(ByteSequence key) { |
| | | final byte[] real=key.toByteArray(); |
| | | // The row is wrapped inside the handler rather than after it, so that null keeps meaning |
| | | // "no such key" and only that: a row whose v is null - which the schema allows, however |
| | | // this backend writes it - has to fail here as it fails in read(), rather than report a |
| | | // key that exists as absent. Both go through valueOfRow(), which is where that failure |
| | | // is named. |
| | | final ByteString value; |
| | | try (final PreparedStatement statement=con.prepareStatement("select v from "+tableName+" where h="+hashParam(con)+" and k=?")){ |
| | | statement.setString(1,key2hash.get(ByteBuffer.wrap(real))); |
| | | statement.setBytes(2,real2db(real)); |
| | | try(final ResultSet rc=executeResultSet(statement)) { |
| | | if (rc.next()) { |
| | | buffer.clear(); |
| | | nextBatchSize=initialBatchSize; |
| | | currentKeyDb=real2db(real); |
| | | currentKey=ByteString.wrap(real); |
| | | currentValue=ByteString.wrap(rc.getBytes("v")); |
| | | defined=true; |
| | | return true; |
| | | } |
| | | } |
| | | value=executeResultSet(statement, batchBound, rc -> rc.next() ? valueOfRow(rc, tableName) : null); |
| | | }catch (SQLException e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | | if (value!=null) { |
| | | buffer.clear(); |
| | | nextBatchSize=initialBatchSize; |
| | | currentKeyDb=real2db(real); |
| | | currentKey=ByteString.wrap(real); |
| | | currentValue=value; |
| | | defined=true; |
| | | return true; |
| | | } |
| | | defined=false; |
| | | return false; |
| | | } |
| | | |
| | | /** |
| | | * Bulk, not operation: with no condition to seek on, this is {@code order by k desc} over |
| | | * the whole table - and on mssql, where {@code k} is a {@code varbinary(max)} that cannot |
| | | * be an index key, a scan and a sort of it. It is also not on a search path: every open of |
| | | * a backend runs it once per base DN, through {@code EntryContainer.getHighestEntryID()}, |
| | | * outside the try/catch of {@code BackendImpl.openBackend()} - a bound of two minutes here |
| | | * would turn a large backend that opens slowly into one that does not open at all. |
| | | */ |
| | | @Override |
| | | public boolean positionToLastKey() { |
| | | if (fetchBatch(null,null,0,true,1)) { |
| | | if (fetchBatch(null,null,0,true,1,StatementBound.BULK)) { |
| | | advanceFromBuffer(); |
| | | return true; |
| | | } |
| | |
| | | return false; |
| | | } |
| | | |
| | | /** |
| | | * The class of the cursor, unlike {@link #positionToLastKey()}, which is bulk however it |
| | | * was opened: an offset comes from the VLV request of a client, so this runs on a search |
| | | * path and has to give the worker thread back - an import has no VLV position to seek to. |
| | | * That a deep offset is served by walking to it - the engines have no other way to answer |
| | | * an {@code offset ?} - is what makes the bound reachable here, and reaching it answers the |
| | | * request with an error rather than parking a thread of the server on it. |
| | | */ |
| | | @Override |
| | | public boolean positionToIndex(int index) { |
| | | if (!buffer.isEmpty()) { // absolute jump: random access, back to small batches |
| | | nextBatchSize=initialBatchSize; |
| | | } |
| | | if (index>=0 && fetchBatch(null,null,index,false,adaptiveBatchSize())) { |
| | | if (index>=0 && fetchBatch(null,null,index,false,adaptiveBatchSize(),batchBound)) { |
| | | advanceFromBuffer(); |
| | | return true; |
| | | } |
| | |
| | | return tree2table.asMap().keySet(); |
| | | } |
| | | |
| | | private final class ImporterImpl implements Importer { |
| | | final class ImporterImpl implements Importer { |
| | | final Connection con; |
| | | final ReadableTransactionImpl txr; |
| | | final WriteableTransactionTransactionImpl txw; |
| | |
| | | |
| | | final Boolean isOpen; |
| | | |
| | | public ImporterImpl() { |
| | | isOpen=getStorageStatus().isWorking(); |
| | | if (!isOpen) { |
| | | try { |
| | | open(AccessMode.READ_WRITE); |
| | | }catch (Exception e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | | /** |
| | | * Both transactions of an import take the bulk class, and with them every statement it |
| | | * issues: phase one writes the trees through {@code put()}, phase two reads them back |
| | | * through {@code read()} and walks them through {@code openCursor()}, and none of that has |
| | | * a client waiting on it. Bounding those as entry reads is not merely strict, it fails work |
| | | * that ran to the end before this bound existed: {@code h} is the primary key on every |
| | | * dialect and the default lock wait is forever on mssql, postgres and oracle, so an upsert |
| | | * of an online import blocked by an LDAP write on the same table sat until the bound of an |
| | | * entry read and then failed the import. |
| | | */ |
| | | ImporterImpl(Connection con, boolean isOpen) { |
| | | // An import writes by definition, so a storage that is not writeable refuses one where the |
| | | // importer is built - which is where it was refused until the write transaction of a read-only |
| | | // storage became one that is granted and checks per operation (#874). Left to that check, an |
| | | // import of such a storage would take a connection out of the pool, begin its transaction and |
| | | // fail at the first tree it clears rather than at its start. |
| | | // What arrives here read-only is a storage that was already open: import-ldif and |
| | | // rebuild-index both close it first, and startImport() opens a closed one READ_WRITE - an |
| | | // import of any storage of this server reopens it that way - so those two arrive writeable. |
| | | if (!accessMode.isWriteable()) { |
| | | throw new ReadOnlyStorageException(); |
| | | } |
| | | try { |
| | | con = getValidatedConnection(); |
| | | }catch (Exception e){ |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | | txr =new ReadableTransactionImpl(con); |
| | | txw =new WriteableTransactionTransactionImpl(con); |
| | | this.con=con; |
| | | this.isOpen=isOpen; |
| | | txr=new ReadableTransactionImpl(con, StatementBound.BULK); |
| | | txw=new WriteableTransactionTransactionImpl(con, StatementBound.BULK); |
| | | } |
| | | |
| | | @Override |
| | |
| | | aborted = true; |
| | | } |
| | | |
| | | // The connection goes back whatever the commit does, and the storage this importer opened |
| | | // is closed whatever the connection does: an importer is closed on the way out of a failed |
| | | // import as readily as a finished one - a clearTree() that reaches the bulk bound is one |
| | | // way there - and a commit that throws on the way would otherwise leave the connection |
| | | // out of the pool for good, holding the transaction and the locks of that import. |
| | | @Override |
| | | public void close() { |
| | | try { |
| | |
| | | return txr.read(treeName, key); |
| | | } |
| | | |
| | | // Bulk like every other statement of an import, by the class of the transaction it comes |
| | | // from: this walks a whole tree with no client waiting on it - phase one of a rebuild-index |
| | | // reads every record of id2entry through this cursor (OnDiskMergeImporter.ID2EntrySource) - |
| | | // and on mssql it walks it unindexed, so a batch of it is a scan and a sort of the table |
| | | // rather than a step along an index. |
| | | @Override |
| | | public SequentialCursor<ByteString, ByteString> openCursor(TreeName treeName) { |
| | | return txr.openCursor(treeName); |
| | |
| | | //import |
| | | @Override |
| | | public Importer startImport() throws ConfigException, StorageRuntimeException { |
| | | return new ImporterImpl(); |
| | | final boolean wasOpen=getStorageStatus().isWorking(); |
| | | if (!wasOpen) { |
| | | try { |
| | | open(AccessMode.READ_WRITE); |
| | | }catch (Exception e) { |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | | } |
| | | final Connection con; |
| | | try { |
| | | con=getValidatedConnection(); |
| | | }catch (Exception e){ |
| | | // and the storage this method opened goes back with it: ImporterImpl.close() is what closes |
| | | // it again when an import opened it, and no importer is going to be built to reach that |
| | | if (!wasOpen) { |
| | | close(); |
| | | } |
| | | throw new StorageRuntimeException(e); |
| | | } |
| | | // outside the catch: the importer of a read-only storage throws ReadOnlyStorageException, |
| | | // which a caller tells apart from any other failure of an import |
| | | boolean built=false; |
| | | try { |
| | | final Importer importer=new ImporterImpl(con, wasOpen); |
| | | built=true; |
| | | return importer; |
| | | }finally { |
| | | // and the connection borrowed above goes back on every path that does not build an |
| | | // importer to hold it: it is the one an import keeps for its whole duration, so leaving it |
| | | // here takes it out of the pool for good, with the transaction it had already begun. A |
| | | // finally rather than a catch, so that it covers what a catch has to name - an Error |
| | | // leaves the pool one connection short exactly as ReadOnlyStorageException did. |
| | | if (!built) { |
| | | try { |
| | | con.close(); |
| | | }catch (SQLException ignored) { |
| | | // the importer was never built; the failure to report is the one on its way out |
| | | } |
| | | // and the storage this method opened goes back with the connection, for the reason the |
| | | // borrow above gives: ImporterImpl.close() is what closes it again when an import |
| | | // opened it, and there is no importer here to reach that |
| | | if (!wasOpen) { |
| | | close(); |
| | | } |
| | | } |
| | | } |
| | | } |
| | | |
| | | //backup |