diff --git a/src/main/java/com/avaje/ebeaninternal/server/core/DefaultServer.java b/src/main/java/com/avaje/ebeaninternal/server/core/DefaultServer.java index 5c85680ae..bdd9d10c1 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/core/DefaultServer.java +++ b/src/main/java/com/avaje/ebeaninternal/server/core/DefaultServer.java @@ -1340,10 +1340,9 @@ public final class DefaultServer implements SpiEbeanServer { Transaction newTxn = createTransaction(); CallableQueryRowCount call = new CallableQueryRowCount(this, copy, newTxn); - FutureTask futureTask = new FutureTask(call); - QueryFutureRowCount queryFuture = new QueryFutureRowCount(copy, futureTask); - backgroundExecutor.execute(futureTask); + QueryFutureRowCount queryFuture = new QueryFutureRowCount(call); + backgroundExecutor.execute(queryFuture.getFutureTask()); return queryFuture; } @@ -1362,11 +1361,9 @@ public final class DefaultServer implements SpiEbeanServer { Transaction newTxn = createTransaction(); CallableQueryIds call = new CallableQueryIds(this, copy, newTxn); - FutureTask> futureTask = new FutureTask>(call); + QueryFutureIds queryFuture = new QueryFutureIds(call); - QueryFutureIds queryFuture = new QueryFutureIds(copy, futureTask); - - backgroundExecutor.execute(futureTask); + backgroundExecutor.execute(queryFuture.getFutureTask()); return queryFuture; } @@ -1376,6 +1373,7 @@ public final class DefaultServer implements SpiEbeanServer { SpiQuery spiQuery = (SpiQuery) query; spiQuery.setFutureFetch(true); + // transfer the persistence content from the transaction if (spiQuery.getPersistenceContext() == null) { if (t != null) { spiQuery.setPersistenceContext(((SpiTransaction) t).getPersistenceContext()); @@ -1387,14 +1385,13 @@ public final class DefaultServer implements SpiEbeanServer { } } + // Create a new transaction solely to execute the findList() at some future time Transaction newTxn = createTransaction(); - CallableQueryList call = new CallableQueryList(this, query, newTxn); + CallableQueryList call = new CallableQueryList(this, spiQuery, newTxn); + QueryFutureList queryFuture = new QueryFutureList(call); + backgroundExecutor.execute(queryFuture.getFutureTask()); - FutureTask> futureTask = new FutureTask>(call); - - backgroundExecutor.execute(futureTask); - - return new QueryFutureList(query, futureTask); + return queryFuture; } public PagingList findPagingList(Query query, Transaction t, int pageSize) { diff --git a/src/main/java/com/avaje/ebeaninternal/server/query/BaseFuture.java b/src/main/java/com/avaje/ebeaninternal/server/query/BaseFuture.java index e2067989e..a0854cae7 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/query/BaseFuture.java +++ b/src/main/java/com/avaje/ebeaninternal/server/query/BaseFuture.java @@ -15,7 +15,7 @@ import java.util.concurrent.TimeoutException; */ public abstract class BaseFuture implements Future { - private final FutureTask futureTask; + protected final FutureTask futureTask; public BaseFuture(FutureTask futureTask) { this.futureTask = futureTask; @@ -43,5 +43,4 @@ public abstract class BaseFuture implements Future { return futureTask.isDone(); } - } diff --git a/src/main/java/com/avaje/ebeaninternal/server/query/CallableQuery.java b/src/main/java/com/avaje/ebeaninternal/server/query/CallableQuery.java index ca56a81f4..b15c07dab 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/query/CallableQuery.java +++ b/src/main/java/com/avaje/ebeaninternal/server/query/CallableQuery.java @@ -1,8 +1,8 @@ package com.avaje.ebeaninternal.server.query; -import com.avaje.ebean.Query; import com.avaje.ebean.Transaction; import com.avaje.ebeaninternal.api.SpiEbeanServer; +import com.avaje.ebeaninternal.api.SpiQuery; /** * Base object for making query execution into Callable's. @@ -13,16 +13,24 @@ import com.avaje.ebeaninternal.api.SpiEbeanServer; */ public abstract class CallableQuery { - protected final Query query; + protected final SpiQuery query; protected final SpiEbeanServer server; - protected final Transaction t; + protected final Transaction transaction; - public CallableQuery(SpiEbeanServer server, Query query, Transaction t) { + public CallableQuery(SpiEbeanServer server, SpiQuery query, Transaction t) { this.server = server; this.query = query; - this.t = t; + this.transaction = t; } + + public SpiQuery getQuery() { + return query; + } + + public Transaction getTransaction() { + return transaction; + } } diff --git a/src/main/java/com/avaje/ebeaninternal/server/query/CallableQueryIds.java b/src/main/java/com/avaje/ebeaninternal/server/query/CallableQueryIds.java index 18b404ea9..8e58af8a8 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/query/CallableQueryIds.java +++ b/src/main/java/com/avaje/ebeaninternal/server/query/CallableQueryIds.java @@ -3,9 +3,9 @@ package com.avaje.ebeaninternal.server.query; import java.util.List; import java.util.concurrent.Callable; -import com.avaje.ebean.Query; import com.avaje.ebean.Transaction; import com.avaje.ebeaninternal.api.SpiEbeanServer; +import com.avaje.ebeaninternal.api.SpiQuery; /** * Represent the fetch Id's query as a Callable. @@ -15,7 +15,7 @@ import com.avaje.ebeaninternal.api.SpiEbeanServer; public class CallableQueryIds extends CallableQuery implements Callable> { - public CallableQueryIds(SpiEbeanServer server, Query query, Transaction t) { + public CallableQueryIds(SpiEbeanServer server, SpiQuery query, Transaction t) { super(server, query, t); } @@ -26,7 +26,11 @@ public class CallableQueryIds extends CallableQuery implements Callable
  • extends CallableQuery implements Callable> { - public CallableQueryList(SpiEbeanServer server, Query query, Transaction t) { + public CallableQueryList(SpiEbeanServer server, SpiQuery query, Transaction t) { super(server, query, t); } @@ -23,7 +23,12 @@ public class CallableQueryList extends CallableQuery implements Callable call() throws Exception { - return server.findList(query, t); + try { + return server.findList(query, transaction); + } finally { + // cleanup the underlying connection + transaction.end(); + } } diff --git a/src/main/java/com/avaje/ebeaninternal/server/query/CallableQueryRowCount.java b/src/main/java/com/avaje/ebeaninternal/server/query/CallableQueryRowCount.java index 55d3a0b17..e64536363 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/query/CallableQueryRowCount.java +++ b/src/main/java/com/avaje/ebeaninternal/server/query/CallableQueryRowCount.java @@ -2,29 +2,36 @@ package com.avaje.ebeaninternal.server.query; import java.util.concurrent.Callable; -import com.avaje.ebean.Query; import com.avaje.ebean.Transaction; import com.avaje.ebeaninternal.api.SpiEbeanServer; +import com.avaje.ebeaninternal.api.SpiQuery; /** * Represent the findRowCount query as a Callable. - * - * @param the entity bean type + * + * @param + * the entity bean type */ public class CallableQueryRowCount extends CallableQuery implements Callable { - - public CallableQueryRowCount(SpiEbeanServer server, Query query, Transaction t) { - super(server, query, t); - } - - /** - * Execute the query returning the row count. - */ - public Integer call() throws Exception { - return server.findRowCountWithCopy(query, t); - } + /** + * Note that the transaction passed in is always a new transaction solely to + * find the row count so it must be cleaned up by this CallableQueryRowCount. + */ + public CallableQueryRowCount(SpiEbeanServer server, SpiQuery query, Transaction t) { + super(server, query, t); + } + + /** + * Execute the query returning the row count. + */ + public Integer call() throws Exception { + try { + return server.findRowCountWithCopy(query, transaction); + } finally { + // cleanup the underlying connection + transaction.end(); + } + } - - } diff --git a/src/main/java/com/avaje/ebeaninternal/server/query/CallableSqlQueryList.java b/src/main/java/com/avaje/ebeaninternal/server/query/CallableSqlQueryList.java index 70a3f4864..1d476ffb9 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/query/CallableSqlQueryList.java +++ b/src/main/java/com/avaje/ebeaninternal/server/query/CallableSqlQueryList.java @@ -17,19 +17,23 @@ public class CallableSqlQueryList implements Callable> { private final EbeanServer server; - private final Transaction t; + private final Transaction transaction; public CallableSqlQueryList(EbeanServer server, SqlQuery query, Transaction t) { this.server = server; this.query = query; - this.t = t; + this.transaction = t; } /** * Execute the query returning the resulting list. */ public List call() throws Exception { - return server.findList(query, t); + try { + return server.findList(query, transaction); + } finally { + transaction.end(); + } } diff --git a/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureIds.java b/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureIds.java index cf9877154..1936c4808 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureIds.java +++ b/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureIds.java @@ -5,30 +5,38 @@ import java.util.concurrent.FutureTask; import com.avaje.ebean.FutureIds; import com.avaje.ebean.Query; -import com.avaje.ebeaninternal.api.SpiQuery; +import com.avaje.ebean.Transaction; /** * Default implementation of FutureIds. */ public class QueryFutureIds extends BaseFuture> implements FutureIds { - private final SpiQuery query; + private final CallableQueryIds call; - public QueryFutureIds(SpiQuery query, FutureTask> futureTask) { - super(futureTask); - this.query = query; + public QueryFutureIds(CallableQueryIds call ) { + super(new FutureTask>(call)); + this.call = call; + } + + public FutureTask> getFutureTask() { + return futureTask; + } + + public Transaction getTransaction() { + return call.transaction; } public Query getQuery() { - return query; + return call.query; } public List getPartialIds() { - return query.getIdList(); + return call.query.getIdList(); } public boolean cancel(boolean mayInterruptIfRunning) { - query.cancel(); + call.query.cancel(); return super.cancel(mayInterruptIfRunning); } diff --git a/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureList.java b/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureList.java index ce7e93c0a..6424b9faa 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureList.java +++ b/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureList.java @@ -5,26 +5,34 @@ import java.util.concurrent.FutureTask; import com.avaje.ebean.FutureList; import com.avaje.ebean.Query; +import com.avaje.ebean.Transaction; /** * Default implementation for FutureList. */ public class QueryFutureList extends BaseFuture> implements FutureList { - private final Query query; + private final CallableQueryList call; + public QueryFutureList(CallableQueryList call) { + super(new FutureTask>(call)); + this.call = call; + } - public QueryFutureList(Query query, FutureTask> futureTask) { - super(futureTask); - this.query = query; + public FutureTask> getFutureTask() { + return futureTask; + } + + public Transaction getTransaction() { + return call.transaction; } public Query getQuery() { - return query; + return call.query; } public boolean cancel(boolean mayInterruptIfRunning) { - query.cancel(); + call.query.cancel(); return super.cancel(mayInterruptIfRunning); } diff --git a/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureRowCount.java b/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureRowCount.java index fee28f74e..a98300c1a 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureRowCount.java +++ b/src/main/java/com/avaje/ebeaninternal/server/query/QueryFutureRowCount.java @@ -4,25 +4,34 @@ import java.util.concurrent.FutureTask; import com.avaje.ebean.FutureRowCount; import com.avaje.ebean.Query; +import com.avaje.ebean.Transaction; /** * Future implementation for the row count query. */ public class QueryFutureRowCount extends BaseFuture implements FutureRowCount { - private final Query query; + private final CallableQueryRowCount call; - public QueryFutureRowCount(Query query, FutureTask futureTask) { - super(futureTask); - this.query = query; + public QueryFutureRowCount(CallableQueryRowCount call ) { + super(new FutureTask(call)); + this.call = call; + } + + public FutureTask getFutureTask() { + return futureTask; + } + + public Transaction getTransaction() { + return call.transaction; } public Query getQuery() { - return query; + return call.query; } public boolean cancel(boolean mayInterruptIfRunning) { - query.cancel(); + call.query.cancel(); return super.cancel(mayInterruptIfRunning); } diff --git a/src/test/java/com/avaje/ebeaninternal/server/query/TestFutureRowCountErrorHandling.java b/src/test/java/com/avaje/ebeaninternal/server/query/TestFutureRowCountErrorHandling.java new file mode 100644 index 000000000..278466bec --- /dev/null +++ b/src/test/java/com/avaje/ebeaninternal/server/query/TestFutureRowCountErrorHandling.java @@ -0,0 +1,102 @@ +package com.avaje.ebeaninternal.server.query; + +import java.util.concurrent.ExecutionException; + +import org.junit.Assert; +import org.junit.Test; + +import com.avaje.ebean.BaseTestCase; +import com.avaje.ebean.Ebean; +import com.avaje.ebean.EbeanServer; +import com.avaje.ebean.FutureIds; +import com.avaje.ebean.FutureList; +import com.avaje.ebean.FutureRowCount; +import com.avaje.ebean.Query; +import com.avaje.ebean.Transaction; +import com.avaje.tests.model.basic.Customer; +import com.avaje.tests.model.basic.ResetBasicData; + +public class TestFutureRowCountErrorHandling extends BaseTestCase { + + @Test + public void testFutureRowCount() throws InterruptedException { + + ResetBasicData.reset(); + + EbeanServer server = Ebean.getServer(null); + + Query query = server.createQuery(Customer.class) + .where().eq("doesNotExist", "this will fail") + .query(); + + FutureRowCount futureRowCount = server.findFutureRowCount(query, null); + + QueryFutureRowCount internalRowCount = (QueryFutureRowCount)futureRowCount; + Transaction t = internalRowCount.getTransaction(); + + try { + futureRowCount.get(); + Assert.assertTrue("never get here as the SQL is invalid",false); + + } catch (ExecutionException e) { + // Confirm the Transaction has been rolled back + Assert.assertFalse("Underlying transaction was rolled back cleanly", t.isActive()); + } + + } + + @Test + public void testFutureIds() throws InterruptedException { + + ResetBasicData.reset(); + + EbeanServer server = Ebean.getServer(null); + + Query query = server.createQuery(Customer.class) + .where().eq("doesNotExist", "this will fail") + .query(); + + FutureIds futureIds = server.findFutureIds(query, null); + + QueryFutureIds internalFuture = (QueryFutureIds)futureIds; + Transaction t = internalFuture.getTransaction(); + + try { + internalFuture.get(); + Assert.assertTrue("never get here as the SQL is invalid",false); + + } catch (ExecutionException e) { + // Confirm the Transaction has been rolled back + Assert.assertFalse("Underlying transaction was rolled back cleanly", t.isActive()); + } + + } + + + @Test + public void testFutureList() throws InterruptedException { + + ResetBasicData.reset(); + + EbeanServer server = Ebean.getServer(null); + + Query query = server.createQuery(Customer.class) + .where().eq("doesNotExist", "this will fail") + .query(); + + FutureList futureList = server.findFutureList(query, null); + + QueryFutureList internalFuture = (QueryFutureList)futureList; + Transaction t = internalFuture.getTransaction(); + + try { + internalFuture.get(); + Assert.assertTrue("never get here as the SQL is invalid",false); + + } catch (ExecutionException e) { + // Confirm the Transaction has been rolled back + Assert.assertFalse("Underlying transaction was rolled back cleanly", t.isActive()); + } + + } +}