Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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")
Comment thread
semen-flamingo marked this conversation as resolved.
public class NatsDeliveryPublisher implements DeliveryPublisher {

private final NatsMessagePublisher natsMessagePublisher;

@Override
public void publish(String subject, Object payload) {
natsMessagePublisher.publish(subject, payload);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -29,7 +28,6 @@ public class ToolInstallationDeliverySpec implements DeliverySpec<ToolInstallati

private static final String SUBJECT_TEMPLATE = "machine.%s.tool-installation";

private final NatsMessagePublisher natsMessagePublisher;
private final DownloadConfigurationMapper downloadConfigurationMapper;
private final LocalFilenameConfigurationMapper localFilenameConfigurationMapper;

Expand Down Expand Up @@ -58,9 +56,8 @@ public DeliveryRequest<ToolInstallationMessage> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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)
Expand All @@ -31,7 +29,6 @@ class ToolInstallationDeliverySpecTest {
private static final String VERSION = "1.2.3";
private static final List<String> INSTALL_ARGS = List.of("--silent");

@Mock private NatsMessagePublisher natsMessagePublisher;
@Mock private DownloadConfigurationMapper downloadConfigurationMapper;
@Mock private LocalFilenameConfigurationMapper localFilenameConfigurationMapper;

Expand Down Expand Up @@ -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");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<DeliveryPublisher> publisher;

public void dispatch(DeliverySeed seed) {
DeliveryType type = seed.type();
Expand All @@ -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);
Comment thread
semen-flamingo marked this conversation as resolved.
log.info("Delivery dispatched: type={} targetId={} machineId={} dispatchId={}",
type, request.getTargetId(), machineId, dispatchId);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
package com.openframe.delivery.dispatch;

public interface DeliveryPublisher {

void publish(String subject, Object payload);
}
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ public interface DeliverySpec<S extends DeliverySeed, P extends DeliveryPayload>

DeliveryRequest<P> request(S seed);

void publish(String machineId, P payload);
String subject(String machineId);

void onFailed(MachineDelivery delivery, DeliveryFailure failure);
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -132,9 +134,10 @@ private void republish(MachineDelivery delivery, Policy policy, Instant now) {
type, delivery.getTargetId(), machineId, attempt, dueAt);
}

private static boolean publish(DeliverySpec<DeliverySeed, DeliveryPayload> spec, String machineId, DeliveryPayload payload) {
private boolean publish(DeliverySpec<DeliverySeed, DeliveryPayload> 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -30,6 +31,8 @@ class DeliveryDispatcherTest {
@Mock private DeliverySpecRegistry registry;
@Mock private DeliveryRecorder recorder;
@Mock private DeliverySpec<TestSeed, TestPayload> spec;
@Mock private ObjectProvider<DeliveryPublisher> publisherProvider;
@Mock private DeliveryPublisher publisher;

@InjectMocks private DeliveryDispatcher dispatcher;

Expand All @@ -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);
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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";
Expand All @@ -67,6 +69,7 @@ class DeliverySweepServiceTest {
@Mock private DeliveryCloser closer;
@Mock private DeliveryMetrics metrics;
@Mock private DeliverySpec<TestSeed, TestPayload> spec;
@Mock private DeliveryPublisher publisher;

@Captor private ArgumentCaptor<TestPayload> payloadCaptor;
@Captor private ArgumentCaptor<Instant> dueAtCaptor;
Expand All @@ -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
Expand All @@ -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))
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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);
}
}
Loading