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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,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);
Comment thread
semen-flamingo marked this conversation as resolved.
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;
}
}
Original file line number Diff line number Diff line change
@@ -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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -4,5 +4,6 @@ public enum DeliveryFailure {
EXHAUSTED,
OFFLINE,
TIMEOUT,
ERROR
ERROR,
AGENT_ERROR
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,17 +20,16 @@ public interface CustomMachineDeliveryRepository {

boolean postponeAfterError(String id, Set<DeliveryStatus> from, Instant dispatchedAt, Instant dueAt);

boolean park(String id, Set<DeliveryStatus> from, Instant dispatchedAt, Instant dueAt);

boolean markAcked(String id, Set<DeliveryStatus> from, Instant ackedAt, Instant dueAt);
boolean markAcked(String id, String dispatchId, Set<DeliveryStatus> from, Instant ackedAt, Instant dueAt);

boolean markDone(String id, Set<DeliveryStatus> from, Instant finishedAt, Instant expiresAt);
boolean markDone(String id, String dispatchId, Set<DeliveryStatus> from, Instant finishedAt, Instant expiresAt);

boolean markCancelled(String id, Set<DeliveryStatus> from, Instant finishedAt, Instant expiresAt);

boolean markCancelled(String id, Set<DeliveryStatus> from, Instant dispatchedAt, Instant finishedAt, Instant expiresAt);

boolean markFailed(String id, Set<DeliveryStatus> from, Instant dispatchedAt, DeliveryFailure failure, Instant finishedAt, Instant expiresAt);
boolean markFailed(String id, String dispatchId, Set<DeliveryStatus> from, DeliveryFailure failure, String error, Instant finishedAt, Instant expiresAt);

long wake(String machineId, Set<DeliveryStatus> from, Instant dueAt);
}
Loading
Loading