From 2fdf810312276f486c5d8b613e079828092c8e3a Mon Sep 17 00:00:00 2001 From: semen-flamingo Date: Wed, 16 Sep 2026 14:46:42 +0200 Subject: [PATCH] feat(delivery): wire ack, success and the tool-installation call site into the delivery engine Second PR of the plan, on top of openframe-machine-delivery: - ScriptExecutionAcknowledgeMessage gains type/targetId; the listener routes any type to DeliveryTracker.acknowledge (PENDING -> ACKED) and keeps the script path for legacy/SCRIPT_SCHEDULE acks. - InstalledAgentService completes TOOL_INSTALLATION rows on installed-agent. - ToolInstallationService dispatches through DeliveryDispatcher with the spec's Seed; ToolInstallationNatsPublisher has no callers left and is deleted. - The spec publishes over core NATS on the same subject; the TOOL_INSTALLATION stream keeps feeding agents that still hold a durable consumer. Still behind openframe.delivery.enabled; with the flag off no row is written and no tracking runs. Co-Authored-By: Claude Fable 5.1 Claude-Session: https://claude.ai/code/session_01BnYUpYgKvuMw6hB5VZSimY --- .../ScriptExecutionAcknowledgeListener.java | 22 +++- .../client/service/InstalledAgentService.java | 4 + ...criptExecutionAcknowledgeListenerTest.java | 114 ++++++++++++++++++ .../service/InstalledAgentServiceTest.java | 103 ++++++++++++++++ .../ToolInstallationDeliverySpec.java | 2 +- .../ToolInstallationNatsPublisher.java | 104 ---------------- .../ScriptExecutionAcknowledgeMessage.java | 5 + .../ToolInstallationDeliverySpecTest.java | 2 +- .../data/service/ToolInstallationService.java | 8 +- 9 files changed, 253 insertions(+), 111 deletions(-) create mode 100644 openframe-client-core/src/test/java/com/openframe/client/listener/rmm/ScriptExecutionAcknowledgeListenerTest.java create mode 100644 openframe-client-core/src/test/java/com/openframe/client/service/InstalledAgentServiceTest.java delete mode 100644 openframe-data-nats/src/main/java/com/openframe/data/nats/publisher/ToolInstallationNatsPublisher.java diff --git a/openframe-client-core/src/main/java/com/openframe/client/listener/rmm/ScriptExecutionAcknowledgeListener.java b/openframe-client-core/src/main/java/com/openframe/client/listener/rmm/ScriptExecutionAcknowledgeListener.java index 7a22c0d443..3f895817b3 100644 --- a/openframe-client-core/src/main/java/com/openframe/client/listener/rmm/ScriptExecutionAcknowledgeListener.java +++ b/openframe-client-core/src/main/java/com/openframe/client/listener/rmm/ScriptExecutionAcknowledgeListener.java @@ -2,6 +2,8 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.openframe.client.service.rmm.ScriptExecutionAcknowledgeService; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.delivery.DeliveryTracker; import com.openframe.data.nats.listener.AbstractJetStreamPushListener; import com.openframe.data.nats.rmm.model.ScriptExecutionAcknowledgeMessage; import io.nats.client.Connection; @@ -17,15 +19,18 @@ public class ScriptExecutionAcknowledgeListener extends AbstractJetStreamPushLis private final ObjectMapper objectMapper; private final ScriptExecutionAcknowledgeService acknowledgeService; + private final DeliveryTracker deliveryTracker; public ScriptExecutionAcknowledgeListener( Connection natsConnection, ObjectMapper objectMapper, - ScriptExecutionAcknowledgeService acknowledgeService + ScriptExecutionAcknowledgeService acknowledgeService, + DeliveryTracker deliveryTracker ) { super(natsConnection); this.objectMapper = objectMapper; this.acknowledgeService = acknowledgeService; + this.deliveryTracker = deliveryTracker; } @Override @@ -58,10 +63,23 @@ protected void handleMessage(Message message) { String payload = new String(message.getData(), StandardCharsets.UTF_8); try { ScriptExecutionAcknowledgeMessage ack = objectMapper.readValue(payload, ScriptExecutionAcknowledgeMessage.class); - acknowledgeService.acknowledge(ack); + if (isDeliveryAck(ack)) { + deliveryTracker.acknowledge(ack.getType(), ack.getTargetId(), ack.getMachineId()); + } + if (isScriptAck(ack)) { + acknowledgeService.acknowledge(ack); + } message.ack(); } catch (Exception e) { log.error("Unexpected error processing execution ack: {}", payload, e); } } + + private static boolean isDeliveryAck(ScriptExecutionAcknowledgeMessage ack) { + return ack.getType() != null && ack.getTargetId() != null; + } + + private static boolean isScriptAck(ScriptExecutionAcknowledgeMessage ack) { + return ack.getType() == null || ack.getType() == DeliveryType.SCRIPT_SCHEDULE; + } } diff --git a/openframe-client-core/src/main/java/com/openframe/client/service/InstalledAgentService.java b/openframe-client-core/src/main/java/com/openframe/client/service/InstalledAgentService.java index 85d1511be4..fae538a0da 100644 --- a/openframe-client-core/src/main/java/com/openframe/client/service/InstalledAgentService.java +++ b/openframe-client-core/src/main/java/com/openframe/client/service/InstalledAgentService.java @@ -4,6 +4,8 @@ import com.openframe.client.exception.MachineNotFoundException; import com.openframe.data.document.installedagents.InstalledAgent; import com.openframe.data.document.tool.ConnectionStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.delivery.DeliveryTracker; import com.openframe.data.repository.device.MachineRepository; import com.openframe.data.repository.installedagents.InstalledAgentRepository; import lombok.RequiredArgsConstructor; @@ -21,6 +23,7 @@ public class InstalledAgentService { private final InstalledAgentRepository installedAgentRepository; private final MachineRepository machineRepository; + private final DeliveryTracker deliveryTracker; @Transactional public void addInstalledAgent(String machineId, String agentType, String version, boolean lastAttempt) { @@ -36,6 +39,7 @@ public void addInstalledAgent(String machineId, String agentType, String version installedAgent -> updateExistingInstalledAgent(installedAgent, version, machineId, agentType), () -> addNewInstalledAgent(machineId, agentType, version) ); + deliveryTracker.complete(DeliveryType.TOOL_INSTALLATION, agentType, machineId); } @Transactional diff --git a/openframe-client-core/src/test/java/com/openframe/client/listener/rmm/ScriptExecutionAcknowledgeListenerTest.java b/openframe-client-core/src/test/java/com/openframe/client/listener/rmm/ScriptExecutionAcknowledgeListenerTest.java new file mode 100644 index 0000000000..64c61e9e28 --- /dev/null +++ b/openframe-client-core/src/test/java/com/openframe/client/listener/rmm/ScriptExecutionAcknowledgeListenerTest.java @@ -0,0 +1,114 @@ +package com.openframe.client.listener.rmm; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.client.service.rmm.ScriptExecutionAcknowledgeService; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.delivery.DeliveryTracker; +import com.openframe.data.nats.rmm.model.ScriptExecutionAcknowledgeMessage; +import io.nats.client.Connection; +import io.nats.client.Message; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class ScriptExecutionAcknowledgeListenerTest { + + private static final String MACHINE_ID = "mach-42"; + private static final String EXECUTION_ID = "exec-1"; + private static final String TOOL_AGENT_ID = "tactical-agent"; + private static final String LEGACY_SCRIPT_ACK = + "{\"executionId\":\"exec-1\",\"machineId\":\"mach-42\",\"scriptIds\":[\"s1\"]}"; + private static final String TOOL_INSTALLATION_ACK = + "{\"type\":\"TOOL_INSTALLATION\",\"targetId\":\"tactical-agent\",\"machineId\":\"mach-42\"}"; + private static final String SCRIPT_SCHEDULE_ACK = + "{\"type\":\"SCRIPT_SCHEDULE\",\"targetId\":\"exec-1\",\"executionId\":\"exec-1\",\"machineId\":\"mach-42\",\"scriptIds\":[\"s1\"]}"; + private static final String MALFORMED = "not json"; + + @Mock private Connection natsConnection; + @Mock private ScriptExecutionAcknowledgeService acknowledgeService; + @Mock private DeliveryTracker deliveryTracker; + @Mock private Message message; + + @Captor private ArgumentCaptor ackCaptor; + + private ScriptExecutionAcknowledgeListener listener; + + @BeforeEach + void setUp() { + listener = new ScriptExecutionAcknowledgeListener(natsConnection, new ObjectMapper(), acknowledgeService, deliveryTracker); + } + + @Test + void handleMessage_legacyScriptAck_scriptServiceOnly() { + // setup + stubPayload(LEGACY_SCRIPT_ACK); + + // execution + listener.handleMessage(message); + + // verifications + verify(acknowledgeService).acknowledge(ackCaptor.capture()); + assertThat(ackCaptor.getValue().getExecutionId()).isEqualTo(EXECUTION_ID); + verifyNoInteractions(deliveryTracker); + verify(message).ack(); + } + + @Test + void handleMessage_toolInstallationAck_trackerOnly() { + // setup + stubPayload(TOOL_INSTALLATION_ACK); + + // execution + listener.handleMessage(message); + + // verifications + verify(deliveryTracker).acknowledge(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, MACHINE_ID); + verifyNoInteractions(acknowledgeService); + verify(message).ack(); + } + + @Test + void handleMessage_scriptScheduleAckWithType_trackerAndScriptService() { + // setup + stubPayload(SCRIPT_SCHEDULE_ACK); + + // execution + listener.handleMessage(message); + + // verifications + verify(deliveryTracker).acknowledge(DeliveryType.SCRIPT_SCHEDULE, EXECUTION_ID, MACHINE_ID); + verify(acknowledgeService).acknowledge(ackCaptor.capture()); + assertThat(ackCaptor.getValue().getExecutionId()).isEqualTo(EXECUTION_ID); + verify(message).ack(); + } + + @Test + void handleMessage_malformedPayload_leftUnacked() { + // setup + stubPayload(MALFORMED); + + // execution + listener.handleMessage(message); + + // verifications + verify(message, never()).ack(); + verifyNoInteractions(acknowledgeService); + verifyNoInteractions(deliveryTracker); + } + + private void stubPayload(String json) { + when(message.getData()).thenReturn(json.getBytes(UTF_8)); + } +} diff --git a/openframe-client-core/src/test/java/com/openframe/client/service/InstalledAgentServiceTest.java b/openframe-client-core/src/test/java/com/openframe/client/service/InstalledAgentServiceTest.java new file mode 100644 index 0000000000..9252801f83 --- /dev/null +++ b/openframe-client-core/src/test/java/com/openframe/client/service/InstalledAgentServiceTest.java @@ -0,0 +1,103 @@ +package com.openframe.client.service; + +import com.openframe.client.exception.MachineNotFoundException; +import com.openframe.delivery.DeliveryTracker; +import com.openframe.data.document.device.Machine; +import com.openframe.data.document.installedagents.InstalledAgent; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.tool.ConnectionStatus; +import com.openframe.data.repository.device.MachineRepository; +import com.openframe.data.repository.installedagents.InstalledAgentRepository; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import java.util.Optional; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class InstalledAgentServiceTest { + + private static final String MACHINE_ID = "mach-42"; + private static final String AGENT_TYPE = "tactical-agent"; + private static final String OLD_VERSION = "1.0.0"; + private static final String VERSION = "1.2.3"; + + @Mock private InstalledAgentRepository installedAgentRepository; + @Mock private MachineRepository machineRepository; + @Mock private DeliveryTracker deliveryTracker; + + @Captor private ArgumentCaptor installedAgentCaptor; + + @InjectMocks private InstalledAgentService service; + + private Machine machine; + private InstalledAgent existing; + + @BeforeEach + void setUp() { + machine = new Machine(); + machine.setMachineId(MACHINE_ID); + existing = new InstalledAgent(); + existing.setMachineId(MACHINE_ID); + existing.setAgentType(AGENT_TYPE); + existing.setVersion(OLD_VERSION); + existing.setStatus(ConnectionStatus.DISCONNECTED); + } + + @Test + void addInstalledAgent_newAgent_savedAndDeliveryCompleted() { + // setup + when(machineRepository.findByMachineId(MACHINE_ID)).thenReturn(Optional.of(machine)); + when(installedAgentRepository.findByMachineIdAndAgentType(MACHINE_ID, AGENT_TYPE)).thenReturn(Optional.empty()); + + // execution + service.addInstalledAgent(MACHINE_ID, AGENT_TYPE, VERSION, false); + + // verifications + verify(installedAgentRepository).save(installedAgentCaptor.capture()); + assertThat(installedAgentCaptor.getValue().getVersion()).isEqualTo(VERSION); + assertThat(installedAgentCaptor.getValue().getStatus()).isEqualTo(ConnectionStatus.CONNECTED); + verify(deliveryTracker).complete(DeliveryType.TOOL_INSTALLATION, AGENT_TYPE, MACHINE_ID); + } + + @Test + void addInstalledAgent_existingAgent_versionUpdatedAndDeliveryCompleted() { + // setup + when(machineRepository.findByMachineId(MACHINE_ID)).thenReturn(Optional.of(machine)); + when(installedAgentRepository.findByMachineIdAndAgentType(MACHINE_ID, AGENT_TYPE)).thenReturn(Optional.of(existing)); + + // execution + service.addInstalledAgent(MACHINE_ID, AGENT_TYPE, VERSION, false); + + // verifications + assertThat(existing.getVersion()).isEqualTo(VERSION); + assertThat(existing.getStatus()).isEqualTo(ConnectionStatus.CONNECTED); + verify(installedAgentRepository).save(existing); + verify(deliveryTracker).complete(DeliveryType.TOOL_INSTALLATION, AGENT_TYPE, MACHINE_ID); + } + + @Test + void addInstalledAgent_unknownMachine_throwsAndDeliveryUntouched() { + // setup + when(machineRepository.findByMachineId(MACHINE_ID)).thenReturn(Optional.empty()); + + // execution + MachineNotFoundException ex = assertThrows(MachineNotFoundException.class, + () -> service.addInstalledAgent(MACHINE_ID, AGENT_TYPE, VERSION, false)); + + // verifications + assertThat(ex.getMessage()).contains(MACHINE_ID); + verifyNoInteractions(deliveryTracker); + } +} diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpec.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpec.java index aaccc9c2d8..39c3fcc9f1 100644 --- a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpec.java +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpec.java @@ -82,7 +82,7 @@ public DeliveryRequest request(Seed seed) { @Override public void publish(String machineId, ToolInstallationMessage payload) { String subject = format(SUBJECT_TEMPLATE, machineId); - natsMessagePublisher.publishPersistent(subject, payload); + natsMessagePublisher.publish(subject, payload); } @Override diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/publisher/ToolInstallationNatsPublisher.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/publisher/ToolInstallationNatsPublisher.java deleted file mode 100644 index f0c47b9ded..0000000000 --- a/openframe-data-nats/src/main/java/com/openframe/data/nats/publisher/ToolInstallationNatsPublisher.java +++ /dev/null @@ -1,104 +0,0 @@ -package com.openframe.data.nats.publisher; - -import com.openframe.data.document.tool.IntegratedTool; -import com.openframe.data.document.toolagent.IntegratedToolAgent; -import com.openframe.data.document.toolagent.ToolAgentAsset; -import com.openframe.data.document.toolagent.ToolAgentAssetSource; -import com.openframe.data.nats.mapper.DownloadConfigurationMapper; -import com.openframe.data.nats.mapper.LocalFilenameConfigurationMapper; -import com.openframe.data.nats.model.ToolInstallationMessage; -import com.openframe.data.nats.publisher.NatsMessagePublisher; -import lombok.RequiredArgsConstructor; -import lombok.extern.slf4j.Slf4j; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.stereotype.Component; - -import java.util.List; -import java.util.stream.Collectors; - -import static java.lang.String.format; - -@Component -@RequiredArgsConstructor -@ConditionalOnProperty("spring.cloud.stream.enabled") -@Slf4j -public class ToolInstallationNatsPublisher { - - private final static String TOPIC_NAME_TEMPLATE = "machine.%s.tool-installation"; - - private final NatsMessagePublisher natsMessagePublisher; - private final DownloadConfigurationMapper downloadConfigurationMapper; - private final LocalFilenameConfigurationMapper localFilenameConfigurationMapper; - - public void publish(String machineId, IntegratedToolAgent toolAgent, IntegratedTool tool) { - publish(machineId, toolAgent, tool, false); - } - - public void publish(String machineId, IntegratedToolAgent toolAgent, IntegratedTool tool, boolean reinstall) { - String topicName = buildTopicName(machineId); - ToolInstallationMessage message = buildMessage(toolAgent, tool, reinstall); - natsMessagePublisher.publishPersistent(topicName, message); - } - - private String buildTopicName(String machineId) { - return format(TOPIC_NAME_TEMPLATE, machineId); - } - - private ToolInstallationMessage buildMessage(IntegratedToolAgent toolAgent, IntegratedTool tool) { - return buildMessage(toolAgent, tool, false); - } - - private ToolInstallationMessage buildMessage(IntegratedToolAgent toolAgent, IntegratedTool tool, boolean reinstall) { - ToolInstallationMessage message = new ToolInstallationMessage(); - message.setToolAgentId(toolAgent.getKey()); - // TODO: need refactoring - message.setToolId(toolAgent.getToolId() == null ? "" : toolAgent.getToolId()); - message.setToolType(tool.getToolType() == null ? "" : tool.getToolType()); - - message.setVersion(toolAgent.getVersion()); - message.setSessionType(toolAgent.getSessionType()); - message.setDownloadConfigurations(downloadConfigurationMapper.map(toolAgent.getDownloadConfigurations(), toolAgent.getVersion())); - message.setAssets(mapAssets(toolAgent.getAssets())); - message.setInstallationCommandArgs(toolAgent.getInstallationCommandArgs()); - message.setUninstallationCommandArgs(toolAgent.getUninstallationCommandArgs()); - message.setRunCommandArgs(toolAgent.getRunCommandArgs()); - message.setToolAgentIdCommandArgs(toolAgent.getAgentToolIdCommandArgs()); - message.setReinstall(reinstall); - return message; - } - - private List mapAssets(List assets) { - if (assets == null) { - return null; - } - return assets.stream() - .map(this::mapAsset) - .collect(Collectors.toList()); - } - - private ToolInstallationMessage.Asset mapAsset(ToolAgentAsset asset) { - ToolInstallationMessage.Asset messageAsset = new ToolInstallationMessage.Asset(); - messageAsset.setId(asset.getId()); - messageAsset.setVersion(asset.getVersion()); - messageAsset.setLocalFilenameConfiguration(localFilenameConfigurationMapper.map(asset.getLocalFilenameConfiguration())); - messageAsset.setDownloadConfigurations( - downloadConfigurationMapper.map(asset.getDownloadConfigurations(), asset.getVersion()) - ); - messageAsset.setSource(mapAssetSource(asset.getSource())); - messageAsset.setPath(asset.getPath()); - messageAsset.setExecutable(asset.isExecutable()); - return messageAsset; - } - - private ToolInstallationMessage.AssetSource mapAssetSource(ToolAgentAssetSource source) { - if (source == null) { - return null; - } - return switch (source) { - case ARTIFACTORY -> ToolInstallationMessage.AssetSource.ARTIFACTORY; - case TOOL_API -> ToolInstallationMessage.AssetSource.TOOL_API; - case GITHUB -> ToolInstallationMessage.AssetSource.GITHUB; - }; - } - -} diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/rmm/model/ScriptExecutionAcknowledgeMessage.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/rmm/model/ScriptExecutionAcknowledgeMessage.java index 9100972645..b908c99e85 100644 --- a/openframe-data-nats/src/main/java/com/openframe/data/nats/rmm/model/ScriptExecutionAcknowledgeMessage.java +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/rmm/model/ScriptExecutionAcknowledgeMessage.java @@ -1,5 +1,6 @@ package com.openframe.data.nats.rmm.model; +import com.openframe.data.document.delivery.DeliveryType; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import lombok.AllArgsConstructor; import lombok.Builder; @@ -19,4 +20,8 @@ public class ScriptExecutionAcknowledgeMessage { private String machineId; private String scheduleId; private List scriptIds; + + // null from agents that predate the delivery engine: treat as a script ack + private DeliveryType type; + private String targetId; } diff --git a/openframe-data-nats/src/test/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpecTest.java b/openframe-data-nats/src/test/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpecTest.java index 1a4c3511f3..51e3ef7957 100644 --- a/openframe-data-nats/src/test/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpecTest.java +++ b/openframe-data-nats/src/test/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpecTest.java @@ -98,6 +98,6 @@ void publish_payload_sentToMachineSubject() { spec.publish(MACHINE_ID, message); // verifications - verify(natsMessagePublisher).publishPersistent("machine.mach-42.tool-installation", message); + verify(natsMessagePublisher).publish("machine.mach-42.tool-installation", message); } } diff --git a/openframe-tool-agent-nats-installation/src/main/java/com/openframe/data/service/ToolInstallationService.java b/openframe-tool-agent-nats-installation/src/main/java/com/openframe/data/service/ToolInstallationService.java index f23f850667..d4295b0c6b 100644 --- a/openframe-tool-agent-nats-installation/src/main/java/com/openframe/data/service/ToolInstallationService.java +++ b/openframe-tool-agent-nats-installation/src/main/java/com/openframe/data/service/ToolInstallationService.java @@ -2,7 +2,8 @@ import com.openframe.data.document.tool.IntegratedTool; import com.openframe.data.document.toolagent.IntegratedToolAgent; -import com.openframe.data.nats.publisher.ToolInstallationNatsPublisher; +import com.openframe.data.nats.delivery.ToolInstallationDeliverySpec; +import com.openframe.delivery.DeliveryDispatcher; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -20,7 +21,7 @@ public class ToolInstallationService { private final IntegratedToolService integratedToolService; private final ToolCommandParamsResolver toolCommandParamsResolver; - private final ToolInstallationNatsPublisher toolInstallationNatsPublisher; + private final DeliveryDispatcher deliveryDispatcher; public void process(String machineId, IntegratedToolAgent toolAgent) { process(machineId, toolAgent, false); @@ -41,7 +42,8 @@ public void process(String machineId, IntegratedToolAgent toolAgent, boolean rei List runCommandArgs = toolAgent.getRunCommandArgs(); toolAgent.setRunCommandArgs(toolCommandParamsResolver.process(toolId, runCommandArgs)); - toolInstallationNatsPublisher.publish(machineId, toolAgent, tool, reinstall); + ToolInstallationDeliverySpec.Seed seed = new ToolInstallationDeliverySpec.Seed(machineId, toolAgent, tool, reinstall); + deliveryDispatcher.dispatch(seed); log.info("Published {} agent installation message for machine {}", toolId, machineId); } catch (Exception e) { // TODO: add fallback mechanism