From 17ed1f53194a464ae5e1eb24d79842b0675adbdb Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 19:19:04 +0900 Subject: [PATCH 1/3] test(reliability): require one recovery clock observation --- ...aultConversionWorkerRecoveryClockTest.java | 91 +++++++++++++++++++ 1 file changed, 91 insertions(+) create mode 100644 src/test/java/com/clearfolio/viewer/service/DefaultConversionWorkerRecoveryClockTest.java diff --git a/src/test/java/com/clearfolio/viewer/service/DefaultConversionWorkerRecoveryClockTest.java b/src/test/java/com/clearfolio/viewer/service/DefaultConversionWorkerRecoveryClockTest.java new file mode 100644 index 00000000..8150ba7b --- /dev/null +++ b/src/test/java/com/clearfolio/viewer/service/DefaultConversionWorkerRecoveryClockTest.java @@ -0,0 +1,91 @@ +package com.clearfolio.viewer.service; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.time.Duration; +import java.time.Instant; +import java.util.Optional; +import java.util.UUID; + +import org.junit.jupiter.api.Test; + +import com.clearfolio.viewer.artifact.InMemoryArtifactStore; +import com.clearfolio.viewer.artifact.PdfBoxArtifactGenerator; +import com.clearfolio.viewer.config.ConversionProperties; +import com.clearfolio.viewer.model.ConversionJob; +import com.clearfolio.viewer.repository.ConversionJobStateStore; +import com.clearfolio.viewer.repository.InMemoryConversionJobRepository; + +/** + * Verifies that startup recovery uses one caller-supplied clock observation. + */ +class DefaultConversionWorkerRecoveryClockTest { + + @Test + void staleProcessingRetryUsesTheRecoveryEvaluationTimestamp() { + InMemoryConversionJobRepository repository = new InMemoryConversionJobRepository(); + CapturingStateStore stateStore = new CapturingStateStore(repository); + ConversionJob staleProcessing = new ConversionJob( + UUID.randomUUID(), + "stale.docx", + "application/octet-stream", + "hash-recovery-clock", + 10L, + 3 + ); + assertTrue(staleProcessing.markProcessing("worker exited")); + repository.save(staleProcessing); + + Instant recoveryNow = Instant.now().plus(Duration.ofDays(1)); + DefaultConversionWorker worker = new DefaultConversionWorker( + repository, + stateStore, + command -> { }, + new InMemoryArtifactStore(), + new PdfBoxArtifactGenerator(), + new ConversionProperties(), + id -> "/artifacts/" + id + ".pdf" + ); + + int recovered = worker.recoverPendingJobs(recoveryNow, Duration.ofSeconds(60)); + + assertEquals(1, recovered); + assertEquals(recoveryNow, stateStore.retryAt); + } + + private static final class CapturingStateStore implements ConversionJobStateStore { + private final ConversionJobStateStore delegate; + private Instant retryAt; + + private CapturingStateStore(ConversionJobStateStore delegate) { + this.delegate = delegate; + } + + @Override + public Optional claimForProcessing(UUID jobId, Instant now) { + return delegate.claimForProcessing(jobId, now); + } + + @Override + public void scheduleRetry(UUID jobId, String message, Instant retryAt) { + this.retryAt = retryAt; + delegate.scheduleRetry(jobId, message, retryAt); + } + + @Override + public void markSucceeded(UUID jobId, String resourcePath, String message) { + delegate.markSucceeded(jobId, resourcePath, message); + } + + @Override + public void markDeadLettered(UUID jobId, String message) { + delegate.markDeadLettered(jobId, message); + } + + @Override + public boolean retryDeadLettered(UUID jobId, String operatorId) { + return delegate.retryDeadLettered(jobId, operatorId); + } + } +} From 5930c53fc0ed8f40efa4b7d9909946924c6875ca Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 19:21:40 +0900 Subject: [PATCH 2/3] fix(reliability): reuse recovery evaluation timestamp --- .../com/clearfolio/viewer/service/DefaultConversionWorker.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main/java/com/clearfolio/viewer/service/DefaultConversionWorker.java b/src/main/java/com/clearfolio/viewer/service/DefaultConversionWorker.java index abe17667..0a73dd87 100644 --- a/src/main/java/com/clearfolio/viewer/service/DefaultConversionWorker.java +++ b/src/main/java/com/clearfolio/viewer/service/DefaultConversionWorker.java @@ -186,7 +186,7 @@ public int recoverPendingJobs(Instant now, Duration processingLeaseTimeout) { List recoverableJobs = repository.findRecoverableJobs(now, staleProcessingBefore); recoverableJobs.forEach(job -> { if (job.getStatus() == ConversionJobStatus.PROCESSING) { - stateStore.scheduleRetry(job.getJobId(), "worker lease expired; retry queued", Instant.now()); + stateStore.scheduleRetry(job.getJobId(), "worker lease expired; retry queued", now); } enqueue(job.getJobId()); }); From 90ed8990c24aa4a5ab761b954387001c32ff5124 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Tue, 11 Aug 2026 20:12:04 +0900 Subject: [PATCH 3/3] test(reliability): align recovery fixture clock --- .../viewer/service/DefaultConversionWorkerTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/test/java/com/clearfolio/viewer/service/DefaultConversionWorkerTest.java b/src/test/java/com/clearfolio/viewer/service/DefaultConversionWorkerTest.java index 88503b48..b950eaed 100644 --- a/src/test/java/com/clearfolio/viewer/service/DefaultConversionWorkerTest.java +++ b/src/test/java/com/clearfolio/viewer/service/DefaultConversionWorkerTest.java @@ -607,7 +607,6 @@ void scheduleRetryImmediatelyRequeuesWhenRetryTimeAlreadyPassed() throws Excepti void recoverPendingJobsRequeuesDueSubmittedAndStaleProcessingJobs() { InMemoryConversionJobRepository repository = new InMemoryConversionJobRepository(); ConversionProperties conversionProperties = new ConversionProperties(); - Instant recoveryNow = Instant.now().plusSeconds(120); ConversionJob dueSubmitted = new ConversionJob( UUID.randomUUID(), @@ -625,7 +624,6 @@ void recoverPendingJobsRequeuesDueSubmittedAndStaleProcessingJobs() { 10L, 3 ); - futureRetry.markRetryScheduled("retry later", recoveryNow.plusSeconds(30)); ConversionJob staleProcessing = new ConversionJob( UUID.randomUUID(), "stale.docx", @@ -635,6 +633,8 @@ void recoverPendingJobsRequeuesDueSubmittedAndStaleProcessingJobs() { 3 ); assertTrue(staleProcessing.markProcessing("worker exited")); + Instant recoveryNow = staleProcessing.getStartedAt().plusNanos(1); + futureRetry.markRetryScheduled("retry later", recoveryNow.plusSeconds(30)); repository.save(dueSubmitted); repository.save(futureRetry); repository.save(staleProcessing); @@ -652,7 +652,7 @@ void recoverPendingJobsRequeuesDueSubmittedAndStaleProcessingJobs() { } ); - int recovered = worker.recoverPendingJobs(recoveryNow, Duration.ofSeconds(60)); + int recovered = worker.recoverPendingJobs(recoveryNow, Duration.ZERO); assertEquals(2, recovered); assertEquals(2, attempts.get());