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 new file mode 100644 index 0000000000..99a5657f5d --- /dev/null +++ b/openframe-client-core/src/main/java/com/openframe/client/listener/delivery/DeliveryResultListener.java @@ -0,0 +1,117 @@ +package com.openframe.client.listener.delivery; + +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; +import com.openframe.delivery.spec.DeliveryRef; +import com.openframe.delivery.track.DeliveryTracker; +import io.nats.client.Connection; +import io.nats.client.Message; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +import java.nio.charset.StandardCharsets; + +import static org.springframework.util.StringUtils.hasText; + +@Slf4j +@Component +public class DeliveryResultListener extends AbstractJetStreamPushListener { + + private static final String REJECTED_INCOMPLETE = "incomplete"; + private static final String REJECTED_MALFORMED = "malformed"; + + private final ObjectMapper objectMapper; + private final NatsTopicMachineIdExtractor machineIdExtractor; + private final DeliveryTracker deliveryTracker; + private final DeliveryMetrics metrics; + + public DeliveryResultListener( + Connection natsConnection, + ObjectMapper objectMapper, + NatsTopicMachineIdExtractor machineIdExtractor, + DeliveryTracker deliveryTracker, + DeliveryMetrics metrics + ) { + super(natsConnection); + this.objectMapper = objectMapper; + this.machineIdExtractor = machineIdExtractor; + this.deliveryTracker = deliveryTracker; + this.metrics = metrics; + } + + @Override + protected String getStreamName() { + return DeliveryResultMessage.STREAM; + } + + @Override + protected String getSubject() { + return DeliveryResultMessage.SUBJECT_FILTER; + } + + @Override + protected String getConsumerName() { + return "delivery-result-processor-v1"; + } + + @Override + protected String getDeliveryGroup() { + return "delivery-result"; + } + + @Override + protected String getDeliverySubject() { + return "machine.delivery.result.delivery"; + } + + @Override + protected void handleMessage(Message message) { + String payload = new String(message.getData(), StandardCharsets.UTF_8); + String subject = message.getSubject(); + try { + String machineId = machineIdExtractor.extract(subject); + DeliveryResultMessage report = objectMapper.readValue(payload, DeliveryResultMessage.class); + if (!isComplete(report)) { + metrics.recordResultRejected(REJECTED_INCOMPLETE); + log.error("Delivery result rejected, agent violates the contract (delivery block incomplete or result unknown): machineId={} payload={}", + machineId, payload); + message.ack(); + return; + } + apply(machineId, report); + message.ack(); + } catch (JsonProcessingException | IllegalArgumentException permanentlyBad) { + metrics.recordResultRejected(REJECTED_MALFORMED); + log.error("Delivery result rejected, malformed: subject={} payload={}", subject, payload, permanentlyBad); + message.ack(); + } catch (Exception e) { + log.error("Unexpected error processing delivery result: {}", payload, e); + } + } + + 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()); + } + } + + private static boolean isComplete(DeliveryResultMessage report) { + DeliveryRef delivery = report.getDelivery(); + return delivery != null + && delivery.getType() != null + && hasText(delivery.getTargetId()) + && hasText(delivery.getDispatchId()) + && report.getResult() != null; + } +} 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 new file mode 100644 index 0000000000..1572ba422a --- /dev/null +++ b/openframe-client-core/src/test/java/com/openframe/client/listener/delivery/DeliveryResultListenerTest.java @@ -0,0 +1,184 @@ +package com.openframe.client.listener.delivery; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.client.service.NatsTopicMachineIdExtractor; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.delivery.metrics.DeliveryMetrics; +import com.openframe.delivery.track.DeliveryTracker; +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.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.mockito.Mockito.doThrow; +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 DeliveryResultListenerTest { + + private static final String MACHINE_ID = "mach-42"; + 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 String ERROR = "download failed"; + private static final String ACKED = + "{\"delivery\":{\"type\":\"TOOL_INSTALLATION\",\"targetId\":\"fleetmdm-agent\",\"dispatchId\":\"d-1\"},\"result\":\"ACKED\"}"; + private static final String DONE = + "{\"delivery\":{\"type\":\"TOOL_INSTALLATION\",\"targetId\":\"fleetmdm-agent\",\"dispatchId\":\"d-1\"},\"result\":\"DONE\"}"; + private static final String FAILED = + "{\"delivery\":{\"type\":\"TOOL_INSTALLATION\",\"targetId\":\"fleetmdm-agent\",\"dispatchId\":\"d-1\"},\"result\":\"FAILED\",\"error\":\"download failed\"}"; + private static final String WITHOUT_DISPATCH_ID = + "{\"delivery\":{\"type\":\"TOOL_INSTALLATION\",\"targetId\":\"fleetmdm-agent\"},\"result\":\"ACKED\"}"; + private static final String UNKNOWN_RESULT = + "{\"delivery\":{\"type\":\"TOOL_INSTALLATION\",\"targetId\":\"fleetmdm-agent\",\"dispatchId\":\"d-1\"},\"result\":\"RETRYING\"}"; + private static final String UNKNOWN_TYPE = + "{\"delivery\":{\"type\":\"HOLOGRAM\",\"targetId\":\"fleetmdm-agent\",\"dispatchId\":\"d-1\"},\"result\":\"ACKED\"}"; + private static final String WITHOUT_DELIVERY = "{\"result\":\"ACKED\"}"; + private static final String MALFORMED = "not json"; + + @Mock private Connection natsConnection; + @Mock private DeliveryTracker deliveryTracker; + @Mock private DeliveryMetrics metrics; + @Mock private Message message; + + private DeliveryResultListener listener; + + @BeforeEach + void setUp() { + listener = new DeliveryResultListener(natsConnection, new ObjectMapper(), new NatsTopicMachineIdExtractor(), deliveryTracker, metrics); + } + + @Test + void handleMessage_acked_trackerAcknowledgesDispatch() { + // setup + stubMessage(ACKED); + + // execution + listener.handleMessage(message); + + // verifications + verify(deliveryTracker).acknowledge(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, MACHINE_ID, DISPATCH_ID); + verify(message).ack(); + } + + @Test + void handleMessage_done_trackerCompletesDispatch() { + // setup + stubMessage(DONE); + + // execution + listener.handleMessage(message); + + // verifications + verify(deliveryTracker).complete(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, MACHINE_ID, DISPATCH_ID); + verify(message).ack(); + } + + @Test + void handleMessage_failed_trackerFailsDispatchWithError() { + // setup + stubMessage(FAILED); + + // execution + listener.handleMessage(message); + + // verifications + verify(deliveryTracker).fail(DeliveryType.TOOL_INSTALLATION, TOOL_AGENT_ID, MACHINE_ID, DISPATCH_ID, ERROR); + verify(message).ack(); + } + + @Test + void handleMessage_withoutDispatchId_rejectedCountedAndAcked() { + // setup + stubMessage(WITHOUT_DISPATCH_ID); + + // execution + listener.handleMessage(message); + + // verifications + verifyNoInteractions(deliveryTracker); + verify(metrics).recordResultRejected("incomplete"); + verify(message).ack(); + } + + @Test + void handleMessage_withoutDeliveryBlock_rejectedCountedAndAcked() { + // setup + stubMessage(WITHOUT_DELIVERY); + + // execution + listener.handleMessage(message); + + // verifications + verifyNoInteractions(deliveryTracker); + verify(metrics).recordResultRejected("incomplete"); + verify(message).ack(); + } + + @Test + void handleMessage_unknownResult_rejectedCountedAndAcked() { + // setup + stubMessage(UNKNOWN_RESULT); + + // execution + listener.handleMessage(message); + + // verifications + verifyNoInteractions(deliveryTracker); + verify(metrics).recordResultRejected("incomplete"); + verify(message).ack(); + } + + @Test + void handleMessage_unknownType_rejectedCountedAndAcked() { + // setup + stubMessage(UNKNOWN_TYPE); + + // execution + listener.handleMessage(message); + + // verifications + verifyNoInteractions(deliveryTracker); + verify(metrics).recordResultRejected("incomplete"); + verify(message).ack(); + } + + @Test + void handleMessage_malformedPayload_rejectedCountedAndAcked() { + // setup + stubMessage(MALFORMED); + + // execution + listener.handleMessage(message); + + // verifications + verifyNoInteractions(deliveryTracker); + verify(metrics).recordResultRejected("malformed"); + verify(message).ack(); + } + + @Test + 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); + + // execution + listener.handleMessage(message); + + // verifications + verify(message, never()).ack(); + } + + private void stubMessage(String json) { + when(message.getSubject()).thenReturn(SUBJECT); + when(message.getData()).thenReturn(json.getBytes(UTF_8)); + } +} diff --git a/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryFailure.java b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryFailure.java index b9afc5802f..0b81c0ac0a 100644 --- a/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryFailure.java +++ b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryFailure.java @@ -4,5 +4,6 @@ public enum DeliveryFailure { EXHAUSTED, OFFLINE, TIMEOUT, - ERROR + ERROR, + AGENT_ERROR } diff --git a/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/MachineDelivery.java b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/MachineDelivery.java index cd334ce889..84b584fbaa 100644 --- a/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/MachineDelivery.java +++ b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/MachineDelivery.java @@ -18,7 +18,6 @@ @AllArgsConstructor @Document(collection = "machine_delivery") @CompoundIndex(name = "machine_delivery_due", def = "{'tenantId': 1, 'status': 1, 'dueAt': 1}") -@CompoundIndex(name = "machine_delivery_machine", def = "{'tenantId': 1, 'machineId': 1}") public class MachineDelivery implements TenantScoped { @Id @@ -32,15 +31,16 @@ public class MachineDelivery implements TenantScoped { private DeliveryStatus status; private int attempts; private int errors; + private String dispatchId; private String payloadJson; private Instant dispatchedAt; private Instant dueAt; - private boolean parked; private Instant ackedAt; private Instant finishedAt; private DeliveryFailure failure; + private String error; @Indexed(name = "machine_delivery_ttl", expireAfterSeconds = 0) private Instant expiresAt; 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 2ee53b7173..8b6c14da93 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 @@ -20,17 +20,16 @@ public interface CustomMachineDeliveryRepository { boolean postponeAfterError(String id, Set from, Instant dispatchedAt, Instant dueAt); - boolean park(String id, Set from, Instant dispatchedAt, Instant dueAt); - boolean markAcked(String id, Set from, Instant ackedAt, Instant dueAt); + boolean markAcked(String id, String dispatchId, Set from, Instant ackedAt, Instant dueAt); - boolean markDone(String id, Set from, Instant finishedAt, Instant expiresAt); + 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); + boolean markFailed(String id, String dispatchId, Set from, DeliveryFailure failure, String error, Instant finishedAt, Instant expiresAt); - long wake(String machineId, Set from, Instant dueAt); } 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 ee10d7ba62..d524096044 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 @@ -25,18 +25,18 @@ public class CustomMachineDeliveryRepositoryImpl extends TenantAwareRepositorySu private static final String FIELD_ID = "_id"; private static final String FIELD_TENANT_ID = "tenantId"; - private static final String FIELD_MACHINE_ID = "machineId"; private static final String FIELD_STATUS = "status"; private static final String FIELD_ATTEMPTS = "attempts"; private static final String FIELD_ERRORS = "errors"; + private static final String FIELD_DISPATCH_ID = "dispatchId"; private static final String FIELD_PAYLOAD_JSON = "payloadJson"; private static final String FIELD_DISPATCHED_AT = "dispatchedAt"; private static final String FIELD_DUE_AT = "dueAt"; - private static final String FIELD_PARKED = "parked"; private static final String FIELD_ACKED_AT = "ackedAt"; private static final String FIELD_FINISHED_AT = "finishedAt"; private static final String FIELD_EXPIRES_AT = "expiresAt"; private static final String FIELD_FAILURE = "failure"; + private static final String FIELD_ERROR = "error"; private static final Sort OLDEST_DUE_FIRST = Sort.by(FIELD_DUE_AT); @@ -55,12 +55,17 @@ public List findDue(DeliveryStatus status, Instant before, int return mongoTemplate.find(query, MachineDelivery.class); } - // $set per field, not a replacement document: a replacement is inserted without the tenant of the scoped filter + // $set per field, not a replacement document: a replacement is inserted without the tenant of the scoped filter; + // null fields are not written, so what a previous dispatch closed the row with has to be unset explicitly @Override public void upsertPending(MachineDelivery delivery) { Document document = new Document(); mongoTemplate.getConverter().write(delivery, document); - Update update = new Update().set(FIELD_TENANT_ID, tenantId()); + Update update = new Update().set(FIELD_TENANT_ID, tenantId()) + .unset(FIELD_ACKED_AT) + .unset(FIELD_FINISHED_AT) + .unset(FIELD_FAILURE) + .unset(FIELD_ERROR); document.forEach((field, value) -> setField(update, field, value)); String id = delivery.getId(); Query byId = new Query(Criteria.where(FIELD_ID).is(id)); @@ -91,27 +96,18 @@ public boolean postponeAfterError(String id, Set from, Instant d } @Override - public boolean park(String id, Set from, Instant dispatchedAt, Instant dueAt) { - Update update = new Update() - .set(FIELD_DUE_AT, dueAt) - .set(FIELD_PARKED, true); - return updateOne(sameDispatch(id, from, dispatchedAt), update); - } - - @Override - public boolean markAcked(String id, Set from, Instant ackedAt, Instant dueAt) { + public boolean markAcked(String id, String dispatchId, Set from, Instant ackedAt, Instant dueAt) { Update update = new Update() .set(FIELD_STATUS, DeliveryStatus.ACKED) .set(FIELD_ACKED_AT, ackedAt) - .set(FIELD_DUE_AT, dueAt) - .set(FIELD_PARKED, false); - return updateOne(stillIn(id, from), update); + .set(FIELD_DUE_AT, dueAt); + return updateOne(thisDispatch(id, from, dispatchId), update); } @Override - public boolean markDone(String id, Set from, Instant finishedAt, Instant expiresAt) { + public boolean markDone(String id, String dispatchId, Set from, Instant finishedAt, Instant expiresAt) { Update update = closed(DeliveryStatus.DONE, finishedAt, expiresAt); - return updateOne(stillIn(id, from), update); + return updateOne(thisDispatch(id, from, dispatchId), update); } @Override @@ -134,16 +130,11 @@ public boolean markFailed(String id, Set from, Instant dispatche } @Override - public long wake(String machineId, Set from, Instant dueAt) { - Criteria parkedRowsOfMachine = Criteria.where(FIELD_MACHINE_ID).is(machineId) - .and(FIELD_STATUS).in(from) - .and(FIELD_PARKED).is(true); - Query query = new Query(parkedRowsOfMachine); - Update update = new Update() - .set(FIELD_DUE_AT, dueAt) - .set(FIELD_PARKED, false); - UpdateResult result = mongoTemplate.updateMulti(query, update, MachineDelivery.class); - return result.getModifiedCount(); + public boolean markFailed(String id, String dispatchId, Set from, DeliveryFailure failure, String error, Instant finishedAt, Instant expiresAt) { + Update update = closed(DeliveryStatus.FAILED, finishedAt, expiresAt) + .set(FIELD_FAILURE, failure) + .set(FIELD_ERROR, error); + return updateOne(thisDispatch(id, from, dispatchId), update); } private static void setField(Update update, String field, Object value) { @@ -169,6 +160,10 @@ private static Criteria sameDispatch(String id, Set from, Instan return stillIn(id, from).and(FIELD_DISPATCHED_AT).is(dispatchedAt); } + private static Criteria thisDispatch(String id, Set from, String dispatchId) { + return stillIn(id, from).and(FIELD_DISPATCH_ID).is(dispatchId); + } + private boolean updateOne(Criteria criteria, Update update) { Query query = new Query(criteria); UpdateResult result = mongoTemplate.updateFirst(query, update, MachineDelivery.class); diff --git a/openframe-data-mongo-sync/src/test/java/com/openframe/data/repository/delivery/impl/CustomMachineDeliveryRepositoryImplTest.java b/openframe-data-mongo-sync/src/test/java/com/openframe/data/repository/delivery/impl/CustomMachineDeliveryRepositoryImplTest.java index b4e9a7f28e..022540e6b6 100644 --- a/openframe-data-mongo-sync/src/test/java/com/openframe/data/repository/delivery/impl/CustomMachineDeliveryRepositoryImplTest.java +++ b/openframe-data-mongo-sync/src/test/java/com/openframe/data/repository/delivery/impl/CustomMachineDeliveryRepositoryImplTest.java @@ -33,6 +33,8 @@ class CustomMachineDeliveryRepositoryImplTest { private static final String TENANT_ID = "tenant-1"; private static final int LIMIT = 500; private static final int ATTEMPTS = 1; + private static final String DISPATCH_ID = "d-1"; + private static final String ERROR = "download failed"; @Mock private TenantAwareMongoTemplate mongoTemplate; @Mock private MongoConverter converter; @@ -84,6 +86,26 @@ void upsertPending_row_setsEveryFieldAndTenantWithinScopedUpsert() { .contains("tenantId=" + TENANT_ID); } + @Test + void upsertPending_rowClosedByPreviousDispatch_closingFieldsUnset() { + // setup + MachineDelivery delivery = MachineDelivery.builder().id(ID).machineId(MACHINE_ID).build(); + when(mongoTemplate.getConverter()).thenReturn(converter); + when(mongoTemplate.tenantId()).thenReturn(TENANT_ID); + + // execution + repository.upsertPending(delivery); + + // verifications + verify(mongoTemplate).upsert(queryCaptor.capture(), updateCaptor.capture(), eq(MachineDelivery.class)); + assertThat(updateCaptor.getValue().getUpdateObject().toString()) + .contains("$unset") + .contains("ackedAt") + .contains("finishedAt") + .contains("failure") + .contains("error"); + } + @Test void markRepublished_sameDispatchAndAttemptStillPending_attemptCountedAndTrue() { // setup @@ -119,6 +141,27 @@ void markRepublished_rowMovedOn_false() { assertThat(republished).isFalse(); } + @Test + void markAcked_unackedRowOfThisDispatch_ackedAndTrue() { + // setup + UpdateResult oneRow = UpdateResult.acknowledged(1, 1L, null); + when(mongoTemplate.updateFirst(queryCaptor.capture(), updateCaptor.capture(), eq(MachineDelivery.class))).thenReturn(oneRow); + + // execution + boolean acked = repository.markAcked(ID, DISPATCH_ID, DeliveryStatus.UNACKED, now, now); + + // verifications + assertThat(acked).isTrue(); + assertThat(queryCaptor.getValue().getQueryObject().toString()) + .contains(ID) + .contains("PENDING") + .contains("dispatchId=" + DISPATCH_ID); + assertThat(updateCaptor.getValue().getUpdateObject().toString()) + .contains("ACKED") + .contains("ackedAt") + .contains("dueAt"); + } + @Test void postponeAfterError_pendingRow_errorCountedAndDueMoved() { // setup @@ -159,22 +202,25 @@ void markFailed_openRowSameDispatch_failureAndExpiryWrittenPayloadDropped() { } @Test - void wake_parkedRowsOfMachine_unparkedAndModifiedCountReturned() { + void markFailed_openRowOfThisDispatch_agentErrorWrittenPayloadDropped() { // setup - UpdateResult twoRows = UpdateResult.acknowledged(2, 2L, null); - when(mongoTemplate.updateMulti(queryCaptor.capture(), updateCaptor.capture(), eq(MachineDelivery.class))).thenReturn(twoRows); + UpdateResult oneRow = UpdateResult.acknowledged(1, 1L, null); + when(mongoTemplate.updateFirst(queryCaptor.capture(), updateCaptor.capture(), eq(MachineDelivery.class))).thenReturn(oneRow); // execution - long woken = repository.wake(MACHINE_ID, DeliveryStatus.UNACKED, now); + boolean failed = repository.markFailed(ID, DISPATCH_ID, DeliveryStatus.OPEN, DeliveryFailure.AGENT_ERROR, ERROR, now, now); // verifications - assertThat(woken).isEqualTo(2); + assertThat(failed).isTrue(); assertThat(queryCaptor.getValue().getQueryObject().toString()) - .contains(MACHINE_ID) + .contains(ID) + .contains("dispatchId=" + DISPATCH_ID) .contains("PENDING") - .contains("parked=true"); + .contains("ACKED"); assertThat(updateCaptor.getValue().getUpdateObject().toString()) - .contains("parked=false") - .contains("dueAt"); + .contains("FAILED") + .contains("AGENT_ERROR") + .contains("error=" + ERROR) + .contains("$unset"); } } diff --git a/openframe-data-nats/pom.xml b/openframe-data-nats/pom.xml index d7ec7edd04..219bb25d05 100644 --- a/openframe-data-nats/pom.xml +++ b/openframe-data-nats/pom.xml @@ -28,6 +28,10 @@ com.openframe.oss openframe-data-mongo-sync + + com.openframe.oss + openframe-machine-delivery + org.springframework.boot diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/DeliveryResult.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/DeliveryResult.java new file mode 100644 index 0000000000..29922384eb --- /dev/null +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/DeliveryResult.java @@ -0,0 +1,7 @@ +package com.openframe.data.nats.delivery; + +public enum DeliveryResult { + ACKED, + DONE, + FAILED +} diff --git a/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/DeliveryResultMessage.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/DeliveryResultMessage.java new file mode 100644 index 0000000000..e0e3d20752 --- /dev/null +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/DeliveryResultMessage.java @@ -0,0 +1,20 @@ +package com.openframe.data.nats.delivery; + +import com.fasterxml.jackson.annotation.JsonFormat; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.openframe.delivery.spec.DeliveryRef; +import lombok.Data; + +@Data +@JsonIgnoreProperties(ignoreUnknown = true) +public class DeliveryResultMessage { + + public static final String STREAM = "DELIVERY_RESULT"; + public static final String SUBJECT_FILTER = "machine.*.delivery.result"; + + // the `delivery` block of the command, copied back by the agent as is + private DeliveryRef delivery; + @JsonFormat(with = JsonFormat.Feature.READ_UNKNOWN_ENUM_VALUES_AS_NULL) + private DeliveryResult result; + private String error; +} 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 new file mode 100644 index 0000000000..058ba4907e --- /dev/null +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySeed.java @@ -0,0 +1,23 @@ +package com.openframe.data.nats.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.tool.IntegratedTool; +import com.openframe.data.document.toolagent.IntegratedToolAgent; +import com.openframe.delivery.spec.DeliverySeed; +import lombok.AllArgsConstructor; +import lombok.Getter; + +@Getter +@AllArgsConstructor +public class ToolInstallationDeliverySeed implements DeliverySeed { + + private final String machineId; + private final IntegratedToolAgent toolAgent; + private final IntegratedTool tool; + private final boolean reinstall; + + @Override + public DeliveryType type() { + return DeliveryType.TOOL_INSTALLATION; + } +} 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 new file mode 100644 index 0000000000..8ec5c643b7 --- /dev/null +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpec.java @@ -0,0 +1,121 @@ +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; +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 com.openframe.delivery.spec.DeliveryRequest; +import com.openframe.delivery.spec.DeliverySpec; +import lombok.RequiredArgsConstructor; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import java.util.List; + +import static java.lang.String.format; +import static java.util.Objects.requireNonNullElse; + +@Component +@RequiredArgsConstructor +@ConditionalOnProperty("spring.cloud.stream.enabled") +public class ToolInstallationDeliverySpec implements DeliverySpec { + + private static final String SUBJECT_TEMPLATE = "machine.%s.tool-installation"; + + private final NatsMessagePublisher natsMessagePublisher; + private final DownloadConfigurationMapper downloadConfigurationMapper; + private final LocalFilenameConfigurationMapper localFilenameConfigurationMapper; + + @Override + public DeliveryType getType() { + return DeliveryType.TOOL_INSTALLATION; + } + + @Override + 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) + .machineId(seed.getMachineId()) + .payload(message) + .build(); + } + + @Override + public void publish(String machineId, ToolInstallationMessage payload) { + String subject = format(SUBJECT_TEMPLATE, machineId); + natsMessagePublisher.publish(subject, payload); + } + + @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(); + message.setToolAgentId(toolAgent.getKey()); + message.setToolId(requireNonNullElse(toolAgent.getToolId(), "")); + message.setToolType(requireNonNullElse(tool.getToolType(), "")); + message.setVersion(version); + message.setSessionType(toolAgent.getSessionType()); + message.setDownloadConfigurations(downloadConfigurationMapper.map(toolAgent.getDownloadConfigurations(), version)); + 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) + .toList(); + } + + private ToolInstallationMessage.Asset mapAsset(ToolAgentAsset asset) { + String version = asset.getVersion(); + ToolInstallationMessage.Asset messageAsset = new ToolInstallationMessage.Asset(); + messageAsset.setId(asset.getId()); + messageAsset.setVersion(version); + messageAsset.setLocalFilenameConfiguration(localFilenameConfigurationMapper.map(asset.getLocalFilenameConfiguration())); + messageAsset.setDownloadConfigurations(downloadConfigurationMapper.map(asset.getDownloadConfigurations(), version)); + messageAsset.setSource(mapAssetSource(asset.getSource())); + messageAsset.setPath(asset.getPath()); + messageAsset.setExecutable(asset.isExecutable()); + return messageAsset; + } + + private static 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/model/ToolInstallationMessage.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/model/ToolInstallationMessage.java index 8d60d93875..79475398c8 100644 --- a/openframe-data-nats/src/main/java/com/openframe/data/nats/model/ToolInstallationMessage.java +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/model/ToolInstallationMessage.java @@ -1,5 +1,9 @@ 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 com.openframe.data.document.toolagent.SessionType; import lombok.Getter; import lombok.Setter; @@ -8,7 +12,10 @@ @Getter @Setter -public class ToolInstallationMessage { +public class ToolInstallationMessage implements DeliveryPayload { + + @JsonInclude(JsonInclude.Include.NON_NULL) + private DeliveryRef delivery; private String toolAgentId; private String toolId; 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 new file mode 100644 index 0000000000..782e522bc0 --- /dev/null +++ b/openframe-data-nats/src/test/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpecTest.java @@ -0,0 +1,103 @@ +package com.openframe.data.nats.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.tool.IntegratedTool; +import com.openframe.data.document.toolagent.IntegratedToolAgent; +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; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +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) +class ToolInstallationDeliverySpecTest { + + private static final String MACHINE_ID = "mach-42"; + private static final String TOOL_AGENT_KEY = "tactical-agent"; + private static final String TOOL_ID = "tactical"; + private static final String TOOL_TYPE = "TACTICAL"; + 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; + + @InjectMocks private ToolInstallationDeliverySpec spec; + + private IntegratedToolAgent toolAgent; + private IntegratedTool tool; + + @BeforeEach + void setUp() { + toolAgent = new IntegratedToolAgent(); + toolAgent.setKey(TOOL_AGENT_KEY); + toolAgent.setToolId(TOOL_ID); + toolAgent.setVersion(VERSION); + toolAgent.setInstallationCommandArgs(INSTALL_ARGS); + tool = new IntegratedTool(); + tool.setToolType(TOOL_TYPE); + } + + @Test + void request_toolAgent_messageBuiltAndTargetIsAgentKey() { + // setup + when(downloadConfigurationMapper.map(null, VERSION)).thenReturn(List.of()); + ToolInstallationDeliverySeed seed = new ToolInstallationDeliverySeed(MACHINE_ID, toolAgent, tool, true); + + // execution + DeliveryRequest request = spec.request(seed); + + // verifications + assertThat(request.getType()).isEqualTo(DeliveryType.TOOL_INSTALLATION); + assertThat(request.getTargetId()).isEqualTo(TOOL_AGENT_KEY); + assertThat(request.getMachineId()).isEqualTo(MACHINE_ID); + ToolInstallationMessage message = request.getPayload(); + assertThat(message.getToolAgentId()).isEqualTo(TOOL_AGENT_KEY); + assertThat(message.getToolId()).isEqualTo(TOOL_ID); + assertThat(message.getToolType()).isEqualTo(TOOL_TYPE); + assertThat(message.getVersion()).isEqualTo(VERSION); + assertThat(message.getInstallationCommandArgs()).isEqualTo(INSTALL_ARGS); + assertThat(message.isReinstall()).isTrue(); + } + + @Test + void request_toolWithoutIdAndType_emptyStringsNotNulls() { + // setup + toolAgent.setToolId(null); + tool.setToolType(null); + when(downloadConfigurationMapper.map(null, VERSION)).thenReturn(List.of()); + ToolInstallationDeliverySeed seed = new ToolInstallationDeliverySeed(MACHINE_ID, toolAgent, tool, false); + + // execution + DeliveryRequest request = spec.request(seed); + + // verifications + assertThat(request.getPayload().getToolId()).isEmpty(); + assertThat(request.getPayload().getToolType()).isEmpty(); + } + + @Test + void publish_payload_sentToMachineSubject() { + // setup + ToolInstallationMessage message = new ToolInstallationMessage(); + + // execution + spec.publish(MACHINE_ID, message); + + // verifications + verify(natsMessagePublisher).publish("machine.mach-42.tool-installation", message); + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/config/DeliveryProperties.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/config/DeliveryProperties.java index bab1be0e42..7d3482cc1b 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/config/DeliveryProperties.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/config/DeliveryProperties.java @@ -14,6 +14,7 @@ import java.util.EnumMap; import java.util.Map; +import static java.lang.Boolean.FALSE; import static java.util.Objects.requireNonNullElse; @Getter @@ -23,9 +24,6 @@ @ConfigurationProperties(prefix = "openframe.delivery") public class DeliveryProperties { - @NotNull - private Boolean enabled; - @Valid @NotNull private Sweep sweep; @@ -37,8 +35,11 @@ public class DeliveryProperties { // deliberately not @Valid: a per-type entry lists only the fields it overrides private Map types = new EnumMap<>(DeliveryType.class); - public boolean isEnabled() { - return enabled; + // a type not listed here is off: every environment switches each type on explicitly + private Map enabled = new EnumMap<>(DeliveryType.class); + + public boolean isEnabled(DeliveryType type) { + return enabled.getOrDefault(type, FALSE); } public Policy resolve(DeliveryType type) { @@ -53,6 +54,9 @@ public Policy resolve(DeliveryType type) { @Setter public static class Sweep { + @NotNull + @Positive + private Long interval; @NotNull @Positive private Integer batchSize; 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 2b41424802..1561c28bd9 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 @@ -1,6 +1,8 @@ package com.openframe.delivery.dispatch; import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.delivery.spec.DeliveryPayload; +import com.openframe.delivery.spec.DeliveryRef; import com.openframe.delivery.spec.DeliveryRequest; import com.openframe.delivery.spec.DeliverySeed; import com.openframe.delivery.spec.DeliverySpec; @@ -9,6 +11,8 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; +import java.util.UUID; + @Slf4j @Component @RequiredArgsConstructor @@ -19,12 +23,17 @@ public class DeliveryDispatcher { public void dispatch(DeliverySeed seed) { DeliveryType type = seed.type(); - DeliverySpec spec = registry.require(type); - DeliveryRequest request = spec.request(seed); + DeliverySpec spec = registry.require(type); + DeliveryRequest request = spec.request(seed); + DeliveryPayload payload = request.getPayload(); + String dispatchId = UUID.randomUUID().toString(); + String targetId = request.getTargetId(); + DeliveryRef delivery = new DeliveryRef(type, targetId, dispatchId); + payload.setDelivery(delivery); recorder.record(request); String machineId = request.getMachineId(); - Object payload = request.getPayload(); spec.publish(machineId, payload); - log.info("Delivery dispatched: type={} targetId={} machineId={}", type, request.getTargetId(), machineId); + 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/DeliveryRecorder.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryRecorder.java index f6b2e755c7..0bb54a30cf 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryRecorder.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/dispatch/DeliveryRecorder.java @@ -8,6 +8,8 @@ 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.DeliveryPayload; +import com.openframe.delivery.spec.DeliveryRef; import com.openframe.delivery.spec.DeliveryRequest; import com.openframe.delivery.track.DeliveryId; import lombok.RequiredArgsConstructor; @@ -26,9 +28,6 @@ public class DeliveryRecorder { private final ObjectMapper objectMapper; public void record(DeliveryRequest request) { - if (!properties.isEnabled()) { - return; - } MachineDelivery delivery = pendingRow(request); repository.upsertPending(delivery); log.info("Delivery recorded: type={} targetId={} machineId={}", @@ -41,7 +40,9 @@ private MachineDelivery pendingRow(DeliveryRequest request) { String targetId = request.getTargetId(); String machineId = request.getMachineId(); String id = DeliveryId.of(type, targetId, machineId); - Object payload = request.getPayload(); + DeliveryPayload payload = request.getPayload(); + DeliveryRef delivery = payload.getDelivery(); + String dispatchId = delivery.getDispatchId(); String payloadJson = toJson(payload); Policy policy = properties.resolve(type); long ackThresholdSeconds = policy.getAckThresholdSeconds(); @@ -51,6 +52,7 @@ private MachineDelivery pendingRow(DeliveryRequest request) { .type(type) .targetId(targetId) .machineId(machineId) + .dispatchId(dispatchId) .status(DeliveryStatus.PENDING) .attempts(0) .errors(0) diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/metrics/DeliveryMetrics.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/metrics/DeliveryMetrics.java index fac9b0b27a..ee766fadb5 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/metrics/DeliveryMetrics.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/metrics/DeliveryMetrics.java @@ -17,6 +17,7 @@ public class DeliveryMetrics { private static final String FAILED_COUNTER = "openframe.delivery.failed"; private static final String PUBLISH_FAILED_COUNTER = "openframe.delivery.publish_failed"; private static final String ROW_ERROR_COUNTER = "openframe.delivery.sweep.row_errors"; + private static final String RESULT_REJECTED_COUNTER = "openframe.delivery.result.rejected"; private static final String SWEEP_TIMER = "openframe.delivery.sweep.duration"; private static final String TAG_TYPE = "type"; private static final String TAG_REASON = "reason"; @@ -38,6 +39,10 @@ public void recordFailed(DeliveryType type, DeliveryFailure failure) { meterRegistry.counter(FAILED_COUNTER, TAG_TYPE, typeTag, TAG_REASON, reasonTag).increment(); } + public void recordResultRejected(String reason) { + meterRegistry.counter(RESULT_REJECTED_COUNTER, TAG_REASON, reason).increment(); + } + public void recordPublishFailed(DeliveryType type) { String typeTag = tagValue(type.name()); meterRegistry.counter(PUBLISH_FAILED_COUNTER, TAG_TYPE, typeTag).increment(); diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryPayload.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryPayload.java new file mode 100644 index 0000000000..ce1722720b --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryPayload.java @@ -0,0 +1,8 @@ +package com.openframe.delivery.spec; + +public interface DeliveryPayload { + + DeliveryRef getDelivery(); + + void setDelivery(DeliveryRef delivery); +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryRef.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryRef.java new file mode 100644 index 0000000000..1f62b71d09 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryRef.java @@ -0,0 +1,19 @@ +package com.openframe.delivery.spec; + +import com.fasterxml.jackson.annotation.JsonFormat; +import com.openframe.data.document.delivery.DeliveryType; +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@NoArgsConstructor +@AllArgsConstructor +public class DeliveryRef { + + // a type this server does not know yet reads as null and the report is rejected instead of poisoning the consumer + @JsonFormat(with = JsonFormat.Feature.READ_UNKNOWN_ENUM_VALUES_AS_NULL) + private DeliveryType type; + private String targetId; + private String dispatchId; +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryRequest.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryRequest.java index 469580feb3..04dbdcd930 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryRequest.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliveryRequest.java @@ -6,7 +6,7 @@ @Getter @Builder -public class DeliveryRequest

{ +public class DeliveryRequest

{ private final DeliveryType type; private final String targetId; private final String machineId; 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 12c60cf908..1e461bf637 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 @@ -4,7 +4,7 @@ import com.openframe.data.document.delivery.DeliveryType; import com.openframe.data.document.delivery.MachineDelivery; -public interface DeliverySpec { +public interface DeliverySpec { DeliveryType getType(); diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySpecRegistry.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySpecRegistry.java index b6fa71e87b..11f85af26d 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySpecRegistry.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/spec/DeliverySpecRegistry.java @@ -33,12 +33,12 @@ public DeliverySpecRegistry(ObjectProvider> specs) { } @SuppressWarnings("unchecked") - public Optional> find(DeliveryType type) { + public Optional> find(DeliveryType type) { DeliverySpec spec = (DeliverySpec) byType.get(type); return Optional.ofNullable(spec); } - public DeliverySpec require(DeliveryType type) { + public DeliverySpec require(DeliveryType type) { Optional> spec = find(type); return spec.orElseThrow(() -> new IllegalArgumentException("No spec registered for delivery type: " + type.name())); } diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliverySweepScheduler.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliverySweepScheduler.java index 09881ff4b4..f6dd706ca5 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliverySweepScheduler.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliverySweepScheduler.java @@ -11,7 +11,7 @@ @Slf4j @Component @RequiredArgsConstructor -@ConditionalOnProperty(name = {"openframe.delivery.enabled", "openframe.delivery.sweep.enabled"}, havingValue = "true") +@ConditionalOnProperty(name = "openframe.delivery.sweep.enabled", havingValue = "true") public class DeliverySweepScheduler { private static final String PASS_RETRY = "retry"; 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 155242cb77..aa7ce270f9 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 @@ -9,8 +9,11 @@ import com.openframe.data.document.delivery.MachineDelivery; import com.openframe.data.repository.delivery.MachineDeliveryRepository; import com.openframe.delivery.config.DeliveryProperties; +import com.openframe.delivery.track.DeliveryCloser; import com.openframe.delivery.config.DeliveryProperties.Policy; +import com.openframe.delivery.config.DeliveryProperties.Sweep; 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; @@ -28,7 +31,7 @@ @Slf4j @Service @RequiredArgsConstructor -@ConditionalOnProperty(name = {"openframe.delivery.enabled", "openframe.delivery.sweep.enabled"}, havingValue = "true") +@ConditionalOnProperty(name = "openframe.delivery.sweep.enabled", havingValue = "true") public class DeliverySweepService { private final MachineDeliveryRepository repository; @@ -91,20 +94,21 @@ private void parkSkipOrFailOffline(MachineDelivery delivery, Policy policy, Inst closer.fail(delivery, DeliveryFailure.OFFLINE, DeliveryStatus.UNACKED, now); return; } - long recheckSeconds = policy.getMaxRetryIntervalSeconds(); - Instant recheckAt = now.plusSeconds(recheckSeconds); + Sweep sweep = properties.getSweep(); + long recheckMillis = sweep.getInterval(); + Instant recheckAt = now.plusMillis(recheckMillis); Instant dueAt = earliest(windowEnd, recheckAt); String id = delivery.getId(); Instant dispatchedAt = delivery.getDispatchedAt(); - repository.park(id, DeliveryStatus.UNACKED, dispatchedAt, dueAt); - log.debug("Delivery parked, machine not online: id={} dueAt={} windowEnd={}", id, dueAt, windowEnd); + repository.postpone(id, DeliveryStatus.UNACKED, dispatchedAt, dueAt); + log.debug("Delivery waits for the machine to come online: id={} dueAt={} windowEnd={}", id, dueAt, windowEnd); } private void republish(MachineDelivery delivery, Policy policy, Instant now) { DeliveryType type = delivery.getType(); - DeliverySpec spec = registry.require(type); - Class payloadClass = spec.getPayloadClass(); - Object payload = readPayload(delivery, payloadClass); + DeliverySpec spec = registry.require(type); + Class payloadClass = spec.getPayloadClass(); + DeliveryPayload payload = readPayload(delivery, payloadClass); String machineId = delivery.getMachineId(); boolean published = publish(spec, machineId, payload); if (!published) { @@ -128,7 +132,7 @@ private void republish(MachineDelivery delivery, Policy policy, Instant now) { type, delivery.getTargetId(), machineId, attempt, dueAt); } - private static boolean publish(DeliverySpec spec, String machineId, Object payload) { + private static boolean publish(DeliverySpec spec, String machineId, DeliveryPayload payload) { try { spec.publish(machineId, payload); return true; diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliveryWatchdogService.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliveryWatchdogService.java index c7ddb6b752..da8dc168ed 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliveryWatchdogService.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliveryWatchdogService.java @@ -6,6 +6,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.track.DeliveryCloser; import com.openframe.delivery.config.DeliveryProperties.Policy; import com.openframe.delivery.metrics.DeliveryMetrics; import lombok.RequiredArgsConstructor; @@ -19,7 +20,7 @@ @Slf4j @Service @RequiredArgsConstructor -@ConditionalOnProperty(name = {"openframe.delivery.enabled", "openframe.delivery.sweep.enabled"}, havingValue = "true") +@ConditionalOnProperty(name = "openframe.delivery.sweep.enabled", havingValue = "true") public class DeliveryWatchdogService { private final MachineDeliveryRepository repository; 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 6e1e7fbb46..0375bba728 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 @@ -15,7 +15,7 @@ @Component @RequiredArgsConstructor -@ConditionalOnProperty(name = {"openframe.delivery.enabled", "openframe.delivery.sweep.enabled"}, havingValue = "true") +@ConditionalOnProperty(name = "openframe.delivery.sweep.enabled", havingValue = "true") public class MachineOnlineStatus { private static final Set GONE = EnumSet.of( diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliveryCloser.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryCloser.java similarity index 73% rename from openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliveryCloser.java rename to openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryCloser.java index 6e065c6705..b4e4c85bd2 100644 --- a/openframe-machine-delivery/src/main/java/com/openframe/delivery/sweep/DeliveryCloser.java +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/track/DeliveryCloser.java @@ -1,4 +1,4 @@ -package com.openframe.delivery.sweep; +package com.openframe.delivery.track; import com.openframe.data.document.delivery.DeliveryFailure; import com.openframe.data.document.delivery.DeliveryStatus; @@ -8,12 +8,12 @@ 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.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.stereotype.Component; import java.time.Instant; @@ -23,7 +23,6 @@ @Slf4j @Component @RequiredArgsConstructor -@ConditionalOnProperty(name = {"openframe.delivery.enabled", "openframe.delivery.sweep.enabled"}, havingValue = "true") public class DeliveryCloser { private final MachineDeliveryRepository repository; @@ -53,6 +52,20 @@ public void fail(MachineDelivery delivery, DeliveryFailure failure, Set from, String reason, Instant now) { DeliveryType type = delivery.getType(); Instant expiresAt = expiresAt(type, now); @@ -65,9 +78,13 @@ public void cancel(MachineDelivery delivery, 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); + Optional> spec = registry.find(type); spec.ifPresentOrElse( registered -> registered.onFailed(delivery, failure), () -> log.warn("No spec registered for delivery type {}, onFailed skipped", type)); 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 1d900a1146..9efa2068f7 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 @@ -18,29 +18,39 @@ public class DeliveryTracker { private final MachineDeliveryRepository repository; private final DeliveryProperties properties; + private final DeliveryCloser closer; - public void acknowledge(DeliveryType type, String targetId, String machineId) { + public void acknowledge(DeliveryType type, String targetId, String machineId, String dispatchId) { String id = DeliveryId.of(type, targetId, machineId); Instant now = Instant.now(); Policy policy = properties.resolve(type); long resultTimeoutSeconds = policy.getResultTimeoutSeconds(); Instant resultDueAt = now.plusSeconds(resultTimeoutSeconds); - boolean acked = repository.markAcked(id, DeliveryStatus.UNACKED, now, resultDueAt); + boolean acked = repository.markAcked(id, dispatchId, DeliveryStatus.UNACKED, now, resultDueAt); if (acked) { - log.info("Delivery ACKED: type={} targetId={} machineId={}", type, targetId, machineId); + log.info("Delivery ACKED: type={} targetId={} machineId={} dispatchId={}", type, targetId, machineId, dispatchId); + } else { + log.debug("Delivery ack ignored, no unacked row for this dispatch: id={} dispatchId={}", id, dispatchId); } } - public void complete(DeliveryType type, String targetId, String machineId) { + public void complete(DeliveryType type, String targetId, String machineId, String dispatchId) { String id = DeliveryId.of(type, targetId, machineId); Instant now = Instant.now(); Instant expiresAt = expiresAt(type, now); - boolean done = repository.markDone(id, DeliveryStatus.COMPLETABLE, now, expiresAt); + boolean done = repository.markDone(id, dispatchId, DeliveryStatus.COMPLETABLE, now, expiresAt); if (done) { - log.info("Delivery DONE: type={} targetId={} machineId={}", type, targetId, machineId); + log.info("Delivery DONE: type={} targetId={} machineId={} dispatchId={}", type, targetId, machineId, dispatchId); + } else { + log.debug("Delivery completion ignored, no row for this dispatch: id={} dispatchId={}", id, dispatchId); } } + public void fail(DeliveryType type, String targetId, String machineId, String dispatchId, 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(); @@ -51,14 +61,6 @@ public void cancel(DeliveryType type, String targetId, String machineId) { } } - public void wake(String machineId) { - Instant now = Instant.now(); - long woken = repository.wake(machineId, DeliveryStatus.UNACKED, now); - if (woken > 0) { - log.info("Delivery rows woken for retry: machineId={} count={}", machineId, woken); - } - } - private Instant expiresAt(DeliveryType type, Instant now) { Policy policy = properties.resolve(type); long ttlSeconds = policy.getTtlSeconds(); diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/config/DeliveryPropertiesTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/config/DeliveryPropertiesTest.java index 4bc81b1bdb..1ef4bfc446 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/config/DeliveryPropertiesTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/config/DeliveryPropertiesTest.java @@ -73,4 +73,39 @@ void resolve_typeOverridesOfflineBehavior_overrideWins() { // verifications assertThat(resolved.getOfflineBehavior()).isEqualTo(DeliveryOfflineBehavior.SKIP); } + @Test + void isEnabled_typeNotListed_false() { + // setup + properties.setEnabled(Map.of()); + + // execution + boolean enabled = properties.isEnabled(DeliveryType.TOOL_INSTALLATION); + + // verifications + assertThat(enabled).isFalse(); + } + + @Test + void isEnabled_typeListedOff_false() { + // setup + properties.setEnabled(Map.of(DeliveryType.TOOL_INSTALLATION, false)); + + // execution + boolean enabled = properties.isEnabled(DeliveryType.TOOL_INSTALLATION); + + // verifications + assertThat(enabled).isFalse(); + } + + @Test + void isEnabled_typeListedOn_true() { + // setup + properties.setEnabled(Map.of(DeliveryType.TOOL_INSTALLATION, true)); + + // execution + boolean enabled = properties.isEnabled(DeliveryType.TOOL_INSTALLATION); + + // verifications + assertThat(enabled).isTrue(); + } } diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/config/DeliveryTestPolicies.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/config/DeliveryTestPolicies.java index 161ffaab64..ba4f55c48b 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/config/DeliveryTestPolicies.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/config/DeliveryTestPolicies.java @@ -11,6 +11,7 @@ public final class DeliveryTestPolicies { public static final int BACKOFF_MULTIPLIER = 2; public static final long MAX_RETRY_INTERVAL = 300L; public static final int BATCH_SIZE = 500; + public static final long SWEEP_INTERVAL_MILLIS = 30_000L; public static final long RECONNECT_WINDOW = 86_400L; public static final long RESULT_TIMEOUT = 600L; public static final long TTL = 604_800L; @@ -29,9 +30,9 @@ public static DeliveryProperties properties() { defaults.setResultTimeoutSeconds(RESULT_TIMEOUT); defaults.setTtlSeconds(TTL); Sweep sweep = new Sweep(); + sweep.setInterval(SWEEP_INTERVAL_MILLIS); sweep.setBatchSize(BATCH_SIZE); DeliveryProperties properties = new DeliveryProperties(); - properties.setEnabled(true); properties.setDefaults(defaults); properties.setSweep(sweep); return properties; 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 47c9c132de..e6569755a7 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 @@ -13,6 +13,7 @@ import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.verify; @@ -50,7 +51,7 @@ void setUp() { } @Test - void dispatch_seed_requestRecordedThenPublishedThroughSpec() { + void dispatch_seed_dispatchIdSetThenRecordedThenPublishedThroughSpec() { // setup doReturn(spec).when(registry).require(DeliveryType.TOOL_INSTALLATION); when(spec.request(seed)).thenReturn(request); @@ -59,6 +60,9 @@ void dispatch_seed_requestRecordedThenPublishedThroughSpec() { dispatcher.dispatch(seed); // verifications + assertThat(payload.getDelivery().getType()).isEqualTo(DeliveryType.TOOL_INSTALLATION); + assertThat(payload.getDelivery().getTargetId()).isEqualTo(TARGET_ID); + assertThat(payload.getDelivery().getDispatchId()).isNotBlank(); verify(recorder).record(request); verify(spec).publish(MACHINE_ID, payload); } diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/dispatch/DeliveryRecorderTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/dispatch/DeliveryRecorderTest.java index f1a5ee635c..97dde8a894 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/dispatch/DeliveryRecorderTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/dispatch/DeliveryRecorderTest.java @@ -5,8 +5,8 @@ import com.openframe.data.document.delivery.DeliveryType; import com.openframe.data.document.delivery.MachineDelivery; import com.openframe.data.repository.delivery.MachineDeliveryRepository; -import com.openframe.delivery.config.DeliveryProperties; import com.openframe.delivery.config.DeliveryTestPolicies; +import com.openframe.delivery.spec.DeliveryRef; import com.openframe.delivery.spec.DeliveryRequest; import com.openframe.delivery.spec.TestPayload; import org.junit.jupiter.api.BeforeEach; @@ -21,13 +21,13 @@ import static com.openframe.delivery.config.DeliveryTestPolicies.TTL; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.verifyNoInteractions; @ExtendWith(MockitoExtension.class) class DeliveryRecorderTest { private static final String MACHINE_ID = "mach-42"; private static final String VALUE = "issued"; + private static final String DISPATCH_ID = "d-1"; @Mock private MachineDeliveryRepository repository; @@ -35,21 +35,20 @@ class DeliveryRecorderTest { private DeliveryRecorder recorder; - private DeliveryProperties properties; private DeliveryRequest request; @BeforeEach void setUp() { TestPayload payload = new TestPayload(); payload.setValue(VALUE); + payload.setDelivery(new DeliveryRef(DeliveryType.CLIENT_UNINSTALL, MACHINE_ID, DISPATCH_ID)); request = DeliveryRequest.builder() .type(DeliveryType.CLIENT_UNINSTALL) .targetId(MACHINE_ID) .machineId(MACHINE_ID) .payload(payload) .build(); - properties = DeliveryTestPolicies.properties(); - recorder = new DeliveryRecorder(repository, properties, new ObjectMapper()); + recorder = new DeliveryRecorder(repository, DeliveryTestPolicies.properties(), new ObjectMapper()); } @Test @@ -66,21 +65,10 @@ void record_request_pendingRowUpsertedDueAfterAckThreshold() { assertThat(saved.getType()).isEqualTo(DeliveryType.CLIENT_UNINSTALL); assertThat(saved.getStatus()).isEqualTo(DeliveryStatus.PENDING); assertThat(saved.getAttempts()).isZero(); - assertThat(saved.getPayloadJson()).contains(VALUE); + assertThat(saved.getDispatchId()).isEqualTo(DISPATCH_ID); + assertThat(saved.getPayloadJson()).contains(VALUE).contains(DISPATCH_ID); assertThat(saved.getDueAt()).isEqualTo(saved.getDispatchedAt().plusSeconds(ACK_THRESHOLD)); assertThat(saved.getExpiresAt()).isEqualTo(saved.getDispatchedAt().plusSeconds(TTL)); assertThat(saved.getErrors()).isZero(); } - - @Test - void record_engineDisabled_nothingWritten() { - // setup - properties.setEnabled(false); - - // execution - recorder.record(request); - - // verifications - verifyNoInteractions(repository); - } } diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/metrics/DeliveryMetricsTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/metrics/DeliveryMetricsTest.java index 8311adc353..9bf199b1db 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/metrics/DeliveryMetricsTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/metrics/DeliveryMetricsTest.java @@ -83,6 +83,19 @@ void recordRowError_called_counterIncremented() { assertThat(counter.count()).isEqualTo(1.0); } + @Test + void recordResultRejected_reason_counterTaggedWithReason() { + // setup + + // execution + metrics.recordResultRejected("incomplete"); + + // verifications + Counter counter = registry.find("openframe.delivery.result.rejected").tags("reason", "incomplete").counter(); + assertThat(counter).isNotNull(); + assertThat(counter.count()).isEqualTo(1.0); + } + @Test void recordFailed_typeAndReason_counterTaggedLowercase() { // setup diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/DeliverySpecRegistryTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/DeliverySpecRegistryTest.java index b638d80bf1..63628e3e73 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/DeliverySpecRegistryTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/DeliverySpecRegistryTest.java @@ -23,7 +23,7 @@ void require_registeredType_specReturned() { DeliverySpecRegistry registry = new DeliverySpecRegistry(specs); // execution - DeliverySpec resolved = registry.require(DeliveryType.TOOL_INSTALLATION); + DeliverySpec resolved = registry.require(DeliveryType.TOOL_INSTALLATION); // verifications assertThat(resolved).isSameAs(spec); @@ -77,7 +77,7 @@ void find_registeredType_specPresent() { DeliverySpecRegistry registry = new DeliverySpecRegistry(specs); // execution - Optional> found = registry.find(DeliveryType.TOOL_INSTALLATION); + Optional> found = registry.find(DeliveryType.TOOL_INSTALLATION); // verifications assertThat(found).get().isSameAs(spec); @@ -90,7 +90,7 @@ void find_unregisteredType_empty() { DeliverySpecRegistry registry = new DeliverySpecRegistry(specs); // execution - Optional> found = registry.find(DeliveryType.CLIENT_UNINSTALL); + Optional> found = registry.find(DeliveryType.CLIENT_UNINSTALL); // verifications assertThat(found).isEmpty(); diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/TestPayload.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/TestPayload.java index c344b62b92..9f6eeccc7a 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/TestPayload.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/spec/TestPayload.java @@ -3,6 +3,8 @@ import lombok.Data; @Data -public class TestPayload { +public class TestPayload implements DeliveryPayload { + + private DeliveryRef delivery; private String value; } 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 81cba4bf96..a61a916249 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.track.DeliveryCloser; import com.openframe.delivery.config.DeliveryTestPolicies; import com.openframe.delivery.metrics.DeliveryMetrics; import com.openframe.delivery.spec.DeliverySpec; @@ -33,6 +34,7 @@ import static com.openframe.delivery.config.DeliveryTestPolicies.MAX_ATTEMPTS; import static com.openframe.delivery.config.DeliveryTestPolicies.MAX_RETRY_INTERVAL; import static com.openframe.delivery.config.DeliveryTestPolicies.RECONNECT_WINDOW; +import static com.openframe.delivery.config.DeliveryTestPolicies.SWEEP_INTERVAL_MILLIS; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; @@ -52,7 +54,7 @@ class DeliverySweepServiceTest { private static final String PAYLOAD_JSON = "{\"value\":\"fleetmdm-agent\"}"; private static final String CORRUPT_JSON = "not-json"; private static final long TWO_DAYS_SECONDS = 172_800L; - private static final long ONE_MINUTE_SECONDS = 60L; + private static final long TEN_SECONDS = 10L; private static final long FIRST_RETRY_DELAY = ACK_THRESHOLD * BACKOFF_MULTIPLIER; private static final long CLOCK_SLACK_SECONDS = 5L; private static final int MANY_ATTEMPTS_ALLOWED = 10; @@ -216,7 +218,7 @@ void retryPending_onlineAttemptsExhausted_failedExhaustedOnlyIfStillUnacked() { } @Test - void retryPending_offlineFarFromWindowEnd_parkedUntilNextRecheck() { + void retryPending_offlineFarFromWindowEnd_postponedToNextSweep() { // setup Instant before = Instant.now(); stubDue(delivery); @@ -226,17 +228,18 @@ void retryPending_offlineFarFromWindowEnd_parkedUntilNextRecheck() { service.retryPending(); // verifications - verify(repository).park(eq(delivery.getId()), eq(DeliveryStatus.UNACKED), eq(dispatchedAt), dueAtCaptor.capture()); + verify(repository).postpone(eq(delivery.getId()), eq(DeliveryStatus.UNACKED), eq(dispatchedAt), dueAtCaptor.capture()); + Instant nextSweep = before.plusMillis(SWEEP_INTERVAL_MILLIS); assertThat(dueAtCaptor.getValue()) - .isAfterOrEqualTo(before.plusSeconds(MAX_RETRY_INTERVAL)) - .isBefore(before.plusSeconds(MAX_RETRY_INTERVAL + CLOCK_SLACK_SECONDS)); + .isAfterOrEqualTo(nextSweep) + .isBefore(nextSweep.plusSeconds(CLOCK_SLACK_SECONDS)); verifyNoInteractions(registry, closer, metrics); } @Test - void retryPending_offlineCloseToWindowEnd_parkedUntilWindowEnd() { + void retryPending_offlineCloseToWindowEnd_postponedToWindowEnd() { // setup - Instant recently = Instant.now().minusSeconds(RECONNECT_WINDOW - ONE_MINUTE_SECONDS); + Instant recently = Instant.now().minusSeconds(RECONNECT_WINDOW - TEN_SECONDS); delivery.setDispatchedAt(recently); stubDue(delivery); stubMachineNotOnline(); @@ -245,7 +248,7 @@ void retryPending_offlineCloseToWindowEnd_parkedUntilWindowEnd() { service.retryPending(); // verifications - verify(repository).park(delivery.getId(), DeliveryStatus.UNACKED, recently, recently.plusSeconds(RECONNECT_WINDOW)); + verify(repository).postpone(delivery.getId(), DeliveryStatus.UNACKED, recently, recently.plusSeconds(RECONNECT_WINDOW)); } @Test @@ -261,7 +264,7 @@ void retryPending_offlineReconnectWindowOver_failedOffline() { // verifications verify(closer).fail(eq(delivery), eq(DeliveryFailure.OFFLINE), eq(DeliveryStatus.UNACKED), any(Instant.class)); - verify(repository, never()).park(eq(delivery.getId()), eq(DeliveryStatus.UNACKED), eq(twoDaysAgo), any(Instant.class)); + verify(repository, never()).postpone(eq(delivery.getId()), eq(DeliveryStatus.UNACKED), eq(twoDaysAgo), any(Instant.class)); } @Test diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliveryWatchdogServiceTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliveryWatchdogServiceTest.java index dae8b5c9c4..a2f5daa077 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliveryWatchdogServiceTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliveryWatchdogServiceTest.java @@ -6,6 +6,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.track.DeliveryCloser; import com.openframe.delivery.config.DeliveryTestPolicies; import com.openframe.delivery.metrics.DeliveryMetrics; import org.junit.jupiter.api.BeforeEach; diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliveryCloserTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryCloserTest.java similarity index 77% rename from openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliveryCloserTest.java rename to openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryCloserTest.java index 760d75840e..a448ba4338 100644 --- a/openframe-machine-delivery/src/test/java/com/openframe/delivery/sweep/DeliveryCloserTest.java +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/track/DeliveryCloserTest.java @@ -1,4 +1,4 @@ -package com.openframe.delivery.sweep; +package com.openframe.delivery.track; import com.openframe.data.document.delivery.DeliveryFailure; import com.openframe.data.document.delivery.DeliveryStatus; @@ -33,6 +33,10 @@ class DeliveryCloserTest { private static final String DELIVERY_ID = "CLIENT_UNINSTALL:openframe-client:mach-42"; private static final String REASON = "machine gone"; + private static final String TARGET_ID = "openframe-client"; + private static final String MACHINE_ID = "mach-42"; + private static final String DISPATCH_ID = "d-1"; + private static final String ERROR = "download failed"; @Mock private MachineDeliveryRepository repository; @Mock private DeliverySpecRegistry registry; @@ -109,6 +113,33 @@ void fail_typeWithoutSpec_rowStillFailedSpecSkipped() { verifyNoInteractions(spec); } + @Test + void failReported_openRowOfThisDispatch_rowFailedMetricCountedSpecNotified() { + // 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 + void failReported_rowOfAnotherDispatchOrClosed_nothingRecorded() { + // setup + when(repository.markFailed(DELIVERY_ID, DISPATCH_ID, DeliveryStatus.OPEN, DeliveryFailure.AGENT_ERROR, ERROR, now, now.plusSeconds(TTL))).thenReturn(false); + + // execution + closer.failReported(DeliveryType.CLIENT_UNINSTALL, TARGET_ID, MACHINE_ID, DISPATCH_ID, ERROR, now); + + // verifications + verifyNoInteractions(metrics, registry, spec); + } + @Test void cancel_rowStillUnackedSameDispatch_rowCancelledWithTtlExpiry() { // setup 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 4bb5e6e1aa..c9adc32d52 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 @@ -28,9 +28,11 @@ class DeliveryTrackerTest { private static final String MACHINE_ID = "mach-42"; private static final String TARGET_ID = "fleetmdm-agent"; private static final String DELIVERY_ID = DeliveryId.of(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); - private static final long TWO_ROWS = 2L; + private static final String DISPATCH_ID = "d-1"; + private static final String ERROR = "download failed"; @Mock private MachineDeliveryRepository repository; + @Mock private DeliveryCloser closer; @Captor private ArgumentCaptor atCaptor; @Captor private ArgumentCaptor untilCaptor; @@ -39,16 +41,16 @@ class DeliveryTrackerTest { @BeforeEach void setUp() { - tracker = new DeliveryTracker(repository, DeliveryTestPolicies.properties()); + tracker = new DeliveryTracker(repository, DeliveryTestPolicies.properties(), closer); } @Test void acknowledge_typedKey_unackedRowMarkedAckedWithResultDeadline() { // setup - when(repository.markAcked(eq(DELIVERY_ID), eq(DeliveryStatus.UNACKED), atCaptor.capture(), untilCaptor.capture())).thenReturn(true); + 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); + tracker.acknowledge(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID, DISPATCH_ID); // verifications Instant ackedAt = atCaptor.getValue(); @@ -58,10 +60,10 @@ void acknowledge_typedKey_unackedRowMarkedAckedWithResultDeadline() { @Test void complete_typedKey_openOrFailedRowMarkedDoneWithTtlExpiry() { // setup - when(repository.markDone(eq(DELIVERY_ID), eq(DeliveryStatus.COMPLETABLE), atCaptor.capture(), untilCaptor.capture())).thenReturn(true); + 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); + tracker.complete(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID, DISPATCH_ID); // verifications Instant finishedAt = atCaptor.getValue(); @@ -82,14 +84,13 @@ void cancel_typedKey_openRowMarkedCancelledWithTtlExpiry() { } @Test - void wake_machineId_unackedRowsOfThatMachineWoken() { + void fail_agentReportedError_closerAsked() { // setup - when(repository.wake(eq(MACHINE_ID), eq(DeliveryStatus.UNACKED), any(Instant.class))).thenReturn(TWO_ROWS); // execution - tracker.wake(MACHINE_ID); + tracker.fail(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID, DISPATCH_ID, ERROR); // verifications - verify(repository).wake(eq(MACHINE_ID), eq(DeliveryStatus.UNACKED), any(Instant.class)); + verify(closer).failReported(eq(DeliveryType.TOOL_INSTALLATION), eq(TARGET_ID), eq(MACHINE_ID), eq(DISPATCH_ID), eq(ERROR), any(Instant.class)); } } diff --git a/openframe-management-service-core/src/main/java/com/openframe/management/initializer/NatsStreamConfigurationInitializer.java b/openframe-management-service-core/src/main/java/com/openframe/management/initializer/NatsStreamConfigurationInitializer.java index 0b4d087b48..4798f17ddf 100644 --- a/openframe-management-service-core/src/main/java/com/openframe/management/initializer/NatsStreamConfigurationInitializer.java +++ b/openframe-management-service-core/src/main/java/com/openframe/management/initializer/NatsStreamConfigurationInitializer.java @@ -1,5 +1,6 @@ package com.openframe.management.initializer; +import com.openframe.data.nats.delivery.DeliveryResultMessage; import com.openframe.data.nats.rmm.model.PackageManagerMissingMessage; import com.openframe.management.service.NatsStreamManagementService; import io.nats.client.api.RetentionPolicy; @@ -84,6 +85,13 @@ public class NatsStreamConfigurationInitializer implements ApplicationRunner { .storageType(StorageType.File) .retentionPolicy(RetentionPolicy.Limits) .build(), + StreamConfiguration.builder() + .name(DeliveryResultMessage.STREAM) + .subjects(List.of(DeliveryResultMessage.SUBJECT_FILTER)) + .storageType(StorageType.File) + .retentionPolicy(RetentionPolicy.Limits) + .maxAge(Duration.ofDays(1)) + .build(), StreamConfiguration.builder() .name(PackageManagerMissingMessage.STREAM) .subjects(List.of(PackageManagerMissingMessage.SUBJECT_FILTER)) 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..71d53fbf52 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,11 @@ import com.openframe.data.document.tool.IntegratedTool; import com.openframe.data.document.toolagent.IntegratedToolAgent; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.nats.delivery.ToolInstallationDeliverySeed; import com.openframe.data.nats.publisher.ToolInstallationNatsPublisher; +import com.openframe.delivery.config.DeliveryProperties; +import com.openframe.delivery.dispatch.DeliveryDispatcher; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -21,6 +25,8 @@ public class ToolInstallationService { private final IntegratedToolService integratedToolService; private final ToolCommandParamsResolver toolCommandParamsResolver; private final ToolInstallationNatsPublisher toolInstallationNatsPublisher; + private final DeliveryProperties deliveryProperties; + private final DeliveryDispatcher deliveryDispatcher; public void process(String machineId, IntegratedToolAgent toolAgent) { process(machineId, toolAgent, false); @@ -41,7 +47,7 @@ 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); + publish(machineId, toolAgent, tool, reinstall); log.info("Published {} agent installation message for machine {}", toolId, machineId); } catch (Exception e) { // TODO: add fallback mechanism @@ -64,4 +70,13 @@ private IntegratedTool getIntegratedToolData(String toolId) { } } + + private void publish(String machineId, IntegratedToolAgent toolAgent, IntegratedTool tool, boolean reinstall) { + if (deliveryProperties.isEnabled(DeliveryType.TOOL_INSTALLATION)) { + ToolInstallationDeliverySeed seed = new ToolInstallationDeliverySeed(machineId, toolAgent, tool, reinstall); + deliveryDispatcher.dispatch(seed); + return; + } + toolInstallationNatsPublisher.publish(machineId, toolAgent, tool, reinstall); + } }