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 64f9d190c7 fix db commit, fixes #8288 (#8480)
64f9d190c7 is described below

commit 64f9d190c77dca34f7657cfbe42a91508f4afe3e
Author: Hans Van Akelyen <[email protected]>
AuthorDate: Mon Sep 21 11:15:04 2026 +0200

    fix db commit, fixes #8288 (#8480)
    
    * fix db commit, fixes #8288
    
    * extra hardening
---
 .../monetdbbulkloader/MonetDbBulkLoader.java       |  17 ++
 .../monetdbbulkloader/MonetDbBulkLoaderTest.java   |  39 +++
 .../transforms/pgbulkloader/PGBulkLoader.java      |  43 ++-
 .../transforms/pgbulkloader/PGBulkLoaderTest.java  | 100 +++++++
 .../SynchronizeAfterMerge.java                     | 100 ++++++-
 .../SynchronizeAfterMergeDisposeTest.java          | 307 +++++++++++++++++++++
 6 files changed, 596 insertions(+), 10 deletions(-)

diff --git 
a/plugins/databases/monetdb/src/main/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoader.java
 
b/plugins/databases/monetdb/src/main/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoader.java
index 9ddb2d8caf..3c2d39e9ee 100644
--- 
a/plugins/databases/monetdb/src/main/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoader.java
+++ 
b/plugins/databases/monetdb/src/main/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoader.java
@@ -141,6 +141,23 @@ public class MonetDbBulkLoader extends 
BaseTransform<MonetDbBulkLoaderMeta, Mone
     return true;
   }
 
