This is an automated email from the ASF dual-hosted git repository.
hansva pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new 525f1ca8e2 Issue #8568 : Do not keep writing execution info through a
closed database connection (#8570)
525f1ca8e2 is described below
commit 525f1ca8e246f88f3ba091227dce3e2bd5f0743c
Author: Matt Casters <[email protected]>
AuthorDate: Thu Sep 24 14:56:39 2026 +0200
Issue #8568 : Do not keep writing execution info through a closed database
connection (#8570)
---
.../caching/BaseCachingExecutionInfoLocation.java | 67 ++++-
.../engines/local/LocalWorkflowEngine.java | 26 +-
.../CachingDatabaseExecutionInfoLocation.java | 294 ++++++++++++++++-----
.../CachingDatabaseExecutionInfoLocationTest.java | 116 ++++++++
4 files changed, 423 insertions(+), 80 deletions(-)
diff --git
a/engine/src/main/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocation.java
b/engine/src/main/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocation.java
index af2d03314e..dc69ef19e7 100644
---
a/engine/src/main/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocation.java
+++
b/engine/src/main/java/org/apache/hop/execution/caching/BaseCachingExecutionInfoLocation.java
@@ -31,6 +31,7 @@ import java.util.Set;
import java.util.Timer;
import java.util.TimerTask;
import java.util.concurrent.atomic.AtomicBoolean;
+import lombok.AccessLevel;
import lombok.Getter;
import lombok.Setter;
import org.apache.commons.lang3.StringUtils;
@@ -84,6 +85,22 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
protected Timer cacheTimer;
+ /**
+ * Set when {@link #close()} has cancelled the timer. A timer task that
already passed {@code
+ * schedule} must not persist after that, and {@link #initialize} starts a
fresh timer.
+ */
+ @Getter(AccessLevel.NONE)
+ @Setter(AccessLevel.NONE)
+ private volatile boolean cacheClosed = true;
+
+ @Getter(AccessLevel.NONE)
+ @Setter(AccessLevel.NONE)
+ private boolean loggedManageCacheError;
+
+ @Getter(AccessLevel.NONE)
+ @Setter(AccessLevel.NONE)
+ private String lastManageCacheError;
+
protected final AtomicBoolean locked;
protected int delay;
@@ -131,17 +148,27 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
//
maxAge = Const.toInt(variables.resolve(maxCacheAge), 86400000);
- // Let's start a timer to manage the cache every second or so
+ // Let's start a timer to manage the cache every second or so.
+ // Cancel any previous timer first: a second initialize() used to leave
the old one running
+ // after close() disconnected the only connection that timer still wrote
through.
//
- cacheTimer = new Timer("Caching execution location timer");
- TimerTask cacheManageTask =
- new TimerTask() {
- @Override
- public void run() {
- manageCache();
- }
- };
- cacheTimer.schedule(cacheManageTask, 1000L, 1000L);
+ synchronized (this) {
+ cacheClosed = false;
+ loggedManageCacheError = false;
+ lastManageCacheError = null;
+ Timer previousTimer = cacheTimer;
+ cacheTimer = null;
+ ExecutorUtil.cleanup(previousTimer);
+ cacheTimer = new Timer("Caching execution location timer", true);
+ TimerTask cacheManageTask =
+ new TimerTask() {
+ @Override
+ public void run() {
+ manageCache();
+ }
+ };
+ cacheTimer.schedule(cacheManageTask, 1000L, 1000L);
+ }
}
@Override
@@ -150,6 +177,9 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
}
protected synchronized void manageCache() {
+ if (cacheClosed) {
+ return;
+ }
try {
// Let's make sure we never run this method in parallel
//
@@ -177,8 +207,16 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
// Remove these entries.
//
tooOld.forEach(id -> cache.remove(id));
+ loggedManageCacheError = false;
+ lastManageCacheError = null;
} catch (Exception e) {
- LogChannel.GENERAL.logError("Error managing file execution information
location cache", e);
+ // A dead JDBC connection used to log a full stack trace from this timer
once a second.
+ String message = e.getMessage();
+ if (!loggedManageCacheError || !StringUtils.equals(message,
lastManageCacheError)) {
+ LogChannel.GENERAL.logError("Error managing execution information
location cache", e);
+ loggedManageCacheError = true;
+ lastManageCacheError = message;
+ }
} finally {
locked.set(false);
}
@@ -186,8 +224,13 @@ public abstract class BaseCachingExecutionInfoLocation
implements IExecutionInfo
@Override
public synchronized void close() throws HopException {
+ // Stop the timer before the final flush. Tasks that are already inside
manageCache hold this
+ // lock and finish first; tasks still queued see cacheClosed and return.
+ cacheClosed = true;
+ Timer timer = cacheTimer;
+ cacheTimer = null;
+ ExecutorUtil.cleanup(timer);
try {
- ExecutorUtil.cleanup(cacheTimer);
for (CacheEntry cacheEntry : cache.values()) {
if (cacheEntry.isDirty()) {
persistCacheEntry(cacheEntry);
diff --git
a/engine/src/main/java/org/apache/hop/workflow/engines/local/LocalWorkflowEngine.java
b/engine/src/main/java/org/apache/hop/workflow/engines/local/LocalWorkflowEngine.java
index 632437e763..951a67310f 100644
---
a/engine/src/main/java/org/apache/hop/workflow/engines/local/LocalWorkflowEngine.java
+++
b/engine/src/main/java/org/apache/hop/workflow/engines/local/LocalWorkflowEngine.java
@@ -260,7 +260,17 @@ public class LocalWorkflowEngine extends Workflow
implements IWorkflowEngine<Wor
startExecutionInfoTimer();
});
- return super.startExecution();
+ try {
+ return super.startExecution();
+ } finally {
+ // Finished listeners are not guaranteed to run to the end: an earlier
listener that throws
+ // skips stopExecutionInfoTimer(), and the cache timer then keeps
writing. Close here too.
+ try {
+ stopExecutionInfoTimer();
+ } catch (Exception e) {
+ log.logError("Error closing execution information location after
workflow execution", e);
+ }
+ }
}
/** This method looks up the execution information location specified in the
run configuration. */
@@ -503,16 +513,20 @@ public class LocalWorkflowEngine extends Workflow
implements IWorkflowEngine<Wor
}
}
- public void stopExecutionInfoTimer() throws HopException {
+ public synchronized void stopExecutionInfoTimer() throws HopException {
ExecutorUtil.cleanup(executionInfoTimer);
+ executionInfoTimer = null;
- if (executionInfoLocation == null) {
+ ExecutionInfoLocation location = executionInfoLocation;
+ // Claim it so the finished listener and the startExecution() finally do
not both flush and
+ // close, and so a second run cannot observe this location while it is
being closed.
+ executionInfoLocation = null;
+ if (location == null || location.getExecutionInfoLocation() == null) {
return;
}
+ IExecutionInfoLocation iLocation = location.getExecutionInfoLocation();
try {
- IExecutionInfoLocation iLocation =
executionInfoLocation.getExecutionInfoLocation();
-
// Register one final last state of the workflow
//
ExecutionState executionState =
@@ -522,7 +536,7 @@ public class LocalWorkflowEngine extends Workflow
implements IWorkflowEngine<Wor
} finally {
// Nothing more needs to be done. We can now close the location.
//
- executionInfoLocation.getExecutionInfoLocation().close();
+ iLocation.close();
}
}
}
diff --git
a/plugins/misc/execution-database/src/main/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocation.java
b/plugins/misc/execution-database/src/main/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocation.java
index b6c413b295..f790666222 100644
---
a/plugins/misc/execution-database/src/main/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocation.java
+++
b/plugins/misc/execution-database/src/main/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocation.java
@@ -18,7 +18,9 @@
package org.apache.hop.execution.database;
import com.fasterxml.jackson.databind.ObjectMapper;
+import java.sql.Connection;
import java.sql.ResultSet;
+import java.sql.SQLException;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.util.ArrayList;
@@ -26,6 +28,7 @@ import java.util.Date;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import lombok.AccessLevel;
import lombok.Getter;
import lombok.Setter;
import org.apache.commons.lang3.StringUtils;
@@ -132,6 +135,16 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
protected transient Database database;
+ /** Guards the live {@link #database}. Not the connection itself: that
object is replaced. */
+ @Getter(AccessLevel.NONE)
+ @Setter(AccessLevel.NONE)
+ private final Object dbLock = new Object();
+
+ /** True after {@link #close()} so a late timer tick cannot open another
connection. */
+ @Getter(AccessLevel.NONE)
+ @Setter(AccessLevel.NONE)
+ private volatile boolean databaseClosed = true;
+
protected String actualConnectionName;
protected String actualSchemaName;
protected String actualTableName;
@@ -157,7 +170,7 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
}
@Override
- public void initialize(IVariables variables, IHopMetadataProvider
metadataProvider)
+ public synchronized void initialize(IVariables variables,
IHopMetadataProvider metadataProvider)
throws HopException {
this.variables = variables;
this.metadataProvider = metadataProvider;
@@ -195,20 +208,25 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
}
}
+ synchronized (dbLock) {
+ databaseClosed = false;
+ discardDatabase();
+ connectDatabase();
+ }
+
try {
- database =
- new Database(
- new LoggingObject("CachingDatabaseExecutionInfoLocation"),
variables, databaseMeta);
- database.connect();
+ super.initialize(variables, metadataProvider);
} catch (Exception e) {
+ synchronized (dbLock) {
+ databaseClosed = true;
+ discardDatabase();
+ }
+ if (e instanceof HopException hopException) {
+ throw hopException;
+ }
throw new HopException(
- "Error connecting to database for execution information location
using connection '"
- + Const.NVL(actualConnectionName, databaseMeta.getName())
- + "'",
- e);
+ "Error starting the caching database execution information
location", e);
}
-
- super.initialize(variables, metadataProvider);
LogChannel.GENERAL.logBasic(
"Caching database execution info location ready: connection="
+ Const.NVL(actualConnectionName, databaseMeta.getName())
@@ -218,23 +236,168 @@ public class CachingDatabaseExecutionInfoLocation
extends BaseCachingExecutionIn
@Override
public synchronized void close() throws HopException {
+ if (databaseClosed) {
+ return;
+ }
try {
super.close();
} finally {
- if (database != null) {
- try {
- database.disconnect();
- } catch (Exception e) {
- LogChannel.GENERAL.logError(
- "Error disconnecting database for execution information
location", e);
+ synchronized (dbLock) {
+ databaseClosed = true;
+ discardDatabase();
+ }
+ }
+ }
+
+ /**
+ * Open an independent auto-commit connection. A workflow transaction
connection is closed when
+ * the workflow stops, which left the cache timer writing through that
closed connection.
+ */
+ private void connectDatabase() throws HopException {
+ Database db =
+ new Database(
+ new LoggingObject("CachingDatabaseExecutionInfoLocation"),
variables, databaseMeta);
+ db.setConnectionGroup(null);
+ try {
+ db.connect();
+ db.setConnectionGroup(null);
+ enableAutoCommit(db);
+ } catch (Exception e) {
+ try {
+ db.setConnectionGroup(null);
+ db.disconnect();
+ } catch (Exception disconnectError) {
+ // The original connect error is the one that matters.
+ }
+ throw new HopException(
+ "Error connecting to database for execution information location
using connection '"
+ + Const.NVL(actualConnectionName, databaseMeta.getName())
+ + "'",
+ e);
+ }
+ database = db;
+ }
+
+ /**
+ * Keep writes out of a workflow transaction. Drivers that reject the change
keep the connection;
+ * the connection group was already cleared, which is what stops the
workflow from closing it.
+ */
+ private static void enableAutoCommit(Database db) {
+ try {
+ Connection connection = db.getConnection();
+ if (connection != null && !connection.getAutoCommit()) {
+ db.setAutoCommit(true);
+ }
+ } catch (Exception e) {
+ LogChannel.GENERAL.logBasic(
+ "Execution information connection stays with the driver's commit
mode: "
+ + e.getMessage());
+ }
+ }
+
+ /** Drop the current connection. A grouped connection is left for the
workflow to close. */
+ private void discardDatabase() {
+ Database current = database;
+ database = null;
+ if (current == null ||
StringUtils.isNotEmpty(current.getConnectionGroup())) {
+ return;
+ }
+ try {
+ current.disconnect();
+ } catch (Exception e) {
+ LogChannel.GENERAL.logError(
+ "Error disconnecting database for execution information location",
e);
+ }
+ }
+
+ private void ensureConnected() throws HopException {
+ if (isConnectionUsable(database) && !databaseClosed) {
+ return;
+ }
+ synchronized (dbLock) {
+ if (databaseClosed) {
+ throw new HopException("Caching database execution information
location is closed");
+ }
+ if (isConnectionUsable(database)) {
+ return;
+ }
+ discardDatabase();
+ connectDatabase();
+ }
+ }
+
+ private static boolean isConnectionUsable(Database db) {
+ if (db == null) {
+ return false;
+ }
+ try {
+ Connection connection = db.getConnection();
+ return connection != null && !connection.isClosed();
+ } catch (SQLException e) {
+ return false;
+ }
+ }
+
+ /**
+ * PostgreSQL reports a client-side close as {@code 08003} / "This
connection has been closed."
+ */
+ static boolean isClosedConnectionFailure(Throwable error) {
+ Throwable current = error;
+ while (current != null) {
+ if (current instanceof SQLException sqlException) {
+ String state = sqlException.getSQLState();
+ if ("08003".equals(state) || "08006".equals(state) ||
"57P01".equals(state)) {
+ return true;
}
- database = null;
}
+ String message = current.getMessage();
+ if (message != null) {
+ String lower = message.toLowerCase();
+ if (lower.contains("connection")
+ && (lower.contains("closed") || lower.contains("broken"))) {
+ return true;
+ }
+ }
+ current = current.getCause();
+ }
+ return false;
+ }
+
+ @FunctionalInterface
+ private interface DatabaseWork<T> {
+ T run() throws Exception;
+ }
+
+ private <T> T callWithDatabase(DatabaseWork<T> work) throws HopException {
+ return callWithDatabase(work, true);
+ }
+
+ private <T> T callWithDatabase(DatabaseWork<T> work, boolean allowRetry)
throws HopException {
+ try {
+ ensureConnected();
+ synchronized (dbLock) {
+ ensureConnected();
+ return work.run();
+ }
+ } catch (Exception e) {
+ if (allowRetry && !databaseClosed && isClosedConnectionFailure(e)) {
+ synchronized (dbLock) {
+ discardDatabase();
+ }
+ return callWithDatabase(work, false);
+ }
+ if (e instanceof HopException hopException) {
+ throw hopException;
+ }
+ throw new HopException("Error accessing the execution information
database", e);
}
}
@Override
protected void persistCacheEntry(CacheEntry cacheEntry) throws HopException {
+ if (databaseClosed) {
+ throw new HopException("Caching database execution information location
is closed");
+ }
try {
mergeChildrenFromDatabase(cacheEntry);
cacheEntry.calculateSummary();
@@ -245,9 +408,11 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
IRowMeta rowMeta = createDataRowMeta();
Object[] data = buildRowData(cacheEntry, json);
- synchronized (database) {
- upsertCacheEntry(rowMeta, data);
- }
+ callWithDatabase(
+ () -> {
+ upsertCacheEntry(rowMeta, data);
+ return null;
+ });
cacheEntry.setDirty(false);
cacheEntry.setLastWritten(new Date());
@@ -370,19 +535,20 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
IRowMeta paramMeta = new RowMeta();
paramMeta.addValueMeta(new ValueMetaString(COL_ID, 100, -1));
- synchronized (database) {
- RowMetaAndData row = database.getOneRow(sql, paramMeta, new Object[]
{executionId});
- if (row == null || row.getData() == null) {
- return null;
- }
- Object jsonObj = row.getData()[0];
- if (jsonObj == null) {
- return null;
- }
- String json = jsonObj.toString();
- ObjectMapper mapper = new ObjectMapper();
- return mapper.readValue(json, CacheEntry.class);
- }
+ return callWithDatabase(
+ () -> {
+ RowMetaAndData row = database.getOneRow(sql, paramMeta, new
Object[] {executionId});
+ if (row == null || row.getData() == null) {
+ return null;
+ }
+ Object jsonObj = row.getData()[0];
+ if (jsonObj == null) {
+ return null;
+ }
+ String json = jsonObj.toString();
+ ObjectMapper mapper = new ObjectMapper();
+ return mapper.readValue(json, CacheEntry.class);
+ });
} catch (Exception e) {
throw new HopException(
"Error loading execution information from database for executionId
'" + executionId + "'",
@@ -404,9 +570,11 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
+ " = ?";
IRowMeta paramMeta = new RowMeta();
paramMeta.addValueMeta(new ValueMetaString(COL_ID, 100, -1));
- synchronized (database) {
- database.execStatement(sql, paramMeta, new Object[]
{cacheEntry.getId()});
- }
+ callWithDatabase(
+ () -> {
+ database.execStatement(sql, paramMeta, new Object[]
{cacheEntry.getId()});
+ return null;
+ });
} catch (Exception e) {
throw new HopException(
"Error deleting execution information from database for id '" +
cacheEntry.getId() + "'",
@@ -447,34 +615,36 @@ public class CachingDatabaseExecutionInfoLocation extends
BaseCachingExecutionIn
sql.append(databaseMeta.getLimitClause(limit));
}
- synchronized (database) {
- ResultSet rs = database.openQuery(sql.toString(), paramMeta,
params.toArray());
- try {
- Object[] row = database.getRow(rs);
- while (row != null) {
- String id = row[0] != null ? row[0].toString() : null;
- Date startDate = null;
- if (row[1] instanceof Date date) {
- startDate = date;
- } else if (row[1] != null) {
- // Timestamp / other
- startDate = (Date) row[1];
- }
- if (id != null) {
- ids.add(new DatedId(id, startDate != null ? startDate : new
Date(0L)));
- if (includeChildren && !activeSelector.isSelectingParents()) {
- CacheEntry entry = loadCacheEntry(id);
- if (entry != null) {
- addChildIds(entry, ids, activeSelector);
+ callWithDatabase(
+ () -> {
+ ResultSet rs = database.openQuery(sql.toString(), paramMeta,
params.toArray());
+ try {
+ Object[] row = database.getRow(rs);
+ while (row != null) {
+ String id = row[0] != null ? row[0].toString() : null;
+ Date startDate = null;
+ if (row[1] instanceof Date date) {
+ startDate = date;
+ } else if (row[1] != null) {
+ // Timestamp / other
+ startDate = (Date) row[1];
+ }
+ if (id != null) {
+ ids.add(new DatedId(id, startDate != null ? startDate : new
Date(0L)));
+ if (includeChildren && !activeSelector.isSelectingParents())
{
+ CacheEntry entry = loadCacheEntry(id);
+ if (entry != null) {
+ addChildIds(entry, ids, activeSelector);
+ }
+ }
}
+ row = database.getRow(rs);
}
+ } finally {
+ database.closeQuery(rs);
}
- row = database.getRow(rs);
- }
- } finally {
- database.closeQuery(rs);
- }
- }
+ return null;
+ });
} catch (Exception e) {
throw new HopException(
"Error finding execution ids from database table " +
getQuotedSchemaTable(), e);
diff --git
a/plugins/misc/execution-database/src/test/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocationTest.java
b/plugins/misc/execution-database/src/test/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocationTest.java
index 022bff0a46..04031320f2 100644
---
a/plugins/misc/execution-database/src/test/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocationTest.java
+++
b/plugins/misc/execution-database/src/test/java/org/apache/hop/execution/database/CachingDatabaseExecutionInfoLocationTest.java
@@ -20,18 +20,25 @@ package org.apache.hop.execution.database;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import java.sql.Connection;
+import java.sql.SQLException;
import java.util.Date;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
+import java.util.Timer;
+import java.util.TimerTask;
import java.util.UUID;
import org.apache.hop.core.HopClientEnvironment;
import org.apache.hop.core.database.Database;
import org.apache.hop.core.database.DatabaseMeta;
import org.apache.hop.core.database.DatabasePluginType;
+import org.apache.hop.core.exception.HopException;
import org.apache.hop.core.logging.LoggingObject;
import org.apache.hop.core.variables.Variables;
import org.apache.hop.databases.h2.H2DatabaseMeta;
@@ -252,6 +259,115 @@ class CachingDatabaseExecutionInfoLocationTest {
assertEquals(newId, ids.get(0));
}
+ @Test
+ void persistReopensAClosedJdbcConnection() throws Exception {
+ String id = UUID.randomUUID().toString();
+ location.persistCacheEntry(sampleEntry(id, "Before",
ExecutionType.Workflow, false, "Running"));
+
+ Connection closed = location.getDatabase().getConnection();
+ closed.close();
+ assertTrue(closed.isClosed());
+
+ location.persistCacheEntry(sampleEntry(id, "Before",
ExecutionType.Workflow, true, "Finished"));
+
+ Connection reopened = location.getDatabase().getConnection();
+ assertNotSame(closed, reopened);
+ assertFalse(reopened.isClosed());
+ assertNull(location.getDatabase().getConnectionGroup());
+ assertTrue(reopened.getAutoCommit());
+
+ CacheEntry loaded = location.loadCacheEntry(id);
+ assertNotNull(loaded);
+ assertTrue(loaded.getExecutionState().isFailed());
+
assertTrue(loaded.getExecutionState().getStatusDescription().startsWith("Finished"));
+ }
+
+ @Test
+ void initializeCancelsThePreviousTimerAndDropsTheClosedConnection() throws
Exception {
+ Timer firstTimer = location.getCacheTimer();
+ assertNotNull(firstTimer);
+ Connection firstConnection = location.getDatabase().getConnection();
+ firstConnection.close();
+
+ location.initialize(variables, metadataProvider);
+
+ assertThrows(
+ IllegalStateException.class,
+ () ->
+ firstTimer.schedule(
+ new TimerTask() {
+ @Override
+ public void run() {
+ // Cancelled timers reject new tasks.
+ }
+ },
+ 10_000L));
+ assertNotSame(firstTimer, location.getCacheTimer());
+ assertFalse(location.getDatabase().getConnection().isClosed());
+ assertNull(location.getDatabase().getConnectionGroup());
+
+ String id = UUID.randomUUID().toString();
+ location.persistCacheEntry(
+ sampleEntry(id, "AfterReinit", ExecutionType.Workflow, false,
"Finished"));
+ assertNotNull(location.loadCacheEntry(id));
+ }
+
+ @Test
+ void closeIsIdempotentAndStopsTheCacheTimer() throws Exception {
+ Timer timer = location.getCacheTimer();
+ String id = UUID.randomUUID().toString();
+ location.persistCacheEntry(
+ sampleEntry(id, "ToClose", ExecutionType.Pipeline, false, "Finished"));
+
+ location.close();
+ location.close();
+
+ assertNull(location.getDatabase());
+ assertThrows(
+ IllegalStateException.class,
+ () ->
+ timer.schedule(
+ new TimerTask() {
+ @Override
+ public void run() {
+ // Cancelled timers reject new tasks.
+ }
+ },
+ 10_000L));
+
+ HopException exception =
+ assertThrows(
+ HopException.class,
+ () ->
+ location.persistCacheEntry(
+ sampleEntry(
+ UUID.randomUUID().toString(),
+ "AfterClose",
+ ExecutionType.Workflow,
+ false,
+ "Running")));
+ assertTrue(
+ exception.getMessage().toLowerCase().contains("closed")
+ || (exception.getCause() != null
+ && exception.getCause().getMessage() != null
+ &&
exception.getCause().getMessage().toLowerCase().contains("closed")));
+ }
+
+ @Test
+ void recognizesClosedConnectionFailures() {
+ assertTrue(
+ CachingDatabaseExecutionInfoLocation.isClosedConnectionFailure(
+ new SQLException("This connection has been closed.", "08003")));
+ assertFalse(
+ CachingDatabaseExecutionInfoLocation.isClosedConnectionFailure(
+ new SQLException("syntax error", "42000")));
+ SQLException nested = new SQLException("An error occurred executing SQL");
+ nested.initCause(new SQLException("This connection has been closed."));
+ assertTrue(
+ CachingDatabaseExecutionInfoLocation.isClosedConnectionFailure(
+ new HopException("wrapper", nested)));
+ }
+
@Test
void buildDdlContainsIndexes() throws Exception {
String ddl = location.buildDdl(variables);