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

SvenO3 pushed a commit to branch fix-adapter-producer-lifecycle
in repository https://gitbox.apache.org/repos/asf/streampipes.git

commit 31e0791d530ed99e5239277f7481004524c7be6e
Author: Sven Oehler <[email protected]>
AuthorDate: Mon Jul 13 14:55:48 2026 +0200

    Disconnect adapter broker producer
---
 .../connect/AdapterWorkerManagement.java           | 53 +++++++++++++++-------
 .../connect/adapter/model/EventCollector.java      | 11 +++--
 .../adapter/model/pipeline/AdapterPipeline.java    | 14 +++++-
 .../elements/SendToBrokerAdapterSink.java          |  7 ++-
 .../management/init/RunningAdapterInstance.java    |  9 ++++
 .../management/init/RunningAdapterInstances.java   | 20 +++++---
 .../connect/AdapterWorkerManagementTest.java       | 41 +++++++++++++++++
 7 files changed, 128 insertions(+), 27 deletions(-)

diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
index ea0df9a89d..4f08e20e7e 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagement.java
@@ -20,7 +20,6 @@ package org.apache.streampipes.extensions.management.connect;
 
 import org.apache.streampipes.commons.exceptions.connect.AdapterException;
 import org.apache.streampipes.connect.transformer.api.TransformationEngines;
-import org.apache.streampipes.extensions.api.connect.StreamPipesAdapter;
 import 
org.apache.streampipes.extensions.api.connect.context.IAdapterRuntimeContext;
 import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
 import 
