diff --git a/docs/email-async/01-async-send-pipeline.md b/docs/email-async/01-async-send-pipeline.md index f37a6055..fd4442d2 100644 --- a/docs/email-async/01-async-send-pipeline.md +++ b/docs/email-async/01-async-send-pipeline.md @@ -148,9 +148,6 @@ parses CSV ([parseCSV.ts](../../js/src/features/emails/api/parseCSV.ts)) and pos committed, the rows are already visible — no `AFTER_COMMIT` event is needed. **There is no automatic enqueue-time trigger** (the accepted tradeoff of #6: if the kick is never issued, the batch waits for the next kick or a restart). - 2. **On startup** — `@EventListener(ApplicationReadyEvent.class)` first resets `PROCESSING → ERROR` - (at-most-once recovery), then calls `trigger()` once (covers rows left `PENDING` before shutdown). This is the - only safety net for a missed kick. - **Drain job** (runs on the executor thread, loops until no rows, then the thread idles): 1. **Claim** up to 50 `PENDING` rows atomically: ```sql diff --git a/src/main/java/org/patinanetwork/patchats/email/EmailDrainer.java b/src/main/java/org/patinanetwork/patchats/email/EmailDrainer.java index 5087994a..b1e29b8a 100644 --- a/src/main/java/org/patinanetwork/patchats/email/EmailDrainer.java +++ b/src/main/java/org/patinanetwork/patchats/email/EmailDrainer.java @@ -13,9 +13,6 @@ import org.patinanetwork.patchats.email.db.repos.EmailRepo; import org.patinanetwork.patchats.email.db.repos.EmailTemplateRepo; import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.boot.context.event.ApplicationReadyEvent; -import org.springframework.context.annotation.Profile; -import org.springframework.context.event.EventListener; import org.springframework.stereotype.Component; /** @@ -27,7 +24,6 @@ */ @Component @Slf4j -@Profile("!ci") public class EmailDrainer { private static final int BATCH_SIZE = 50; @@ -113,17 +109,4 @@ private void sendOne(final Email email, final Map templateC emailRepo.markError(email.getId(), ex.getMessage()); } } - - /** - * On boot: reset orphaned {@code PROCESSING} rows to {@code ERROR} (at-most-once recovery, decision #9), then kick - * one drain to cover rows left {@code PENDING} before shutdown — the only safety net for a missed kick. - */ - @EventListener(ApplicationReadyEvent.class) - public void onApplicationReady() { - final int reset = emailRepo.resetProcessingToError(); - if (reset > 0) { - log.warn("Reset {} orphaned PROCESSING email(s) to ERROR on startup", reset); - } - trigger(); - } } diff --git a/src/main/java/org/patinanetwork/patchats/email/db/repos/EmailRepo.java b/src/main/java/org/patinanetwork/patchats/email/db/repos/EmailRepo.java index 926fb40e..328fba11 100644 --- a/src/main/java/org/patinanetwork/patchats/email/db/repos/EmailRepo.java +++ b/src/main/java/org/patinanetwork/patchats/email/db/repos/EmailRepo.java @@ -25,12 +25,6 @@ public interface EmailRepo { /** Marks a row {@code ERROR} with the given message (no retry — terminal, decision #8). */ void markError(UUID id, String errorMessage); - /** - * At-most-once crash recovery (decision #9): flips any orphaned {@code PROCESSING} rows to {@code ERROR}. Returns - * the number of rows reset. - */ - int resetProcessingToError(); - /** Per-status row counts for one batch ({@code GROUP BY status}); statuses with no rows are absent. */ Map countByStatus(UUID requestId); diff --git a/src/main/java/org/patinanetwork/patchats/email/db/repos/EmailSqlRepo.java b/src/main/java/org/patinanetwork/patchats/email/db/repos/EmailSqlRepo.java index 49e8acd7..0e6eca59 100644 --- a/src/main/java/org/patinanetwork/patchats/email/db/repos/EmailSqlRepo.java +++ b/src/main/java/org/patinanetwork/patchats/email/db/repos/EmailSqlRepo.java @@ -117,17 +117,6 @@ public void markError(final UUID id, final String errorMessage) { .update(); } - @Override - public int resetProcessingToError() { - return jdbc.sql(""" - UPDATE emails - SET status = 'ERROR', - error_message = 'Reset on startup: orphaned PROCESSING row (at-most-once recovery)', - updated_at = now() - WHERE status = 'PROCESSING' - """).update(); - } - @Override public Map countByStatus(final UUID requestId) { final Map counts = new EnumMap<>(EmailStatus.class); diff --git a/src/test/java/org/patinanetwork/patchats/email/EmailDrainerTest.java b/src/test/java/org/patinanetwork/patchats/email/EmailDrainerTest.java index a9161ff0..efa9e19b 100644 --- a/src/test/java/org/patinanetwork/patchats/email/EmailDrainerTest.java +++ b/src/test/java/org/patinanetwork/patchats/email/EmailDrainerTest.java @@ -128,17 +128,6 @@ void renderFailureIsTerminalError() { verify(emailRepo, never()).markSent(any()); } - @Test - void startupResetsProcessingToErrorThenDrains() { - when(emailRepo.resetProcessingToError()).thenReturn(2); - when(emailRepo.claimBatch(50)).thenReturn(List.of()); - - drainer(SYNC).onApplicationReady(); - - verify(emailRepo).resetProcessingToError(); - verify(emailRepo).claimBatch(50); // the startup kick still runs a drain pass - } - @Test void overlappingTriggersCoalesceToOneDrain() { // A manual executor that captures jobs without running them, to observe submission count.