This is an automated email from the ASF dual-hosted git repository.

mattcasters 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 c482981821 Issue #8546 : Keep pipeline metrics updating after a Kafka 
consumer restart (#8560)
c482981821 is described below

commit c482981821f62c738e0b6d7081e5baa08e4a51a8
Author: Matt Casters <[email protected]>
AuthorDate: Mon Sep 28 12:27:25 2026 +0200

    Issue #8546 : Keep pipeline metrics updating after a Kafka consumer restart 
(#8560)
    
    * Issue #8546 : Keep pipeline metrics updating after a Kafka consumer 
restart
    
    A finished listener from the previous engine could stop the new run's
    canvas and metrics timers, so the grid stayed at 0 until the next stop.
    Kafka also logged each batch and then cleared the sub-pipeline counters.
    
    * Issue #8546 : Address review on metrics timers and restart
    
    Run again when the engine is not running, including a successful Beam
    run and a failed preparation. Do GUI work outside the session lock, and
    cancel a refresh timer once its engine is no longer on screen.
---
 .../pages/pipeline/transforms/kafkaconsumer.adoc   |   2 +
 .../pipeline/SingleThreadedPipelineExecutor.java   |  54 -------
 .../kafka/consumer/KafkaConsumerInput.java         |  17 ++-
 .../kafka/consumer/KafkaConsumerInputTest.java     |   9 ++
 .../hopgui/file/pipeline/HopGuiPipelineGraph.java  | 161 +++++++++++++-------
 .../delegates/HopGuiPipelineGridDelegate.java      |  56 ++++---
 .../ui/hopgui/file/shared/ExecutionGuiSession.java | 163 +++++++++++++++++++++
 .../hopgui/file/workflow/HopGuiWorkflowGraph.java  | 104 ++++++++-----
 .../file/shared/ExecutionGuiSessionTest.java       | 122 +++++++++++++++
 9 files changed, 512 insertions(+), 176 deletions(-)

diff --git 
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
 
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
index 33baddc302..11bf38d132 100644
--- 
a/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
+++ 
b/docs/hop-user-manual/modules/ROOT/pages/pipeline/transforms/kafkaconsumer.adoc
@@ -27,6 +27,8 @@ The Kafka consumer transform pulls streaming data from Kafka. 
It runs a sub-pipe
 
 This sub-pipeline must start with an Injector transform.
 
+The Input metric on the Kafka consumer is the number of records consumed so 
far in the run. Read and Written stay at 0 unless result fields are returned 
from the sub-pipeline. Open the running sub-pipeline to see each transform's 
own metrics. Those totals accumulate until the consumer stops.
+
 You can define the number of messages to accept for processing, as well as the 
specific data formats to stream activity data and system metrics.
 
 You can set up this transform to collect monitored events, track user 
consumption of data streams, and monitor alerts.
diff --git 
a/engine/src/main/java/org/apache/hop/pipeline/SingleThreadedPipelineExecutor.java
 
b/engine/src/main/java/org/apache/hop/pipeline/SingleThreadedPipelineExecutor.java
index e0b0d51259..5257f2017e 100644
--- 
a/engine/src/main/java/org/apache/hop/pipeline/SingleThreadedPipelineExecutor.java
+++ 
b/engine/src/main/java/org/apache/hop/pipeline/SingleThreadedPipelineExecutor.java
@@ -25,8 +25,6 @@ import org.apache.hop.core.IRowSet;
 import org.apache.hop.core.Result;
 import org.apache.hop.core.exception.HopException;
 import org.apache.hop.core.logging.ILogChannel;
-import org.apache.hop.i18n.BaseMessages;
-import org.apache.hop.pipeline.transform.BaseTransform;
 import org.apache.hop.pipeline.transform.TransformMetaDataCombi;
 import org.apache.hop.pipeline.transform.stream.IStream;
 
@@ -45,12 +43,9 @@ public class SingleThreadedPipelineExecutor {
   private List<List<IStream>> transformInfoStreams;
   private List<List<IRowSet>> transformInfoRowSets;
   private ILogChannel log;
-  private static final Class<?> PKG = SingleThreadedPipelineExecutor.class;
   private static final String CONST_SEPARATOR =
       "-------------------------------------------------------";
 
-  @Getter @Setter private boolean clearingMetricsPerIteration = true;
-
   public SingleThreadedPipelineExecutor(Pipeline pipeline) {
     initializeObject(pipeline, false);
   }
@@ -405,55 +400,6 @@ public class SingleThreadedPipelineExecutor {
     return nrDone < transforms.size() && !pipeline.isStopped();
   }
 
-  public void buildExecutionSummary() {
-
-    for (TransformMetaDataCombi combi : transforms) {
-      // Summarize execution results
-      long li = combi.transform.getLinesInput();
-      long lo = combi.transform.getLinesOutput();
-      long lr = combi.transform.getLinesRead();
-      long lw = combi.transform.getLinesWritten();
-      long lu = combi.transform.getLinesUpdated();
-      long lj = combi.transform.getLinesRejected();
-      long e = combi.transform.getErrors();
-
-      ILogChannel tLog = combi.transform.getLogChannel();
-
-      if (li > 0 || lo > 0 || lr > 0 || lw > 0 || lu > 0 || lj > 0 || e > 0) {
-        tLog.logBasic(
-            BaseMessages.getString(
-                PKG,
-                "SingleThreadedPipeline.Log.SummaryInfo",
-                String.valueOf(li),
-                String.valueOf(lo),
-                String.valueOf(lr),
-                String.valueOf(lw),
-                String.valueOf(lu),
-                String.valueOf(e + lj)));
-      } else {
-        tLog.logDetailed(
-            BaseMessages.getString(
-                PKG,
-                "SingleThreadedPipeline.Log.SummaryInfo",
-                String.valueOf(li),
-                String.valueOf(lo),
-                String.valueOf(lr),
-                String.valueOf(lw),
-                String.valueOf(lu),
-                String.valueOf(e + lj)));
-      }
-      if (clearingMetricsPerIteration) {
-        ((BaseTransform<?, ?>) combi.transform).setLinesInput(0);
-        ((BaseTransform<?, ?>) combi.transform).setLinesOutput(0);
-        ((BaseTransform<?, ?>) combi.transform).setLinesWritten(0);
-        ((BaseTransform<?, ?>) combi.transform).setLinesRead(0);
-        ((BaseTransform<?, ?>) combi.transform).setLinesSkipped(0);
-        ((BaseTransform<?, ?>) combi.transform).setLinesUpdated(0);
-        combi.transform.setLinesRejected(0);
-      }
-    }
-  }
-
   protected int getTotalRows(List<IRowSet> rowSets) {
     int total = 0;
     for (IRowSet rowSet : rowSets) {
diff --git 
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
 
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
index 59cf04e850..3636b6a4f3 100644
--- 
a/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
+++ 
b/plugins/transforms/kafka/src/main/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInput.java
@@ -19,6 +19,7 @@ package org.apache.hop.pipeline.transforms.kafka.consumer;
 
 import java.time.Duration;
 import java.util.ArrayList;
+import java.util.Date;
 import java.util.Properties;
 import java.util.Set;
 import java.util.stream.Collectors;
@@ -210,8 +211,6 @@ public class KafkaConsumerInput
         // If the conditions for error handling are not met init 
SingleThreadedExecutor normally
         data.executor = new SingleThreadedPipelineExecutor(kafkaPipeline);
       }
-      data.executor.setClearingMetricsPerIteration(
-          StringUtils.isEmpty(meta.getExecutionInformationLocation()));
 
       // Initialize the sub-pipeline
       //
@@ -351,6 +350,9 @@ public class KafkaConsumerInput
         } else {
           // Grab the records...
           //
+          if (getFirstRowReadDate() == null) {
+            setFirstRowReadDate(new Date());
+          }
           for (ConsumerRecord<Object, Object> record : records) {
             Object[] outputRow = processMessageAsRow(record);
             data.rowProducer.putRow(data.outputRowMeta, outputRow);
@@ -361,7 +363,7 @@ public class KafkaConsumerInput
           }
           data.lastRecordTime = System.currentTimeMillis();
           if (isBasic()) {
-            logBasic("Number of rows read: " + 
data.rowProducer.getRowSet().size());
+            logBasic(batchLogMessage(records.count(), getLinesInput()));
           }
           // Pass them to the single threaded transformation and do an 
iteration...
           //
@@ -402,7 +404,6 @@ public class KafkaConsumerInput
           // "removing" failing items from the kafka queue
           //
           data.consumer.commitAsync();
-          data.executor.buildExecutionSummary();
           if (errorHandlingConditionIsSatisfied()) {
             data.incomingRowsBuffer.clear();
           }
@@ -454,6 +455,14 @@ public class KafkaConsumerInput
     return true;
   }
 
+  /**
+   * One batch in the parent log. This is not the single-threaded executor's 
"Finished processing"
+   * line, and it is not the injector buffer size.
+   */
+  static String batchLogMessage(int batchRecords, long totalInput) {
+    return "Kafka consumer batch of " + batchRecords + " record(s), cumulative 
input " + totalInput;
+  }
+
   /**
    * True when a max consume duration is configured and the wall clock since 
transform start has
    * reached it. {@code maxConsumeDurationMs <= 0} means no limit.
diff --git 
a/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
 
b/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
index d2fb5e186e..ac520e67f0 100644
--- 
a/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
+++ 
b/plugins/transforms/kafka/src/test/java/org/apache/hop/pipeline/transforms/kafka/consumer/KafkaConsumerInputTest.java
@@ -46,6 +46,15 @@ class KafkaConsumerInputTest {
     assertEquals(100L, KafkaConsumerInput.pollTimeoutMs(true, 2_000L, 0L, 0L, 
0L));
   }
 
+  @Test
+  void batchLogReportsBatchSizeAndCumulativeInput() {
+    String message = KafkaConsumerInput.batchLogMessage(12, 40L);
+    assertTrue(message.contains("12"));
+    assertTrue(message.contains("40"));
+    assertFalse(message.toLowerCase().contains("finished processing"));
+    assertFalse(message.toLowerCase().contains("rows read"));
+  }
+
   @Test
   void pollTimeoutIsCappedToRemainingConsumeDuration() {
     // A deadline uses a short poll so empty topics re-check the clock instead 
of blocking in
diff --git 
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/HopGuiPipelineGraph.java
 
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/HopGuiPipelineGraph.java
index edf70fa058..634efa57c4 100644
--- 
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/HopGuiPipelineGraph.java
+++ 
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/HopGuiPipelineGraph.java
@@ -34,7 +34,6 @@ import java.util.Map;
 import java.util.Objects;
 import java.util.Set;
 import java.util.Timer;
-import java.util.TimerTask;
 import java.util.UUID;
 import lombok.Getter;
 import lombok.Setter;
@@ -200,6 +199,7 @@ import 
org.apache.hop.ui.hopgui.file.pipeline.extension.HopGuiPipelineGraphExten
 import 
org.apache.hop.ui.hopgui.file.pipeline.extension.PipelineRenamedExtension;
 import org.apache.hop.ui.hopgui.file.shared.CanvasToolTip;
 import org.apache.hop.ui.hopgui.file.shared.DrillDownGuiPlugin;
+import org.apache.hop.ui.hopgui.file.shared.ExecutionGuiSession;
 import org.apache.hop.ui.hopgui.file.shared.HopGuiAbstractGraph;
 import org.apache.hop.ui.hopgui.file.shared.HopGuiGraphSnapshotUndo;
 import org.apache.hop.ui.hopgui.file.shared.HopGuiTooltipExtension;
@@ -514,6 +514,9 @@ public class HopGuiPipelineGraph extends HopGuiAbstractGraph
 
   Timer redrawTimer;
 
+  /** Ties redraw and metrics timers to the engine currently shown. */
+  private final ExecutionGuiSession executionGuiSession = new 
ExecutionGuiSession();
+
   @Setter private HopPipelineFileType<PipelineMeta> fileType;
   private boolean doubleClick;
 
@@ -6269,6 +6272,8 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
     if (handlePipelineMetaChanges(pipelineMeta)) {
 
       // If the pipeline is not running, start the pipeline...
+      // Stopped counts as not running. Beam leaves a successful run 
unfinished, and a failed
+      // preparation only stops the engine, so "not finished" would block the 
next Run.
       //
       if (!isRunning()) {
         try {
@@ -6303,12 +6308,12 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
           //
           pipelineMeta.clearCaches();
 
-          pipeline =
+          setDisplayedPipeline(
               PipelineEngineFactory.createPipelineEngine(
                   variables,
                   variables.resolve(pipelineRunConfigurationName),
                   hopGui.getMetadataProvider(),
-                  pipelineMeta);
+                  pipelineMeta));
           DrillDownGuiPlugin.bindToHopGui(pipeline, hopGui.getId());
 
           // Set the variables from the execution configuration
@@ -6350,7 +6355,7 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
           }
 
         } catch (HopException e) {
-          pipeline = null;
+          setDisplayedPipeline(null);
           new ErrorDialog(
               hopShell(),
               BaseMessages.getString(PKG, 
"PipelineLog.Dialog.ErrorOpeningPipeline.Title"),
@@ -6382,10 +6387,13 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
 
           updateGui();
 
-          // Update the GUI at the end of the pipeline
+          // Update the GUI at the end of the pipeline. Ignore the event when 
this engine is no
+          // longer the one on screen (a restart overlapped its shutdown). The 
work stays outside
+          // the session lock: it can take the graph lock or open a dialog.
           //
-          pipeline.addExecutionFinishedListener(e -> pipelineFinished());
-          pipeline.addExecutionStoppedListener(e -> pipelineStopped());
+          final IPipelineEngine<PipelineMeta> engine = pipeline;
+          engine.addExecutionFinishedListener(this::runPipelineFinished);
+          engine.addExecutionStoppedListener(this::runPipelineStopped);
         }
       } else {
         modalMessageDialog(
@@ -6503,7 +6511,8 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
         // Create a new pipeline to execution
         //
         pipelineMeta.clearCaches();
-        pipeline = new LocalPipelineEngine(pipelineMeta, variables, 
hopGui.getLoggingObject());
+        setDisplayedPipeline(
+            new LocalPipelineEngine(pipelineMeta, variables, 
hopGui.getLoggingObject()));
         DrillDownGuiPlugin.bindToHopGui(pipeline, hopGui.getId());
         pipeline.setPreview(true);
         pipeline.setVariable(IPipelineEngine.PIPELINE_IN_PREVIEW_MODE, "Y");
@@ -6708,16 +6717,16 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
   }
 
   private synchronized void startThreads() {
+    final IPipelineEngine<PipelineMeta> engine = pipeline;
+    if (engine == null) {
+      return;
+    }
     try {
       // Add a listener to the pipeline.
       // If the pipeline is done, we want to do the end processing, etc.
+      // A listener from an engine that is no longer on screen must not stop 
the new run's timers.
       //
-      pipeline.addExecutionFinishedListener(
-          p -> {
-            checkPipelineEnded();
-            checkErrorVisuals();
-            stopRedrawTimer();
-          });
+      engine.addExecutionFinishedListener(this::runPipelineMetricsFinished);
 
       hopGui
           .getDisplay()
@@ -6726,14 +6735,13 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
                   new Thread(
                           () -> {
                             try {
-                              pipeline.startThreads();
-                              pipeline.waitUntilFinished();
+                              engine.startThreads();
+                              engine.waitUntilFinished();
                             } catch (Exception e) {
-                              pipeline
+                              engine
                                   .getLogChannel()
                                   .logError("Error starting transform 
threads", e);
-                              checkErrorVisuals();
-                              stopRedrawTimer();
+                              stopPipelineTimersIfCurrent(engine);
                             }
                           })
                       .start());
@@ -6742,13 +6750,8 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
 
       updateGui();
     } catch (Exception e) {
-      if (pipeline != null) {
-        pipeline.getLogChannel().logError("Error starting transform threads", 
e);
-      } else {
-        log.logError("Error starting transform threads", e);
-      }
-      checkErrorVisuals();
-      stopRedrawTimer();
+      engine.getLogChannel().logError("Error starting transform threads", e);
+      stopPipelineTimersIfCurrent(engine);
     }
   }
 
@@ -6763,8 +6766,10 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
       return;
     }
 
-    // Set the pipeline instance
-    this.pipeline = runningPipeline;
+    // Set the pipeline instance. Adopt it before timers and listeners so a 
previous engine's
+    // finished event cannot own the refresh anymore.
+    //
+    setDisplayedPipeline(runningPipeline);
 
     // Add all the execution result tabs (logging, metrics, etc.)
     addAllTabs();
@@ -6794,14 +6799,9 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
 
     if (isRunning) {
       // Add listeners for when the pipeline finishes (only if still running)
-      pipeline.addExecutionFinishedListener(
-          p -> {
-            checkPipelineEnded();
-            checkErrorVisuals();
-            stopRedrawTimer();
-          });
+      pipeline.addExecutionFinishedListener(this::runPipelineMetricsFinished);
 
-      pipeline.addExecutionStoppedListener(e -> pipelineStopped());
+      pipeline.addExecutionStoppedListener(this::runPipelineStopped);
 
       // Start the redraw timer to continuously update the GUI
       startRedrawTimer();
@@ -6817,28 +6817,79 @@ public class HopGuiPipelineGraph extends 
HopGuiAbstractGraph
   }
 
   private void startRedrawTimer() {
-
-    redrawTimer = new Timer("HopGuiPipelineGraph: redraw timer");
-    TimerTask timtask =
-        new TimerTask() {
-          @Override
-          public void run() {
-            if (!hopDisplay().isDisposed()) {
-              hopDisplay()
-                  .asyncExec(
-                      () -> {
-                        if (!HopGuiPipelineGraph.this.canvas.isDisposed()
-                            && perspective.isActive()
-                            && HopGuiPipelineGraph.this.isVisible()) {
-                          HopGuiPipelineGraph.this.canvas.redraw();
-                          HopGuiPipelineGraph.this.updateGui();
-                        }
-                      });
-            }
+    ExecutionGuiSession.Snapshot snapshot = executionGuiSession.current();
+    if (!(snapshot.engine() instanceof IPipelineEngine<?> engine)) {
+      return;
+    }
+    executionGuiSession.scheduleWhileCurrent(
+        snapshot,
+        "HopGuiPipelineGraph: redraw timer",
+        ConstUi.INTERVAL_MS_PIPELINE_CANVAS_REFRESH,
+        () -> !engine.isFinished(),
+        timer -> {
+          ExecutorUtil.cleanup(redrawTimer);
+          redrawTimer = timer;
+        },
+        () -> {
+          if (hopDisplay().isDisposed()) {
+            return;
           }
-        };
+          hopDisplay()
+              .asyncExec(
+                  () -> {
+                    if (!executionGuiSession.isCurrent(engine, 
snapshot.generation())) {
+                      return;
+                    }
+                    if (!HopGuiPipelineGraph.this.canvas.isDisposed()
+                        && perspective.isActive()
+                        && HopGuiPipelineGraph.this.isVisible()) {
+                      HopGuiPipelineGraph.this.canvas.redraw();
+                      HopGuiPipelineGraph.this.updateGui();
+                    }
+                  });
+        });
+  }
+
+  /**
+   * Publish {@code engine} as the pipeline on screen, together with the 
session the timers check.
+   */
+  private void setDisplayedPipeline(IPipelineEngine<PipelineMeta> engine) {
+    executionGuiSession.adopt(engine, () -> this.pipeline = engine);
+  }
+
+  public ExecutionGuiSession getExecutionGuiSession() {
+    return executionGuiSession;
+  }
+
+  /** Extension-point finish hook. Runs outside the session lock because it 
may open a dialog. */
+  private void runPipelineFinished(IPipelineEngine<PipelineMeta> finished) {
+    if (executionGuiSession.isCurrentEngine(finished)) {
+      pipelineFinished();
+    }
+  }
 
-    redrawTimer.schedule(timtask, 0L, 
ConstUi.INTERVAL_MS_PIPELINE_CANVAS_REFRESH);
+  private void runPipelineStopped(IPipelineEngine<PipelineMeta> stopped) {
+    if (executionGuiSession.isCurrentEngine(stopped)) {
+      pipelineStopped();
+    }
+  }
+
+  /** Metrics refresh for a pipeline that just finished. The timer stop is the 
only locked part. */
+  private void runPipelineMetricsFinished(IPipelineEngine<PipelineMeta> 
finished) {
+    if (!executionGuiSession.isCurrentEngine(finished)) {
+      return;
+    }
+    checkPipelineEnded();
+    checkErrorVisuals();
+    executionGuiSession.stopIfCurrent(finished, this::stopRedrawTimer);
+  }
+
+  private void stopPipelineTimersIfCurrent(IPipelineEngine<PipelineMeta> 
engine) {
+    if (!executionGuiSession.isCurrentEngine(engine)) {
+      return;
+    }
+    checkErrorVisuals();
+    executionGuiSession.stopIfCurrent(engine, this::stopRedrawTimer);
   }
 
   protected void stopRedrawTimer() {
diff --git 
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/delegates/HopGuiPipelineGridDelegate.java
 
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/delegates/HopGuiPipelineGridDelegate.java
index bdfeab9325..f22b1aedaa 100644
--- 
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/delegates/HopGuiPipelineGridDelegate.java
+++ 
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/pipeline/delegates/HopGuiPipelineGridDelegate.java
@@ -29,7 +29,6 @@ import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.Timer;
-import java.util.TimerTask;
 import java.util.concurrent.locks.ReentrantLock;
 import java.util.function.Consumer;
 import lombok.Getter;
@@ -49,6 +48,7 @@ import org.apache.hop.pipeline.PipelineMeta;
 import org.apache.hop.pipeline.engine.EngineMetrics;
 import org.apache.hop.pipeline.engine.IEngineComponent;
 import org.apache.hop.pipeline.engine.IEngineMetric;
+import org.apache.hop.pipeline.engine.IPipelineEngine;
 import org.apache.hop.pipeline.transform.ITransform;
 import org.apache.hop.pipeline.transform.TransformMeta;
 import org.apache.hop.pipeline.transform.TransformStatus;
@@ -63,6 +63,7 @@ import org.apache.hop.ui.hopgui.ToolbarFacade;
 import org.apache.hop.ui.hopgui.file.IHopFileTypeHandler;
 import org.apache.hop.ui.hopgui.file.pipeline.HopGuiPipelineGraph;
 import org.apache.hop.ui.hopgui.file.pipeline.PipelineMetricDisplayUtil;
+import org.apache.hop.ui.hopgui.file.shared.ExecutionGuiSession;
 import org.eclipse.swt.SWT;
 import org.eclipse.swt.custom.CLabel;
 import org.eclipse.swt.custom.CTabItem;
@@ -327,31 +328,40 @@ public class HopGuiPipelineGridDelegate {
   }
 
   public void startRefreshMetricsTimer() {
-    if (refreshMetricsTimer != null) {
+    // A new run replaces any timer left by the previous one. Returning early 
when a timer already
+    // existed dropped the new run on the floor if the previous finish then 
cancelled that timer.
+    //
+    ExecutionGuiSession session = pipelineGraph.getExecutionGuiSession();
+    ExecutionGuiSession.Snapshot snapshot = session.current();
+    if (!(snapshot.engine() instanceof IPipelineEngine<?> engine)) {
       return;
     }
-
-    // Timer updates the view every UPDATE_TIME_VIEW interval
-    refreshMetricsTimer = new Timer("HopGuiPipelineGraph: " + 
pipelineGraph.getMeta().getName());
-
-    TimerTask refreshMetricsTimerTask =
-        new TimerTask() {
-          @Override
-          public void run() {
-            if (!hopGui.getDisplay().isDisposed()) {
-              
hopGui.getDisplay().asyncExec(HopGuiPipelineGridDelegate.this::refreshView);
-              if (pipelineGraph.getPipeline() != null
-                  && (pipelineGraph.getPipeline().isFinished()
-                      || pipelineGraph.getPipeline().isStopped())
-                  && !pipelineGraph.getPipeline().isReadyToStart()) {
-                ExecutorUtil.cleanup(refreshMetricsTimer, UPDATE_TIME_VIEW + 
10);
-                refreshMetricsTimer = null;
-              }
-            }
+    session.scheduleWhileCurrent(
+        snapshot,
+        "HopGuiPipelineGraph: " + pipelineGraph.getMeta().getName(),
+        UPDATE_TIME_VIEW,
+        null,
+        timer -> {
+          ExecutorUtil.cleanup(refreshMetricsTimer);
+          refreshMetricsTimer = timer;
+        },
+        () -> {
+          if (hopGui.getDisplay().isDisposed()) {
+            return;
           }
-        };
-
-    refreshMetricsTimer.schedule(refreshMetricsTimerTask, 0L, 
UPDATE_TIME_VIEW);
+          hopGui
+              .getDisplay()
+              .asyncExec(
+                  () -> {
+                    if (!session.isCurrent(engine, snapshot.generation())) {
+                      return;
+                    }
+                    refreshView();
+                  });
+          if ((engine.isFinished() || engine.isStopped()) && 
!engine.isReadyToStart()) {
+            session.stopIfCurrent(engine, this::stopRefreshMetricsTimer);
+          }
+        });
   }
 
   public void stopRefreshMetricsTimer() {
diff --git 
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSession.java
 
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSession.java
new file mode 100644
index 0000000000..12b4a51f7e
--- /dev/null
+++ 
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSession.java
@@ -0,0 +1,163 @@
+/*
+ * 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.ui.hopgui.file.shared;
+
+import java.util.Timer;
+import java.util.TimerTask;
+import java.util.function.BooleanSupplier;
+import java.util.function.Consumer;
+
+/**
+ * Binds GUI refresh timers to the execution currently on screen.
+ *
+ * <p>{@link #adopt(Object, Runnable)} and {@link #stopIfCurrent(Object, 
Runnable)} share one lock,
+ * so a finished listener from the previous engine cannot stop the timers of 
the engine that
+ * replaced it. {@code stop} and anything run from {@link 
#scheduleWhileCurrent} must not take the
+ * graph lock or wait on the GUI thread. The GUI thread takes this lock while 
adopting an engine,
+ * and a listener that waits for the GUI while holding it deadlocks.
+ */
+public final class ExecutionGuiSession {
+
+  /** One engine and the generation assigned to it, read under the session 
lock. */
+  public record Snapshot(Object engine, int generation) {}
+
+  private final Object lock = new Object();
+  private Object engine;
+  private int generation;
+
+  /**
+   * Show {@code engine}. {@code bind} runs under the lock before the session 
publishes the engine,
+   * so readers either see the previous pair or this one.
+   */
+  public Snapshot adopt(Object engine, Runnable bind) {
+    synchronized (lock) {
+      if (bind != null) {
+        bind.run();
+      }
+      this.engine = engine;
+      return new Snapshot(engine, ++generation);
+    }
+  }
+
+  public Snapshot adopt(Object engine) {
+    return adopt(engine, null);
+  }
+
+  /** Engine and generation as of one lock acquisition. */
+  public Snapshot current() {
+    synchronized (lock) {
+      return new Snapshot(engine, generation);
+    }
+  }
+
+  public boolean isCurrent(Object engine, int generation) {
+    synchronized (lock) {
+      return engine != null && this.engine == engine && this.generation == 
generation;
+    }
+  }
+
+  /** True when {@code engine} is the one on screen, ignoring generation. */
+  public boolean isCurrentEngine(Object engine) {
+    synchronized (lock) {
+      return engine != null && this.engine == engine;
+    }
+  }
+
+  /**
+   * Run {@code action} only while {@code engine} is still the one adopted at 
{@code generation}.
+   * The lock is held for the whole action.
+   */
+  public boolean runIfCurrent(Object engine, int generation, Runnable action) {
+    synchronized (lock) {
+      if (engine == null || this.engine != engine || this.generation != 
generation) {
+        return false;
+      }
+      if (action != null) {
+        action.run();
+      }
+      return true;
+    }
+  }
+
+  /**
+   * Run {@code stop} only if {@code engine} is still the one on screen. The 
lock is held across the
+   * check and {@code stop}. {@code stop} must not take the graph lock or wait 
on the GUI thread.
+   */
+  public boolean stopIfCurrent(Object engine, Runnable stop) {
+    synchronized (lock) {
+      if (engine == null || this.engine != engine) {
+        return false;
+      }
+      if (stop != null) {
+        stop.run();
+      }
+      return true;
+    }
+  }
+
+  /**
+   * Replace the caller's timer with one that runs {@code onTick} while {@code 
snapshot} is current.
+   * A tick that finds a newer engine cancels the timer, so a listener that no 
longer owns the
+   * screen does not have to. {@code allow}, {@code replace} and {@code 
onTick} run on the timer
+   * thread or under the session lock and must not wait on the GUI thread.
+   *
+   * @param allow checked under the lock before the timer is stored. When it 
returns false, nothing
+   *     is scheduled.
+   * @param replace stores the new timer and drops the previous one. Called 
under the lock.
+   * @return false when nothing was scheduled
+   */
+  public boolean scheduleWhileCurrent(
+      Snapshot snapshot,
+      String threadName,
+      long periodMs,
+      BooleanSupplier allow,
+      Consumer<Timer> replace,
+      Runnable onTick) {
+    if (snapshot == null || snapshot.engine() == null || replace == null || 
onTick == null) {
+      return false;
+    }
+    Timer timer = new Timer(threadName);
+    TimerTask task =
+        new TimerTask() {
+          @Override
+          public void run() {
+            if (!isCurrent(snapshot.engine(), snapshot.generation())) {
+              timer.cancel();
+              return;
+            }
+            onTick.run();
+          }
+        };
+    boolean[] started = {false};
+    runIfCurrent(
+        snapshot.engine(),
+        snapshot.generation(),
+        () -> {
+          if (allow != null && !allow.getAsBoolean()) {
+            return;
+          }
+          replace.accept(timer);
+          timer.schedule(task, 0L, periodMs);
+          started[0] = true;
+        });
+    if (!started[0]) {
+      timer.cancel();
+    }
+    return started[0];
+  }
+}
diff --git 
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/workflow/HopGuiWorkflowGraph.java
 
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/workflow/HopGuiWorkflowGraph.java
index f3c38ac91a..a45c2cca1c 100644
--- 
a/ui/src/main/java/org/apache/hop/ui/hopgui/file/workflow/HopGuiWorkflowGraph.java
+++ 
b/ui/src/main/java/org/apache/hop/ui/hopgui/file/workflow/HopGuiWorkflowGraph.java
@@ -29,7 +29,6 @@ import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.Timer;
-import java.util.TimerTask;
 import java.util.UUID;
 import lombok.Getter;
 import lombok.Setter;
@@ -139,6 +138,7 @@ import 
org.apache.hop.ui.hopgui.file.delegates.HopGuiNoteLinkSupport;
 import org.apache.hop.ui.hopgui.file.delegates.HopGuiNotePadDelegate;
 import org.apache.hop.ui.hopgui.file.shared.CanvasToolTip;
 import org.apache.hop.ui.hopgui.file.shared.DrillDownGuiPlugin;
+import org.apache.hop.ui.hopgui.file.shared.ExecutionGuiSession;
 import org.apache.hop.ui.hopgui.file.shared.HopGuiAbstractGraph;
 import org.apache.hop.ui.hopgui.file.shared.HopGuiGraphSnapshotUndo;
 import org.apache.hop.ui.hopgui.file.shared.HopGuiTooltipExtension;
@@ -445,6 +445,9 @@ public class HopGuiWorkflowGraph extends HopGuiAbstractGraph
 
   private Timer redrawTimer;
 
+  /** Ties the canvas redraw timer to the workflow currently shown. */
+  private final ExecutionGuiSession executionGuiSession = new 
ExecutionGuiSession();
+
   public HopGuiWorkflowGraph(
       Composite parent,
       final HopGui hopGui,
@@ -4218,7 +4221,7 @@ public class HopGuiWorkflowGraph extends 
HopGuiAbstractGraph
   }
 
   public synchronized void setWorkflow(IWorkflowEngine<WorkflowMeta> workflow) 
{
-    this.workflow = workflow;
+    executionGuiSession.adopt(workflow, () -> this.workflow = workflow);
   }
 
   public void paintControl(PaintEvent e) {
@@ -5423,13 +5426,13 @@ public class HopGuiWorkflowGraph extends 
HopGuiAbstractGraph
           ExtensionPointHandler.callExtensionPoint(
               log, variables, 
HopExtensionPoint.HopGuiWorkflowMetaExecutionStart.id, workflowMeta);
 
-          workflow =
+          setWorkflow(
               WorkflowEngineFactory.createWorkflowEngine(
                   variables,
                   
variables.resolve(executionConfiguration.getRunConfiguration()),
                   hopGui.getMetadataProvider(),
                   runWorkflowMeta,
-                  hopGuiLoggingObject);
+                  hopGuiLoggingObject));
 
           workflow.setLogLevel(executionConfiguration.getLogLevel());
           
workflow.setGatheringMetrics(executionConfiguration.isGatheringMetrics());
@@ -5485,6 +5488,13 @@ public class HopGuiWorkflowGraph extends 
HopGuiAbstractGraph
           }
 
           log.logBasic(BaseMessages.getString(PKG, 
"WorkflowLog.Log.StartingWorkflow"));
+
+          // Listeners before the thread and the timer. A workflow that 
finishes in between would
+          // otherwise leave the redraw timer running.
+          //
+          workflow.addExecutionFinishedListener(this::onWorkflowFinished);
+          workflow.addExecutionStoppedListener(this::onWorkflowStopped);
+
           workflowThread = new Thread(() -> workflow.startExecution());
           workflowThread.start();
 
@@ -5494,11 +5504,6 @@ public class HopGuiWorkflowGraph extends 
HopGuiAbstractGraph
           startRedrawTimer();
 
           updateGui();
-
-          // Attach a listener to notify us that the workflow has finished.
-          //
-          workflow.addExecutionFinishedListener(e -> 
HopGuiWorkflowGraph.this.workflowFinished());
-          workflow.addExecutionStoppedListener(e -> 
HopGuiWorkflowGraph.this.workflowStopped());
           // Show the execution results views
           //
           addAllTabs();
@@ -5508,7 +5513,7 @@ public class HopGuiWorkflowGraph extends 
HopGuiAbstractGraph
               BaseMessages.getString(PKG, 
"WorkflowLog.Dialog.CanNotOpenWorkflow.Title"),
               BaseMessages.getString(PKG, 
"WorkflowLog.Dialog.CanNotOpenWorkflow.Message"),
               e);
-          workflow = null;
+          setWorkflow(null);
         }
       } else {
         MessageBox m = new MessageBox(hopShell(), SWT.OK | SWT.ICON_WARNING);
@@ -5551,27 +5556,36 @@ public class HopGuiWorkflowGraph extends 
HopGuiAbstractGraph
 
   /** This gets called at the very end, when everything is done. */
   protected void workflowFinished() {
-    // Do a final check to see if it all ended...
-    //
-    if (workflow != null && workflow.isInitialized() && workflow.isFinished()) 
{
+    onWorkflowFinished(workflow);
+  }
+
+  private void onWorkflowFinished(IWorkflowEngine<WorkflowMeta> finished) {
+    if (!executionGuiSession.isCurrentEngine(finished)) {
+      return;
+    }
+    if (finished.isInitialized() && finished.isFinished()) {
       log.logBasic(
           BaseMessages.getString(PKG, "WorkflowLog.Log.WorkflowHasEnded", 
workflowMeta.getName()));
     }
-
-    stopRedrawTimer();
-
     updateGui();
+    executionGuiSession.stopIfCurrent(finished, this::stopRedrawTimer);
   }
 
   protected void workflowStopped() {
-    if (workflow != null && workflow.isInitialized() && workflow.isStopped()) {
+    onWorkflowStopped(workflow);
+  }
+
+  private void onWorkflowStopped(IWorkflowEngine<WorkflowMeta> stopped) {
+    if (!executionGuiSession.isCurrentEngine(stopped)) {
+      return;
+    }
+    if (stopped.isInitialized() && stopped.isStopped()) {
       log.logBasic(
           BaseMessages.getString(
               PKG, "WorkflowLog.Log.ProcessingOfWorkflowStopped", 
workflowMeta.getName()));
     }
-
-    stopRedrawTimer();
     updateGui();
+    executionGuiSession.stopIfCurrent(stopped, this::stopRedrawTimer);
   }
 
   @Override
@@ -6068,9 +6082,9 @@ public class HopGuiWorkflowGraph extends 
HopGuiAbstractGraph
 
     if (isRunning) {
       // Add listeners for when the workflow finishes (only if still running)
-      workflow.addExecutionFinishedListener(e -> 
HopGuiWorkflowGraph.this.workflowFinished());
+      workflow.addExecutionFinishedListener(this::onWorkflowFinished);
 
-      workflow.addExecutionStoppedListener(e -> 
HopGuiWorkflowGraph.this.workflowStopped());
+      workflow.addExecutionStoppedListener(this::onWorkflowStopped);
 
       // Start the redraw timer to continuously update the GUI
       startRedrawTimer();
@@ -6085,26 +6099,36 @@ public class HopGuiWorkflowGraph extends 
HopGuiAbstractGraph
   }
 
   private void startRedrawTimer() {
-    redrawTimer = new Timer("WorkflowGraph auto refresh: " + 
workflow.getWorkflowName());
-    TimerTask timerTask =
-        new TimerTask() {
-          @Override
-          public void run() {
-            if (!hopDisplay().isDisposed()) {
-              hopDisplay()
-                  .asyncExec(
-                      () -> {
-                        if (!HopGuiWorkflowGraph.this.canvas.isDisposed()
-                            && perspective.isActive()
-                            && HopGuiWorkflowGraph.this.isVisible()) {
-                          updateGui();
-                        }
-                      });
-            }
+    ExecutionGuiSession.Snapshot snapshot = executionGuiSession.current();
+    if (!(snapshot.engine() instanceof IWorkflowEngine<?> engine)) {
+      return;
+    }
+    executionGuiSession.scheduleWhileCurrent(
+        snapshot,
+        "WorkflowGraph auto refresh: " + engine.getWorkflowName(),
+        ConstUi.INTERVAL_MS_PIPELINE_CANVAS_REFRESH,
+        () -> !engine.isFinished(),
+        timer -> {
+          ExecutorUtil.cleanup(redrawTimer);
+          redrawTimer = timer;
+        },
+        () -> {
+          if (hopDisplay().isDisposed()) {
+            return;
           }
-        };
-
-    redrawTimer.schedule(timerTask, 0L, 
ConstUi.INTERVAL_MS_PIPELINE_CANVAS_REFRESH);
+          hopDisplay()
+              .asyncExec(
+                  () -> {
+                    if (!executionGuiSession.isCurrent(engine, 
snapshot.generation())) {
+                      return;
+                    }
+                    if (!HopGuiWorkflowGraph.this.canvas.isDisposed()
+                        && perspective.isActive()
+                        && HopGuiWorkflowGraph.this.isVisible()) {
+                      updateGui();
+                    }
+                  });
+        });
   }
 
   protected void stopRedrawTimer() {
diff --git 
a/ui/src/test/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSessionTest.java
 
b/ui/src/test/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSessionTest.java
new file mode 100644
index 0000000000..9417ed76c2
--- /dev/null
+++ 
b/ui/src/test/java/org/apache/hop/ui/hopgui/file/shared/ExecutionGuiSessionTest.java
@@ -0,0 +1,122 @@
+/*
+ * 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.ui.hopgui.file.shared;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.junit.jupiter.api.Test;
+
+class ExecutionGuiSessionTest {
+
+  @Test
+  void stopIfCurrentRunsForTheEngineOnScreen() {
+    ExecutionGuiSession session = new ExecutionGuiSession();
+    Object engine = new Object();
+    ExecutionGuiSession.Snapshot snapshot = session.adopt(engine);
+    AtomicBoolean stopped = new AtomicBoolean(false);
+
+    assertTrue(session.stopIfCurrent(engine, () -> stopped.set(true)));
+    assertTrue(stopped.get());
+    assertTrue(session.isCurrent(snapshot.engine(), snapshot.generation()));
+    assertTrue(session.isCurrentEngine(engine));
+  }
+
+  @Test
+  void stopIfCurrentIgnoresTheEngineThatWasReplaced() {
+    ExecutionGuiSession session = new ExecutionGuiSession();
+    Object first = new Object();
+    Object second = new Object();
+    session.adopt(first);
+    int secondGeneration = session.adopt(second).generation();
+    AtomicBoolean stopped = new AtomicBoolean(false);
+
+    assertFalse(session.stopIfCurrent(first, () -> stopped.set(true)));
+    assertFalse(stopped.get());
+    assertTrue(session.isCurrent(second, secondGeneration));
+    assertFalse(session.runIfCurrent(first, secondGeneration, () -> 
stopped.set(true)));
+    assertFalse(stopped.get());
+  }
+
+  @Test
+  void adoptFromAnotherThreadWaitsUntilTheStopperReleasesTheLock() throws 
Exception {
+    ExecutionGuiSession session = new ExecutionGuiSession();
+    Object first = new Object();
+    Object second = new Object();
+    Object third = new Object();
+    session.adopt(first);
+
+    CountDownLatch holding = new CountDownLatch(1);
+    CountDownLatch release = new CountDownLatch(1);
+    AtomicInteger secondGeneration = new AtomicInteger();
+    AtomicBoolean failed = new AtomicBoolean(false);
+    Thread stopper =
+        new Thread(
+            () -> {
+              try {
+                session.stopIfCurrent(
+                    first,
+                    () -> {
+                      secondGeneration.set(session.adopt(second).generation());
+                      holding.countDown();
+                      try {
+                        if (!release.await(5, TimeUnit.SECONDS)) {
+                          throw new AssertionError("stopper was not released");
+                        }
+                      } catch (InterruptedException e) {
+                        Thread.currentThread().interrupt();
+                        throw new AssertionError(e);
+                      }
+                    });
+              } catch (AssertionError e) {
+                failed.set(true);
+              }
+            });
+    stopper.start();
+    assertTrue(holding.await(5, TimeUnit.SECONDS));
+
+    CountDownLatch reachedAdopt = new CountDownLatch(1);
+    CountDownLatch adopted = new CountDownLatch(1);
+    AtomicInteger thirdGeneration = new AtomicInteger();
+    Thread adopter =
+        new Thread(
+            () -> {
+              reachedAdopt.countDown();
+              thirdGeneration.set(session.adopt(third).generation());
+              adopted.countDown();
+            });
+    adopter.start();
+    assertTrue(reachedAdopt.await(5, TimeUnit.SECONDS));
+    // Do not call into the session here: the stopper still holds its lock and 
is waiting for
+    // release. The adopted latch stays open until that lock is released.
+    assertFalse(adopted.await(300, TimeUnit.MILLISECONDS));
+
+    release.countDown();
+    assertTrue(adopted.await(5, TimeUnit.SECONDS));
+    assertTrue(thirdGeneration.get() > secondGeneration.get());
+    assertTrue(session.isCurrent(third, thirdGeneration.get()));
+    assertFalse(session.isCurrent(second, secondGeneration.get()));
+    stopper.join(5_000);
+    assertFalse(stopper.isAlive());
+    assertFalse(failed.get());
+  }
+}

Reply via email to