org.apache.streampipes.extensions.management.connect.adapter.model.EventCollector;
@@ -65,10 +64,6 @@ public class AdapterWorkerManagement {
       if (adapter.isPresent()) {
         var newAdapterInstance = 
adapter.get().declareConfig().getSupplier().get();
         validateScriptLanguage(adapterDescription);
-        runningAdapterInstances.addAdapter(
-            adapterDescription.getElementId(),
-            newAdapterInstance,
-            adapterDescription);
 
         // This method allows adapters to modify the adapter description prior 
to invocation.
         // It is particularly useful for adapters like FileReplayAdapter that 
need to manipulate timestamp values
@@ -79,8 +74,19 @@ public class AdapterWorkerManagement {
         var extractor = AdapterParameterExtractor.from(adapterDescription, 
registeredParsers);
         var runtimeContext = 
makeRuntimeContext(adapterDescription.getElementId());
         var eventCollector = EventCollector.from(adapterDescription, 
runtimeContext);
-
-        newAdapterInstance.onAdapterStarted(extractor, eventCollector, 
runtimeContext);
+        runningAdapterInstances.addAdapter(
+            adapterDescription.getElementId(),
+            newAdapterInstance,
+            adapterDescription,
+            eventCollector);
+
+        try {
+          newAdapterInstance.onAdapterStarted(extractor, eventCollector, 
runtimeContext);
+        } catch (AdapterException | RuntimeException e) {
+          
runningAdapterInstances.removeAdapter(adapterDescription.getElementId());
+          closeEventCollector(eventCollector);
+          throw e;
+        }
       } else {
         var errorMessage = "Adapter with id %s could not be 
found".formatted(adapterDescription.getAppId());
         LOG.error(errorMessage);
@@ -97,14 +103,19 @@ public class AdapterWorkerManagement {
 
     adapterTransitionRegistry.registerStopping(elementId);
     try {
-      StreamPipesAdapter adapter = 
RunningAdapterInstances.INSTANCE.removeAdapter(elementId);
-
-      if (adapter != null) {
-
-        var registeredParsers = adapter.declareConfig().getSupportedParsers();
-        var extractor = AdapterParameterExtractor.from(adapterDescription, 
registeredParsers);
-        var runtimeContext = makeRuntimeContext(elementId);
-        adapter.onAdapterStopped(extractor, runtimeContext);
+      var runningAdapter = runningAdapterInstances.removeAdapter(elementId);
+
+      if (runningAdapter != null) {
+        var adapter = runningAdapter.adapter();
+
+        try {
+          var registeredParsers = 
adapter.declareConfig().getSupportedParsers();
+          var extractor = AdapterParameterExtractor.from(adapterDescription, 
registeredParsers);
+          var runtimeContext = makeRuntimeContext(elementId);
+          adapter.onAdapterStopped(extractor, runtimeContext);
+        } finally {
+          closeEventCollector(runningAdapter.eventCollector());
+        }
       }
 
       resetMonitoring(elementId);
@@ -113,7 +124,7 @@ public class AdapterWorkerManagement {
     }
   }
 
-  private IAdapterRuntimeContext makeRuntimeContext(String adapterInstanceId) {
+  protected IAdapterRuntimeContext makeRuntimeContext(String 
adapterInstanceId) {
     return new AdapterContextGenerator().makeRuntimeContext(adapterInstanceId);
   }
 
@@ -131,4 +142,14 @@ public class AdapterWorkerManagement {
   private void resetMonitoring(String elementId) {
     SpMonitoringManager.INSTANCE.reset(elementId);
   }
+
+  private void closeEventCollector(EventCollector eventCollector) {
+    try {
+      if (eventCollector != null) {
+        eventCollector.close();
+      }
+    } catch (RuntimeException e) {
+      LOG.error("Could not close adapter event collector", e);
+    }
+  }
 }
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/model/EventCollector.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/model/EventCollector.java
index 13929db693..802d5c08f4 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/model/EventCollector.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/model/EventCollector.java
@@ -26,7 +26,7 @@ import 
org.apache.streampipes.model.connect.adapter.AdapterDescription;
 
 import java.util.Map;
 
-public class EventCollector implements IEventCollector {
+public class EventCollector implements IEventCollector, AutoCloseable {
 
   private final AdapterPipeline adapterPipeline;
   private final IAdapterRuntimeContext runtimeContext;
@@ -37,8 +37,8 @@ public class EventCollector implements IEventCollector {
     this.runtimeContext = runtimeContext;
   }
 
-  public static IEventCollector from(AdapterDescription adapterDescription,
-                                     IAdapterRuntimeContext runtimeContext) {
+  public static EventCollector from(AdapterDescription adapterDescription,
+                                    IAdapterRuntimeContext runtimeContext) {
     var adapterPipeline = new 
AdapterPipelineGenerator().generatePipeline(adapterDescription);
     return new EventCollector(adapterPipeline, runtimeContext);
   }
@@ -51,4 +51,9 @@ public class EventCollector implements IEventCollector {
       runtimeContext.getLogger().error(e);
     }
   }
+
+  @Override
+  public void close() {
+    adapterPipeline.close();
+  }
 }
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/model/pipeline/AdapterPipeline.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/model/pipeline/AdapterPipeline.java
index 4c924fd632..c02fd47acb 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/model/pipeline/AdapterPipeline.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/model/pipeline/AdapterPipeline.java
@@ -18,6 +18,7 @@
 
 package 
org.apache.streampipes.extensions.management.connect.adapter.model.pipeline;
 
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
 import 
org.apache.streampipes.connect.shared.preprocessing.elements.ScriptTransformationPipelineElement;
 import org.apache.streampipes.connect.transformer.api.Context;
 import org.apache.streampipes.extensions.api.connect.IAdapterPipeline;
@@ -29,7 +30,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.function.Function;
 
-public class AdapterPipeline implements IAdapterPipeline {
+public class AdapterPipeline implements IAdapterPipeline, AutoCloseable {
 
   private List<IAdapterPipelineElement> pipelineElements;
   private IAdapterPipelineElement pipelineSink;
@@ -101,4 +102,15 @@ public class AdapterPipeline implements IAdapterPipeline {
   public EventSchema getResultingEventSchema() {
     return resultingEventSchema;
   }
+
+  @Override
+  public void close() {
+    if (pipelineSink instanceof AutoCloseable closeable) {
+      try {
+        closeable.close();
+      } catch (Exception e) {
+        throw new SpRuntimeException(e);
+      }
+    }
+  }
 }
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
index 04910de464..4ef6ff2386 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
@@ -33,7 +33,7 @@ import 
org.apache.streampipes.model.grounding.TransportProtocol;
 
 import java.util.Map;
 
-public class SendToBrokerAdapterSink implements IAdapterPipelineElement {
+public class SendToBrokerAdapterSink implements IAdapterPipelineElement, 
AutoCloseable {
 
   protected AdapterDescription adapterDescription;
   protected SpDataFormatDefinition dataFormatDefinition;
@@ -85,6 +85,11 @@ public class SendToBrokerAdapterSink implements 
IAdapterPipelineElement {
     producer.publish(event);
   }
 
+  @Override
+  public void close() {
+    producer.disconnect();
+  }
+
   public void modifyProtocolForDebugging(TransportProtocol protocol) {
     protocol.setBrokerHostname("localhost");
     if (protocol instanceof KafkaTransportProtocol) {
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstance.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstance.java
new file mode 100644
index 0000000000..8a0a394863
--- /dev/null
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstance.java
@@ -0,0 +1,9 @@
+package org.apache.streampipes.extensions.management.init;
+
+import org.apache.streampipes.extensions.api.connect.StreamPipesAdapter;
+import 
org.apache.streampipes.extensions.management.connect.adapter.model.EventCollector;
+
+public record RunningAdapterInstance(StreamPipesAdapter adapter,
+                                     EventCollector eventCollector) {
+}
+
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstances.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstances.java
index 40c9787c30..959bed6777 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstances.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstances.java
@@ -19,6 +19,7 @@
 package org.apache.streampipes.extensions.management.init;
 
 import org.apache.streampipes.extensions.api.connect.StreamPipesAdapter;
+import 
org.apache.streampipes.extensions.management.connect.adapter.model.EventCollector;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 
 import java.util.Collection;
@@ -30,22 +31,29 @@ public enum RunningAdapterInstances {
 
   private final Map<String, StreamPipesAdapter> runningAdapterInstances = new 
HashMap<>();
   private final Map<String, AdapterDescription> 
runningAdapterDescriptionInstances = new HashMap<>();
+  private final Map<String, EventCollector> runningAdapterCollectors = new 
HashMap<>();
 
-  public void addAdapter(String elementId, StreamPipesAdapter adapter, 
AdapterDescription adapterDescription) {
+  public void addAdapter(String elementId,
+                         StreamPipesAdapter adapter,
+                         AdapterDescription adapterDescription,
+                         EventCollector collector) {
     runningAdapterInstances.put(elementId, adapter);
     runningAdapterDescriptionInstances.put(elementId, adapterDescription);
+    runningAdapterCollectors.put(elementId, collector);
   }
 
-  public StreamPipesAdapter removeAdapter(String elementId) {
-    StreamPipesAdapter result = runningAdapterInstances.get(elementId);
+  public RunningAdapterInstance removeAdapter(String elementId) {
+    var result = new RunningAdapterInstance(
+        runningAdapterInstances.get(elementId),
+        runningAdapterCollectors.get(elementId)
+    );
     runningAdapterInstances.remove(elementId);
     runningAdapterDescriptionInstances.remove(elementId);
-    return result;
+    runningAdapterCollectors.remove(elementId);
+    return result.adapter() != null ? result : null;
   }
 
   public Collection<AdapterDescription> getAllRunningAdapterDescriptions() {
     return this.runningAdapterDescriptionInstances.values();
   }
-
-
 }
diff --git 
a/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
 
b/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
index 418162f3f9..9aabc40edd 100644
--- 
a/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
+++ 
b/streampipes-extensions-management/src/test/java/org/apache/streampipes/extensions/management/connect/AdapterWorkerManagementTest.java
@@ -19,18 +19,27 @@
 package org.apache.streampipes.extensions.management.connect;
 
 import org.apache.streampipes.commons.exceptions.connect.AdapterException;
+import org.apache.streampipes.extensions.api.connect.IAdapterConfiguration;
+import org.apache.streampipes.extensions.api.connect.StreamPipesAdapter;
+import 
org.apache.streampipes.extensions.api.connect.context.IAdapterRuntimeContext;
+import 
org.apache.streampipes.extensions.management.connect.adapter.model.EventCollector;
 import org.apache.streampipes.extensions.management.init.IDeclarersSingleton;
+import 
org.apache.streampipes.extensions.management.init.RunningAdapterInstances;
 import org.apache.streampipes.model.health.AdapterInstanceState;
 import org.apache.streampipes.sdk.builder.adapter.AdapterConfigurationBuilder;
 
 import org.junit.jupiter.api.Test;
 
+import java.util.List;
 import java.util.Optional;
+import java.util.UUID;
 
 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.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 public class AdapterWorkerManagementTest {
@@ -61,4 +70,36 @@ public class AdapterWorkerManagementTest {
     
assertTrue(adapterTransitionRegistry.getTransitioningAdapterInstanceStates().isEmpty());
 
   }
+
+  @Test
+  public void stopAdapterClosesEventCollectorWhenAdapterStopFails() throws 
AdapterException {
+    var elementId = "adapter-id-" + UUID.randomUUID();
+    var adapterDescription = AdapterConfigurationBuilder
+        .create("id", 0,  null)
+        .build();
+    adapterDescription.setElementId(elementId);
+
+    var adapter = mock(StreamPipesAdapter.class);
+    var adapterConfig = mock(IAdapterConfiguration.class);
+    var eventCollector = mock(EventCollector.class);
+
+    when(adapter.declareConfig()).thenReturn(adapterConfig);
+    when(adapterConfig.getSupportedParsers()).thenReturn(List.of());
+    doThrow(new AdapterException("stop failed"))
+        .when(adapter)
+        .onAdapterStopped(any(), any());
+
+    RunningAdapterInstances.INSTANCE.addAdapter(elementId, adapter, 
adapterDescription, eventCollector);
+    var adapterWorkerManagement = new AdapterWorkerManagement(
+        RunningAdapterInstances.INSTANCE, null, new 
AdapterTransitionRegistry()) {
+      @Override
+      protected IAdapterRuntimeContext makeRuntimeContext(String 
adapterInstanceId) {
+        return mock(IAdapterRuntimeContext.class);
+      }
+    };
+
+    assertThrows(AdapterException.class, () -> 
adapterWorkerManagement.stopAdapter(adapterDescription));
+
+    verify(eventCollector).close();
+  }
 }

Reply via email to