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 new file mode 100644 index 0000000000..7095d22f8a --- /dev/null +++ b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryFailure.java @@ -0,0 +1,7 @@ +package com.openframe.data.document.delivery; + +public enum DeliveryFailure { + EXHAUSTED, + OFFLINE, + TIMEOUT +} diff --git a/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryStatus.java b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryStatus.java new file mode 100644 index 0000000000..9c05515151 --- /dev/null +++ b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryStatus.java @@ -0,0 +1,8 @@ +package com.openframe.data.document.delivery; + +public enum DeliveryStatus { + PENDING, + ACKED, + DONE, + FAILED +} diff --git a/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryType.java b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryType.java new file mode 100644 index 0000000000..3bd4c20354 --- /dev/null +++ b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/DeliveryType.java @@ -0,0 +1,9 @@ +package com.openframe.data.document.delivery; + +public enum DeliveryType { + SCRIPT_SCHEDULE, + TOOL_INSTALLATION, + TOOL_UPDATE, + CLIENT_UPDATE, + CLIENT_UNINSTALL +} 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 new file mode 100644 index 0000000000..8f2b487cc8 --- /dev/null +++ b/openframe-data-mongo-common/src/main/java/com/openframe/data/document/delivery/MachineDelivery.java @@ -0,0 +1,48 @@ +package com.openframe.data.document.delivery; + +import com.openframe.data.document.rmm.schedule.ScheduleOfflineBehavior; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; +import org.springframework.data.annotation.Id; +import org.springframework.data.mongodb.core.index.CompoundIndex; +import org.springframework.data.mongodb.core.index.Indexed; +import org.springframework.data.mongodb.core.mapping.Document; + +import java.time.Instant; + +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +@Document(collection = "machine_delivery") +@CompoundIndex(name = "machine_delivery_sweep", def = "{'type': 1, 'status': 1, 'lastAttemptAt': 1}") +public class MachineDelivery { + + @Id + private String id; + + private DeliveryType type; + private String targetId; + private String machineId; + private String tenantId; + + private DeliveryStatus status; + private int attempts; + private String payloadJson; + + private Instant dispatchedAt; + private Instant lastAttemptAt; + private Instant ackedAt; + private Instant finishedAt; + + private DeliveryFailure failure; + private String error; + + private ScheduleOfflineBehavior offlineBehavior; + private Long reconnectWindowSeconds; + + @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/MachineDeliveryRepository.java b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/MachineDeliveryRepository.java new file mode 100644 index 0000000000..934d078eb3 --- /dev/null +++ b/openframe-data-mongo-sync/src/main/java/com/openframe/data/repository/delivery/MachineDeliveryRepository.java @@ -0,0 +1,18 @@ +package com.openframe.data.repository.delivery; + +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import org.springframework.data.mongodb.repository.MongoRepository; +import org.springframework.stereotype.Repository; + +import java.time.Instant; +import java.util.List; + +@Repository +public interface MachineDeliveryRepository extends MongoRepository { + + List findByTypeAndStatusAndLastAttemptAtBefore(DeliveryType type, DeliveryStatus status, Instant before); + + List findByTypeAndStatusAndAckedAtBefore(DeliveryType type, DeliveryStatus status, Instant before); +} 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/ToolInstallationDeliverySpec.java b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpec.java new file mode 100644 index 0000000000..aaccc9c2d8 --- /dev/null +++ b/openframe-data-nats/src/main/java/com/openframe/data/nats/delivery/ToolInstallationDeliverySpec.java @@ -0,0 +1,143 @@ +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.DeliveryRequest; +import com.openframe.delivery.DeliverySeed; +import com.openframe.delivery.DeliverySpec; +import lombok.AllArgsConstructor; +import lombok.Getter; +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; + + @Getter + @AllArgsConstructor + public static class Seed 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; + } + } + + @Override + public DeliveryType getType() { + return DeliveryType.TOOL_INSTALLATION; + } + + @Override + public Class getSeedClass() { + return Seed.class; + } + + @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(Seed 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.publishPersistent(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/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..1a4c3511f3 --- /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.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()); + ToolInstallationDeliverySpec.Seed seed = new ToolInstallationDeliverySpec.Seed(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()); + ToolInstallationDeliverySpec.Seed seed = new ToolInstallationDeliverySpec.Seed(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).publishPersistent("machine.mach-42.tool-installation", message); + } +} diff --git a/openframe-machine-delivery/pom.xml b/openframe-machine-delivery/pom.xml new file mode 100644 index 0000000000..5426ab31be --- /dev/null +++ b/openframe-machine-delivery/pom.xml @@ -0,0 +1,46 @@ + + + 4.0.0 + + + com.openframe.oss + openframe-oss-lib + ${revision} + + + openframe-machine-delivery + jar + OpenFrame Machine Delivery + Application-level delivery of commands to machines: per-type specs, delivery rows, sweep, watchdog, ack and completion tracking. Transport is supplied by the specs, the engine has no NATS dependency. + + + + com.openframe.oss + openframe-core + + + com.openframe.oss + openframe-data-mongo-sync + + + org.springframework.boot + spring-boot-starter-validation + + + net.javacrumbs.shedlock + shedlock-spring + 5.10.2 + + + io.micrometer + micrometer-core + + + org.springframework.boot + spring-boot-starter-test + test + + + diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryDispatcher.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryDispatcher.java new file mode 100644 index 0000000000..43fd0210ef --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryDispatcher.java @@ -0,0 +1,29 @@ +package com.openframe.delivery; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +@Slf4j +@Component +@RequiredArgsConstructor +public class DeliveryDispatcher { + + private final DeliverySpecRegistry registry; + private final DeliveryRecorder recorder; + + public void dispatch(DeliverySeed seed) { + DeliverySpec spec = registry.require(seed.type()); + dispatch(spec, seed); + } + + private void dispatch(DeliverySpec spec, DeliverySeed seed) { + Class seedClass = spec.getSeedClass(); + S typedSeed = seedClass.cast(seed); + DeliveryRequest

request = spec.request(typedSeed); + recorder.record(request); + spec.publish(request.getMachineId(), request.getPayload()); + log.info("Delivery dispatched: type={} targetId={} machineId={}", + request.getType(), request.getTargetId(), request.getMachineId()); + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryFailureRecorder.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryFailureRecorder.java new file mode 100644 index 0000000000..50abe845bc --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryFailureRecorder.java @@ -0,0 +1,44 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryFailure; +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +import com.openframe.delivery.DeliveryProperties.Policy; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import java.time.Instant; + +@Slf4j +@Component +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "true") +public class DeliveryFailureRecorder { + + private final MachineDeliveryRepository repository; + private final DeliverySpecRegistry registry; + private final DeliveryProperties properties; + private final DeliveryMetrics metrics; + + public void fail(MachineDelivery delivery, DeliveryFailure failure, Instant now) { + DeliveryType type = delivery.getType(); + Policy policy = properties.resolve(type); + long ttlSeconds = policy.getTtlSeconds(); + + delivery.setStatus(DeliveryStatus.FAILED); + delivery.setFailure(failure); + delivery.setFinishedAt(now); + delivery.setExpiresAt(now.plusSeconds(ttlSeconds)); + repository.save(delivery); + + metrics.recordFailed(type, failure); + DeliverySpec spec = registry.require(type); + spec.onFailed(delivery, failure); + log.warn("Delivery FAILED: type={} targetId={} machineId={} attempts={} reason={}", + type, delivery.getTargetId(), delivery.getMachineId(), delivery.getAttempts(), failure); + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryId.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryId.java new file mode 100644 index 0000000000..0e1b50c43c --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryId.java @@ -0,0 +1,15 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; + +final class DeliveryId { + + private static final String SEPARATOR = ":"; + + private DeliveryId() { + } + + static String of(DeliveryType type, String targetId, String machineId) { + return type.name() + SEPARATOR + targetId + SEPARATOR + machineId; + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryMetrics.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryMetrics.java new file mode 100644 index 0000000000..75a7750e3d --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryMetrics.java @@ -0,0 +1,38 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryFailure; +import com.openframe.data.document.delivery.DeliveryType; +import io.micrometer.core.instrument.MeterRegistry; +import lombok.RequiredArgsConstructor; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import java.util.Locale; + +@Component +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "true") +public class DeliveryMetrics { + + private static final String RETRIED_COUNTER = "openframe.delivery.retried"; + private static final String FAILED_COUNTER = "openframe.delivery.failed"; + private static final String TAG_TYPE = "type"; + private static final String TAG_REASON = "reason"; + + private final MeterRegistry meterRegistry; + + public void recordRetried(DeliveryType type) { + String typeTag = tagValue(type.name()); + meterRegistry.counter(RETRIED_COUNTER, TAG_TYPE, typeTag).increment(); + } + + public void recordFailed(DeliveryType type, DeliveryFailure failure) { + String typeTag = tagValue(type.name()); + String reasonTag = tagValue(failure.name()); + meterRegistry.counter(FAILED_COUNTER, TAG_TYPE, typeTag, TAG_REASON, reasonTag).increment(); + } + + private static String tagValue(String enumName) { + return enumName.toLowerCase(Locale.ROOT); + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryProperties.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryProperties.java new file mode 100644 index 0000000000..1bfb89ec26 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryProperties.java @@ -0,0 +1,79 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.document.rmm.schedule.ScheduleOfflineBehavior; +import jakarta.validation.Valid; +import jakarta.validation.constraints.NotNull; +import lombok.Getter; +import lombok.Setter; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.stereotype.Component; +import org.springframework.validation.annotation.Validated; + +import java.util.EnumMap; +import java.util.Map; + +import static java.util.Objects.requireNonNullElse; + +@Getter +@Setter +@Validated +@Component +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "true") +@ConfigurationProperties(prefix = "openframe.delivery") +public class DeliveryProperties { + + @Valid + @NotNull + private Policy defaults; + + // deliberately not @Valid: a per-type entry lists only the fields it overrides + private Map types = new EnumMap<>(DeliveryType.class); + + public Policy resolve(DeliveryType type) { + Policy override = types.get(type); + if (override == null) { + return defaults; + } + return override.mergeOver(defaults); + } + + public Policy resolve(MachineDelivery delivery) { + Policy typePolicy = resolve(delivery.getType()); + Policy rowOverride = new Policy(); + rowOverride.setOfflineBehavior(delivery.getOfflineBehavior()); + rowOverride.setReconnectWindowSeconds(delivery.getReconnectWindowSeconds()); + return rowOverride.mergeOver(typePolicy); + } + + @Getter + @Setter + public static class Policy { + + @NotNull + private Long ackThresholdSeconds; + @NotNull + private Integer maxAttempts; + @NotNull + private ScheduleOfflineBehavior offlineBehavior; + @NotNull + private Long reconnectWindowSeconds; + @NotNull + private Long resultTimeoutSeconds; + @NotNull + private Long ttlSeconds; + + Policy mergeOver(Policy base) { + Policy merged = new Policy(); + merged.ackThresholdSeconds = requireNonNullElse(ackThresholdSeconds, base.ackThresholdSeconds); + merged.maxAttempts = requireNonNullElse(maxAttempts, base.maxAttempts); + merged.offlineBehavior = requireNonNullElse(offlineBehavior, base.offlineBehavior); + merged.reconnectWindowSeconds = requireNonNullElse(reconnectWindowSeconds, base.reconnectWindowSeconds); + merged.resultTimeoutSeconds = requireNonNullElse(resultTimeoutSeconds, base.resultTimeoutSeconds); + merged.ttlSeconds = requireNonNullElse(ttlSeconds, base.ttlSeconds); + return merged; + } + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryRecorder.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryRecorder.java new file mode 100644 index 0000000000..1012a20274 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryRecorder.java @@ -0,0 +1,6 @@ +package com.openframe.delivery; + +public interface DeliveryRecorder { + + void record(DeliveryRequest request); +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryRequest.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryRequest.java new file mode 100644 index 0000000000..5857761286 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryRequest.java @@ -0,0 +1,17 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.rmm.schedule.ScheduleOfflineBehavior; +import lombok.Builder; +import lombok.Getter; + +@Getter +@Builder +public class DeliveryRequest

{ + private final DeliveryType type; + private final String targetId; + private final String machineId; + private final P payload; + private final ScheduleOfflineBehavior offlineBehavior; + private final Long reconnectWindowSeconds; +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySeed.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySeed.java new file mode 100644 index 0000000000..e0f41cfc08 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySeed.java @@ -0,0 +1,8 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; + +public interface DeliverySeed { + + DeliveryType type(); +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySpec.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySpec.java new file mode 100644 index 0000000000..0048147262 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySpec.java @@ -0,0 +1,20 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryFailure; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; + +public interface DeliverySpec { + + DeliveryType getType(); + + Class getSeedClass(); + + Class

getPayloadClass(); + + DeliveryRequest

request(S seed); + + void publish(String machineId, P payload); + + void onFailed(MachineDelivery delivery, DeliveryFailure failure); +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySpecRegistry.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySpecRegistry.java new file mode 100644 index 0000000000..a95ee96c7c --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySpecRegistry.java @@ -0,0 +1,42 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.stereotype.Component; + +import java.util.Map; +import java.util.Set; +import java.util.TreeSet; + +import static java.util.function.Function.identity; +import static java.util.stream.Collectors.toUnmodifiableMap; + +@Slf4j +@Component +public class DeliverySpecRegistry { + + private final Map> byType; + + // ObjectProvider, not List: a service with zero specs on the classpath must still boot. + // toUnmodifiableMap throws IllegalStateException on a duplicate type — the wanted fail-fast. + public DeliverySpecRegistry(ObjectProvider> specs) { + this.byType = specs.stream() + .collect(toUnmodifiableMap(DeliverySpec::getType, identity())); + Set registered = byType.keySet(); + Set sortedTypes = new TreeSet<>(registered); + log.info("Registered {} delivery spec(s): {}", byType.size(), sortedTypes); + } + + public DeliverySpec require(DeliveryType type) { + DeliverySpec spec = byType.get(type); + if (spec == null) { + throw new IllegalArgumentException("No spec registered for delivery type: " + type.name()); + } + return spec; + } + + public Set types() { + return byType.keySet(); + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySweepScheduler.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySweepScheduler.java new file mode 100644 index 0000000000..801ac34a95 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySweepScheduler.java @@ -0,0 +1,43 @@ +package com.openframe.delivery; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import net.javacrumbs.shedlock.spring.annotation.SchedulerLock; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +@Slf4j +@Component +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.delivery.sweep.enabled", havingValue = "true") +public class DeliverySweepScheduler { + + private final DeliverySweepService sweepService; + private final DeliveryWatchdogService watchdogService; + + @Scheduled(fixedDelayString = "${openframe.delivery.sweep.interval}") + @SchedulerLock(name = "deliverySweep", + lockAtMostFor = "${openframe.delivery.sweep.lock-at-most-for}", + lockAtLeastFor = "${openframe.delivery.sweep.lock-at-least-for}") + public void tick() { + retryPending(); + reapAcked(); + } + + private void retryPending() { + try { + sweepService.retryPending(); + } catch (Exception e) { + log.error("Delivery sweep failed", e); + } + } + + private void reapAcked() { + try { + watchdogService.reapAcked(); + } catch (Exception e) { + log.error("Delivery watchdog failed", e); + } + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySweepService.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySweepService.java new file mode 100644 index 0000000000..e29309b110 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliverySweepService.java @@ -0,0 +1,127 @@ +package com.openframe.delivery; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.data.document.delivery.DeliveryFailure; +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.document.rmm.schedule.ScheduleOfflineBehavior; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +import com.openframe.delivery.DeliveryProperties.Policy; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Service; + +import java.time.Instant; +import java.util.List; +import java.util.Set; + +import static java.util.stream.Collectors.toSet; + +@Slf4j +@Service +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "true") +public class DeliverySweepService { + + private final MachineDeliveryRepository repository; + private final MachineOnlineStatus machineOnlineStatus; + private final DeliverySpecRegistry registry; + private final DeliveryProperties properties; + private final DeliveryFailureRecorder failureRecorder; + private final DeliveryMetrics metrics; + private final ObjectMapper objectMapper; + + public void retryPending() { + Instant now = Instant.now(); + Set types = registry.types(); + types.forEach(type -> retryPending(type, now)); + } + + private void retryPending(DeliveryType type, Instant now) { + Policy policy = properties.resolve(type); + Instant threshold = now.minusSeconds(policy.getAckThresholdSeconds()); + List unacked = repository.findByTypeAndStatusAndLastAttemptAtBefore(type, DeliveryStatus.PENDING, threshold); + if (unacked.isEmpty()) { + return; + } + Set machineIds = unacked.stream().map(MachineDelivery::getMachineId).collect(toSet()); + Set offlineMachineIds = machineOnlineStatus.offline(machineIds); + unacked.forEach(delivery -> retryOrFail(delivery, offlineMachineIds, now)); + } + + private void retryOrFail(MachineDelivery delivery, Set offlineMachineIds, Instant now) { + Policy policy = properties.resolve(delivery); + if (isOffline(delivery, offlineMachineIds)) { + waitOrFailOffline(delivery, policy, now); + return; + } + if (hasAttemptsLeft(delivery, policy)) { + republish(delivery, now); + return; + } + failureRecorder.fail(delivery, DeliveryFailure.EXHAUSTED, now); + } + + private void waitOrFailOffline(MachineDelivery delivery, Policy policy, Instant now) { + if (shouldSkipOffline(policy)) { + failureRecorder.fail(delivery, DeliveryFailure.OFFLINE, now); + return; + } + if (isReconnectWindowOver(delivery, policy, now)) { + failureRecorder.fail(delivery, DeliveryFailure.OFFLINE, now); + return; + } + log.debug("Delivery waits for machine to come online: type={} targetId={} machineId={}", + delivery.getType(), delivery.getTargetId(), delivery.getMachineId()); + } + + private void republish(MachineDelivery delivery, Instant now) { + DeliveryType type = delivery.getType(); + DeliverySpec spec = registry.require(type); + publish(spec, delivery); + + int attempt = delivery.getAttempts() + 1; + delivery.setAttempts(attempt); + delivery.setLastAttemptAt(now); + repository.save(delivery); + metrics.recordRetried(type); + log.info("Delivery re-published: type={} targetId={} machineId={} attempt={}", + type, delivery.getTargetId(), delivery.getMachineId(), attempt); + } + + private void publish(DeliverySpec spec, MachineDelivery delivery) { + Class

payloadClass = spec.getPayloadClass(); + P payload = readPayload(delivery, payloadClass); + String machineId = delivery.getMachineId(); + spec.publish(machineId, payload); + } + + private

P readPayload(MachineDelivery delivery, Class

payloadClass) { + String payloadJson = delivery.getPayloadJson(); + try { + return objectMapper.readValue(payloadJson, payloadClass); + } catch (JsonProcessingException e) { + throw new IllegalStateException("Corrupt delivery payload: " + delivery.getId(), e); + } + } + + private static boolean isOffline(MachineDelivery delivery, Set offlineMachineIds) { + return offlineMachineIds.contains(delivery.getMachineId()); + } + + private static boolean hasAttemptsLeft(MachineDelivery delivery, Policy policy) { + return delivery.getAttempts() < policy.getMaxAttempts(); + } + + private static boolean shouldSkipOffline(Policy policy) { + return policy.getOfflineBehavior() == ScheduleOfflineBehavior.SKIP; + } + + private static boolean isReconnectWindowOver(MachineDelivery delivery, Policy policy, Instant now) { + Instant deadline = delivery.getDispatchedAt().plusSeconds(policy.getReconnectWindowSeconds()); + return now.isAfter(deadline); + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryTracker.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryTracker.java new file mode 100644 index 0000000000..46bfda1ab9 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryTracker.java @@ -0,0 +1,10 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; + +public interface DeliveryTracker { + + void acknowledge(DeliveryType type, String targetId, String machineId); + + void complete(DeliveryType type, String targetId, String machineId); +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryWatchdogService.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryWatchdogService.java new file mode 100644 index 0000000000..23b258a478 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/DeliveryWatchdogService.java @@ -0,0 +1,39 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryFailure; +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +import com.openframe.delivery.DeliveryProperties.Policy; +import lombok.RequiredArgsConstructor; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Service; + +import java.time.Instant; +import java.util.List; +import java.util.Set; + +@Service +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "true") +public class DeliveryWatchdogService { + + private final MachineDeliveryRepository repository; + private final DeliverySpecRegistry registry; + private final DeliveryProperties properties; + private final DeliveryFailureRecorder failureRecorder; + + public void reapAcked() { + Instant now = Instant.now(); + Set types = registry.types(); + types.forEach(type -> reapAcked(type, now)); + } + + private void reapAcked(DeliveryType type, Instant now) { + Policy policy = properties.resolve(type); + Instant threshold = now.minusSeconds(policy.getResultTimeoutSeconds()); + List silent = repository.findByTypeAndStatusAndAckedAtBefore(type, DeliveryStatus.ACKED, threshold); + silent.forEach(delivery -> failureRecorder.fail(delivery, DeliveryFailure.TIMEOUT, now)); + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/MachineOnlineStatus.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/MachineOnlineStatus.java new file mode 100644 index 0000000000..b3ba543c38 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/MachineOnlineStatus.java @@ -0,0 +1,26 @@ +package com.openframe.delivery; + +import com.openframe.data.document.device.DeviceStatus; +import com.openframe.data.document.device.Machine; +import com.openframe.data.repository.device.MachineRepository; +import lombok.RequiredArgsConstructor; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import java.util.List; +import java.util.Set; + +import static java.util.stream.Collectors.toSet; + +@Component +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "true") +public class MachineOnlineStatus { + + private final MachineRepository machineRepository; + + public Set offline(Set machineIds) { + List offline = machineRepository.findByMachineIdInAndStatus(machineIds, DeviceStatus.OFFLINE); + return offline.stream().map(Machine::getMachineId).collect(toSet()); + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/MongoDeliveryRecorder.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/MongoDeliveryRecorder.java new file mode 100644 index 0000000000..a937f2ddfa --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/MongoDeliveryRecorder.java @@ -0,0 +1,59 @@ +package com.openframe.delivery; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +import java.time.Instant; + +@Slf4j +@Component +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "true") +public class MongoDeliveryRecorder implements DeliveryRecorder { + + private final MachineDeliveryRepository repository; + private final ObjectMapper objectMapper; + + @Override + public void record(DeliveryRequest request) { + MachineDelivery delivery = pendingRow(request); + repository.save(delivery); + log.info("Delivery recorded: type={} targetId={} machineId={}", + request.getType(), request.getTargetId(), request.getMachineId()); + } + + private MachineDelivery pendingRow(DeliveryRequest request) { + Instant now = Instant.now(); + String id = DeliveryId.of(request.getType(), request.getTargetId(), request.getMachineId()); + String payloadJson = toJson(request.getPayload()); + return MachineDelivery.builder() + .id(id) + .type(request.getType()) + .targetId(request.getTargetId()) + .machineId(request.getMachineId()) + .status(DeliveryStatus.PENDING) + .attempts(0) + .payloadJson(payloadJson) + .dispatchedAt(now) + .lastAttemptAt(now) + .offlineBehavior(request.getOfflineBehavior()) + .reconnectWindowSeconds(request.getReconnectWindowSeconds()) + .build(); + } + + private String toJson(Object payload) { + try { + return objectMapper.writeValueAsString(payload); + } catch (JsonProcessingException e) { + String payloadType = payload.getClass().getSimpleName(); + throw new IllegalArgumentException("Delivery payload is not serializable: " + payloadType, e); + } + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/MongoDeliveryTracker.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/MongoDeliveryTracker.java new file mode 100644 index 0000000000..2da4561bc5 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/MongoDeliveryTracker.java @@ -0,0 +1,69 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +import com.openframe.delivery.DeliveryProperties.Policy; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Service; + +import java.time.Instant; + +@Slf4j +@Service +@RequiredArgsConstructor +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "true") +public class MongoDeliveryTracker implements DeliveryTracker { + + private final MachineDeliveryRepository repository; + private final DeliveryProperties properties; + + @Override + public void acknowledge(DeliveryType type, String targetId, String machineId) { + String id = DeliveryId.of(type, targetId, machineId); + repository.findById(id) + .filter(MongoDeliveryTracker::isPending) + .ifPresent(this::markAcked); + } + + @Override + public void complete(DeliveryType type, String targetId, String machineId) { + String id = DeliveryId.of(type, targetId, machineId); + repository.findById(id) + .filter(MongoDeliveryTracker::isOpen) + .ifPresent(this::markDone); + } + + private void markAcked(MachineDelivery delivery) { + Instant now = Instant.now(); + delivery.setStatus(DeliveryStatus.ACKED); + delivery.setAckedAt(now); + repository.save(delivery); + log.info("Delivery ACKED: type={} targetId={} machineId={}", + delivery.getType(), delivery.getTargetId(), delivery.getMachineId()); + } + + private void markDone(MachineDelivery delivery) { + Instant now = Instant.now(); + Policy policy = properties.resolve(delivery.getType()); + long ttlSeconds = policy.getTtlSeconds(); + delivery.setStatus(DeliveryStatus.DONE); + delivery.setFinishedAt(now); + delivery.setExpiresAt(now.plusSeconds(ttlSeconds)); + repository.save(delivery); + log.info("Delivery DONE: type={} targetId={} machineId={} attempts={}", + delivery.getType(), delivery.getTargetId(), delivery.getMachineId(), delivery.getAttempts()); + } + + private static boolean isPending(MachineDelivery delivery) { + return delivery.getStatus() == DeliveryStatus.PENDING; + } + + private static boolean isOpen(MachineDelivery delivery) { + DeliveryStatus status = delivery.getStatus(); + return status == DeliveryStatus.PENDING || status == DeliveryStatus.ACKED; + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/NoopDeliveryRecorder.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/NoopDeliveryRecorder.java new file mode 100644 index 0000000000..1f26787029 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/NoopDeliveryRecorder.java @@ -0,0 +1,13 @@ +package com.openframe.delivery; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +@Component +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "false", matchIfMissing = true) +public class NoopDeliveryRecorder implements DeliveryRecorder { + + @Override + public void record(DeliveryRequest request) { + } +} diff --git a/openframe-machine-delivery/src/main/java/com/openframe/delivery/NoopDeliveryTracker.java b/openframe-machine-delivery/src/main/java/com/openframe/delivery/NoopDeliveryTracker.java new file mode 100644 index 0000000000..7cc2b96429 --- /dev/null +++ b/openframe-machine-delivery/src/main/java/com/openframe/delivery/NoopDeliveryTracker.java @@ -0,0 +1,18 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Component; + +@Component +@ConditionalOnProperty(name = "openframe.delivery.enabled", havingValue = "false", matchIfMissing = true) +public class NoopDeliveryTracker implements DeliveryTracker { + + @Override + public void acknowledge(DeliveryType type, String targetId, String machineId) { + } + + @Override + public void complete(DeliveryType type, String targetId, String machineId) { + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryDispatcherTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryDispatcherTest.java new file mode 100644 index 0000000000..977dbc349d --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryDispatcherTest.java @@ -0,0 +1,74 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +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 static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class DeliveryDispatcherTest { + + private static final String MACHINE_ID = "mach-42"; + private static final String TARGET_ID = "tactical-agent"; + + @Mock private DeliverySpecRegistry registry; + @Mock private DeliveryRecorder recorder; + @Mock private DeliverySpec spec; + + @InjectMocks private DeliveryDispatcher dispatcher; + + private TestSeed seed; + private TestPayload payload; + private DeliveryRequest request; + + @BeforeEach + void setUp() { + seed = new TestSeed(MACHINE_ID); + payload = new TestPayload(); + request = DeliveryRequest.builder() + .type(DeliveryType.TOOL_INSTALLATION) + .targetId(TARGET_ID) + .machineId(MACHINE_ID) + .payload(payload) + .build(); + } + + @Test + void dispatch_seed_requestRecordedThenPublishedThroughSpec() { + // setup + doReturn(spec).when(registry).require(DeliveryType.TOOL_INSTALLATION); + when(spec.getSeedClass()).thenReturn(TestSeed.class); + when(spec.request(seed)).thenReturn(request); + + // execution + dispatcher.dispatch(seed); + + // verifications + verify(recorder).record(request); + verify(spec).publish(MACHINE_ID, payload); + } + + @Test + void dispatch_unregisteredType_throwsWithoutRecording() { + // setup + when(registry.require(DeliveryType.TOOL_INSTALLATION)) + .thenThrow(new IllegalArgumentException("No spec registered for delivery type: TOOL_INSTALLATION")); + + // execution + IllegalArgumentException ex = assertThrows(IllegalArgumentException.class, () -> dispatcher.dispatch(seed)); + + // verifications + assertThat(ex.getMessage()).contains("TOOL_INSTALLATION"); + verifyNoInteractions(recorder); + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryFailureRecorderTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryFailureRecorderTest.java new file mode 100644 index 0000000000..d7f2a706a7 --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryFailureRecorderTest.java @@ -0,0 +1,62 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryFailure; +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +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 java.time.Instant; + +import static com.openframe.delivery.DeliveryTestPolicies.TTL; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.verify; + +@ExtendWith(MockitoExtension.class) +class DeliveryFailureRecorderTest { + + @Mock private MachineDeliveryRepository repository; + @Mock private DeliverySpecRegistry registry; + @Mock private DeliveryMetrics metrics; + @Mock private DeliverySpec spec; + + private DeliveryFailureRecorder recorder; + + private MachineDelivery delivery; + + @BeforeEach + void setUp() { + delivery = MachineDelivery.builder() + .type(DeliveryType.CLIENT_UNINSTALL) + .status(DeliveryStatus.PENDING) + .attempts(5) + .build(); + DeliveryProperties properties = DeliveryTestPolicies.properties(); + recorder = new DeliveryFailureRecorder(repository, registry, properties, metrics); + } + + @Test + void fail_exhausted_rowFailedMetricCountedSpecNotified() { + // setup + Instant now = Instant.now(); + doReturn(spec).when(registry).require(DeliveryType.CLIENT_UNINSTALL); + + // execution + recorder.fail(delivery, DeliveryFailure.EXHAUSTED, now); + + // verifications + assertThat(delivery.getStatus()).isEqualTo(DeliveryStatus.FAILED); + assertThat(delivery.getFailure()).isEqualTo(DeliveryFailure.EXHAUSTED); + assertThat(delivery.getFinishedAt()).isEqualTo(now); + assertThat(delivery.getExpiresAt()).isEqualTo(now.plusSeconds(TTL)); + verify(repository).save(delivery); + verify(metrics).recordFailed(DeliveryType.CLIENT_UNINSTALL, DeliveryFailure.EXHAUSTED); + verify(spec).onFailed(delivery, DeliveryFailure.EXHAUSTED); + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryPropertiesTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryPropertiesTest.java new file mode 100644 index 0000000000..0703955da4 --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryPropertiesTest.java @@ -0,0 +1,113 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.document.rmm.schedule.ScheduleOfflineBehavior; +import com.openframe.delivery.DeliveryProperties.Policy; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.Map; + +import static com.openframe.delivery.DeliveryTestPolicies.ACK_THRESHOLD; +import static com.openframe.delivery.DeliveryTestPolicies.MAX_ATTEMPTS; +import static com.openframe.delivery.DeliveryTestPolicies.RECONNECT_WINDOW; +import static com.openframe.delivery.DeliveryTestPolicies.RESULT_TIMEOUT; +import static com.openframe.delivery.DeliveryTestPolicies.TTL; +import static org.assertj.core.api.Assertions.assertThat; + +class DeliveryPropertiesTest { + + private static final int UNINSTALL_MAX_ATTEMPTS = 5; + private static final long ROW_RECONNECT_WINDOW = 3_600L; + + private DeliveryProperties properties; + + @BeforeEach + void setUp() { + properties = DeliveryTestPolicies.properties(); + } + + @Test + void resolve_typeWithoutOverride_defaultsReturned() { + // setup + Policy defaults = properties.getDefaults(); + + // execution + Policy resolved = properties.resolve(DeliveryType.TOOL_INSTALLATION); + + // verifications + assertThat(resolved).isSameAs(defaults); + } + + @Test + void resolve_typeOverridesOneField_otherFieldsFromDefaults() { + // setup + Policy uninstall = new Policy(); + uninstall.setMaxAttempts(UNINSTALL_MAX_ATTEMPTS); + properties.setTypes(Map.of(DeliveryType.CLIENT_UNINSTALL, uninstall)); + + // execution + Policy resolved = properties.resolve(DeliveryType.CLIENT_UNINSTALL); + + // verifications + assertThat(resolved.getMaxAttempts()).isEqualTo(UNINSTALL_MAX_ATTEMPTS); + assertThat(resolved.getAckThresholdSeconds()).isEqualTo(ACK_THRESHOLD); + assertThat(resolved.getOfflineBehavior()).isEqualTo(ScheduleOfflineBehavior.RETRY_ON_RECONNECT); + assertThat(resolved.getReconnectWindowSeconds()).isEqualTo(RECONNECT_WINDOW); + assertThat(resolved.getResultTimeoutSeconds()).isEqualTo(RESULT_TIMEOUT); + assertThat(resolved.getTtlSeconds()).isEqualTo(TTL); + } + + @Test + void resolve_typeOverridesOfflineBehavior_overrideWins() { + // setup + Policy scripts = new Policy(); + scripts.setOfflineBehavior(ScheduleOfflineBehavior.SKIP); + properties.setTypes(Map.of(DeliveryType.SCRIPT_SCHEDULE, scripts)); + + // execution + Policy resolved = properties.resolve(DeliveryType.SCRIPT_SCHEDULE); + + // verifications + assertThat(resolved.getOfflineBehavior()).isEqualTo(ScheduleOfflineBehavior.SKIP); + } + + @Test + void resolve_rowWithoutOverrides_typePolicyReturned() { + // setup + Policy uninstall = new Policy(); + uninstall.setMaxAttempts(UNINSTALL_MAX_ATTEMPTS); + properties.setTypes(Map.of(DeliveryType.CLIENT_UNINSTALL, uninstall)); + MachineDelivery delivery = MachineDelivery.builder().type(DeliveryType.CLIENT_UNINSTALL).build(); + + // execution + Policy resolved = properties.resolve(delivery); + + // verifications + assertThat(resolved.getMaxAttempts()).isEqualTo(UNINSTALL_MAX_ATTEMPTS); + assertThat(resolved.getOfflineBehavior()).isEqualTo(ScheduleOfflineBehavior.RETRY_ON_RECONNECT); + assertThat(resolved.getReconnectWindowSeconds()).isEqualTo(RECONNECT_WINDOW); + } + + @Test + void resolve_rowOverridesOfflineFields_rowWinsOverType() { + // setup + Policy scripts = new Policy(); + scripts.setOfflineBehavior(ScheduleOfflineBehavior.SKIP); + properties.setTypes(Map.of(DeliveryType.SCRIPT_SCHEDULE, scripts)); + MachineDelivery delivery = MachineDelivery.builder() + .type(DeliveryType.SCRIPT_SCHEDULE) + .offlineBehavior(ScheduleOfflineBehavior.RETRY_ON_RECONNECT) + .reconnectWindowSeconds(ROW_RECONNECT_WINDOW) + .build(); + + // execution + Policy resolved = properties.resolve(delivery); + + // verifications + assertThat(resolved.getOfflineBehavior()).isEqualTo(ScheduleOfflineBehavior.RETRY_ON_RECONNECT); + assertThat(resolved.getReconnectWindowSeconds()).isEqualTo(ROW_RECONNECT_WINDOW); + assertThat(resolved.getMaxAttempts()).isEqualTo(MAX_ATTEMPTS); + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliverySpecRegistryTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliverySpecRegistryTest.java new file mode 100644 index 0000000000..727ced7c99 --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliverySpecRegistryTest.java @@ -0,0 +1,81 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.ObjectProvider; + +import java.util.stream.Stream; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class DeliverySpecRegistryTest { + + @Test + void require_registeredType_specReturned() { + // setup + DeliverySpec spec = spec(DeliveryType.TOOL_INSTALLATION); + ObjectProvider> specs = provider(spec); + DeliverySpecRegistry registry = new DeliverySpecRegistry(specs); + + // execution + DeliverySpec resolved = registry.require(DeliveryType.TOOL_INSTALLATION); + + // verifications + assertThat(resolved).isSameAs(spec); + } + + @Test + void require_unregisteredType_throwsIllegalArgument() { + // setup + ObjectProvider> specs = provider(); + DeliverySpecRegistry registry = new DeliverySpecRegistry(specs); + + // execution + verifications + assertThatThrownBy(() -> registry.require(DeliveryType.CLIENT_UNINSTALL)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("CLIENT_UNINSTALL"); + } + + @Test + void constructor_duplicateType_throwsIllegalState() { + // setup + DeliverySpec first = spec(DeliveryType.TOOL_INSTALLATION); + DeliverySpec second = spec(DeliveryType.TOOL_INSTALLATION); + ObjectProvider> duplicates = provider(first, second); + + // execution + verifications + assertThatThrownBy(() -> new DeliverySpecRegistry(duplicates)) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("TOOL_INSTALLATION"); + } + + @Test + void types_twoSpecs_bothTypesListed() { + // setup + DeliverySpec install = spec(DeliveryType.TOOL_INSTALLATION); + DeliverySpec uninstall = spec(DeliveryType.CLIENT_UNINSTALL); + ObjectProvider> specs = provider(install, uninstall); + + // execution + DeliverySpecRegistry registry = new DeliverySpecRegistry(specs); + + // verifications + assertThat(registry.types()).containsExactlyInAnyOrder(DeliveryType.TOOL_INSTALLATION, DeliveryType.CLIENT_UNINSTALL); + } + + @SuppressWarnings("unchecked") + private static ObjectProvider> provider(DeliverySpec... specs) { + ObjectProvider> provider = mock(ObjectProvider.class); + when(provider.stream()).thenReturn(Stream.of(specs)); + return provider; + } + + private static DeliverySpec spec(DeliveryType type) { + DeliverySpec spec = mock(DeliverySpec.class); + when(spec.getType()).thenReturn(type); + return spec; + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliverySweepServiceTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliverySweepServiceTest.java new file mode 100644 index 0000000000..bffd6c591e --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliverySweepServiceTest.java @@ -0,0 +1,182 @@ +package com.openframe.delivery; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.data.document.delivery.DeliveryFailure; +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.document.rmm.schedule.ScheduleOfflineBehavior; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import java.time.Instant; +import java.util.List; +import java.util.Set; + +import static com.openframe.delivery.DeliveryTestPolicies.ACK_THRESHOLD; +import static com.openframe.delivery.DeliveryTestPolicies.MAX_ATTEMPTS; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doReturn; +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 DeliverySweepServiceTest { + + private static final String MACHINE_ID = "mach-42"; + private static final String TARGET_ID = "tactical-agent"; + private static final String PAYLOAD_JSON = "{\"value\":\"tactical-agent\"}"; + private static final long TWO_DAYS_SECONDS = 172_800L; + + @Mock private MachineDeliveryRepository repository; + @Mock private MachineOnlineStatus machineOnlineStatus; + @Mock private DeliverySpecRegistry registry; + @Mock private DeliveryFailureRecorder failureRecorder; + @Mock private DeliveryMetrics metrics; + @Mock private DeliverySpec spec; + + @Captor private ArgumentCaptor payloadCaptor; + + private DeliverySweepService service; + + private MachineDelivery delivery; + + @BeforeEach + void setUp() { + Instant dispatchedAt = Instant.now().minusSeconds(ACK_THRESHOLD * 2); + delivery = MachineDelivery.builder() + .id(DeliveryId.of(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID)) + .type(DeliveryType.TOOL_INSTALLATION) + .targetId(TARGET_ID) + .machineId(MACHINE_ID) + .status(DeliveryStatus.PENDING) + .attempts(0) + .payloadJson(PAYLOAD_JSON) + .dispatchedAt(dispatchedAt) + .lastAttemptAt(dispatchedAt) + .build(); + DeliveryProperties properties = DeliveryTestPolicies.properties(); + service = new DeliverySweepService(repository, machineOnlineStatus, registry, properties, failureRecorder, metrics, new ObjectMapper()); + } + + @Test + void retryPending_onlineWithAttemptsLeft_republishedAndAttemptCounted() { + // setup + stubOneStuckDelivery(); + stubMachineOnline(); + doReturn(spec).when(registry).require(DeliveryType.TOOL_INSTALLATION); + when(spec.getPayloadClass()).thenReturn(TestPayload.class); + + // execution + service.retryPending(); + + // verifications + verify(spec).publish(eq(MACHINE_ID), payloadCaptor.capture()); + assertThat(payloadCaptor.getValue().getValue()).isEqualTo(TARGET_ID); + assertThat(delivery.getAttempts()).isEqualTo(1); + assertThat(delivery.getStatus()).isEqualTo(DeliveryStatus.PENDING); + verify(repository).save(delivery); + verify(metrics).recordRetried(DeliveryType.TOOL_INSTALLATION); + verifyNoInteractions(failureRecorder); + } + + @Test + void retryPending_onlineAttemptsExhausted_failedExhausted() { + // setup + delivery.setAttempts(MAX_ATTEMPTS); + stubOneStuckDelivery(); + stubMachineOnline(); + + // execution + service.retryPending(); + + // verifications + verify(failureRecorder).fail(eq(delivery), eq(DeliveryFailure.EXHAUSTED), any(Instant.class)); + verify(repository, never()).save(delivery); + verifyNoInteractions(metrics); + } + + @Test + void retryPending_offlineInsideReconnectWindow_leftPendingWithoutAttempt() { + // setup + stubOneStuckDelivery(); + stubMachineOffline(); + + // execution + service.retryPending(); + + // verifications + assertThat(delivery.getAttempts()).isZero(); + verify(repository, never()).save(delivery); + verifyNoInteractions(failureRecorder); + verifyNoInteractions(metrics); + } + + @Test + void retryPending_offlineReconnectWindowOver_failedOffline() { + // setup + Instant twoDaysAgo = Instant.now().minusSeconds(TWO_DAYS_SECONDS); + delivery.setDispatchedAt(twoDaysAgo); + stubOneStuckDelivery(); + stubMachineOffline(); + + // execution + service.retryPending(); + + // verifications + verify(failureRecorder).fail(eq(delivery), eq(DeliveryFailure.OFFLINE), any(Instant.class)); + } + + @Test + void retryPending_offlineWithRowSkipOverride_failedOfflineImmediately() { + // setup + delivery.setOfflineBehavior(ScheduleOfflineBehavior.SKIP); + stubOneStuckDelivery(); + stubMachineOffline(); + + // execution + service.retryPending(); + + // verifications + verify(failureRecorder).fail(eq(delivery), eq(DeliveryFailure.OFFLINE), any(Instant.class)); + } + + @Test + void retryPending_nothingStuck_onlineStatusNotQueried() { + // setup + when(registry.types()).thenReturn(Set.of(DeliveryType.TOOL_INSTALLATION)); + when(repository.findByTypeAndStatusAndLastAttemptAtBefore(eq(DeliveryType.TOOL_INSTALLATION), eq(DeliveryStatus.PENDING), any(Instant.class))) + .thenReturn(List.of()); + + // execution + service.retryPending(); + + // verifications + verifyNoInteractions(machineOnlineStatus); + verifyNoInteractions(failureRecorder); + } + + private void stubOneStuckDelivery() { + when(registry.types()).thenReturn(Set.of(DeliveryType.TOOL_INSTALLATION)); + when(repository.findByTypeAndStatusAndLastAttemptAtBefore(eq(DeliveryType.TOOL_INSTALLATION), eq(DeliveryStatus.PENDING), any(Instant.class))) + .thenReturn(List.of(delivery)); + } + + private void stubMachineOnline() { + when(machineOnlineStatus.offline(Set.of(MACHINE_ID))).thenReturn(Set.of()); + } + + private void stubMachineOffline() { + when(machineOnlineStatus.offline(Set.of(MACHINE_ID))).thenReturn(Set.of(MACHINE_ID)); + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryTestPolicies.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryTestPolicies.java new file mode 100644 index 0000000000..b7c3a1867b --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryTestPolicies.java @@ -0,0 +1,29 @@ +package com.openframe.delivery; + +import com.openframe.data.document.rmm.schedule.ScheduleOfflineBehavior; +import com.openframe.delivery.DeliveryProperties.Policy; + +final class DeliveryTestPolicies { + + static final long ACK_THRESHOLD = 30L; + static final int MAX_ATTEMPTS = 3; + static final long RECONNECT_WINDOW = 86_400L; + static final long RESULT_TIMEOUT = 600L; + static final long TTL = 604_800L; + + private DeliveryTestPolicies() { + } + + static DeliveryProperties properties() { + Policy defaults = new Policy(); + defaults.setAckThresholdSeconds(ACK_THRESHOLD); + defaults.setMaxAttempts(MAX_ATTEMPTS); + defaults.setOfflineBehavior(ScheduleOfflineBehavior.RETRY_ON_RECONNECT); + defaults.setReconnectWindowSeconds(RECONNECT_WINDOW); + defaults.setResultTimeoutSeconds(RESULT_TIMEOUT); + defaults.setTtlSeconds(TTL); + DeliveryProperties properties = new DeliveryProperties(); + properties.setDefaults(defaults); + return properties; + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryWatchdogServiceTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryWatchdogServiceTest.java new file mode 100644 index 0000000000..dd2ed28d96 --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/DeliveryWatchdogServiceTest.java @@ -0,0 +1,74 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryFailure; +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +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 java.time.Instant; +import java.util.List; +import java.util.Set; + +import static com.openframe.delivery.DeliveryTestPolicies.RESULT_TIMEOUT; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class DeliveryWatchdogServiceTest { + + @Mock private MachineDeliveryRepository repository; + @Mock private DeliverySpecRegistry registry; + @Mock private DeliveryFailureRecorder failureRecorder; + + private DeliveryWatchdogService service; + + private MachineDelivery silentDelivery; + + @BeforeEach + void setUp() { + silentDelivery = MachineDelivery.builder() + .type(DeliveryType.CLIENT_UNINSTALL) + .status(DeliveryStatus.ACKED) + .ackedAt(Instant.now().minusSeconds(RESULT_TIMEOUT * 2)) + .build(); + DeliveryProperties properties = DeliveryTestPolicies.properties(); + service = new DeliveryWatchdogService(repository, registry, properties, failureRecorder); + } + + @Test + void reapAcked_ackedOlderThanResultTimeout_failedTimeout() { + // setup + when(registry.types()).thenReturn(Set.of(DeliveryType.CLIENT_UNINSTALL)); + when(repository.findByTypeAndStatusAndAckedAtBefore(eq(DeliveryType.CLIENT_UNINSTALL), eq(DeliveryStatus.ACKED), any(Instant.class))) + .thenReturn(List.of(silentDelivery)); + + // execution + service.reapAcked(); + + // verifications + verify(failureRecorder).fail(eq(silentDelivery), eq(DeliveryFailure.TIMEOUT), any(Instant.class)); + } + + @Test + void reapAcked_nothingSilent_noFailure() { + // setup + when(registry.types()).thenReturn(Set.of(DeliveryType.CLIENT_UNINSTALL)); + when(repository.findByTypeAndStatusAndAckedAtBefore(eq(DeliveryType.CLIENT_UNINSTALL), eq(DeliveryStatus.ACKED), any(Instant.class))) + .thenReturn(List.of()); + + // execution + service.reapAcked(); + + // verifications + verifyNoInteractions(failureRecorder); + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/MachineOnlineStatusTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/MachineOnlineStatusTest.java new file mode 100644 index 0000000000..c0a6e6e210 --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/MachineOnlineStatusTest.java @@ -0,0 +1,42 @@ +package com.openframe.delivery; + +import com.openframe.data.document.device.DeviceStatus; +import com.openframe.data.document.device.Machine; +import com.openframe.data.repository.device.MachineRepository; +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 java.util.Set; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class MachineOnlineStatusTest { + + private static final String ONLINE_ID = "mach-1"; + private static final String OFFLINE_ID = "mach-2"; + + @Mock private MachineRepository machineRepository; + + @InjectMocks private MachineOnlineStatus status; + + @Test + void offline_mixedMachines_onlyOfflineIdsReturned() { + // setup + Machine offline = new Machine(); + offline.setMachineId(OFFLINE_ID); + Set asked = Set.of(ONLINE_ID, OFFLINE_ID); + when(machineRepository.findByMachineIdInAndStatus(asked, DeviceStatus.OFFLINE)).thenReturn(List.of(offline)); + + // execution + Set result = status.offline(asked); + + // verifications + assertThat(result).containsExactly(OFFLINE_ID); + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/MongoDeliveryRecorderTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/MongoDeliveryRecorderTest.java new file mode 100644 index 0000000000..2eb870b77c --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/MongoDeliveryRecorderTest.java @@ -0,0 +1,66 @@ +package com.openframe.delivery; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.document.rmm.schedule.ScheduleOfflineBehavior; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.verify; + +@ExtendWith(MockitoExtension.class) +class MongoDeliveryRecorderTest { + + private static final String MACHINE_ID = "mach-42"; + private static final String VALUE = "issued"; + + @Mock private MachineDeliveryRepository repository; + + @Captor private ArgumentCaptor deliveryCaptor; + + private MongoDeliveryRecorder recorder; + + private DeliveryRequest request; + + @BeforeEach + void setUp() { + TestPayload payload = new TestPayload(); + payload.setValue(VALUE); + request = DeliveryRequest.builder() + .type(DeliveryType.CLIENT_UNINSTALL) + .targetId(MACHINE_ID) + .machineId(MACHINE_ID) + .payload(payload) + .offlineBehavior(ScheduleOfflineBehavior.RETRY_ON_RECONNECT) + .build(); + recorder = new MongoDeliveryRecorder(repository, new ObjectMapper()); + } + + @Test + void record_request_pendingRowSaved() { + // setup + + // execution + recorder.record(request); + + // verifications + verify(repository).save(deliveryCaptor.capture()); + MachineDelivery saved = deliveryCaptor.getValue(); + assertThat(saved.getId()).isEqualTo("CLIENT_UNINSTALL:mach-42:mach-42"); + assertThat(saved.getType()).isEqualTo(DeliveryType.CLIENT_UNINSTALL); + assertThat(saved.getStatus()).isEqualTo(DeliveryStatus.PENDING); + assertThat(saved.getAttempts()).isZero(); + assertThat(saved.getPayloadJson()).contains(VALUE); + assertThat(saved.getOfflineBehavior()).isEqualTo(ScheduleOfflineBehavior.RETRY_ON_RECONNECT); + assertThat(saved.getDispatchedAt()).isEqualTo(saved.getLastAttemptAt()); + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/MongoDeliveryTrackerTest.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/MongoDeliveryTrackerTest.java new file mode 100644 index 0000000000..fe9a8c66b1 --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/MongoDeliveryTrackerTest.java @@ -0,0 +1,114 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryStatus; +import com.openframe.data.document.delivery.DeliveryType; +import com.openframe.data.document.delivery.MachineDelivery; +import com.openframe.data.repository.delivery.MachineDeliveryRepository; +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 java.util.Optional; + +import static com.openframe.delivery.DeliveryTestPolicies.TTL; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +class MongoDeliveryTrackerTest { + + private static final String MACHINE_ID = "mach-42"; + private static final String TARGET_ID = "tactical-agent"; + private static final String DELIVERY_ID = DeliveryId.of(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); + + @Mock private MachineDeliveryRepository repository; + + private MongoDeliveryTracker tracker; + + private MachineDelivery delivery; + + @BeforeEach + void setUp() { + delivery = MachineDelivery.builder() + .id(DELIVERY_ID) + .type(DeliveryType.TOOL_INSTALLATION) + .targetId(TARGET_ID) + .machineId(MACHINE_ID) + .status(DeliveryStatus.PENDING) + .build(); + tracker = new MongoDeliveryTracker(repository, DeliveryTestPolicies.properties()); + } + + @Test + void acknowledge_pendingRow_ackedWithTimestamp() { + // setup + when(repository.findById(DELIVERY_ID)).thenReturn(Optional.of(delivery)); + + // execution + tracker.acknowledge(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); + + // verifications + assertThat(delivery.getStatus()).isEqualTo(DeliveryStatus.ACKED); + assertThat(delivery.getAckedAt()).isNotNull(); + verify(repository).save(delivery); + } + + @Test + void acknowledge_alreadyAcked_untouched() { + // setup + delivery.setStatus(DeliveryStatus.ACKED); + when(repository.findById(DELIVERY_ID)).thenReturn(Optional.of(delivery)); + + // execution + tracker.acknowledge(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); + + // verifications + verify(repository, never()).save(delivery); + } + + @Test + void acknowledge_unknownRow_noop() { + // setup + when(repository.findById(DELIVERY_ID)).thenReturn(Optional.empty()); + + // execution + tracker.acknowledge(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); + + // verifications + verify(repository, never()).save(delivery); + } + + @Test + void complete_ackedRow_doneWithExpiry() { + // setup + delivery.setStatus(DeliveryStatus.ACKED); + when(repository.findById(DELIVERY_ID)).thenReturn(Optional.of(delivery)); + + // execution + tracker.complete(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); + + // verifications + assertThat(delivery.getStatus()).isEqualTo(DeliveryStatus.DONE); + assertThat(delivery.getFinishedAt()).isNotNull(); + assertThat(delivery.getExpiresAt()).isEqualTo(delivery.getFinishedAt().plusSeconds(TTL)); + verify(repository).save(delivery); + } + + @Test + void complete_failedRow_untouched() { + // setup + delivery.setStatus(DeliveryStatus.FAILED); + when(repository.findById(DELIVERY_ID)).thenReturn(Optional.of(delivery)); + + // execution + tracker.complete(DeliveryType.TOOL_INSTALLATION, TARGET_ID, MACHINE_ID); + + // verifications + assertThat(delivery.getStatus()).isEqualTo(DeliveryStatus.FAILED); + verify(repository, never()).save(delivery); + } +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/TestPayload.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/TestPayload.java new file mode 100644 index 0000000000..f2b4e46f28 --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/TestPayload.java @@ -0,0 +1,8 @@ +package com.openframe.delivery; + +import lombok.Data; + +@Data +class TestPayload { + private String value; +} diff --git a/openframe-machine-delivery/src/test/java/com/openframe/delivery/TestSeed.java b/openframe-machine-delivery/src/test/java/com/openframe/delivery/TestSeed.java new file mode 100644 index 0000000000..d5b16ddd58 --- /dev/null +++ b/openframe-machine-delivery/src/test/java/com/openframe/delivery/TestSeed.java @@ -0,0 +1,16 @@ +package com.openframe.delivery; + +import com.openframe.data.document.delivery.DeliveryType; +import lombok.AllArgsConstructor; +import lombok.Getter; + +@Getter +@AllArgsConstructor +class TestSeed implements DeliverySeed { + private final String machineId; + + @Override + public DeliveryType type() { + return DeliveryType.TOOL_INSTALLATION; + } +} diff --git a/pom.xml b/pom.xml index 8226a9d8f1..3edacc7464 100644 --- a/pom.xml +++ b/pom.xml @@ -44,6 +44,7 @@ openframe-core-crypto openframe-data-mongo-common openframe-data-mongo-sync + openframe-machine-delivery openframe-data-mongo-reactive openframe-data-redis openframe-presence @@ -142,6 +143,11 @@ openframe-notification-core ${revision} + + com.openframe.oss + openframe-machine-delivery + ${revision} + com.openframe.oss openframe-data-kafka