diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/NatsDeliveryPublisher.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/NatsDeliveryPublisher.java new file mode 100644 index 0000000000..af6b4a097e --- /dev/null +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/NatsDeliveryPublisher.java @@ -0,0 +1,20 @@ +package com.openframe.data.nats.delivery; + +import com.openframe.data.nats.publisher.NatsMessagePublisher; +import com.openframe.delivery.dispatch.DeliveryPublisher; +import lombok.RequiredArgsConstructor; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +@Component +@RequiredArgsConstructor +@ConditionalOnProperty("spring.cloud.stream.enabled") +public class NatsDeliveryPublisher implements DeliveryPublisher { + + private final NatsMessagePublisher natsMessagePublisher; + + @Override + public void publish(String subject, Object payload) { + natsMessagePublisher.publish(subject, payload); + } +} 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 8ec5c643b7..1d2c4f26eb 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 @@ -10,7 +10,6 @@ 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 com.openframe.delivery.spec.DeliveryRequest; import com.openframe.delivery.spec.DeliverySpec; import lombok.RequiredArgsConstructor; @@ -29,7 +28,6 @@ public class ToolInstallationDeliverySpec implements DeliverySpec request(ToolInstallationDelivery } @Override - public void publish(String machineId, ToolInstallationMessage payload) { - String subject = format(SUBJECT_TEMPLATE, machineId); - natsMessagePublisher.publish(subject, payload); + public String subject(String machineId) { + return format(SUBJECT_TEMPLATE, machineId); } @Override 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 782e522bc0..5590cde41d 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 @@ -6,7 +6,6 @@ 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 com.openframe.delivery.spec.DeliveryRequest; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -18,7 +17,6 @@ import java.util.List; import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @ExtendWith(MockitoExtension.class) @@ -31,7 +29,6 @@ class ToolInstallationDeliverySpecTest { private static final String VERSION = "1.2.3"; private static final List INSTALL_ARGS = List.of("--silent"); - @Mock private NatsMessagePublisher natsMessagePublisher; @Mock private DownloadConfigurationMapper downloadConfigurationMapper; @Mock private LocalFilenameConfigurationMapper localFilenameConfigurationMapper; @@ -90,14 +87,11 @@ void request_toolWithoutIdAndType_emptyStringsNotNulls() { } @Test - void publish_payload_sentToMachineSubject() { - // setup - ToolInstallationMessage message = new ToolInstallationMessage(); - + void subject_machineId_machineToolInstallationSubject() { // execution - spec.publish(MACHINE_ID, message); + String subject = spec.subject(MACHINE_ID); // verifications - verify(natsMessagePublisher).publish("machine.mach-42.tool-installation", message); + assertThat(subject).isEqualTo("machine.mach-42.tool-installation"); } } diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryDispatcher.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryDispatcher.java index 1561c28bd9..86916f8781 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryDispatcher.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryDispatcher.java @@ -9,6 +9,7 @@ import com.openframe.delivery.spec.DeliverySpecRegistry; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.stereotype.Component; import java.util.UUID; @@ -20,6 +21,8 @@ public class DeliveryDispatcher { private final DeliverySpecRegistry registry; private final DeliveryRecorder recorder; + // ObjectProvider: the dispatcher boots in services without NATS, where no publisher bean exists + private final ObjectProvider publisher; public void dispatch(DeliverySeed seed) { DeliveryType type = seed.type(); @@ -32,7 +35,8 @@ public void dispatch(DeliverySeed seed) { payload.setDelivery(delivery); recorder.record(request); String machineId = request.getMachineId(); - spec.publish(machineId, payload); + String subject = spec.subject(machineId); + publisher.getObject().publish(subject, payload); log.info("Delivery dispatched: type={} targetId={} machineId={} dispatchId={}", type, request.getTargetId(), machineId, dispatchId); } diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryPublisher.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryPublisher.java new file mode 100644 index 0000000000..ce9f54c0c0 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryPublisher.java @@ -0,0 +1,6 @@ +package com.openframe.delivery.dispatch; + +public interface DeliveryPublisher { + + void publish(String subject, Object payload); +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySpec.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySpec.java index 1e461bf637..159cc67e30 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySpec.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySpec.java @@ -12,7 +12,7 @@ public interface DeliverySpec DeliveryRequest

request(S seed); - void publish(String machineId, P payload); + String subject(String machineId); void onFailed(MachineDelivery delivery, DeliveryFailure failure); } diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliverySweepService.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliverySweepService.java index aa7ce270f9..33faadec19 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliverySweepService.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliverySweepService.java @@ -12,6 +12,7 @@ import com.openframe.delivery.track.DeliveryCloser; import com.openframe.delivery.config.DeliveryProperties.Policy; import com.openframe.delivery.config.DeliveryProperties.Sweep; +import com.openframe.delivery.dispatch.DeliveryPublisher; import com.openframe.delivery.metrics.DeliveryMetrics; import com.openframe.delivery.spec.DeliveryPayload; import com.openframe.delivery.spec.DeliverySeed; @@ -40,6 +41,7 @@ public class DeliverySweepService { private final DeliveryProperties properties; private final DeliveryCloser closer; private final DeliveryMetrics metrics; + private final DeliveryPublisher publisher; private final ObjectMapper objectMapper; public void retryPending() { @@ -132,9 +134,10 @@ private void republish(MachineDelivery delivery, Policy policy, Instant now) { type, delivery.getTargetId(), machineId, attempt, dueAt); } - private static boolean publish(DeliverySpec spec, String machineId, DeliveryPayload payload) { + private boolean publish(DeliverySpec spec, String machineId, DeliveryPayload payload) { try { - spec.publish(machineId, payload); + String subject = spec.subject(machineId); + publisher.publish(subject, payload); return true; } catch (RuntimeException e) { log.error("Delivery publish failed: machineId={}", machineId, e); diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/dispatch/DeliveryDispatcherTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/dispatch/DeliveryDispatcherTest.java index e6569755a7..8be953d96c 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/dispatch/DeliveryDispatcherTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/dispatch/DeliveryDispatcherTest.java @@ -12,6 +12,7 @@ import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.beans.factory.ObjectProvider; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -30,6 +31,8 @@ class DeliveryDispatcherTest { @Mock private DeliverySpecRegistry registry; @Mock private DeliveryRecorder recorder; @Mock private DeliverySpec spec; + @Mock private ObjectProvider publisherProvider; + @Mock private DeliveryPublisher publisher; @InjectMocks private DeliveryDispatcher dispatcher; @@ -51,10 +54,12 @@ void setUp() { } @Test - void dispatch_seed_dispatchIdSetThenRecordedThenPublishedThroughSpec() { + void dispatch_seed_dispatchIdSetThenRecordedThenPublishedToSpecSubject() { // setup doReturn(spec).when(registry).require(DeliveryType.TOOL_INSTALLATION); when(spec.request(seed)).thenReturn(request); + when(spec.subject(MACHINE_ID)).thenReturn("machine.mach-42.test"); + when(publisherProvider.getObject()).thenReturn(publisher); // execution dispatcher.dispatch(seed); @@ -64,7 +69,7 @@ void dispatch_seed_dispatchIdSetThenRecordedThenPublishedThroughSpec() { assertThat(payload.getDelivery().getTargetId()).isEqualTo(TARGET_ID); assertThat(payload.getDelivery().getDispatchId()).isNotBlank(); verify(recorder).record(request); - verify(spec).publish(MACHINE_ID, payload); + verify(publisher).publish("machine.mach-42.test", payload); } @Test diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliverySweepServiceTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliverySweepServiceTest.java index a61a916249..80a1e4dc52 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliverySweepServiceTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliverySweepServiceTest.java @@ -8,6 +8,7 @@ import com.openframe.data.document.delivery.MachineDelivery; import com.openframe.data.repository.delivery.MachineDeliveryRepository; import com.openframe.delivery.config.DeliveryProperties; +import com.openframe.delivery.dispatch.DeliveryPublisher; import com.openframe.delivery.track.DeliveryCloser; import com.openframe.delivery.config.DeliveryTestPolicies; import com.openframe.delivery.metrics.DeliveryMetrics; @@ -50,6 +51,7 @@ class DeliverySweepServiceTest { private static final String MACHINE_ID = "mach-42"; private static final String OTHER_MACHINE_ID = "mach-43"; + private static final String SUBJECT = "machine.mach-42.test"; private static final String TARGET_ID = "fleetmdm-agent"; private static final String PAYLOAD_JSON = "{\"value\":\"fleetmdm-agent\"}"; private static final String CORRUPT_JSON = "not-json"; @@ -67,6 +69,7 @@ class DeliverySweepServiceTest { @Mock private DeliveryCloser closer; @Mock private DeliveryMetrics metrics; @Mock private DeliverySpec spec; + @Mock private DeliveryPublisher publisher; @Captor private ArgumentCaptor payloadCaptor; @Captor private ArgumentCaptor dueAtCaptor; @@ -82,7 +85,7 @@ void setUp() { dispatchedAt = Instant.now().minusSeconds(ACK_THRESHOLD * 2); delivery = row(MACHINE_ID, PAYLOAD_JSON); properties = DeliveryTestPolicies.properties(); - service = new DeliverySweepService(repository, machineOnlineStatus, registry, properties, closer, metrics, new ObjectMapper()); + service = new DeliverySweepService(repository, machineOnlineStatus, registry, properties, closer, metrics, publisher, new ObjectMapper()); } @Test @@ -98,7 +101,7 @@ void retryPending_onlineWithAttemptsLeft_publishedThenCountedWithBackoff() { service.retryPending(); // verifications - verify(spec).publish(eq(MACHINE_ID), payloadCaptor.capture()); + verify(publisher).publish(eq(SUBJECT), payloadCaptor.capture()); assertThat(payloadCaptor.getValue().getValue()).isEqualTo(TARGET_ID); assertThat(dueAtCaptor.getValue()) .isAfterOrEqualTo(before.plusSeconds(FIRST_RETRY_DELAY)) @@ -134,7 +137,7 @@ void retryPending_publishThrows_attemptNotCountedRowPostponedWithoutErrorCount() stubDue(delivery); stubMachineOnline(); stubSpec(); - doThrow(new IllegalStateException("nats down")).when(spec).publish(eq(MACHINE_ID), any(TestPayload.class)); + doThrow(new IllegalStateException("nats down")).when(publisher).publish(eq(SUBJECT), any(TestPayload.class)); // execution service.retryPending(); @@ -197,7 +200,7 @@ void retryPending_rowMovedOnWhilePublishing_publishedButNotCounted() { service.retryPending(); // verifications - verify(spec).publish(eq(MACHINE_ID), any(TestPayload.class)); + verify(publisher).publish(eq(SUBJECT), any(TestPayload.class)); verify(metrics, never()).recordRetried(DeliveryType.TOOL_INSTALLATION); verifyNoInteractions(closer); } @@ -313,9 +316,9 @@ void retryPending_oneRowCorrupt_corruptCountedAndPostponedOtherRepublished() { // verifications verify(repository).postponeAfterError(eq(corrupt.getId()), eq(DeliveryStatus.UNACKED), eq(dispatchedAt), any(Instant.class)); - verify(spec).publish(eq(MACHINE_ID), payloadCaptor.capture()); + verify(publisher).publish(eq(SUBJECT), payloadCaptor.capture()); assertThat(payloadCaptor.getValue().getValue()).isEqualTo(TARGET_ID); - verify(spec, never()).publish(eq(OTHER_MACHINE_ID), any(TestPayload.class)); + verify(spec, never()).subject(OTHER_MACHINE_ID); verify(metrics).recordRetried(DeliveryType.TOOL_INSTALLATION); verify(metrics).recordRowError(); } @@ -369,5 +372,6 @@ private void stubMachineGone() { private void stubSpec() { doReturn(spec).when(registry).require(DeliveryType.TOOL_INSTALLATION); when(spec.getPayloadClass()).thenReturn(TestPayload.class); + when(spec.subject(MACHINE_ID)).thenReturn(SUBJECT); } }