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);

Reply via email to