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(); + } }
