diff --git a/openframe-api-service-core/src/main/java/com/openframe/api/service/ForceClientUninstallService.java b/openframe-api-service-core/src/main/java/com/openframe/api/service/ForceClientUninstallService.java index ab5bb25613..b2e4bb926e 100644 --- a/openframe-api-service-core/src/main/java/com/openframe/api/service/ForceClientUninstallService.java +++ b/openframe-api-service-core/src/main/java/com/openframe/api/service/ForceClientUninstallService.java @@ -4,9 +4,13 @@ import com.openframe.api.dto.force.response.ForceAgentStatus; import com.openframe.api.dto.force.response.ForceClientUninstallResponse; import com.openframe.api.dto.force.response.ForceClientUninstallResponseItem; +import com.openframe.data.document.delivery.DeliveryType; import com.openframe.data.document.device.Machine; +import com.openframe.data.nats.delivery.ClientUninstallDeliverySeed; import com.openframe.data.nats.publisher.ClientUninstallNatsPublisher; import com.openframe.data.repository.device.MachineRepository; +import com.openframe.delivery.config.DeliveryProperties; +import com.openframe.delivery.dispatch.DeliveryDispatcher; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -25,6 +29,8 @@ public class ForceClientUninstallService { private final ClientUninstallNatsPublisher clientUninstallNatsPublisher; private final MachineRepository machineRepository; + private final DeliveryProperties deliveryProperties; + private final DeliveryDispatcher deliveryDispatcher; public ForceClientUninstallResponse process(ForceClientUninstallRequest request) { List machineIds = request.getMachineIds(); @@ -57,7 +63,11 @@ private ForceClientUninstallResponseItem processMachine(String machineId) { return buildResponseItem(machineId, ForceAgentStatus.FAILED); } - clientUninstallNatsPublisher.publish(machineId); + if (deliveryProperties.isEnabled(DeliveryType.CLIENT_UNINSTALL)) { + deliveryDispatcher.dispatch(new ClientUninstallDeliverySeed(machineId)); + } else { + clientUninstallNatsPublisher.publish(machineId); + } markPendingDeletion(machine); diff --git a/openframe-api-service-core/src/test/java/com/openframe/api/service/ForceClientUninstallServiceTest.java b/openframe-api-service-core/src/test/java/com/openframe/api/service/ForceClientUninstallServiceTest.java new file mode 100644 index 0000000000..ef5ddd5685 --- /dev/null +++ b/openframe-api-service-core/src/test/java/com/openframe/api/service/ForceClientUninstallServiceTest.java @@ -0,0 +1,90 @@ +package com.openframe.api.service; + +import com.openframe.api.dto.force.request.ForceClientUninstallRequest; +import com.openframe.api.dto.force.response.ForceAgentStatus; +import com.openframe.api.dto.force.response.ForceClientUninstallResponse; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.device.DeviceStatus; +import com.openframe.data.document.device.Machine; +import com.openframe.data.nats.delivery.ClientUninstallDeliverySeed; +import com.openframe.data.nats.publisher.ClientUninstallNatsPublisher; +import com.openframe.data.repository.device.MachineRepository; +import com.openframe.delivery.config.DeliveryProperties; +import com.openframe.delivery.dispatch.DeliveryDispatcher; +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.List; +import java.util.Optional; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class ForceClientUninstallServiceTest { + + private static final String MACHINE_ID = "mach-42"; + + @Mock private ClientUninstallNatsPublisher clientUninstallNatsPublisher; + @Mock private MachineRepository machineRepository; + @Mock private DeliveryProperties deliveryProperties; + @Mock private DeliveryDispatcher deliveryDispatcher; + + @Captor private ArgumentCaptor seedCaptor; + + @InjectMocks private ForceClientUninstallService service; + + private Machine machine; + private ForceClientUninstallRequest request; + + @BeforeEach + void setUp() { + machine = new Machine(); + machine.setMachineId(MACHINE_ID); + machine.setStatus(DeviceStatus.ONLINE); + request = new ForceClientUninstallRequest(); + request.setMachineIds(List.of(MACHINE_ID)); + when(machineRepository.findByMachineId(MACHINE_ID)).thenReturn(Optional.of(machine)); + } + + @Test + void process_flagOff_publishedToJetStreamAndMarkedPendingDeletion() { + // setup + when(deliveryProperties.isEnabled(DeliveryType.CLIENT_UNINSTALL)).thenReturn(false); + + // execution + ForceClientUninstallResponse response = service.process(request); + + // verifications + verify(clientUninstallNatsPublisher).publish(MACHINE_ID); + assertThat(machine.getStatus()).isEqualTo(DeviceStatus.PENDING_DELETION); + verify(machineRepository).save(machine); + verifyNoInteractions(deliveryDispatcher); + assertThat(response.getItems().get(0).getStatus()).isEqualTo(ForceAgentStatus.PROCESSED); + } + + @Test + void process_flagOn_dispatchedThroughEngineAndMarkedPendingDeletion() { + // setup + when(deliveryProperties.isEnabled(DeliveryType.CLIENT_UNINSTALL)).thenReturn(true); + + // execution + ForceClientUninstallResponse response = service.process(request); + + // verifications + verify(deliveryDispatcher).dispatch(seedCaptor.capture()); + assertThat(seedCaptor.getValue().getMachineId()).isEqualTo(MACHINE_ID); + assertThat(machine.getStatus()).isEqualTo(DeviceStatus.PENDING_DELETION); + verify(machineRepository).save(machine); + verifyNoInteractions(clientUninstallNatsPublisher); + assertThat(response.getItems().get(0).getStatus()).isEqualTo(ForceAgentStatus.PROCESSED); + } +} diff --git a/openframe-client-core/src/main/java/com/openframe/client/listener/delivery/DeliveryResultListener.java b/openframe-client-core/src/main/java/com/openframe/client/listener/delivery/DeliveryResultListener.java index 99a5657f5d..4469c7f259 100644 --- a/openframe-client-core/src/main/java/com/openframe/client/listener/delivery/DeliveryResultListener.java +++ b/openframe-client-core/src/main/java/com/openframe/client/listener/delivery/DeliveryResultListener.java @@ -3,7 +3,6 @@ import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.openframe.client.service.NatsTopicMachineIdExtractor; -import com.openframe.data.document.delivery.DeliveryType; import com.openframe.data.nats.delivery.DeliveryResultMessage; import com.openframe.data.nats.listener.AbstractJetStreamPushListener; import com.openframe.delivery.metrics.DeliveryMetrics; @@ -96,13 +95,10 @@ protected void handleMessage(Message message) { private void apply(String machineId, DeliveryResultMessage report) { DeliveryRef delivery = report.getDelivery(); - DeliveryType type = delivery.getType(); - String targetId = delivery.getTargetId(); - String dispatchId = delivery.getDispatchId(); switch (report.getResult()) { - case ACKED -> deliveryTracker.acknowledge(type, targetId, machineId, dispatchId); - case DONE -> deliveryTracker.complete(type, targetId, machineId, dispatchId); - case FAILED -> deliveryTracker.fail(type, targetId, machineId, dispatchId, report.getError()); + case ACKED -> deliveryTracker.acknowledge(delivery, machineId); + case DONE -> deliveryTracker.done(delivery, machineId); + case FAILED -> deliveryTracker.fail(delivery, machineId, report.getError()); } } diff --git a/openframe-client-core/src/test/java/com/openframe/client/listener/delivery/DeliveryResultListenerTest.java b/openframe-client-core/src/test/java/com/openframe/client/listener/delivery/DeliveryResultListenerTest.java index 1572ba422a..fb00df3b40 100644 --- a/openframe-client-core/src/test/java/com/openframe/client/listener/delivery/DeliveryResultListenerTest.java +++ b/openframe-client-core/src/test/java/com/openframe/client/listener/delivery/DeliveryResultListenerTest.java @@ -4,6 +4,7 @@ import com.openframe.client.service.NatsTopicMachineIdExtractor; import com.openframe.data.document.delivery.DeliveryType; import com.openframe.delivery.metrics.DeliveryMetrics; +import com.openframe.delivery.spec.DeliveryRef; import com.openframe.delivery.track.DeliveryTracker; import io.nats.client.Connection; import io.nats.client.Message; @@ -27,6 +28,7 @@ class DeliveryResultListenerTest { private static final String SUBJECT = "machine.mach-42.delivery.result"; private static final String TOOL_AGENT_ID = "fleetmdm-agent"; private static final String DISPATCH_ID = "d-1"; + private static final DeliveryRef REF = new DeliveryRef(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, DISPATCH_ID); private static final String ERROR = "download failed"; private static final String ACKED = "{\"delivery\":{\"type\":\"TOOL_INSTALLATION\",\"targetId\":\"fleetmdm-agent\",\"dispatchId\":\"d-1\"},\"result\":\"ACKED\"}"; @@ -64,7 +66,7 @@ void handleMessage_acked_trackerAcknowledgesDispatch() { listener.handleMessage(message); // verifications - verify(deliveryTracker).acknowledge(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, MACHINE_ID, DISPATCH_ID); + verify(deliveryTracker).acknowledge(REF, MACHINE_ID); verify(message).ack(); } @@ -77,7 +79,7 @@ void handleMessage_done_trackerCompletesDispatch() { listener.handleMessage(message); // verifications - verify(deliveryTracker).complete(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, MACHINE_ID, DISPATCH_ID); + verify(deliveryTracker).done(REF, MACHINE_ID); verify(message).ack(); } @@ -90,7 +92,7 @@ void handleMessage_failed_trackerFailsDispatchWithError() { listener.handleMessage(message); // verifications - verify(deliveryTracker).fail(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, MACHINE_ID, DISPATCH_ID, ERROR); + verify(deliveryTracker).fail(REF, MACHINE_ID, ERROR); verify(message).ack(); } @@ -168,7 +170,7 @@ void handleMessage_malformedPayload_rejectedCountedAndAcked() { void handleMessage_trackerThrows_leftUnackedForRedelivery() { // setup stubMessage(ACKED); - doThrow(new IllegalStateException("mongo down")).when(deliveryTracker).acknowledge(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, MACHINE_ID, DISPATCH_ID); + doThrow(new IllegalStateException("mongo down")).when(deliveryTracker).acknowledge(REF, MACHINE_ID); // execution listener.handleMessage(message); diff --git a/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/CustomMachineDeliveryRepository.java b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/CustomMachineDeliveryRepository.java index 8b6c14da93..50ed6ef145 100644 --- a/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/CustomMachineDeliveryRepository.java +++ b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/CustomMachineDeliveryRepository.java @@ -25,8 +25,6 @@ public interface CustomMachineDeliveryRepository { boolean markDone(String id, String dispatchId, Set from, Instant finishedAt, Instant expiresAt); - boolean markCancelled(String id, Set from, Instant finishedAt, Instant expiresAt); - boolean markCancelled(String id, Set from, Instant dispatchedAt, Instant finishedAt, Instant expiresAt); boolean markFailed(String id, Set from, Instant dispatchedAt, DeliveryFailure failure, Instant finishedAt, Instant expiresAt); diff --git a/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/impl/CustomMachineDeliveryRepositoryImpl.java b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/impl/CustomMachineDeliveryRepositoryImpl.java index d524096044..faddca2802 100644 --- a/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/impl/CustomMachineDeliveryRepositoryImpl.java +++ b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/impl/CustomMachineDeliveryRepositoryImpl.java @@ -110,12 +110,6 @@ public boolean markDone(String id, String dispatchId, Set from, return updateOne(thisDispatch(id, from, dispatchId), update); } - @Override - public boolean markCancelled(String id, Set from, Instant finishedAt, Instant expiresAt) { - Update update = closed(DeliveryStatus.CANCELLED, finishedAt, expiresAt); - return updateOne(stillIn(id, from), update); - } - @Override public boolean markCancelled(String id, Set from, Instant dispatchedAt, Instant finishedAt, Instant expiresAt) { Update update = closed(DeliveryStatus.CANCELLED, finishedAt, expiresAt); diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySeed.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySeed.java new file mode 100644 index 0000000000..8620339a2c --- /dev/null +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySeed.java @@ -0,0 +1,25 @@ +package com.openframe.data.nats.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.delivery.spec.DeliverySeed; +import lombok.AllArgsConstructor; +import lombok.Getter; + +@Getter +@AllArgsConstructor +public class ClientUninstallDeliverySeed implements DeliverySeed { + + private static final String TARGET_ID = "openframe-client"; + + private final String machineId; + + @Override + public DeliveryType getType() { + return DeliveryType.CLIENT_UNINSTALL; + } + + @Override + public String getTargetId() { + return TARGET_ID; + } +} diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySpec.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySpec.java new file mode 100644 index 0000000000..d33fe9bcf6 --- /dev/null +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySpec.java @@ -0,0 +1,55 @@ +package com.openframe.data.nats.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.device.DeviceStatus; +import com.openframe.data.nats.model.ClientUninstallMessage; +import com.openframe.delivery.spec.DeliveryRequest; +import com.openframe.delivery.spec.DeliverySpec; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import java.time.Instant; +import java.util.EnumSet; +import java.util.Set; + +import static java.lang.String.format; + +@Component +@ConditionalOnProperty("spring.cloud.stream.enabled") +public class ClientUninstallDeliverySpec implements DeliverySpec { + + private static final String SUBJECT_TEMPLATE = "machine.%s.client-uninstall"; + + @Override + public DeliveryType getType() { + return DeliveryType.CLIENT_UNINSTALL; + } + + @Override + public Class getPayloadClass() { + return ClientUninstallMessage.class; + } + + @Override + public DeliveryRequest request(ClientUninstallDeliverySeed seed) { + ClientUninstallMessage message = new ClientUninstallMessage(); + message.setIssuedAt(Instant.now().toString()); + return DeliveryRequest.builder() + .type(DeliveryType.CLIENT_UNINSTALL) + .targetId(seed.getTargetId()) + .machineId(seed.getMachineId()) + .payload(message) + .build(); + } + + @Override + public String subject(String machineId) { + return format(SUBJECT_TEMPLATE, machineId); + } + + // the uninstall is the one command a machine marked for deletion is still waiting for + @Override + public Set getDeliverableStatuses() { + return EnumSet.of(DeviceStatus.ONLINE, DeviceStatus.OFFLINE, DeviceStatus.PENDING, DeviceStatus.PENDING_DELETION); + } +} diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySeed.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySeed.java index 058ba4907e..9522532cee 100644 --- a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySeed.java +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySeed.java @@ -17,7 +17,12 @@ public class ToolInstallationDeliverySeed implements DeliverySeed { private final boolean reinstall; @Override - public DeliveryType type() { + public DeliveryType getType() { return DeliveryType.TOOL_INSTALLATION; } + + @Override + public String getTargetId() { + return toolAgent.getKey(); + } } 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 1d2c4f26eb..14a7233f24 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 @@ -1,8 +1,6 @@ package com.openframe.data.nats.delivery; -import com.openframe.data.document.delivery.DeliveryFailure; import com.openframe.data.document.delivery.DeliveryType; -import com.openframe.data.document.delivery.MachineDelivery; import com.openframe.data.document.tool.IntegratedTool; import com.openframe.data.document.toolagent.IntegratedToolAgent; import com.openframe.data.document.toolagent.ToolAgentAsset; @@ -41,15 +39,13 @@ public Class getPayloadClass() { return ToolInstallationMessage.class; } - // targetId must equal the agentType the agent sends in installed-agent, or complete() never finds the row @Override public DeliveryRequest request(ToolInstallationDeliverySeed seed) { IntegratedToolAgent toolAgent = seed.getToolAgent(); ToolInstallationMessage message = buildMessage(toolAgent, seed.getTool(), seed.isReinstall()); - String targetId = toolAgent.getKey(); return DeliveryRequest.builder() .type(DeliveryType.TOOL_INSTALLATION) - .targetId(targetId) + .targetId(seed.getTargetId()) .machineId(seed.getMachineId()) .payload(message) .build(); @@ -60,11 +56,6 @@ public String subject(String machineId) { return format(SUBJECT_TEMPLATE, machineId); } - @Override - public void onFailed(MachineDelivery delivery, DeliveryFailure failure) { - // intentionally empty: a failed install leaves nothing to compensate - } - private ToolInstallationMessage buildMessage(IntegratedToolAgent toolAgent, IntegratedTool tool, boolean reinstall) { String version = toolAgent.getVersion(); ToolInstallationMessage message = new ToolInstallationMessage(); diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/model/ClientUninstallMessage.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/model/ClientUninstallMessage.java index 13ffcd8d57..f9f32e15a4 100644 --- a/openframe-data-nats/src/main/java/com/openframe/data/nats/model/ClientUninstallMessage.java +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/model/ClientUninstallMessage.java @@ -1,9 +1,15 @@ package com.openframe.data.nats.model; +import com.fasterxml.jackson.annotation.JsonInclude; +import com.openframe.delivery.spec.DeliveryPayload; +import com.openframe.delivery.spec.DeliveryRef; import lombok.Data; @Data -public class ClientUninstallMessage { +public class ClientUninstallMessage implements DeliveryPayload { + + @JsonInclude(JsonInclude.Include.NON_NULL) + private DeliveryRef delivery; /** * When the command was issued (ISO-8601 instant). Lets the agent ignore diff --git a/openframe-data-nats/src/test/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySpecTest.java b/openframe-data-nats/src/test/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySpecTest.java new file mode 100644 index 0000000000..d9acdcc035 --- /dev/null +++ b/openframe-data-nats/src/test/java/com/openframe/data/nats/delivery/ClientUninstallDeliverySpecTest.java @@ -0,0 +1,53 @@ +package com.openframe.data.nats.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.device.DeviceStatus; +import com.openframe.data.nats.model.ClientUninstallMessage; +import com.openframe.delivery.spec.DeliveryRequest; +import org.junit.jupiter.api.Test; + +import java.util.Set; + +import static org.assertj.core.api.Assertions.assertThat; + +class ClientUninstallDeliverySpecTest { + + private static final String MACHINE_ID = "mach-42"; + + private final ClientUninstallDeliverySpec spec = new ClientUninstallDeliverySpec(); + + @Test + void request_seed_messageStampedAndTargetIsTheClient() { + // setup + ClientUninstallDeliverySeed seed = new ClientUninstallDeliverySeed(MACHINE_ID); + + // execution + DeliveryRequest request = spec.request(seed); + + // verifications + assertThat(request.getType()).isEqualTo(DeliveryType.CLIENT_UNINSTALL); + assertThat(request.getTargetId()).isEqualTo("openframe-client"); + assertThat(request.getMachineId()).isEqualTo(MACHINE_ID); + assertThat(request.getPayload().getIssuedAt()).isNotBlank(); + assertThat(request.getPayload().getDelivery()).isNull(); + } + + @Test + void getDeliverableStatuses_machineMarkedForDeletion_stillReceivesTheUninstall() { + // execution + Set statuses = spec.getDeliverableStatuses(); + + // verifications + assertThat(statuses).contains(DeviceStatus.PENDING_DELETION, DeviceStatus.ONLINE, DeviceStatus.OFFLINE); + assertThat(statuses).doesNotContain(DeviceStatus.DELETED, DeviceStatus.ARCHIVED, DeviceStatus.DECOMMISSIONED); + } + + @Test + void subject_machineId_machineClientUninstallSubject() { + // execution + String subject = spec.subject(MACHINE_ID); + + // verifications + assertThat(subject).isEqualTo("machine.mach-42.client-uninstall"); + } +} 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 86916f8781..9e9a5d646e 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 @@ -25,7 +25,7 @@ public class DeliveryDispatcher { private final ObjectProvider publisher; public void dispatch(DeliverySeed seed) { - DeliveryType type = seed.type(); + DeliveryType type = seed.getType(); DeliverySpec spec = registry.require(type); DeliveryRequest request = spec.request(seed); DeliveryPayload payload = request.getPayload(); diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySeed.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySeed.java index 85c907aa23..7da02223f5 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySeed.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySeed.java @@ -4,5 +4,9 @@ public interface DeliverySeed { - DeliveryType type(); + DeliveryType getType(); + + String getTargetId(); + + String getMachineId(); } 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 159cc67e30..6931869b11 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 @@ -1,8 +1,10 @@ package com.openframe.delivery.spec; -import com.openframe.data.document.delivery.DeliveryFailure; import com.openframe.data.document.delivery.DeliveryType; -import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.document.device.DeviceStatus; + +import java.util.EnumSet; +import java.util.Set; public interface DeliverySpec { @@ -14,5 +16,8 @@ public interface DeliverySpec String subject(String machineId); - void onFailed(MachineDelivery delivery, DeliveryFailure failure); + // a machine that never connected is still waiting for its first commands; one that is gone or leaving gets none + default Set getDeliverableStatuses() { + return EnumSet.of(DeviceStatus.ONLINE, DeviceStatus.OFFLINE, DeviceStatus.PENDING); + } } 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 33faadec19..b2aca42223 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 @@ -7,6 +7,7 @@ import com.openframe.data.document.delivery.DeliveryStatus; import com.openframe.data.document.delivery.DeliveryType; import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.document.device.DeviceStatus; import com.openframe.data.repository.delivery.MachineDeliveryRepository; import com.openframe.delivery.config.DeliveryProperties; import com.openframe.delivery.track.DeliveryCloser; @@ -25,6 +26,7 @@ import java.time.Instant; import java.util.List; +import java.util.Map; import java.util.Set; import static java.util.stream.Collectors.toSet; @@ -52,14 +54,14 @@ public void retryPending() { return; } Set machineIds = due.stream().map(MachineDelivery::getMachineId).collect(toSet()); - Set gone = machineOnlineStatus.gone(machineIds); + Map statuses = machineOnlineStatus.statuses(machineIds); Set online = machineOnlineStatus.online(machineIds); - due.forEach(delivery -> retryOne(delivery, gone, online, now)); + due.forEach(delivery -> retryOne(delivery, statuses, online, now)); } - private void retryOne(MachineDelivery delivery, Set gone, Set online, Instant now) { + private void retryOne(MachineDelivery delivery, Map statuses, Set online, Instant now) { try { - retryOrClose(delivery, gone, online, now); + retryOrClose(delivery, statuses, online, now); } catch (Exception e) { metrics.recordRowError(); countErrorAndBackOff(delivery, now); @@ -67,13 +69,14 @@ private void retryOne(MachineDelivery delivery, Set gone, Set on } } - private void retryOrClose(MachineDelivery delivery, Set gone, Set online, Instant now) { + private void retryOrClose(MachineDelivery delivery, Map statuses, Set online, Instant now) { String machineId = delivery.getMachineId(); - if (gone.contains(machineId)) { - closer.cancel(delivery, DeliveryStatus.UNACKED, "machine gone", now); + DeliveryType type = delivery.getType(); + DeviceStatus status = statuses.get(machineId); + if (status != null && !registry.require(type).getDeliverableStatuses().contains(status)) { + closer.cancel(delivery, DeliveryStatus.UNACKED, "machine in status " + status + " takes no " + type, now); return; } - DeliveryType type = delivery.getType(); Policy policy = properties.resolve(type); if (!online.contains(machineId)) { parkSkipOrFailOffline(delivery, policy, now); diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/MachineOnlineStatus.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/MachineOnlineStatus.java index f556fe8e8e..11a7cffdc7 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/MachineOnlineStatus.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/MachineOnlineStatus.java @@ -8,10 +8,11 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Component; -import java.util.EnumSet; import java.util.List; +import java.util.Map; import java.util.Set; +import static java.util.stream.Collectors.toMap; import static java.util.stream.Collectors.toSet; @Component @@ -19,22 +20,18 @@ @ConditionalOnProperty(name = "openframe.delivery.sweep.enabled", havingValue = "true") public class MachineOnlineStatus { - private static final Set GONE = EnumSet.of( - DeviceStatus.DELETED, DeviceStatus.ARCHIVED, DeviceStatus.DECOMMISSIONED); - private final MachineRepository machineRepository; public Set online(Set machineIds) { List online = machineRepository.findByMachineIdInAndTelemetryStatus(machineIds, TelemetryStatus.ONLINE); - return ids(online); - } - - public Set gone(Set machineIds) { - List gone = machineRepository.findByMachineIdInAndStatusIn(machineIds, GONE); - return ids(gone); + return online.stream().map(Machine::getMachineId).collect(toSet()); } - private static Set ids(List machines) { - return machines.stream().map(Machine::getMachineId).collect(toSet()); + // whether a machine may still receive a command is the spec's call; a machine the repository does not know stays absent + public Map statuses(Set machineIds) { + List machines = machineRepository.findByMachineIdIn(machineIds); + return machines.stream() + .filter(machine -> machine.getStatus() != null) + .collect(toMap(Machine::getMachineId, Machine::getStatus, (first, second) -> first)); } } diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryCloser.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryCloser.java index b4e4c85bd2..6976c83cdb 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryCloser.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryCloser.java @@ -8,16 +8,11 @@ import com.openframe.delivery.config.DeliveryProperties; import com.openframe.delivery.config.DeliveryProperties.Policy; import com.openframe.delivery.metrics.DeliveryMetrics; -import com.openframe.delivery.spec.DeliveryPayload; -import com.openframe.delivery.spec.DeliverySeed; -import com.openframe.delivery.spec.DeliverySpec; -import com.openframe.delivery.spec.DeliverySpecRegistry; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import java.time.Instant; -import java.util.Optional; import java.util.Set; @Slf4j @@ -26,7 +21,6 @@ public class DeliveryCloser { private final MachineDeliveryRepository repository; - private final DeliverySpecRegistry registry; private final DeliveryProperties properties; private final DeliveryMetrics metrics; @@ -47,7 +41,6 @@ public void fail(MachineDelivery delivery, DeliveryFailure failure, Set from, String re } } - private void notifyAgentError(MachineDelivery delivery) { - notifySpec(delivery, DeliveryFailure.AGENT_ERROR); - } - - private void notifySpec(MachineDelivery delivery, DeliveryFailure failure) { - DeliveryType type = delivery.getType(); - Optional> spec = registry.find(type); - spec.ifPresentOrElse( - registered -> registered.onFailed(delivery, failure), - () -> log.warn("No spec registered for delivery type {}, onFailed skipped", type)); - } - private Instant expiresAt(DeliveryType type, Instant now) { Policy policy = properties.resolve(type); long ttlSeconds = policy.getTtlSeconds(); diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryTracker.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryTracker.java index 9efa2068f7..9a8d74559b 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryTracker.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryTracker.java @@ -5,6 +5,7 @@ import com.openframe.data.repository.delivery.MachineDeliveryRepository; import com.openframe.delivery.config.DeliveryProperties; import com.openframe.delivery.config.DeliveryProperties.Policy; +import com.openframe.delivery.spec.DeliveryRef; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -20,7 +21,11 @@ public class DeliveryTracker { private final DeliveryProperties properties; private final DeliveryCloser closer; - public void acknowledge(DeliveryType type, String targetId, String machineId, String dispatchId) { + // ref = the delivery block the agent copied back from the command; the row is looked up by that exact dispatch + public void acknowledge(DeliveryRef ref, String machineId) { + DeliveryType type = ref.getType(); + String targetId = ref.getTargetId(); + String dispatchId = ref.getDispatchId(); String id = DeliveryId.of(type, targetId, machineId); Instant now = Instant.now(); Policy policy = properties.resolve(type); @@ -34,7 +39,10 @@ public void acknowledge(DeliveryType type, String targetId, String machineId, St } } - public void complete(DeliveryType type, String targetId, String machineId, String dispatchId) { + public void done(DeliveryRef ref, String machineId) { + DeliveryType type = ref.getType(); + String targetId = ref.getTargetId(); + String dispatchId = ref.getDispatchId(); String id = DeliveryId.of(type, targetId, machineId); Instant now = Instant.now(); Instant expiresAt = expiresAt(type, now); @@ -46,19 +54,9 @@ public void complete(DeliveryType type, String targetId, String machineId, Strin } } - public void fail(DeliveryType type, String targetId, String machineId, String dispatchId, String error) { + public void fail(DeliveryRef ref, String machineId, String error) { Instant now = Instant.now(); - closer.failReported(type, targetId, machineId, dispatchId, error, now); - } - - public void cancel(DeliveryType type, String targetId, String machineId) { - String id = DeliveryId.of(type, targetId, machineId); - Instant now = Instant.now(); - Instant expiresAt = expiresAt(type, now); - boolean cancelled = repository.markCancelled(id, DeliveryStatus.OPEN, now, expiresAt); - if (cancelled) { - log.info("Delivery CANCELLED: type={} targetId={} machineId={}", type, targetId, machineId); - } + closer.failReported(ref.getType(), ref.getTargetId(), machineId, ref.getDispatchId(), error, now); } private Instant expiresAt(DeliveryType type, Instant now) { diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/TestSeed.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/TestSeed.java index c6d8052f50..b6e9ee87e4 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/TestSeed.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/TestSeed.java @@ -11,7 +11,12 @@ public class TestSeed implements DeliverySeed { private final String machineId; @Override - public DeliveryType type() { + public DeliveryType getType() { return DeliveryType.TOOL_INSTALLATION; } + + @Override + public String getTargetId() { + return "fleetmdm-agent"; + } } 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 80a1e4dc52..4a95df07fc 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 @@ -6,6 +6,7 @@ import com.openframe.data.document.delivery.DeliveryStatus; import com.openframe.data.document.delivery.DeliveryType; import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.document.device.DeviceStatus; import com.openframe.data.repository.delivery.MachineDeliveryRepository; import com.openframe.delivery.config.DeliveryProperties; import com.openframe.delivery.dispatch.DeliveryPublisher; @@ -26,6 +27,8 @@ import org.mockito.junit.jupiter.MockitoExtension; import java.time.Instant; +import java.util.EnumSet; +import java.util.Map; import java.util.List; import java.util.Set; @@ -41,6 +44,7 @@ import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; @@ -217,7 +221,7 @@ void retryPending_onlineAttemptsExhausted_failedExhaustedOnlyIfStillUnacked() { // verifications verify(closer).fail(eq(delivery), eq(DeliveryFailure.EXHAUSTED), eq(DeliveryStatus.UNACKED), any(Instant.class)); - verifyNoInteractions(registry, metrics); + verifyNoInteractions(metrics); } @Test @@ -236,7 +240,7 @@ void retryPending_offlineFarFromWindowEnd_postponedToNextSweep() { assertThat(dueAtCaptor.getValue()) .isAfterOrEqualTo(nextSweep) .isBefore(nextSweep.plusSeconds(CLOCK_SLACK_SECONDS)); - verifyNoInteractions(registry, closer, metrics); + verifyNoInteractions(closer, metrics); } @Test @@ -283,7 +287,7 @@ void retryPending_offlineWithSkipBehavior_cancelledNotFailed() { // verifications verify(closer).cancel(eq(delivery), eq(DeliveryStatus.UNACKED), any(String.class), any(Instant.class)); verify(closer, never()).fail(eq(delivery), any(DeliveryFailure.class), eq(DeliveryStatus.UNACKED), any(Instant.class)); - verifyNoInteractions(registry, metrics); + verifyNoInteractions(metrics); } @Test @@ -297,7 +301,36 @@ void retryPending_machineDeleted_cancelled() { // verifications verify(closer).cancel(eq(delivery), eq(DeliveryStatus.UNACKED), any(String.class), any(Instant.class)); - verifyNoInteractions(registry, metrics); + verifyNoInteractions(metrics); + } + + @Test + void retryPending_machineLeavingAndTypeDoesNotReachIt_cancelled() { + // setup + stubDue(delivery); + stubMachine(DeviceStatus.PENDING_DELETION, false); + + // execution + service.retryPending(); + + // verifications + verify(closer).cancel(eq(delivery), eq(DeliveryStatus.UNACKED), any(String.class), any(Instant.class)); + verifyNoInteractions(metrics); + } + + @Test + void retryPending_machineLeavingAndTypeReachesIt_waitsForOnlineInstead() { + // setup + stubDue(delivery); + stubMachine(DeviceStatus.PENDING_DELETION, false); + when(spec.getDeliverableStatuses()).thenReturn(EnumSet.complementOf(EnumSet.of(DeviceStatus.DELETED))); + + // execution + service.retryPending(); + + // verifications + verify(closer, never()).cancel(eq(delivery), eq(DeliveryStatus.UNACKED), any(String.class), any(Instant.class)); + verify(repository).postpone(eq(delivery.getId()), eq(DeliveryStatus.UNACKED), eq(dispatchedAt), any(Instant.class)); } @Test @@ -306,8 +339,9 @@ void retryPending_oneRowCorrupt_corruptCountedAndPostponedOtherRepublished() { MachineDelivery corrupt = row(OTHER_MACHINE_ID, CORRUPT_JSON); stubDue(corrupt, delivery); Set both = Set.of(OTHER_MACHINE_ID, MACHINE_ID); - when(machineOnlineStatus.gone(both)).thenReturn(Set.of()); + when(machineOnlineStatus.statuses(both)).thenReturn(Map.of(OTHER_MACHINE_ID, DeviceStatus.ONLINE, MACHINE_ID, DeviceStatus.ONLINE)); when(machineOnlineStatus.online(both)).thenReturn(both); + when(spec.getDeliverableStatuses()).thenReturn(EnumSet.of(DeviceStatus.ONLINE, DeviceStatus.OFFLINE, DeviceStatus.PENDING)); stubSpec(); when(repository.markRepublished(eq(delivery.getId()), eq(DeliveryStatus.UNACKED), eq(dispatchedAt), eq(NO_ATTEMPTS), any(Instant.class))).thenReturn(true); @@ -355,18 +389,22 @@ private void stubDue(MachineDelivery... rows) { } private void stubMachineOnline() { - when(machineOnlineStatus.gone(Set.of(MACHINE_ID))).thenReturn(Set.of()); - when(machineOnlineStatus.online(Set.of(MACHINE_ID))).thenReturn(Set.of(MACHINE_ID)); + stubMachine(DeviceStatus.ONLINE, true); } private void stubMachineNotOnline() { - when(machineOnlineStatus.gone(Set.of(MACHINE_ID))).thenReturn(Set.of()); - when(machineOnlineStatus.online(Set.of(MACHINE_ID))).thenReturn(Set.of()); + stubMachine(DeviceStatus.OFFLINE, false); } private void stubMachineGone() { - when(machineOnlineStatus.gone(Set.of(MACHINE_ID))).thenReturn(Set.of(MACHINE_ID)); - when(machineOnlineStatus.online(Set.of(MACHINE_ID))).thenReturn(Set.of()); + stubMachine(DeviceStatus.DELETED, false); + } + + private void stubMachine(DeviceStatus status, boolean online) { + when(machineOnlineStatus.statuses(Set.of(MACHINE_ID))).thenReturn(Map.of(MACHINE_ID, status)); + when(machineOnlineStatus.online(Set.of(MACHINE_ID))).thenReturn(online ? Set.of(MACHINE_ID) : Set.of()); + lenient().doReturn(spec).when(registry).require(DeliveryType.TOOL_INSTALLATION); + lenient().when(spec.getDeliverableStatuses()).thenReturn(EnumSet.of(DeviceStatus.ONLINE, DeviceStatus.OFFLINE, DeviceStatus.PENDING)); } private void stubSpec() { diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/MachineOnlineStatusTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/MachineOnlineStatusTest.java index f96d77b9be..47698b4d45 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/MachineOnlineStatusTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/MachineOnlineStatusTest.java @@ -10,8 +10,8 @@ import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; -import java.util.EnumSet; import java.util.List; +import java.util.Map; import java.util.Set; import static org.assertj.core.api.Assertions.assertThat; @@ -21,46 +21,43 @@ class MachineOnlineStatusTest { private static final String ONLINE_ID = "mach-1"; - private static final String OFFLINE_ID = "mach-2"; - private static final String DELETED_ID = "mach-3"; - private static final Set GONE = EnumSet.of( - DeviceStatus.DELETED, DeviceStatus.ARCHIVED, DeviceStatus.DECOMMISSIONED); + private static final String LEAVING_ID = "mach-2"; + private static final String UNKNOWN_ID = "mach-3"; + private static final Set ALL = Set.of(ONLINE_ID, LEAVING_ID, UNKNOWN_ID); @Mock private MachineRepository machineRepository; @InjectMocks private MachineOnlineStatus status; @Test - void online_mixedMachines_onlyOnlineIdsReturned() { + void online_machinesTalkingToUs_returned() { // setup - Set asked = Set.of(ONLINE_ID, OFFLINE_ID); - List found = List.of(machine(ONLINE_ID)); - when(machineRepository.findByMachineIdInAndTelemetryStatus(asked, TelemetryStatus.ONLINE)).thenReturn(found); + when(machineRepository.findByMachineIdInAndTelemetryStatus(ALL, TelemetryStatus.ONLINE)) + .thenReturn(List.of(machine(ONLINE_ID, DeviceStatus.PENDING_DELETION))); // execution - Set online = status.online(asked); + Set online = status.online(ALL); // verifications assertThat(online).containsExactly(ONLINE_ID); } @Test - void gone_mixedMachines_onlyRemovedIdsReturned() { + void statuses_knownMachines_mappedUnknownAbsent() { // setup - Set asked = Set.of(ONLINE_ID, DELETED_ID); - List found = List.of(machine(DELETED_ID)); - when(machineRepository.findByMachineIdInAndStatusIn(asked, GONE)).thenReturn(found); + when(machineRepository.findByMachineIdIn(ALL)).thenReturn(List.of(machine(ONLINE_ID, DeviceStatus.ONLINE), machine(LEAVING_ID, DeviceStatus.PENDING_DELETION))); // execution - Set gone = status.gone(asked); + Map statuses = status.statuses(ALL); // verifications - assertThat(gone).containsExactly(DELETED_ID); + assertThat(statuses).containsOnly(Map.entry(ONLINE_ID, DeviceStatus.ONLINE), Map.entry(LEAVING_ID, DeviceStatus.PENDING_DELETION)); } - private static Machine machine(String machineId) { + private static Machine machine(String machineId, DeviceStatus deviceStatus) { Machine machine = new Machine(); machine.setMachineId(machineId); + machine.setStatus(deviceStatus); return machine; } } diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryCloserTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryCloserTest.java index a448ba4338..8bd88368ec 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryCloserTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryCloserTest.java @@ -8,10 +8,6 @@ import com.openframe.delivery.config.DeliveryProperties; import com.openframe.delivery.config.DeliveryTestPolicies; import com.openframe.delivery.metrics.DeliveryMetrics; -import com.openframe.delivery.spec.DeliverySpec; -import com.openframe.delivery.spec.DeliverySpecRegistry; -import com.openframe.delivery.spec.TestPayload; -import com.openframe.delivery.spec.TestSeed; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -19,11 +15,9 @@ import org.mockito.junit.jupiter.MockitoExtension; import java.time.Instant; -import java.util.Optional; import static com.openframe.delivery.config.DeliveryTestPolicies.TTL; import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; @@ -39,9 +33,7 @@ class DeliveryCloserTest { private static final String ERROR = "download failed"; @Mock private MachineDeliveryRepository repository; - @Mock private DeliverySpecRegistry registry; @Mock private DeliveryMetrics metrics; - @Mock private DeliverySpec spec; private DeliveryCloser closer; @@ -56,20 +48,22 @@ void setUp() { delivery = MachineDelivery.builder() .id(DELIVERY_ID) .type(DeliveryType.CLIENT_UNINSTALL) + .targetId(TARGET_ID) + .machineId(MACHINE_ID) + .dispatchId(DISPATCH_ID) .status(DeliveryStatus.PENDING) .attempts(5) .dispatchedAt(dispatchedAt) .build(); DeliveryProperties properties = DeliveryTestPolicies.properties(); - closer = new DeliveryCloser(repository, registry, properties, metrics); + closer = new DeliveryCloser(repository, properties, metrics); } @Test - void fail_rowStillUnacked_rowFailedMetricCountedSpecNotified() { + void fail_rowStillUnacked_rowFailedMetricCounted() { // setup Instant expiresAt = now.plusSeconds(TTL); when(repository.markFailed(DELIVERY_ID, DeliveryStatus.UNACKED, dispatchedAt, DeliveryFailure.EXHAUSTED, now, expiresAt)).thenReturn(true); - doReturn(Optional.of(spec)).when(registry).find(DeliveryType.CLIENT_UNINSTALL); // execution closer.fail(delivery, DeliveryFailure.EXHAUSTED, DeliveryStatus.UNACKED, now); @@ -80,7 +74,6 @@ void fail_rowStillUnacked_rowFailedMetricCountedSpecNotified() { assertThat(delivery.getFinishedAt()).isEqualTo(now); assertThat(delivery.getExpiresAt()).isEqualTo(expiresAt); verify(metrics).recordFailed(DeliveryType.CLIENT_UNINSTALL, DeliveryFailure.EXHAUSTED); - verify(spec).onFailed(delivery, DeliveryFailure.EXHAUSTED); } @Test @@ -94,38 +87,19 @@ void fail_rowAckedMeanwhile_nothingRecorded() { // verifications assertThat(delivery.getStatus()).isEqualTo(DeliveryStatus.PENDING); - verifyNoInteractions(metrics, registry, spec); + verifyNoInteractions(metrics); } @Test - void fail_typeWithoutSpec_rowStillFailedSpecSkipped() { - // setup - Instant expiresAt = now.plusSeconds(TTL); - when(repository.markFailed(DELIVERY_ID, DeliveryStatus.AWAITING_RESULT, dispatchedAt, DeliveryFailure.TIMEOUT, now, expiresAt)).thenReturn(true); - when(registry.find(DeliveryType.CLIENT_UNINSTALL)).thenReturn(Optional.empty()); - - // execution - closer.fail(delivery, DeliveryFailure.TIMEOUT, DeliveryStatus.AWAITING_RESULT, now); - - // verifications - assertThat(delivery.getStatus()).isEqualTo(DeliveryStatus.FAILED); - verify(metrics).recordFailed(DeliveryType.CLIENT_UNINSTALL, DeliveryFailure.TIMEOUT); - verifyNoInteractions(spec); - } - - @Test - void failReported_openRowOfThisDispatch_rowFailedMetricCountedSpecNotified() { + void failReported_openRowOfThisDispatch_rowFailedMetricCounted() { // setup when(repository.markFailed(DELIVERY_ID, DISPATCH_ID, DeliveryStatus.OPEN, DeliveryFailure.AGENT_ERROR, ERROR, now, now.plusSeconds(TTL))).thenReturn(true); - when(repository.findById(DELIVERY_ID)).thenReturn(Optional.of(delivery)); - doReturn(Optional.of(spec)).when(registry).find(DeliveryType.CLIENT_UNINSTALL); // execution closer.failReported(DeliveryType.CLIENT_UNINSTALL, TARGET_ID, MACHINE_ID, DISPATCH_ID, ERROR, now); // verifications verify(metrics).recordFailed(DeliveryType.CLIENT_UNINSTALL, DeliveryFailure.AGENT_ERROR); - verify(spec).onFailed(delivery, DeliveryFailure.AGENT_ERROR); } @Test @@ -137,7 +111,7 @@ void failReported_rowOfAnotherDispatchOrClosed_nothingRecorded() { closer.failReported(DeliveryType.CLIENT_UNINSTALL, TARGET_ID, MACHINE_ID, DISPATCH_ID, ERROR, now); // verifications - verifyNoInteractions(metrics, registry, spec); + verifyNoInteractions(metrics); } @Test @@ -151,6 +125,6 @@ void cancel_rowStillUnackedSameDispatch_rowCancelledWithTtlExpiry() { // verifications verify(repository).markCancelled(DELIVERY_ID, DeliveryStatus.UNACKED, dispatchedAt, now, expiresAt); - verifyNoInteractions(registry, metrics); + verifyNoInteractions(metrics); } } diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryTrackerTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryTrackerTest.java index c9adc32d52..e7189e97d3 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryTrackerTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryTrackerTest.java @@ -4,6 +4,7 @@ import com.openframe.data.document.delivery.DeliveryType; import com.openframe.data.repository.delivery.MachineDeliveryRepository; import com.openframe.delivery.config.DeliveryTestPolicies; +import com.openframe.delivery.spec.DeliveryRef; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -30,6 +31,7 @@ class DeliveryTrackerTest { private static final String DELIVERY_ID = DeliveryId.of(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); private static final String DISPATCH_ID = "d-1"; private static final String ERROR = "download failed"; + private static final DeliveryRef REF = new DeliveryRef(DeliveryType.TOOL_INSTALLATION, TARGET_ID, DISPATCH_ID); @Mock private MachineDeliveryRepository repository; @Mock private DeliveryCloser closer; @@ -45,12 +47,12 @@ void setUp() { } @Test - void acknowledge_typedKey_unackedRowMarkedAckedWithResultDeadline() { + void acknowledge_ref_unackedRowMarkedAckedWithResultDeadline() { // setup when(repository.markAcked(eq(DELIVERY_ID), eq(DISPATCH_ID), eq(DeliveryStatus.UNACKED), atCaptor.capture(), untilCaptor.capture())).thenReturn(true); // execution - tracker.acknowledge(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID, DISPATCH_ID); + tracker.acknowledge(REF, MACHINE_ID); // verifications Instant ackedAt = atCaptor.getValue(); @@ -58,25 +60,12 @@ void acknowledge_typedKey_unackedRowMarkedAckedWithResultDeadline() { } @Test - void complete_typedKey_openOrFailedRowMarkedDoneWithTtlExpiry() { + void done_ref_openOrFailedRowMarkedDoneWithTtlExpiry() { // setup when(repository.markDone(eq(DELIVERY_ID), eq(DISPATCH_ID), eq(DeliveryStatus.COMPLETABLE), atCaptor.capture(), untilCaptor.capture())).thenReturn(true); // execution - tracker.complete(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID, DISPATCH_ID); - - // verifications - Instant finishedAt = atCaptor.getValue(); - assertThat(untilCaptor.getValue()).isEqualTo(finishedAt.plusSeconds(TTL)); - } - - @Test - void cancel_typedKey_openRowMarkedCancelledWithTtlExpiry() { - // setup - when(repository.markCancelled(eq(DELIVERY_ID), eq(DeliveryStatus.OPEN), atCaptor.capture(), untilCaptor.capture())).thenReturn(true); - - // execution - tracker.cancel(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); + tracker.done(REF, MACHINE_ID); // verifications Instant finishedAt = atCaptor.getValue(); @@ -85,10 +74,8 @@ void cancel_typedKey_openRowMarkedCancelledWithTtlExpiry() { @Test void fail_agentReportedError_closerAsked() { - // setup - // execution - tracker.fail(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID, DISPATCH_ID, ERROR); + tracker.fail(REF, MACHINE_ID, ERROR); // verifications verify(closer).failReported(eq(DeliveryType.TOOL_INSTALLATION), eq(TARGET_ID), eq(MACHINE_ID), eq(DISPATCH_ID), eq(ERROR), any(Instant.class));