diff --git a/ebean-api/src/main/java/io/ebean/CancelableQuery.java b/ebean-api/src/main/java/io/ebean/CancelableQuery.java new file mode 100644 index 000000000..fb3100bd2 --- /dev/null +++ b/ebean-api/src/main/java/io/ebean/CancelableQuery.java @@ -0,0 +1,19 @@ +package io.ebean; + +/** + * Defines a cancelable query. + *

+ * Typically holds a representation of the PreparedStatement to perform the + * actual cancel. + *

+ */ +public interface CancelableQuery { + + /** + * Cancel the query. + *

+ * For JDBC this translates to calling cancel on the PreparedStatement. + *

+ */ + void cancel(); +} diff --git a/ebean-api/src/main/java/io/ebean/DtoQuery.java b/ebean-api/src/main/java/io/ebean/DtoQuery.java index 95ad6eff5..69ed034ef 100644 --- a/ebean-api/src/main/java/io/ebean/DtoQuery.java +++ b/ebean-api/src/main/java/io/ebean/DtoQuery.java @@ -38,7 +38,7 @@ import java.util.stream.Stream; * * } */ -public interface DtoQuery { +public interface DtoQuery extends CancelableQuery { /** * Execute the query returning a list. diff --git a/ebean-api/src/main/java/io/ebean/Query.java b/ebean-api/src/main/java/io/ebean/Query.java index db7bf8cb9..0dc86905e 100644 --- a/ebean-api/src/main/java/io/ebean/Query.java +++ b/ebean-api/src/main/java/io/ebean/Query.java @@ -177,7 +177,7 @@ import java.util.stream.Stream; * * @param the type of Entity bean this query will fetch. */ -public interface Query { +public interface Query extends CancelableQuery { /** * The lock type (strength) to use with query FOR UPDATE row locking. @@ -291,15 +291,6 @@ public interface Query { */ UpdateQuery asUpdate(); - /** - * Cancel the query execution if supported by the underlying database and - * driver. - *

- * This must be called from a different thread to the query executor. - *

- */ - void cancel(); - /** * Return a copy of the query. *

diff --git a/ebean-api/src/main/java/io/ebean/SqlQuery.java b/ebean-api/src/main/java/io/ebean/SqlQuery.java index 46b26c90a..9a878f233 100644 --- a/ebean-api/src/main/java/io/ebean/SqlQuery.java +++ b/ebean-api/src/main/java/io/ebean/SqlQuery.java @@ -37,7 +37,7 @@ import java.util.function.Predicate; * * } */ -public interface SqlQuery extends Serializable { +public interface SqlQuery extends Serializable, CancelableQuery { /** * Execute the query returning a list. diff --git a/ebean-api/src/main/java/io/ebean/util/JdbcClose.java b/ebean-api/src/main/java/io/ebean/util/JdbcClose.java index 92b3f7e28..b72f59c4d 100644 --- a/ebean-api/src/main/java/io/ebean/util/JdbcClose.java +++ b/ebean-api/src/main/java/io/ebean/util/JdbcClose.java @@ -66,4 +66,17 @@ public class JdbcClose { logger.warn("Error on connection rollback", e); } } + + /** + * Cancels the statement + */ + public static void cancel(Statement stmt) { + try { + if (stmt != null) { + stmt.cancel(); + } + } catch (SQLException e) { + logger.warn("Error on cancelling statement", e); + } + } } diff --git a/ebean-core/src/main/java/io/ebeaninternal/api/SpiCancelableQuery.java b/ebean-core/src/main/java/io/ebeaninternal/api/SpiCancelableQuery.java new file mode 100644 index 000000000..10a53b9a7 --- /dev/null +++ b/ebean-core/src/main/java/io/ebeaninternal/api/SpiCancelableQuery.java @@ -0,0 +1,26 @@ +package io.ebeaninternal.api; + +import javax.persistence.PersistenceException; + +import io.ebean.CancelableQuery; + +/** + * Cancellable query, that has a delegate. + * + * @author Roland Praml, FOCONIS AG + * + */ +public interface SpiCancelableQuery extends CancelableQuery { + + /** + * Checks if the query was cancelled. + * @throws PersistenceException if query was cancelled. + */ + void checkCancelled(); + + /** + * Set the underlying cancelable query (with the PreparedStatement). + */ + void setCancelableQuery(CancelableQuery cancelableQuery); + +} diff --git a/ebean-core/src/main/java/io/ebeaninternal/api/SpiQuery.java b/ebean-core/src/main/java/io/ebeaninternal/api/SpiQuery.java index 4686a3644..12942e7d6 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/api/SpiQuery.java +++ b/ebean-core/src/main/java/io/ebeaninternal/api/SpiQuery.java @@ -17,7 +17,6 @@ import io.ebeaninternal.server.core.SpiOrmQueryRequest; import io.ebeaninternal.server.deploy.BeanDescriptor; import io.ebeaninternal.server.deploy.BeanPropertyAssocMany; import io.ebeaninternal.server.deploy.TableJoin; -import io.ebeaninternal.server.query.CancelableQuery; import io.ebeaninternal.server.querydefn.NaturalKeyBindParam; import io.ebeaninternal.server.querydefn.OrmQueryDetail; import io.ebeaninternal.server.querydefn.OrmQueryProperties; @@ -31,7 +30,7 @@ import java.util.Set; /** * Object Relational query - Internal extension to Query object. */ -public interface SpiQuery extends Query, SpiQueryFetch, TxnProfileEventCodes { +public interface SpiQuery extends Query, SpiQueryFetch, TxnProfileEventCodes, SpiCancelableQuery { enum Mode { NORMAL(false), LAZYLOAD_MANY(false), LAZYLOAD_BEAN(true), REFRESH_BEAN(true); @@ -847,16 +846,6 @@ public interface SpiQuery extends Query, SpiQueryFetch, TxnProfileEventCod */ ReadEvent getFutureFetchAudit(); - /** - * Set the underlying cancelable query (with the PreparedStatement). - */ - void setCancelableQuery(CancelableQuery cancelableQuery); - - /** - * Return true if this query has been cancelled. - */ - boolean isCancelled(); - /** * Return the base table to use if user defined on the query. */ diff --git a/ebean-core/src/main/java/io/ebeaninternal/api/SpiSqlBinding.java b/ebean-core/src/main/java/io/ebeaninternal/api/SpiSqlBinding.java index 849ad0cb0..62b93d34c 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/api/SpiSqlBinding.java +++ b/ebean-core/src/main/java/io/ebeaninternal/api/SpiSqlBinding.java @@ -3,7 +3,7 @@ package io.ebeaninternal.api; /** * SQL query binding (for SqlQuery and DtoQuery). */ -public interface SpiSqlBinding { +public interface SpiSqlBinding extends SpiCancelableQuery { /** * Return the named or positioned parameters. diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/AbstractSqlQueryRequest.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/AbstractSqlQueryRequest.java index 8a5e9214d..67844c631 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/AbstractSqlQueryRequest.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/AbstractSqlQueryRequest.java @@ -1,5 +1,6 @@ package io.ebeaninternal.server.core; +import io.ebean.CancelableQuery; import io.ebean.Transaction; import io.ebean.util.JdbcClose; import io.ebeaninternal.api.*; @@ -12,11 +13,14 @@ import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.ResultSet; import java.sql.SQLException; +import java.util.concurrent.locks.ReentrantLock; + +import javax.persistence.PersistenceException; /** * Wraps the objects involved in executing a SQL / Relational Query. */ -public abstract class AbstractSqlQueryRequest { +public abstract class AbstractSqlQueryRequest implements CancelableQuery { protected final SpiSqlBinding query; @@ -36,6 +40,8 @@ public abstract class AbstractSqlQueryRequest { protected long startNano; + private final ReentrantLock lock = new ReentrantLock(); + /** * Create the BeanFindRequest. */ @@ -43,6 +49,7 @@ public abstract class AbstractSqlQueryRequest { this.server = server; this.query = query; this.transaction = (SpiTransaction) t; + this.query.setCancelableQuery(this); } /** @@ -137,24 +144,30 @@ public abstract class AbstractSqlQueryRequest { } protected void executeAsSql(Binder binder) throws SQLException { - prepareSql(); - Connection conn = transaction.getInternalConnection(); - pstmt = conn.prepareStatement(sql); - if (query.getTimeout() > 0) { - pstmt.setQueryTimeout(query.getTimeout()); + lock.lock(); + try { + query.checkCancelled(); + prepareSql(); + Connection conn = transaction.getInternalConnection(); + pstmt = conn.prepareStatement(sql); + if (query.getTimeout() > 0) { + pstmt.setQueryTimeout(query.getTimeout()); + } + if (query.getBufferFetchSizeHint() > 0) { + pstmt.setFetchSize(query.getBufferFetchSizeHint()); + } + BindParams bindParams = query.getBindParams(); + if (!bindParams.isEmpty()) { + this.bindLog = binder.bind(bindParams, pstmt, conn); + } + if (isLogSql()) { + transaction.logSql(Str.add(TrimLogSql.trim(sql), "; --bind(", bindLog, ")")); + } + } finally { + lock.unlock(); } - if (query.getBufferFetchSizeHint() > 0) { - pstmt.setFetchSize(query.getBufferFetchSizeHint()); - } - BindParams bindParams = query.getBindParams(); - if (!bindParams.isEmpty()) { - this.bindLog = binder.bind(bindParams, pstmt, conn); - } - if (isLogSql()) { - transaction.logSql(Str.add(TrimLogSql.trim(sql), "; --bind(", bindLog, ")")); - } - setResultSet(pstmt.executeQuery(), null); + query.checkCancelled(); } /** @@ -164,4 +177,13 @@ public abstract class AbstractSqlQueryRequest { return sql; } + @Override + public void cancel() { + lock.lock(); + try { + JdbcClose.cancel(pstmt); + } finally { + lock.unlock(); + } + } } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/DefaultServer.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/DefaultServer.java index 231f0b718..87580f90b 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/DefaultServer.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/DefaultServer.java @@ -1382,7 +1382,7 @@ public final class DefaultServer implements SpiServer, SpiEbeanServer { @Nonnull @Override public FutureList findFutureList(Query query, Transaction t) { - SpiQuery spiQuery = (SpiQuery) query; + SpiQuery spiQuery = (SpiQuery) query.copy(); spiQuery.setFutureFetch(true); // FutureList query always run in it's own persistence content spiQuery.setPersistenceContext(new DefaultPersistenceContext()); diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/DtoQueryRequest.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/DtoQueryRequest.java index f4cdbc236..2b1a55f6e 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/DtoQueryRequest.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/DtoQueryRequest.java @@ -53,6 +53,7 @@ public final class DtoQueryRequest extends AbstractSqlQueryRequest { ormQuery.setType(type); ormQuery.setManualId(); + query.setCancelableQuery(ormQuery); // execute the underlying ORM query returning the ResultSet SpiResultSet result = server.findResultSet(ormQuery, transaction); this.pstmt = result.getStatement(); @@ -117,6 +118,7 @@ public final class DtoQueryRequest extends AbstractSqlQueryRequest { } public boolean next() throws SQLException { + query.checkCancelled(); return dataReader.next(); } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java index 5548c88e3..0fb83a18f 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/OrmQueryRequest.java @@ -1,6 +1,7 @@ package io.ebeaninternal.server.core; import io.ebean.CacheMode; +import io.ebean.CancelableQuery; import io.ebean.OrderBy; import io.ebean.PersistenceContextScope; import io.ebean.QueryIterator; @@ -35,7 +36,6 @@ import io.ebeaninternal.server.deploy.DeployPropertyParserMap; import io.ebeaninternal.server.el.ElPropertyValue; import io.ebeaninternal.server.loadcontext.DLoadContext; import io.ebeaninternal.server.query.CQueryPlan; -import io.ebeaninternal.server.query.CancelableQuery; import io.ebeaninternal.server.transaction.DefaultPersistenceContext; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java b/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java index 9a0ceb153..a5130ef22 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/core/RelationalQueryRequest.java @@ -139,6 +139,7 @@ public final class RelationalQueryRequest extends AbstractSqlQueryRequest { @Override public boolean next() throws SQLException { + query.checkCancelled(); if (!resultSet.next()) { return false; } else { diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQuery.java b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQuery.java index 0fb7e9dcf..e7674a31c 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQuery.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQuery.java @@ -1,5 +1,6 @@ package io.ebeaninternal.server.query; +import io.ebean.CancelableQuery; import io.ebean.QueryIterator; import io.ebean.Version; import io.ebean.bean.*; @@ -139,8 +140,6 @@ public class CQuery implements DbReadContext, CancelableQuery, SpiProfileTran */ private PreparedStatement pstmt; - private boolean cancelled; - private String bindLog; private final CQueryPlan queryPlan; @@ -272,15 +271,7 @@ public class CQuery implements DbReadContext, CancelableQuery, SpiProfileTran public void cancel() { lock.lock(); try { - this.cancelled = true; - if (pstmt != null) { - try { - logger.debug("Cancelling query"); - pstmt.cancel(); - } catch (SQLException e) { - throw new PersistenceException("Error cancelling query", e); - } - } + JdbcClose.cancel(pstmt); } finally { lock.unlock(); } @@ -312,9 +303,8 @@ public class CQuery implements DbReadContext, CancelableQuery, SpiProfileTran ResultSet prepareResultSet(boolean forwardOnlyHint) throws SQLException { lock.lock(); try { - if (cancelled) { - throw new SQLException("Query cancelled"); - } + // cancelled before we started + query.checkCancelled(); startNano = System.nanoTime(); SpiTransaction t = request.getTransaction(); profileOffset = t.profileOffset(); @@ -342,10 +332,12 @@ public class CQuery implements DbReadContext, CancelableQuery, SpiProfileTran pstmt.setFetchSize(query.getBufferFetchSizeHint()); } bindLog = predicates.bind(queryPlan.bindEncryptedProperties(pstmt, conn)); - return pstmt.executeQuery(); } finally { lock.unlock(); } + ResultSet ret = pstmt.executeQuery(); + query.checkCancelled(); + return ret; } /** @@ -491,7 +483,8 @@ public class CQuery implements DbReadContext, CancelableQuery, SpiProfileTran boolean hasNext() throws SQLException { lock.lock(); try { - if (noMoreRows || cancelled) { + query.checkCancelled(); + if (noMoreRows) { return false; } if (hasNextCache) { diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryEngine.java b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryEngine.java index c632c9aeb..a1a211c0a 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryEngine.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryEngine.java @@ -66,11 +66,13 @@ public class CQueryEngine { public int delete(OrmQueryRequest request) { CQueryUpdate query = queryBuilder.buildUpdateQuery(true, request); + request.setCancelableQuery(query); return executeUpdate(request, query); } public int update(OrmQueryRequest request) { CQueryUpdate query = queryBuilder.buildUpdateQuery(false, request); + request.setCancelableQuery(query); return executeUpdate(request, query); } @@ -97,6 +99,7 @@ public class CQueryEngine { public List findSingleAttributeList(OrmQueryRequest request) { CQueryFetchSingleAttribute rcQuery = queryBuilder.buildFetchAttributeQuery(request); + request.setCancelableQuery(rcQuery); return findAttributeList(request, rcQuery); } @@ -151,6 +154,7 @@ public class CQueryEngine { public List findIds(OrmQueryRequest request) { CQueryFetchSingleAttribute rcQuery = queryBuilder.buildFetchIdsQuery(request); + request.setCancelableQuery(rcQuery); return findAttributeList(request, rcQuery); } @@ -164,6 +168,7 @@ public class CQueryEngine { public int findCount(OrmQueryRequest request) { CQueryRowCount rcQuery = queryBuilder.buildRowCountQuery(request); + request.setCancelableQuery(rcQuery); try { int count = rcQuery.findCount(); @@ -235,8 +240,10 @@ public class CQueryEngine { } catch (SQLException e) { try { + PersistenceException pex = cquery.createPersistenceException(e); + // create exception before closing connection cquery.close(); - throw cquery.createPersistenceException(e); + throw pex; } finally { request.rollbackTransIfRequired(); } @@ -259,6 +266,7 @@ public class CQueryEngine { // order by lower sys period desc query.order().desc(sysPeriodLower); CQuery cquery = queryBuilder.buildQuery(request); + request.setCancelableQuery(cquery); try { cquery.prepareBindExecuteQuery(); if (request.isLogSql()) { @@ -327,6 +335,7 @@ public class CQueryEngine { */ public SpiResultSet findResultSet(OrmQueryRequest request) { CQuery cquery = queryBuilder.buildQuery(request); + request.setCancelableQuery(cquery); try { boolean fwdOnly; if (request.isFindIterate()) { @@ -411,6 +420,7 @@ public class CQueryEngine { EntityBean bean = null; CQuery cquery = queryBuilder.buildQuery(request); + request.setCancelableQuery(cquery); try { cquery.prepareBindExecuteQuery(); diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryFetchSingleAttribute.java b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryFetchSingleAttribute.java index 9fd7278cc..b68813bd4 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryFetchSingleAttribute.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryFetchSingleAttribute.java @@ -1,5 +1,6 @@ package io.ebeaninternal.server.query; +import io.ebean.CancelableQuery; import io.ebean.CountedValue; import io.ebean.core.type.ScalarDataReader; import io.ebean.util.JdbcClose; @@ -18,11 +19,12 @@ import java.sql.SQLException; import java.util.ArrayList; import java.util.List; import java.util.Set; +import java.util.concurrent.locks.ReentrantLock; /** * Base compiled query request for single attribute queries. */ -class CQueryFetchSingleAttribute implements SpiProfileTransactionEvent { +class CQueryFetchSingleAttribute implements SpiProfileTransactionEvent, CancelableQuery { private static final Logger logger = LoggerFactory.getLogger(CQueryFetchSingleAttribute.class); @@ -65,6 +67,8 @@ class CQueryFetchSingleAttribute implements SpiProfileTransactionEvent { private final boolean containsCounts; private long profileOffset; + + private final ReentrantLock lock = new ReentrantLock(); /** * Create the Sql select based on the request. @@ -147,21 +151,27 @@ class CQueryFetchSingleAttribute implements SpiProfileTransactionEvent { } private void prepareExecute() throws SQLException { - - SpiTransaction t = getTransaction(); - profileOffset = t.profileOffset(); - Connection conn = t.getInternalConnection(); - pstmt = conn.prepareStatement(sql); - - if (query.getBufferFetchSizeHint() > 0) { - pstmt.setFetchSize(query.getBufferFetchSizeHint()); + lock.lock(); + try { + query.checkCancelled(); + SpiTransaction t = getTransaction(); + profileOffset = t.profileOffset(); + Connection conn = t.getInternalConnection(); + pstmt = conn.prepareStatement(sql); + + if (query.getBufferFetchSizeHint() > 0) { + pstmt.setFetchSize(query.getBufferFetchSizeHint()); + } + if (query.getTimeout() > 0) { + pstmt.setQueryTimeout(query.getTimeout()); + } + + bindLog = predicates.bind(pstmt, conn); + } finally { + lock.unlock(); } - if (query.getTimeout() > 0) { - pstmt.setQueryTimeout(query.getTimeout()); - } - - bindLog = predicates.bind(pstmt, conn); dataReader = new RsetDataReader(request.getDataTimeZone(), pstmt.executeQuery()); + query.checkCancelled(); } /** @@ -194,4 +204,14 @@ class CQueryFetchSingleAttribute implements SpiProfileTransactionEvent { Set getDependentTables() { return queryPlan.getDependentTables(); } + + @Override + public void cancel() { + lock.lock(); + try { + JdbcClose.cancel(pstmt); + } finally { + lock.unlock(); + } + } } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryRowCount.java b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryRowCount.java index 6a3d40a59..248fa9058 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryRowCount.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryRowCount.java @@ -1,5 +1,6 @@ package io.ebeaninternal.server.query; +import io.ebean.CancelableQuery; import io.ebean.util.JdbcClose; import io.ebeaninternal.api.SpiProfileTransactionEvent; import io.ebeaninternal.api.SpiQuery; @@ -13,11 +14,12 @@ import java.sql.PreparedStatement; import java.sql.ResultSet; import java.sql.SQLException; import java.util.Set; +import java.util.concurrent.locks.ReentrantLock; /** * Executes the select row count query. */ -class CQueryRowCount implements SpiProfileTransactionEvent { +class CQueryRowCount implements SpiProfileTransactionEvent, CancelableQuery { private final CQueryPlan queryPlan; @@ -57,6 +59,8 @@ class CQueryRowCount implements SpiProfileTransactionEvent { private int rowCount; private long profileOffset; + + private final ReentrantLock lock = new ReentrantLock(); /** * Create the Sql select based on the request. @@ -110,14 +114,22 @@ class CQueryRowCount implements SpiProfileTransactionEvent { SpiTransaction t = getTransaction(); profileOffset = t.profileOffset(); Connection conn = t.getInternalConnection(); - pstmt = conn.prepareStatement(sql); + lock.lock(); + try { + query.checkCancelled(); + pstmt = conn.prepareStatement(sql); - if (query.getTimeout() > 0) { - pstmt.setQueryTimeout(query.getTimeout()); + if (query.getTimeout() > 0) { + pstmt.setQueryTimeout(query.getTimeout()); + } + + bindLog = predicates.bind(pstmt, conn); + } finally { + lock.unlock(); } - - bindLog = predicates.bind(pstmt, conn); rset = pstmt.executeQuery(); + query.checkCancelled(); + if (!rset.next()) { throw new PersistenceException("Expecting 1 row but got none?"); } @@ -161,4 +173,14 @@ class CQueryRowCount implements SpiProfileTransactionEvent { Set getDependentTables() { return queryPlan.getDependentTables(); } + + @Override + public void cancel() { + lock.lock(); + try { + JdbcClose.cancel(pstmt); + } finally { + lock.unlock(); + } + } } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryUpdate.java b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryUpdate.java index 76858c271..eec5d2e6f 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryUpdate.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/query/CQueryUpdate.java @@ -1,5 +1,6 @@ package io.ebeaninternal.server.query; +import io.ebean.CancelableQuery; import io.ebean.util.JdbcClose; import io.ebeaninternal.api.SpiProfileTransactionEvent; import io.ebeaninternal.api.SpiQuery; @@ -10,11 +11,12 @@ import io.ebeaninternal.server.deploy.BeanDescriptor; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.SQLException; +import java.util.concurrent.locks.ReentrantLock; /** - * Executes the delete query. + * Executes the update query. */ -class CQueryUpdate implements SpiProfileTransactionEvent { +class CQueryUpdate implements SpiProfileTransactionEvent, CancelableQuery { private final CQueryPlan queryPlan; @@ -45,6 +47,8 @@ class CQueryUpdate implements SpiProfileTransactionEvent { private long profileOffset; + private final ReentrantLock lock = new ReentrantLock(); + /** * Create the Sql select based on the request. */ @@ -82,15 +86,22 @@ class CQueryUpdate implements SpiProfileTransactionEvent { SpiTransaction t = getTransaction(); profileOffset = t.profileOffset(); Connection conn = t.getInternalConnection(); - pstmt = conn.prepareStatement(sql); + lock.lock(); + try { + query.checkCancelled(); + pstmt = conn.prepareStatement(sql); - if (query.getTimeout() > 0) { - pstmt.setQueryTimeout(query.getTimeout()); + if (query.getTimeout() > 0) { + pstmt.setQueryTimeout(query.getTimeout()); + } + + bindLog = predicates.bind(pstmt, conn); + } finally { + lock.unlock(); } - - bindLog = predicates.bind(pstmt, conn); rowCount = pstmt.executeUpdate(); - + query.checkCancelled(); + long executionTimeMicros = (System.nanoTime() - startNano) / 1000L; request.slowQueryCheck(executionTimeMicros, rowCount); if (queryPlan.executionTime(executionTimeMicros)) { @@ -122,4 +133,14 @@ class CQueryUpdate implements SpiProfileTransactionEvent { .profileStream() .addQueryEvent(query.profileEventId(), profileOffset, desc.getName(), rowCount, query.getProfileId()); } + + @Override + public void cancel() { + lock.lock(); + try { + JdbcClose.cancel(pstmt); + } finally { + lock.unlock(); + } + } } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/query/DtoQueryEngine.java b/ebean-core/src/main/java/io/ebeaninternal/server/query/DtoQueryEngine.java index 46de28b8b..526d36dee 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/query/DtoQueryEngine.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/query/DtoQueryEngine.java @@ -29,7 +29,7 @@ public class DtoQueryEngine { } return rows; - } catch (Throwable e) { + } catch (SQLException e) { throw new PersistenceException(errMsg(e.getMessage(), request.getSql()), e); } finally { request.close(); @@ -51,7 +51,7 @@ public class DtoQueryEngine { while (request.next()) { consumer.accept(request.readNextBean()); } - } catch (Exception e) { + } catch (SQLException e) { throw new PersistenceException(errMsg(e.getMessage(), request.getSql()), e); } finally { request.close(); @@ -88,9 +88,8 @@ public class DtoQueryEngine { break; } } - } catch (Exception e) { + } catch (SQLException e) { throw new PersistenceException(errMsg(e.getMessage(), request.getSql()), e); - } finally { request.close(); } diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/AbstractQuery.java b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/AbstractQuery.java new file mode 100644 index 000000000..62fd3cdf0 --- /dev/null +++ b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/AbstractQuery.java @@ -0,0 +1,56 @@ +package io.ebeaninternal.server.querydefn; + +import java.util.concurrent.locks.ReentrantLock; + +import javax.persistence.PersistenceException; + +import io.ebean.CancelableQuery; +import io.ebeaninternal.api.SpiCancelableQuery; + +/** + * Common code for Dto/Orm/RelationalQuery + * + * @author Roland Praml, FOCONIS AG + * + */ +public class AbstractQuery implements SpiCancelableQuery { + + private boolean cancelled; + + private CancelableQuery cancelableQuery; + + private final ReentrantLock lock = new ReentrantLock(); + + @Override + public void cancel() { + lock.lock(); + try { + if (!cancelled) { + cancelled = true; + if (cancelableQuery != null) { + cancelableQuery.cancel(); + } + } + } finally { + lock.unlock(); + } + } + + @Override + public void checkCancelled() { + if (cancelled) { + throw new PersistenceException("Query was cancelled"); + } + } + + @Override + public void setCancelableQuery(CancelableQuery cancelableQuery) { + lock.lock(); + try { + checkCancelled(); + this.cancelableQuery = cancelableQuery; + } finally { + lock.unlock(); + } + } +} diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultDtoQuery.java b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultDtoQuery.java index 9afcb961f..c2cb02198 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultDtoQuery.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultDtoQuery.java @@ -21,7 +21,7 @@ import java.util.stream.Stream; /** * Default implementation of DtoQuery. */ -public class DefaultDtoQuery implements SpiDtoQuery { +public class DefaultDtoQuery extends AbstractQuery implements SpiDtoQuery { private final SpiEbeanServer server; diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultOrmQuery.java b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultOrmQuery.java index 71fe9ff68..0eb5e619c 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultOrmQuery.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultOrmQuery.java @@ -58,7 +58,6 @@ import io.ebeaninternal.server.deploy.TableJoin; import io.ebeaninternal.server.expression.DefaultExpressionList; import io.ebeaninternal.server.expression.IdInExpression; import io.ebeaninternal.server.expression.SimpleExpression; -import io.ebeaninternal.server.query.CancelableQuery; import io.ebeaninternal.server.query.NativeSqlQueryPlanKey; import io.ebeaninternal.server.rawsql.SpiRawSql; import io.ebeaninternal.server.transaction.ExternalJdbcTransaction; @@ -81,7 +80,7 @@ import java.util.stream.Stream; /** * Default implementation of an Object Relational query. */ -public class DefaultOrmQuery implements SpiQuery { +public class DefaultOrmQuery extends AbstractQuery implements SpiQuery { private static final String DEFAULT_QUERY_NAME = "default"; @@ -113,10 +112,6 @@ public class DefaultOrmQuery implements SpiQuery { private ProfilingListener profilingListener; - private boolean cancelled; - - private CancelableQuery cancelableQuery; - private Type type; private String label; @@ -856,6 +851,7 @@ public class DefaultOrmQuery implements SpiQuery { copy.parentNode = parentNode; copy.forUpdate = forUpdate; copy.rawSql = rawSql; + setCancelableQuery(copy); // required to cancel findId query return copy; } @@ -2046,16 +2042,6 @@ public class DefaultOrmQuery implements SpiQuery { return futureFetchAudit; } - @Override - public void setCancelableQuery(CancelableQuery cancelableQuery) { - lock.lock(); - try { - this.cancelableQuery = cancelableQuery; - } finally { - lock.unlock(); - } - } - @Override public Query setBaseTable(String baseTable) { this.baseTable = baseTable; @@ -2083,28 +2069,6 @@ public class DefaultOrmQuery implements SpiQuery { return rootTableAlias != null ? rootTableAlias : defaultAlias; } - @Override - public void cancel() { - lock.lock(); - try { - if (!cancelled && cancelableQuery != null) { - cancelled = true; - cancelableQuery.cancel(); - } - } finally { - lock.unlock(); - } - } - - @Override - public boolean isCancelled() { - lock.lock(); - try { - return cancelled; - } finally { - lock.unlock(); - } - } @Override public Set validate() { diff --git a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java index 182b1cfd4..6afe54a1f 100644 --- a/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java +++ b/ebean-core/src/main/java/io/ebeaninternal/server/querydefn/DefaultRelationalQuery.java @@ -17,7 +17,7 @@ import java.util.function.Predicate; /** * Default implementation of SQuery - SQL Query. */ -public class DefaultRelationalQuery implements SpiSqlQuery { +public class DefaultRelationalQuery extends AbstractQuery implements SpiSqlQuery { private static final long serialVersionUID = -1098305779779591068L; diff --git a/ebean-core/src/test/java/org/tests/model/basic/MyEBasicConfigStartup.java b/ebean-core/src/test/java/org/tests/model/basic/MyEBasicConfigStartup.java index 72cb3e689..c27c5e79b 100644 --- a/ebean-core/src/test/java/org/tests/model/basic/MyEBasicConfigStartup.java +++ b/ebean-core/src/test/java/org/tests/model/basic/MyEBasicConfigStartup.java @@ -59,19 +59,19 @@ public class MyEBasicConfigStartup implements ServerConfigStartup { @Override public void inserted(Object bean) { insertCount.incrementAndGet(); - System.out.println("-- EBasic inserted " + ((EBasic) bean).getId()); + // System.out.println("-- EBasic inserted " + ((EBasic) bean).getId()); } @Override public void updated(Object bean, Set updatedProperties) { updateCount.incrementAndGet(); - System.out.println("-- EBasic updated " + ((EBasic) bean).getId() + " updatedProperties: " + updatedProperties); + // System.out.println("-- EBasic updated " + ((EBasic) bean).getId() + " updatedProperties: " + updatedProperties); } @Override public void deleted(Object bean) { deleteCount.incrementAndGet(); - System.out.println("-- EBasic deleted " + ((EBasic) bean).getId()); + // System.out.println("-- EBasic deleted " + ((EBasic) bean).getId()); } } diff --git a/ebean-core/src/test/java/org/tests/query/cancel/EBasicDto.java b/ebean-core/src/test/java/org/tests/query/cancel/EBasicDto.java new file mode 100644 index 000000000..49b6b157c --- /dev/null +++ b/ebean-core/src/test/java/org/tests/query/cancel/EBasicDto.java @@ -0,0 +1,28 @@ +package org.tests.query.cancel; + +import org.tests.model.basic.EBasic.Status; + +/** + * DTO for Ebasic Queries. + */ +public class EBasicDto { + private Integer id; + + private Status status; + + public Integer getId() { + return id; + } + + public void setId(Integer id) { + this.id = id; + } + + public Status getStatus() { + return status; + } + + public void setStatus(Status status) { + this.status = status; + } +} \ No newline at end of file diff --git a/ebean-core/src/test/java/org/tests/query/cancel/SlowDownEBasic.java b/ebean-core/src/test/java/org/tests/query/cancel/SlowDownEBasic.java new file mode 100644 index 000000000..68d6c038b --- /dev/null +++ b/ebean-core/src/test/java/org/tests/query/cancel/SlowDownEBasic.java @@ -0,0 +1,57 @@ +package org.tests.query.cancel; + +import java.sql.Connection; +import java.sql.SQLException; +import java.sql.Statement; + +import org.h2.api.Trigger; + +import io.ebean.DB; +import io.ebean.Transaction; + +/** + * Class to artificially slow down selects on 'e_basic' table + */ +public class SlowDownEBasic implements Trigger { + + private static int wait; + + private static boolean triggerInstalled; + + @Override + public void init(final Connection conn, final String schemaName, final String triggerName, final String tableName, + final boolean before, final int type) { + } + + @Override + public void fire(final Connection conn, final Object[] oldRow, final Object[] newRow) { + try { + Thread.sleep(wait); + } catch (InterruptedException e) { + // nop + } + } + + @Override + public void close() { + } + + @Override + public void remove() { + } + + + public static void setSelectWaitMillis(final int wait) throws SQLException { + SlowDownEBasic.wait = wait; + if (triggerInstalled) { + return; + } + triggerInstalled = true; + try (Transaction txn = DB.beginTransaction(); Statement stmt = txn.getConnection().createStatement()) { + + stmt.execute("CREATE TRIGGER SLOW_DOWN_E_BASIC BEFORE SELECT ON e_basic " + "CALL \"" + + SlowDownEBasic.class.getName() + "\""); + txn.commit(); + } + } +} diff --git a/ebean-core/src/test/java/org/tests/query/cancel/SqlQueryCancelTest.java b/ebean-core/src/test/java/org/tests/query/cancel/SqlQueryCancelTest.java new file mode 100644 index 000000000..d3dda7443 --- /dev/null +++ b/ebean-core/src/test/java/org/tests/query/cancel/SqlQueryCancelTest.java @@ -0,0 +1,362 @@ +package org.tests.query.cancel; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.sql.SQLException; +import java.util.UUID; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.function.Consumer; +import java.util.function.Function; + +import javax.persistence.PersistenceException; + +import org.junit.BeforeClass; +import org.junit.Test; +import org.tests.model.basic.EBasic; + +import io.ebean.BaseTestCase; +import io.ebean.DB; +import io.ebean.DtoQuery; +import io.ebean.Query; +import io.ebean.QueryIterator; +import io.ebean.SqlQuery; +import io.ebean.annotation.ForPlatform; +import io.ebean.annotation.Platform; + +/** + * Tests, if all kind of queries are cancelable. There are two ways how to + * cancel a query:
+ * At begin: + * + *

+ * query = DB.find(...)
+ * query.cancel();
+ * query.findList();
+ * 
+ * + * The query was caneled before executing. In this case we do hit the DB driver + *
+ *
+ * During run: + * + *
+ * // Thread 1:              Thread 2
+ * query = DB.find(...)
+ * query.findList();
+ *    ...finding
+ *     ...finding            query.cancel();
+ *      ...JDBC-Exception
+ * 
+ * + * The test tries to simulate a slow query by installing the + * {@link SlowDownEBasic} 'SELECT' trigger. The trigger can be configured to + * wait 3 * timing ms and a second thread will cancel the query in + * timing ms. + * + * in this case, we expect a JDBC exception from the driver.
+ *
+ * NOTE:
+ * H2 checks the cancel flag in org.h2.command.Prepared::setCurrentRowNumber + * only every 128th row. So we need at least 128 models and we cannot check + * queries like findCount or findOne, because they only return one row. + * + * @author Roland Praml, FOCONIS AG + * + */ +public class SqlQueryCancelTest extends BaseTestCase { + + private int timing = 10; + + @BeforeClass + public static void setupTestData() throws SQLException { + for (int i = 0; i < 128; i++) { + EBasic model = new EBasic("Basic " + i); + DB.save(model); + } + SlowDownEBasic.setSelectWaitMillis(0); + } + + @Test + public void cancelSqlQueryAtBegin() throws SQLException { + doCancelSqlAtBegin(SqlQuery::findList); + doCancelSqlAtBegin(SqlQuery::findOne); + doCancelSqlAtBegin(q -> q.findEach(e -> {})); + doCancelSqlAtBegin(q -> q.findEachWhile(e -> true)); + } + + @ForPlatform(Platform.H2) + @Test + public void cancelSqlDuringRun() throws SQLException { + + doCancelSqlDuringRun(SqlQuery::findList); + // doCancelSqlDuringRun(q -> q.setMaxRows(1).findOne()); + // findOne cannot be tested, as H2 does the cancel check every 128 rows only + doCancelSqlDuringRun(q -> q.findEach(e -> {})); + doCancelSqlDuringRun(q -> q.findEachWhile(e -> true)); + } + + + @Test + public void cancelOrmQueryAtBegin() throws SQLException { + doCancelOrmAtBegin(Query::findCount); + doCancelOrmAtBegin(Query::findFutureCount); + // We cannot test 'findCount' due H2 restrictions + doCancelOrmAtBegin(Query::findFutureIds); + doCancelOrmAtBegin(Query::findFutureList); + doCancelOrmAtBegin(Query::findIds); + doCancelOrmAtBegin(Query::findIterate); + doCancelOrmAtBegin(Query::findList); + doCancelOrmAtBegin(Query::findMap); + doCancelOrmAtBegin(Query::findOne); + doCancelOrmAtBegin(q -> q.setMaxRows(1000).findPagedList().getList()); // untested + doCancelOrmAtBegin(Query::findSet); + doCancelOrmAtBegin(Query::findSingleAttribute); + doCancelOrmAtBegin(Query::findSingleAttributeList); + doCancelOrmAtBegin(Query::findStream); + // testDuringRun(Query::findVersions); + // EBasic has no history support, but it should work if @History is added + doCancelOrmAtBegin(q -> q.findEach(e -> {})); + doCancelOrmAtBegin(q -> q.findEachWhile(e -> true)); + } + + @ForPlatform(Platform.H2) + @Test + public void cancelOrmDuringRun() throws Throwable { + // doCancelOrmDuringRun(Query::findCount); + // testDuringRunFuture(Query::findFutureCount); + // We cannot test 'findCount' due H2 restrictions + doCancelOrmFutureDuringRun(Query::findFutureIds); + doCancelOrmFutureDuringRun(Query::findFutureList); + doCancelOrmDuringRun(Query::findIds); + doCancelOrmDuringRun(Query::findIterate); + doCancelOrmDuringRun(Query::findList); + doCancelOrmDuringRun(Query::findMap); + // doCancelOrmDuringRun(q -> q.setMaxRows(1).findOne()); + // findOne cannot be tested, as H2 does the cancel check every 128 rows only + doCancelOrmDuringRun(q -> q.setMaxRows(1000).findPagedList().getList()); // untested + doCancelOrmDuringRun(Query::findSet); + doCancelOrmDuringRun(Query::findSingleAttribute); + doCancelOrmDuringRun(Query::findSingleAttributeList); + doCancelOrmDuringRun(Query::findStream); + // testDuringRun(Query::findVersions); + // EBasic has no history support, but it should work if @History is added + doCancelOrmDuringRun(q -> q.findEach(e -> {})); + doCancelOrmDuringRun(q -> q.findEachWhile(e -> true)); + } + + @Test + public void cancelOrmDuringIterate() throws SQLException { + + Query query = DB.find(EBasic.class); + + QueryIterator iter = query.findIterate(); + assertThat(iter.hasNext()).isTrue(); + query.cancel(); + assertThat(iter.next()).isNotNull(); + + // We might have 100 entities in a buffer. So we must iterate through all. + assertThatThrownBy(() -> { + while(iter.hasNext()) iter.next(); + }) + .isInstanceOf(PersistenceException.class) + .hasMessageContaining("Query was cancelled"); + } + + @Test + public void cancelOrmDtoQueryAtBegin() throws SQLException { + + doCancelOrmDtoAtBegin(DtoQuery::findIterate); + doCancelOrmDtoAtBegin(DtoQuery::findList); + doCancelOrmDtoAtBegin(DtoQuery::findOne); + doCancelOrmDtoAtBegin(DtoQuery::findStream); + doCancelOrmDtoAtBegin(q -> q.findEach(e -> {})); + doCancelOrmDtoAtBegin(q -> q.findEachWhile(e -> true)); + } + + @ForPlatform(Platform.H2) + @Test + public void cancelOrmDtoDuringRun() throws SQLException { + + doCancelOrmDtoDuringRun(DtoQuery::findIterate); + doCancelOrmDtoDuringRun(DtoQuery::findList); + // doCancelOrmDtoDuringRun(q -> q.setMaxRows(1).findOne()); + // findOne cannot be tested, as H2 does the cancel check every 128 rows only + doCancelOrmDtoDuringRun(DtoQuery::findStream); + doCancelOrmDtoDuringRun(q -> q.findEach(e -> {})); + doCancelOrmDtoDuringRun(q -> q.findEachWhile(e -> true)); + } + + @Test + public void cancelOrmDtoDuringIterate() throws SQLException { + + DtoQuery query = DB.find(EBasic.class).select("id,status").asDto(EBasicDto.class); + + QueryIterator iter = query.findIterate(); + assertThat(iter.hasNext()).isTrue(); + query.cancel(); + assertThat(iter.next()).isNotNull(); + + // We might have 100 entities in a buffer. So we must iterate through all. + assertThatThrownBy(() -> { + while(iter.hasNext()) iter.next(); + }) + .isInstanceOf(PersistenceException.class) + .hasMessageContaining("Query was cancelled"); + } + + @Test + public void cancelSqlDtoQueryAtBegin() throws SQLException { + + doCancelSqlDtoAtBegin(DtoQuery::findIterate); + doCancelSqlDtoAtBegin(DtoQuery::findList); + doCancelSqlDtoAtBegin(DtoQuery::findOne); + doCancelSqlDtoAtBegin(DtoQuery::findStream); + doCancelSqlDtoAtBegin(q -> q.findEach(e -> {})); + doCancelSqlDtoAtBegin(q -> q.findEachWhile(e -> true)); + } + + @ForPlatform(Platform.H2) + @Test + public void cancelSqlDtoDuringRun() throws SQLException { + + //doCancelSqlDtoDuringRun(DtoQuery::findIterate); + doCancelSqlDtoDuringRun(DtoQuery::findList); + // doCancelSqlDtoDuringRun(q -> q.setMaxRows(1).findOne()); + // findOne cannot be tested, as H2 does the cancel check every 128 rows only + doCancelSqlDtoDuringRun(DtoQuery::findStream); + doCancelSqlDtoDuringRun(q -> q.findEach(e -> {})); + doCancelSqlDtoDuringRun(q -> q.findEachWhile(e -> true)); + } + + @Test + public void cancelSqlDtoDuringIterate() throws SQLException { + + DtoQuery query = DB.findDto(EBasicDto.class, "select id, status from e_basic"); + + QueryIterator iter = query.findIterate(); + assertThat(iter.hasNext()).isTrue(); + query.cancel(); + assertThat(iter.next()).isNotNull(); + + // We might have 100 entities in a buffer. So we must iterate through all. + assertThatThrownBy(() -> { + while(iter.hasNext()) iter.next(); + }) + .isInstanceOf(PersistenceException.class) + .hasMessageContaining("Query was cancelled"); + } + + private void doCancelSqlAtBegin(Consumer test) throws SQLException { + SqlQuery query = DB.sqlQuery("select * from e_basic"); + query.cancel(); + assertThatThrownBy(() -> test.accept(query)) + .isInstanceOf(PersistenceException.class) + .hasMessageContaining("Query was cancelled"); + } + + private void doCancelSqlDuringRun(Consumer test) throws SQLException { + SqlQuery warmup = DB.sqlQuery("select * from e_basic"); + test.accept(warmup); + SqlQuery query = DB.sqlQuery("select * from e_basic"); + executeDelayed(query::cancel); + assertThatThrownBy(() -> test.accept(query)) + .isInstanceOf(PersistenceException.class) + .hasCauseInstanceOf(org.h2.jdbc.JdbcSQLTimeoutException.class); + } + + private void doCancelOrmAtBegin(Consumer> test) throws SQLException { + Query query = DB.find(EBasic.class); + query.cancel(); + assertThatThrownBy(() -> test.accept(query)) + .isInstanceOf(PersistenceException.class) + .hasMessageContaining("Query was cancelled"); + } + + private void doCancelOrmDuringRun(Consumer> test) throws SQLException { + Query warmup = DB.find(EBasic.class); + test.accept(warmup); + Query warmup2 = DB.find(EBasic.class); + test.accept(warmup2); + Query query = DB.find(EBasic.class); + executeDelayed(query::cancel); + assertThatThrownBy(() -> test.accept(query)) + .isInstanceOf(PersistenceException.class) + .hasCauseInstanceOf(org.h2.jdbc.JdbcSQLTimeoutException.class); + } + + private void doCancelOrmFutureDuringRun(Function, Future> test) throws SQLException, InterruptedException, ExecutionException { + Query warmup = DB.find(EBasic.class); + test.apply(warmup).get(); + + Query query = DB.find(EBasic.class); + executeDelayed(query::cancel); + assertThatThrownBy(() -> { + try { + test.apply(query).get(); + } catch (ExecutionException ee) { + throw ee.getCause(); + } + }) + .isInstanceOf(PersistenceException.class) + .hasCauseInstanceOf(org.h2.jdbc.JdbcSQLTimeoutException.class); + } + + private void doCancelOrmDtoAtBegin(Consumer> test) throws SQLException { + DtoQuery query = DB.find(EBasic.class).select("id,status").asDto(EBasicDto.class); + query.cancel(); + assertThatThrownBy(() -> test.accept(query)) + .isInstanceOf(PersistenceException.class) + .hasMessageContaining("Query was cancelled"); + } + + private void doCancelOrmDtoDuringRun(Consumer> test) throws SQLException { + DtoQuery warmup = DB.find(EBasic.class).select("id,status").asDto(EBasicDto.class); + test.accept(warmup); + + DtoQuery query = DB.find(EBasic.class).select("id,status").asDto(EBasicDto.class); + executeDelayed(query::cancel); + assertThatThrownBy(() -> test.accept(query)) + .isInstanceOf(PersistenceException.class) + .hasCauseInstanceOf(org.h2.jdbc.JdbcSQLTimeoutException.class); + } + + private void doCancelSqlDtoAtBegin(Consumer> test) throws SQLException { + DtoQuery query = DB.findDto(EBasicDto.class, "select id, status from e_basic"); + query.cancel(); + assertThatThrownBy(() -> test.accept(query)) + .isInstanceOf(PersistenceException.class) + .hasMessageContaining("Query was cancelled"); + } + + private void doCancelSqlDtoDuringRun(Consumer> test) throws SQLException { + DtoQuery warmup = DB.findDto(EBasicDto.class, "select id, status from e_basic"); + test.accept(warmup); + + DtoQuery query = DB.findDto(EBasicDto.class, "select id, status from e_basic"); + executeDelayed(query::cancel); + assertThatThrownBy(() -> test.accept(query)) + .isInstanceOf(PersistenceException.class) + .hasCauseInstanceOf(org.h2.jdbc.JdbcSQLTimeoutException.class); + } + + private void executeDelayed(Runnable r) throws SQLException { + // We modify the DB here. Otherwise we may hit an internal H2 cache, if the + // same query is performed. Queries from the cache cannot be canceled. + EBasic makeDbDirty = new EBasic("Basic " + UUID.randomUUID()); + DB.save(makeDbDirty); + SlowDownEBasic.setSelectWaitMillis(timing * 3); + new Thread(() -> { + try { + Thread.sleep(timing); + r.run(); + SlowDownEBasic.setSelectWaitMillis(0); + } catch (Exception e) { + e.printStackTrace(); + } + }).start(); + } + +}