+  /**
+   * The end-of-input branch of {@link #processRow()} flushes the buffer and 
closes the MonetDB
+   * socket, but an error in the middle of the stream (the catch below) and a 
stop that breaks the
+   * run loop mid-row never reach it - the socket, and the server-side load 
session it holds, then
+   * stay open. Close it here as a backstop; {@link MapiSocket#close()} is 
null-guarded and
+   * idempotent, so a normal, already-closed load is left untouched. See <a
+   * href="https://github.com/apache/hop/issues/8288";>issue 8288</a>.
+   */
+  @Override
+  public void dispose() {
+    if (data.mserver != null) {
+      data.mserver.close();
+      data.mserver = null;
+    }
+    super.dispose();
+  }
+
   @Override
   public boolean processRow() throws HopException {
     try {
diff --git 
a/plugins/databases/monetdb/src/test/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoaderTest.java
 
b/plugins/databases/monetdb/src/test/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoaderTest.java
index 1d8c019797..c4f03f457e 100644
--- 
a/plugins/databases/monetdb/src/test/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoaderTest.java
+++ 
b/plugins/databases/monetdb/src/test/java/org/apache/hop/pipeline/transforms/monetdbbulkloader/MonetDbBulkLoaderTest.java
@@ -20,6 +20,8 @@ package org.apache.hop.pipeline.transforms.monetdbbulkloader;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
 
 import org.apache.hop.core.HopEnvironment;
 import org.apache.hop.core.plugins.PluginRegistry;
@@ -30,6 +32,7 @@ import 
org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
 import org.apache.hop.pipeline.transform.TransformMeta;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
+import org.monetdb.mcl.net.MapiSocket;
 
 /** Test for MonetDbBulkLoader (excluding dialog). */
 class MonetDbBulkLoaderTest {
@@ -84,6 +87,42 @@ class MonetDbBulkLoaderTest {
     assertNull(MonetDbBulkLoader.hexFieldForMonetDbCopy(meta, null));
   }
 
+  private MonetDbBulkLoader newLoader(MonetDbBulkLoaderData data) {
+    PipelineMeta pipelineMeta = new PipelineMeta();
+    TransformMeta transformMeta = new TransformMeta("test", new 
MonetDbBulkLoaderMeta());
+    pipelineMeta.addTransform(transformMeta);
+    MonetDbBulkLoaderMeta meta = new MonetDbBulkLoaderMeta();
+    Pipeline pipeline = new LocalPipelineEngine(pipelineMeta);
+    return new MonetDbBulkLoader(transformMeta, meta, data, 0, pipelineMeta, 
pipeline);
+  }
+
+  /**
+   * An error or a stop in the middle of the stream skips the end-of-input 
close. dispose() has to
+   * close the MonetDB socket so the load session does not stay open on the 
server. Issue 8288.
+   */
+  @Test
+  void disposeClosesTheOpenSocket() {
+    MonetDbBulkLoaderData data = new MonetDbBulkLoaderData();
+    MapiSocket mserver = mock(MapiSocket.class);
+    data.mserver = mserver;
+
+    newLoader(data).dispose();
+
+    verify(mserver).close();
+    assertNull(data.mserver);
+  }
+
+  /** A transform that never opened a socket must dispose cleanly. */
+  @Test
+  void disposeSurvivesWithoutASocket() {
+    MonetDbBulkLoaderData data = new MonetDbBulkLoaderData();
+    data.mserver = null;
+
+    newLoader(data).dispose();
+
+    assertNull(data.mserver);
+  }
+
   @Test
   void testEscapeOsPathSpacesOnWindows() {
     PipelineMeta pipelineMeta = new PipelineMeta();
diff --git 
a/plugins/databases/postgresql/src/main/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoader.java
 
b/plugins/databases/postgresql/src/main/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoader.java
index 16a4291888..68b06877f2 100644
--- 
a/plugins/databases/postgresql/src/main/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoader.java
+++ 
b/plugins/databases/postgresql/src/main/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoader.java
@@ -26,6 +26,7 @@ package org.apache.hop.pipeline.transforms.pgbulkloader;
 //
 
 import com.google.common.annotations.VisibleForTesting;
+import java.io.IOException;
 import java.math.BigDecimal;
 import java.nio.charset.Charset;
 import java.sql.Connection;
@@ -229,7 +230,11 @@ public class PGBulkLoader extends 
BaseTransform<PGBulkLoaderMeta, PGBulkLoaderDa
           pgCopyOut.flush();
           pgCopyOut.endCopy();
           pgCopyOut.close();
-          data.db.getConnection().close();
+          pgCopyOut = null;
+        }
+        if (data != null && data.db != null) {
+          data.db.disconnect();
+          data.db = null;
         }
 
         return false;
@@ -481,6 +486,42 @@ public class PGBulkLoader extends 
BaseTransform<PGBulkLoaderMeta, PGBulkLoaderDa
     }
   }
 
+  /**
+   * The end-of-input branch of {@link #processRow()} finishes the COPY and 
closes the connection. A
+   * stop or an error in the middle of the stream never reaches it, leaving 
the COPY and its
+   * connection open on the server - and with them the locks the load holds. 
Abort the copy here
+   * ({@code endCopy()} would commit the partial rows, {@code cancelCopy()} 
discards them) and
+   * release the connection. See <a 
href="https://github.com/apache/hop/issues/8288";>issue 8288</a>.
+   */
+  @Override
+  public void dispose() {
+    try {
+      if (pgCopyOut != null && pgCopyOut.isActive()) {
+        pgCopyOut.cancelCopy();
+      }
+    } catch (SQLException e) {
+      logError("Error cancelling the COPY command while stopping the 
transform", e);
+    } finally {
+      try {
+        // Only close a copy that is no longer active. A still-active copy 
here means cancelCopy()
+        // above threw (a broken connection, the likely case), and pgjdbc's 
close() runs endCopy()
+        // on an active copy - which would commit the very rows we are trying 
to discard. The
+        // disconnect() below tears the connection down regardless.
+        if (pgCopyOut != null && !pgCopyOut.isActive()) {
+          pgCopyOut.close();
+        }
+      } catch (IOException e) {
+        logError("Error closing the COPY output stream", e);
+      }
+      pgCopyOut = null;
+      if (data.db != null) {
+        data.db.disconnect();
+        data.db = null;
+      }
+    }
+    super.dispose();
+  }
+
   protected void verifyDatabaseConnection() throws HopException {
     // Confirming Database Connection is defined.
     if (meta.getConnection() == null) {
diff --git 
a/plugins/databases/postgresql/src/test/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoaderTest.java
 
b/plugins/databases/postgresql/src/test/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoaderTest.java
index 716fab3dbe..598a1b3d97 100644
--- 
a/plugins/databases/postgresql/src/test/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoaderTest.java
+++ 
b/plugins/databases/postgresql/src/test/java/org/apache/hop/pipeline/transforms/pgbulkloader/PGBulkLoaderTest.java
@@ -27,12 +27,16 @@ import static org.junit.jupiter.api.Assertions.fail;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.doNothing;
 import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
+import java.lang.reflect.Field;
 import java.nio.charset.StandardCharsets;
+import java.sql.SQLException;
 import java.util.ArrayList;
 import org.apache.hop.core.HopClientEnvironment;
 import org.apache.hop.core.database.Database;
@@ -54,6 +58,7 @@ import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.RegisterExtension;
+import org.postgresql.copy.PGCopyOutputStream;
 
 class PGBulkLoaderTest {
   @RegisterExtension
@@ -246,6 +251,101 @@ class PGBulkLoaderTest {
     }
   }
 
+  private PGBulkLoader disposableLoader(PGBulkLoaderData data, 
PGCopyOutputStream copyOut)
+      throws Exception {
+    PGBulkLoader loader =
+        spy(
+            new PGBulkLoader(
+                transformMockHelper.transformMeta,
+                transformMockHelper.iTransformMeta,
+                data,
+                0,
+                transformMockHelper.pipelineMeta,
+                transformMockHelper.pipeline));
+    Field field = PGBulkLoader.class.getDeclaredField("pgCopyOut");
+    field.setAccessible(true);
+    field.set(loader, copyOut);
+    return loader;
+  }
+
+  /**
+   * A stop or an error leaves the COPY open. dispose() must abort it - not 
endCopy(), which would
+   * commit the partial rows - and release the connection. Issue 8288.
+   */
+  @Test
+  void disposeCancelsAnActiveCopyAndDisconnects() throws Exception {
+    PGBulkLoaderData data = new PGBulkLoaderData();
+    Database db = mock(Database.class);
+    data.db = db;
+    PGCopyOutputStream copyOut = mock(PGCopyOutputStream.class);
+    // Active on entry (so the copy is cancelled), inactive afterwards (a 
successful cancel), so the
+    // stream is then closed without an endCopy() commit.
+    when(copyOut.isActive()).thenReturn(true, false);
+
+    PGBulkLoader loader = disposableLoader(data, copyOut);
+    loader.dispose();
+
+    verify(copyOut).cancelCopy();
+    verify(copyOut, never()).endCopy();
+    verify(copyOut).close();
+    verify(db).disconnect();
+    assertNull(data.db);
+  }
+
+  /**
+   * The normal end-of-input path already finished and closed the COPY. 
dispose() then only has to
+   * release the connection, and must not touch the already-closed copy.
+   */
+  @Test
+  void disposeDoesNotCancelAnAlreadyFinishedCopy() throws Exception {
+    PGBulkLoaderData data = new PGBulkLoaderData();
+    Database db = mock(Database.class);
+    data.db = db;
+    PGCopyOutputStream copyOut = mock(PGCopyOutputStream.class);
+    when(copyOut.isActive()).thenReturn(false);
+
+    PGBulkLoader loader = disposableLoader(data, copyOut);
+    loader.dispose();
+
+    verify(copyOut, never()).cancelCopy();
+    verify(copyOut).close();
+    verify(db).disconnect();
+  }
+
+  /**
+   * If cancelCopy() throws (a broken connection mid-load, the likely case), 
dispose() must not fall
+   * through to close() - pgjdbc's close() runs endCopy() on a still-active 
copy, committing the
+   * very rows we are discarding. The connection is torn down instead. Issue 
8288 review follow-up.
+   */
+  @Test
+  void disposeDoesNotCommitWhenCancelFailsOnAnActiveCopy() throws Exception {
+    PGBulkLoaderData data = new PGBulkLoaderData();
+    Database db = mock(Database.class);
+    data.db = db;
+    PGCopyOutputStream copyOut = mock(PGCopyOutputStream.class);
+    when(copyOut.isActive()).thenReturn(true);
+    doThrow(new SQLException("connection reset")).when(copyOut).cancelCopy();
+
+    PGBulkLoader loader = disposableLoader(data, copyOut);
+    loader.dispose();
+
+    verify(copyOut).cancelCopy();
+    verify(copyOut, never()).endCopy();
+    verify(copyOut, never()).close();
+    verify(db).disconnect();
+    assertNull(data.db);
+  }
+
+  /** A transform stopped before the first row opened neither the copy nor the 
connection. */
+  @Test
+  void disposeSurvivesWithoutACopyOrConnection() throws Exception {
+    PGBulkLoaderData data = new PGBulkLoaderData();
+    data.db = null;
+    PGBulkLoader loader = disposableLoader(data, null);
+
+    loader.dispose();
+  }
+
   private static PGBulkLoaderMeta getPgBulkLoaderMock(String DbNameOverride)
       throws HopXmlException {
     PGBulkLoaderMeta pgBulkLoaderMetaMock = mock(PGBulkLoaderMeta.class);
diff --git 
a/plugins/transforms/synchronizeaftermerge/src/main/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMerge.java
 
b/plugins/transforms/synchronizeaftermerge/src/main/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMerge.java
index 70e1a17b68..5b2506aae3 100644
--- 
a/plugins/transforms/synchronizeaftermerge/src/main/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMerge.java
+++ 
b/plugins/transforms/synchronizeaftermerge/src/main/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMerge.java
@@ -1083,10 +1083,74 @@ public class SynchronizeAfterMerge
     return databaseMeta;
   }
 
+  /**
+   * A single-threaded (streaming) pipeline never sends the end-of-input 
signal that {@link
+   * #processRow()} flushes on; it calls this after every batch of rows 
instead. Commit what is
+   * pending now, so it does not sit uncommitted - and, on databases like 
Oracle, locked - until the
+   * stream ends. See <a 
href="https://github.com/apache/hop/issues/8288";>issue 8288</a>.
+   */
+  @Override
+  public void batchComplete() throws HopException {
+    if (data.db == null || data.db.getConnection() == null) {
+      return;
+    }
+    try {
+      emptyBatchBuffer(false);
+    } catch (HopDatabaseBatchException be) {
+      // The statements stay open for the next batch, so drop what failed 
before going on. The rest
+      // is the recovery a failure in the middle of the stream gets.
+      for (PreparedStatement statement : data.preparedStatements.values()) {
+        data.db.clearBatch(statement);
+      }
+      if (getTransformMeta().isDoingErrorHandling()) {
+        data.db.commit(true);
+        processBatchException(be.toString(), be.getUpdateCounts(), 
be.getExceptionsList());
+      } else {
+        data.db.rollback();
+        throw new HopException(
+            BaseMessages.getString(PKG, 
"SynchronizeAfterMerge.Error.UpdatingBatch"), be);
+      }
+    } catch (SQLException e) {
+      throw new HopDatabaseException("Unexpected error committing the database 
connection.", e);
+    }
+  }
+
+  /**
+   * The end-of-input path in {@link #processRow()} has normally flushed and 
disconnected by now. It
+   * is skipped when the transform is stopped or fails in the middle of a row, 
and a single-threaded
+   * (streaming) pipeline never sends end-of-input at all. In both cases the 
connection stayed open,
+   * and with it the uncommitted transaction and every row lock it holds. See 
<a
+   * href="https://github.com/apache/hop/issues/8288";>issue 8288</a>.
+   *
+   * <p>A graceful stop leaves {@code getErrors() == 0}, so the pending batch 
is committed rather
+   * than rolled back - deliberately, and in line with Table Output, Update 
and Delete: this
+   * transform already commits every {@code commitSize} rows, so committing 
the final partial batch
+   * on a stop keeps the same all-or-a-multiple-of-commitSize contract. Only a 
real error rolls
+   * back. Note that {@link #emptyBatchBuffer(boolean)} still calls {@code 
putRow} for the committed
+   * rows; on a stop {@code putRow} is a no-op (nothing reads downstream 
anyway), while the rows are
+   * safely in the table - the same behaviour Table Output has.
+   */
+  @Override
+  public void dispose() {
+    if (data.db != null) {
+      if (data.db.getConnection() != null) {
+        if (getErrors() > 0) {
+          // The transform failed: nothing that is still pending may reach the 
table.
+          rollback();
+          data.db.disconnect();
+        } else {
+          finishTransform();
+        }
+      }
+      data.db = null;
+    }
+    super.dispose();
+  }
+
   private void finishTransform() {
     if (data.db != null && data.db.getConnection() != null) {
       try {
-        finishTransformEmptyBatchBuffer();
+        emptyBatchBuffer(true);
       } catch (HopDatabaseBatchException be) {
         finishTransformErrorHandling(be);
       } catch (Exception dbe) {
@@ -1098,11 +1162,7 @@ public class SynchronizeAfterMerge
         setOutputDone();
 
         if (getErrors() > 0) {
-          try {
-            data.db.rollback();
-          } catch (HopDatabaseException e) {
-            logError("Unexpected error rolling back the database connection.", 
e);
-          }
+          rollback();
         }
 
         data.db.disconnect();
@@ -1110,7 +1170,21 @@ public class SynchronizeAfterMerge
     }
   }
 
-  private void finishTransformEmptyBatchBuffer()
+  private void rollback() {
+    try {
+      data.db.rollback();
+    } catch (HopDatabaseException e) {
+      logError("Unexpected error rolling back the database connection.", e);
+    }
+  }
+
+  /**
+   * Execute and commit the pending batch of every prepared statement and pass 
the buffered rows on.
+   *
+   * @param closeStatements true at the end of the transform; false between 
the batches of a
+   *     single-threaded pipeline, where the statements are reused.
+   */
+  private void emptyBatchBuffer(boolean closeStatements)
       throws SQLException, HopDatabaseException, HopTransformException, 
HopValueException {
     if (!data.db.getConnection().isClosed()) {
       for (String schemaTable : data.preparedStatements.keySet()) {
@@ -1121,9 +1195,17 @@ public class SynchronizeAfterMerge
           batchCounter = 0;
         }
 
-        PreparedStatement insertStatement = 
data.preparedStatements.get(schemaTable);
+        // Between batches there is nothing to commit for a statement that 
took no rows this batch.
+        // Skip it; at final completion we still fall through so 
emptyAndCommit closes the
+        // statement.
+        if (!closeStatements && batchCounter == 0) {
+          continue;
+        }
+
+        PreparedStatement statement = data.preparedStatements.get(schemaTable);
 
-        data.db.emptyAndCommit(insertStatement, data.batchMode, batchCounter);
+        data.db.emptyAndCommit(statement, data.batchMode, batchCounter, 
closeStatements);
+        data.commitCounterMap.put(schemaTable, 0);
       }
       for (int i = 0; i < data.batchBuffer.size(); i++) {
         Object[] row = data.batchBuffer.get(i);
diff --git 
a/plugins/transforms/synchronizeaftermerge/src/test/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMergeDisposeTest.java
 
b/plugins/transforms/synchronizeaftermerge/src/test/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMergeDisposeTest.java
new file mode 100644
index 0000000000..e5cc731a25
--- /dev/null
+++ 
b/plugins/transforms/synchronizeaftermerge/src/test/java/org/apache/hop/pipeline/transforms/synchronizeaftermerge/SynchronizeAfterMergeDisposeTest.java
@@ -0,0 +1,307 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *       http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hop.pipeline.transforms.synchronizeaftermerge;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+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 static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.nullable;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.inOrder;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.verify;
+
+import java.sql.BatchUpdateException;
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.hop.core.database.Database;
+import org.apache.hop.core.exception.HopDatabaseBatchException;
+import org.apache.hop.core.exception.HopException;
+import org.apache.hop.core.row.IRowMeta;
+import org.apache.hop.core.row.RowMeta;
+import org.apache.hop.core.row.value.ValueMetaString;
+import org.apache.hop.pipeline.PipelineMeta;
+import org.apache.hop.pipeline.engines.local.LocalPipelineEngine;
+import org.apache.hop.pipeline.transform.TransformMeta;
+import org.apache.hop.pipeline.transform.TransformPartitioningMeta;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.InOrder;
+
+/**
+ * Releasing the database connection when the end of the input never comes.
+ *
+ * <p>The transform used to flush and disconnect only from the end-of-input 
branch of {@code
+ * processRow()}. A stop or a failure in the middle of a row skips that 
branch, and a
+ * single-threaded (streaming) pipeline never reaches it at all: the 
connection then stayed open for
+ * as long as the JVM lived, and with it the uncommitted transaction and the 
row locks it held on
+ * the database. See <a href="https://github.com/apache/hop/issues/8288";>issue 
8288</a>.
+ *
+ * <p>The invariants pinned here: {@code dispose()} always ends with {@code 
disconnect()} when a
+ * connection is open, commits pending work unless the transform failed, and 
{@code batchComplete()}
+ * commits between batches without closing the statements the next batch 
reuses.
+ */
+class SynchronizeAfterMergeDisposeTest {
+
+  private static final String INSERT_KEY = "\"T\"" + 
SynchronizeAfterMerge.CONST_INSERT;
+  private static final String UPDATE_KEY = "\"T\"" + 
SynchronizeAfterMerge.CONST_UPDATE;
+
+  private SynchronizeAfterMerge transform;
+  private SynchronizeAfterMergeData data;
+  private TransformMeta transformMeta;
+  private Database db;
+  private PreparedStatement insertStatement;
+  private PreparedStatement updateStatement;
+  private List<Object[]> emitted;
+  private List<String> rejected;
+
+  @BeforeEach
+  void setUp() throws Exception {
+    SynchronizeAfterMergeMeta meta = mock(SynchronizeAfterMergeMeta.class);
+    transformMeta = mock(TransformMeta.class);
+    doReturn("transform").when(transformMeta).getName();
+    doReturn(mock(TransformPartitioningMeta.class))
+        .when(transformMeta)
+        .getTargetTransformPartitioningMeta();
+    doReturn(meta).when(transformMeta).getTransform();
+
+    PipelineMeta pipelineMeta = mock(PipelineMeta.class);
+    doReturn(transformMeta).when(pipelineMeta).findTransform(anyString());
+
+    db = mock(Database.class);
+    Connection connection = mock(Connection.class);
+    doReturn(connection).when(db).getConnection();
+
+    insertStatement = mock(PreparedStatement.class);
+    updateStatement = mock(PreparedStatement.class);
+
+    data = new SynchronizeAfterMergeData();
+    data.db = db;
+    data.batchMode = true;
+    data.insertValue = "insert";
+    data.indexOfOperationOrderField = 1;
+    data.inputRowMeta = new RowMeta();
+    data.inputRowMeta.addValueMeta(new ValueMetaString("name"));
+    data.inputRowMeta.addValueMeta(new ValueMetaString("operation"));
+    data.outputRowMeta = data.inputRowMeta;
+    data.preparedStatements.put(INSERT_KEY, insertStatement);
+    data.preparedStatements.put(UPDATE_KEY, updateStatement);
+    data.commitCounterMap.put(INSERT_KEY, 3);
+    data.commitCounterMap.put(UPDATE_KEY, 2);
+    data.batchBuffer = new ArrayList<>();
+
+    transform =
+        spy(
+            new SynchronizeAfterMerge(
+                transformMeta, meta, data, 1, pipelineMeta, spy(new 
LocalPipelineEngine())));
+    doReturn(transformMeta).when(transform).getTransformMeta();
+    doReturn(false).when(transform).isRowLevel();
+    doNothing().when(transform).logDetailed(anyString());
+    doNothing().when(transform).logError(anyString());
+    doNothing().when(transform).logError(anyString(), any(Throwable.class));
+
+    emitted = new ArrayList<>();
+    rejected = new ArrayList<>();
+    doAnswer(
+            inv -> {
+              emitted.add(inv.getArgument(1));
+              return null;
+            })
+        .when(transform)
+        .putRow(any(IRowMeta.class), any());
+    doAnswer(
+            inv -> {
+              rejected.add(inv.getArgument(3));
+              return null;
+            })
+        .when(transform)
+        .putError(
+            any(IRowMeta.class),
+            any(),
+            anyLong(),
+            anyString(),
+            nullable(String.class),
+            anyString());
+  }
+
+  private void bufferRows(int count) {
+    for (int i = 0; i < count; i++) {
+      data.batchBuffer.add(new Object[] {"row" + i, "insert"});
+    }
+  }
+
+  /** The stop-in-the-middle-of-a-row and the streaming case: nothing flushed 
us before. */
+  @Test
+  void disposeFlushesCommitsAndDisconnectsWhenEndOfInputNeverCame() throws 
Exception {
+    bufferRows(5);
+
+    transform.dispose();
+
+    // Both statements are flushed and closed. Their relative order is a 
Hashtable artifact, so it
+    // is
+    // not asserted; what matters is that a flush precedes the disconnect.
+    verify(db).emptyAndCommit(updateStatement, true, 2, true);
+    InOrder order = inOrder(db);
+    order.verify(db).emptyAndCommit(insertStatement, true, 3, true);
+    order.verify(db).disconnect();
+    verify(db, never()).rollback();
+    assertEquals(5, emitted.size(), "the buffered rows leave on the output 
before we disconnect");
+    assertTrue(data.batchBuffer.isEmpty());
+    assertNull(data.db, "the database handle is released for the garbage 
collector");
+  }
+
+  /** After a failure nothing that is still pending may reach the table. */
+  @Test
+  void disposeRollsBackAndDisconnectsAfterAFailure() throws Exception {
+    bufferRows(5);
+    transform.setErrors(1);
+
+    transform.dispose();
+
+    InOrder order = inOrder(db);
+    order.verify(db).rollback();
+    order.verify(db).disconnect();
+    verify(db, never()).emptyAndCommit(any(), anyBoolean(), anyInt(), 
anyBoolean());
+    assertTrue(emitted.isEmpty(), "no row is reported as written after a 
rollback");
+  }
+
+  /** The normal end-of-input path already disconnected: dispose() must not do 
it twice. */
+  @Test
+  void disposeIsANoOpOnceTheConnectionIsGone() throws Exception {
+    doReturn(null).when(db).getConnection();
+
+    transform.dispose();
+
+    verify(db, never()).disconnect();
+    verify(db, never()).rollback();
+    assertNull(data.db);
+  }
+
+  /** A transform whose init() never got as far as a connection. */
+  @Test
+  void disposeSurvivesAMissingDatabase() throws Exception {
+    data.db = null;
+
+    transform.dispose();
+
+    verify(db, never()).disconnect();
+  }
+
+  /** Between batches the statements are reused: commit, but leave them open. 
*/
+  @Test
+  void batchCompleteCommitsWithoutClosingTheStatements() throws Exception {
+    bufferRows(5);
+
+    transform.batchComplete();
+
+    verify(db).emptyAndCommit(insertStatement, true, 3, false);
+    verify(db).emptyAndCommit(updateStatement, true, 2, false);
+    verify(db, never()).disconnect();
+    assertEquals(0, data.commitCounterMap.get(INSERT_KEY), "the batch counter 
starts over");
+    assertEquals(0, data.commitCounterMap.get(UPDATE_KEY), "the batch counter 
starts over");
+    assertEquals(5, emitted.size());
+    assertTrue(data.batchBuffer.isEmpty());
+  }
+
+  /**
+   * Between batches a statement that took no rows must not trigger a no-op 
commit; only the
+   * statements with pending work are flushed. Issue 8288 review follow-up.
+   */
+  @Test
+  void batchCompleteSkipsStatementsWithAnEmptyBatch() throws Exception {
+    data.commitCounterMap.put(UPDATE_KEY, 0);
+    bufferRows(5);
+
+    transform.batchComplete();
+
+    verify(db).emptyAndCommit(insertStatement, true, 3, false);
+    verify(db, never()).emptyAndCommit(eq(updateStatement), anyBoolean(), 
anyInt(), anyBoolean());
+  }
+
+  @Test
+  void batchCompleteIsANoOpWithoutAConnection() throws Exception {
+    doReturn(null).when(db).getConnection();
+
+    transform.batchComplete();
+
+    verify(db, never()).emptyAndCommit(any(), anyBoolean(), anyInt(), 
anyBoolean());
+  }
+
+  /**
+   * A batch that fails between batches gets the same recovery as one that 
fails mid-stream: the
+   * failed batches are dropped so the next batch does not re-run them, what 
went through is
+   * committed and every buffered row leaves on one stream or the other.
+   */
+  @Test
+  void batchCompleteWithErrorHandlingRoutesTheFailedBatchAndKeepsGoing() 
throws Exception {
+    doReturn(true).when(transformMeta).isDoingErrorHandling();
+    bufferRows(5);
+    BatchUpdateException cause =
+        new BatchUpdateException("boom", new int[] {1, 1, 
java.sql.Statement.EXECUTE_FAILED});
+    HopDatabaseBatchException failure =
+        Database.createHopDatabaseBatchException("Error updating batch", 
cause);
+    doThrow(failure)
+        .when(db)
+        .emptyAndCommit(eq(insertStatement), anyBoolean(), anyInt(), 
eq(false));
+
+    transform.batchComplete();
+
+    verify(db).clearBatch(insertStatement);
+    verify(db).clearBatch(updateStatement);
+    verify(db).commit(true);
+    verify(db, never()).rollback();
+    verify(db, never()).disconnect();
+    assertEquals(5, emitted.size() + rejected.size(), "every buffered row 
leaves exactly once");
+    assertTrue(data.batchBuffer.isEmpty());
+  }
+
+  /** Without error handling a failed batch is a failed transform. */
+  @Test
+  void batchCompleteWithoutErrorHandlingRollsBackAndFails() throws Exception {
+    doReturn(false).when(transformMeta).isDoingErrorHandling();
+    bufferRows(5);
+    HopDatabaseBatchException failure =
+        Database.createHopDatabaseBatchException(
+            "Error updating batch", new BatchUpdateException("boom", new 
int[0]));
+    doThrow(failure)
+        .when(db)
+        .emptyAndCommit(eq(insertStatement), anyBoolean(), anyInt(), 
eq(false));
+
+    assertThrows(HopException.class, () -> transform.batchComplete());
+
+    verify(db).clearBatch(insertStatement);
+    verify(db).clearBatch(updateStatement);
+    verify(db).rollback();
+    verify(db, never()).commit(anyBoolean());
+    verify(db, never()).disconnect();
+  }
+}

Reply via email to