From 15c22d4bf65266c997963a8f127245322bf51dbb Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Wed, 12 Aug 2026 11:35:04 +0200 Subject: [PATCH 01/13] first draft for concurrency using an approach at the database level --- readme.md | 24 ++- .../EmailServiceConfiguration.java | 32 +++- .../application/CleanupAttachmentFiles.java | 42 +++--- .../lib/application/ManageEmail.java | 63 +++++++- .../lib/application/QueryEmail.java | 30 ++-- .../lib/application/SendScheduledEmails.java | 44 +++--- ...itional-spring-configuration-metadata.json | 6 + .../CleanupAttachmentFilesTest.java | 142 ++++++++++++++++++ .../application/SendScheduledEmailsTest.java | 138 +++++++++++++++++ .../database/factory/EmailFactory.java | 22 +++ 10 files changed, 478 insertions(+), 65 deletions(-) create mode 100644 src/test/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFilesTest.java create mode 100644 src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java diff --git a/readme.md b/readme.md index 19c5434..2dffdd3 100644 --- a/readme.md +++ b/readme.md @@ -74,11 +74,25 @@ public class App { The following configuration options are available: -| Name | Default | Description | -|----------------------------------------|-------------|-----------------------------------------------------------------------| -| `lib.emailservice.migrations.enabled` | true | Enables database migrations. | -| `lib.emailservice.scheduling.enabled` | true | Enables the scheduler sending the emails. | -| `lib.emailservice.scheduling.interval` | 30000 | Specifies the milliseconds delay between runs of the scheduler. | +| Name | Default | Description | +|-------------------------------------------------|---------|----------------------------------------------------------------------------------------------------------------| +| `aboutbits.emailservice.migrations.enabled` | true | Enables database migrations. | +| `aboutbits.emailservice.scheduling.enabled` | true | Enables the scheduler sending the emails. | +| `aboutbits.emailservice.scheduling.cleanup.enabled` | true | Enables cleanup of attachment files after sending. | +| `aboutbits.emailservice.scheduling.interval` | 30000 | Specifies the milliseconds delay between runs of the scheduler. | +| `aboutbits.emailservice.scheduling.batch-size` | 50 | Maximum number of emails a single pod claims per scheduler pass (see [Multi-pod deployments](#multi-pod-deployments)). | + +## Multi-pod deployments + +The scheduler is safe to run in every pod concurrently. Each ready email is +claimed by exactly one pod using Postgres row-level `SELECT ... FOR UPDATE SKIP LOCKED`, +so pods work on disjoint rows in parallel and no email is ever sent twice by +different pods. The same guarantee applies to the attachment cleanup scheduler. + +`aboutbits.emailservice.scheduling.batch-size` caps the number of emails one +pod processes per pass. With the default 30s interval and 50 emails per pass, +a single pod can take up to 100 emails per minute; For higher-throughput deployments increase the batch size or +lower the interval. ## Local development: diff --git a/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java b/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java index 535054e..27ed84e 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java @@ -16,6 +16,7 @@ import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; import jakarta.persistence.EntityManager; import org.jspecify.annotations.NullMarked; +import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.AutoConfigurationPackage; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -45,25 +46,44 @@ public EmailMapper emailMapper() { } @Bean - public QueryEmail queryEmail(EmailRepository emailRepository, EmailMapper emailMapper, EntityManager entityManager) { + public QueryEmail queryEmail( + EmailRepository emailRepository, + EmailMapper emailMapper, + EntityManager entityManager + ) { return new QueryEmail(emailRepository, emailMapper, entityManager); } @Bean - public ManageEmail manageEmail(EmailRepository emailRepository, JavaMailSender javaMailSender, AttachmentDataSource attachmentDataSource, EmailMapper emailMapper) { + public ManageEmail manageEmail( + EmailRepository emailRepository, + JavaMailSender javaMailSender, + AttachmentDataSource attachmentDataSource, + EmailMapper emailMapper + ) { return new ManageEmail(emailRepository, javaMailSender, attachmentDataSource, emailMapper); } @Bean @ConditionalOnProperty(value = "aboutbits.emailservice.scheduling.enabled", matchIfMissing = true) - public SendScheduledEmails sendScheduledEmails(QueryEmail queryEmail, ManageEmail manageEmail, List callbacks) { - return new SendScheduledEmails(queryEmail, manageEmail, callbacks); + public SendScheduledEmails sendScheduledEmails( + QueryEmail queryEmail, + ManageEmail manageEmail, + List callbacks, + @Value("${aboutbits.emailservice.scheduling.batch-size:50}") int batchSize + ) { + return new SendScheduledEmails(queryEmail, manageEmail, callbacks, batchSize); } @Bean @ConditionalOnProperty(value = "aboutbits.emailservice.scheduling.cleanup.enabled", matchIfMissing = true) - public CleanupAttachmentFiles cleanupAttachments(QueryEmail queryEmail, ManageEmail manageEmail, List callbacks) { - return new CleanupAttachmentFiles(queryEmail, manageEmail, callbacks); + public CleanupAttachmentFiles cleanupAttachments( + QueryEmail queryEmail, + ManageEmail manageEmail, + List callbacks, + @Value("${aboutbits.emailservice.scheduling.batch-size:50}") int batchSize + ) { + return new CleanupAttachmentFiles(queryEmail, manageEmail, callbacks, batchSize); } @Bean diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFiles.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFiles.java index df7209d..057e4e9 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFiles.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFiles.java @@ -2,16 +2,14 @@ import it.aboutbits.springboot.emailservice.lib.AttachmentCleanerCallback; -import it.aboutbits.springboot.emailservice.lib.exception.AttachmentException; -import lombok.RequiredArgsConstructor; import lombok.extern.log4j.Log4j2; import org.jspecify.annotations.NullMarked; import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.transaction.annotation.Transactional; import java.time.Duration; import java.util.List; -@RequiredArgsConstructor @Log4j2 @NullMarked public class CleanupAttachmentFiles { @@ -20,35 +18,39 @@ public class CleanupAttachmentFiles { private final QueryEmail queryEmail; private final ManageEmail manageEmail; private final List callbacks; + private final int batchSize; private long lastInfoLogMillis = System.currentTimeMillis(); private long silentRuns = 0; private boolean firstRun = true; + public CleanupAttachmentFiles( + QueryEmail queryEmail, + ManageEmail manageEmail, + List callbacks, + int batchSize + ) { + this.queryEmail = queryEmail; + this.manageEmail = manageEmail; + this.callbacks = callbacks; + this.batchSize = batchSize; + } + @Scheduled(initialDelayString = "${aboutbits.emailservice.scheduling.interval:30000}", fixedDelayString = "${aboutbits.emailservice.scheduling.interval:30000}") - void cleanupAttachments() { + @Transactional + void claimAndCleanupAttachments() { logStartOfPass(); - var emailsToCleanup = queryEmail.readyToCleanup(); - - var countCleaned = 0; - var countError = 0; - for (var email : emailsToCleanup) { - try { - manageEmail.cleanupAttachments(email); - countCleaned++; - } catch (AttachmentException e) { - countError++; - } - } + var ids = queryEmail.claimReadyToCleanupIds(batchSize); + var outcome = manageEmail.cleanupBatch(ids); - logEndOfPass(countCleaned, countError); + logEndOfPass(outcome.cleaned(), outcome.errors()); for (var callback : callbacks) { callback.report(new AttachmentCleanerCallback.Report( - emailsToCleanup.size(), - countCleaned, - countError + outcome.total(), + outcome.cleaned(), + outcome.errors() )); } } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java index 3346aac..de43b83 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java @@ -18,6 +18,7 @@ import org.springframework.mail.MailException; import org.springframework.mail.javamail.JavaMailSender; import org.springframework.mail.javamail.MimeMessageHelper; +import org.springframework.transaction.annotation.Transactional; import org.springframework.validation.annotation.Validated; import java.io.IOException; @@ -31,6 +32,18 @@ @Slf4j @NullMarked public class ManageEmail { + record SendBatchOutcome(int sent, int errors) { + int total() { + return sent + errors; + } + } + + record CleanupBatchOutcome(int cleaned, int errors) { + int total() { + return cleaned + errors; + } + } + private final EmailRepository emailRepository; private final JavaMailSender mailSender; private final AttachmentDataSource attachmentDataSource; @@ -40,7 +53,7 @@ public ManageEmail( EmailRepository emailRepository, JavaMailSender mailSender, AttachmentDataSource attachmentDataSource, - final EmailMapper emailMapper + EmailMapper emailMapper ) { this.emailRepository = emailRepository; this.mailSender = mailSender; @@ -79,6 +92,54 @@ public EmailDto sendOrFail(@Valid EmailParameter parameter) throws EmailExceptio return emailMapper.toDto(savedEmail); } + // Sends the emails identified by the given ids and persists SENT/ERROR state for each + @Transactional + SendBatchOutcome sendBatch(List ids) { + if (ids.isEmpty()) { + return new SendBatchOutcome(0, 0); + } + var emails = emailRepository.findByIdIn(ids); + var sent = 0; + var errors = 0; + for (var email : emails) { + try { + var updated = send(email); + if (updated.hasFailed()) { + errors++; + } else { + sent++; + } + } catch (RuntimeException e) { + // A single misbehaving email should not roll back the whole batch's committed state + // (which would risk duplicate delivery for siblings that already left the SMTP relay). + log.error("Unexpected failure while sending email: {}", email.getId(), e); + errors++; + } + } + return new SendBatchOutcome(sent, errors); + } + + // Releases attachment payloads for the given emails and marks them cleaned + @Transactional + CleanupBatchOutcome cleanupBatch(List ids) { + if (ids.isEmpty()) { + return new CleanupBatchOutcome(0, 0); + } + var emails = emailRepository.findByIdIn(ids); + var cleaned = 0; + var errors = 0; + for (var email : emails) { + try { + cleanupAttachments(email); + cleaned++; + } catch (AttachmentException | RuntimeException e) { + log.warn("Failed to cleanup attachments for email: {}", email.getId(), e); + errors++; + } + } + return new CleanupBatchOutcome(cleaned, errors); + } + Email send(Email email) { if (EmailState.SENT.equals(email.getState())) { return email; diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java index 7c43588..268e279 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java @@ -4,8 +4,8 @@ import it.aboutbits.springboot.emailservice.lib.EmailDto; import it.aboutbits.springboot.emailservice.lib.EmailState; import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; -import it.aboutbits.springboot.emailservice.lib.model.Email; import jakarta.persistence.EntityManager; +import jakarta.persistence.LockModeType; import lombok.RequiredArgsConstructor; import org.jspecify.annotations.NullMarked; import org.springframework.data.domain.Page; @@ -21,6 +21,10 @@ @RequiredArgsConstructor @NullMarked public class QueryEmail { + // Hibernate convention for "jakarta.persistence.lock.timeout": + // -2 translates to "SKIP LOCKED" at the database layer (matches "org.hibernate.Timeouts.SKIP_LOCKED_MILLI") + private static final int SKIP_LOCKED_TIMEOUT = -2; + private static final String LOCK_TIMEOUT_HINT = "jakarta.persistence.lock.timeout"; private final EmailRepository emailRepository; private final EmailMapper emailMapper; private final EntityManager entityManager; @@ -42,29 +46,33 @@ public List byIds(Collection ids) { return emailMapper.toDto(emailRepository.findByIdIn(ids)); } - List readyToSend() { - var entityGraph = entityManager.getEntityGraph("email_service_emails-entity-graph"); + List claimReadyToSendIds(int limit) { return entityManager.createQuery( """ - SELECT e from Email e WHERE e.scheduledAt < :scheduledBefore AND e.state IN ( + select e.id from Email e where e.scheduledAt < :scheduledBefore and e.state in ( it.aboutbits.springboot.emailservice.lib.EmailState.PENDING, it.aboutbits.springboot.emailservice.lib.EmailState.ERROR ) - """, Email.class + order by e.scheduledAt + """, Long.class ) .setParameter("scheduledBefore", OffsetDateTime.now()) - .setHint("jakarta.persistence.fetchgraph", entityGraph) + .setLockMode(LockModeType.PESSIMISTIC_WRITE) + .setHint(LOCK_TIMEOUT_HINT, SKIP_LOCKED_TIMEOUT) + .setMaxResults(limit) .getResultList(); } - List readyToCleanup() { - var entityGraph = entityManager.getEntityGraph("email_service_emails-entity-graph"); + List claimReadyToCleanupIds(int limit) { return entityManager.createQuery( """ - SELECT e from Email e WHERE e.attachmentsCleaned=false AND e.state=it.aboutbits.springboot.emailservice.lib.EmailState.SENT - """, Email.class + select e.id from Email e where e.attachmentsCleaned=false and e.state=it.aboutbits.springboot.emailservice.lib.EmailState.SENT + order by e.updatedAt + """, Long.class ) - .setHint("jakarta.persistence.fetchgraph", entityGraph) + .setLockMode(LockModeType.PESSIMISTIC_WRITE) + .setHint(LOCK_TIMEOUT_HINT, SKIP_LOCKED_TIMEOUT) + .setMaxResults(limit) .getResultList(); } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java index f1789d0..988a5de 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java @@ -2,15 +2,14 @@ import it.aboutbits.springboot.emailservice.lib.EmailSchedulerCallback; -import lombok.RequiredArgsConstructor; import lombok.extern.log4j.Log4j2; import org.jspecify.annotations.NullMarked; import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.transaction.annotation.Transactional; import java.time.Duration; import java.util.List; -@RequiredArgsConstructor @Log4j2 @NullMarked public class SendScheduledEmails { @@ -19,38 +18,39 @@ public class SendScheduledEmails { private final QueryEmail queryEmail; private final ManageEmail manageEmail; private final List callbacks; + private final int batchSize; private long lastInfoLogMillis = System.currentTimeMillis(); private long silentRuns = 0; private boolean firstRun = true; + public SendScheduledEmails( + QueryEmail queryEmail, + ManageEmail manageEmail, + List callbacks, + int batchSize + ) { + this.queryEmail = queryEmail; + this.manageEmail = manageEmail; + this.callbacks = callbacks; + this.batchSize = batchSize; + } + @Scheduled(initialDelayString = "${aboutbits.emailservice.scheduling.interval:30000}", fixedDelayString = "${aboutbits.emailservice.scheduling.interval:30000}") - void sendEmails() { + @Transactional + void claimAndSendEmails() { logStartOfPass(); - var emailsToSend = queryEmail.readyToSend(); - - var countSent = 0; - var countError = 0; - for (var email : emailsToSend) { - var updatedEmail = manageEmail.send(email); - switch (updatedEmail.getState()) { - case ERROR -> countError++; - case SENT -> countSent++; - default -> log.warn( - JOB_DESCRIPTION + " | Job produced an invalid notification result state: {}.", - updatedEmail.getState().name() - ); - } - } + var ids = queryEmail.claimReadyToSendIds(batchSize); + var outcome = manageEmail.sendBatch(ids); - logEndOfPass(countSent, countError); + logEndOfPass(outcome.sent(), outcome.errors()); for (var callback : callbacks) { callback.report(new EmailSchedulerCallback.Report( - emailsToSend.size(), - countSent, - countError + outcome.total(), + outcome.sent(), + outcome.errors() )); } } diff --git a/src/main/resources/META-INF/additional-spring-configuration-metadata.json b/src/main/resources/META-INF/additional-spring-configuration-metadata.json index 6653b27..3e5a51b 100644 --- a/src/main/resources/META-INF/additional-spring-configuration-metadata.json +++ b/src/main/resources/META-INF/additional-spring-configuration-metadata.json @@ -23,6 +23,12 @@ "type": "java.lang.Long", "description": "Specifies the milliseconds delay between runs of the scheduler.", "defaultValue": 30000 + }, + { + "name": "aboutbits.emailservice.scheduling.batch-size", + "type": "java.lang.Integer", + "description": "Maximum number of emails one pod claims per scheduler pass.", + "defaultValue": 50 } ] } diff --git a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFilesTest.java b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFilesTest.java new file mode 100644 index 0000000..d30e4e1 --- /dev/null +++ b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFilesTest.java @@ -0,0 +1,142 @@ +package it.aboutbits.springboot.emailservice.lib.application; + +import it.aboutbits.springboot.emailservice.lib.AttachmentDataSource; +import it.aboutbits.springboot.emailservice.lib.EmailState; +import it.aboutbits.springboot.emailservice.lib.exception.AttachmentException; +import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; +import it.aboutbits.springboot.emailservice.lib.model.Email; +import it.aboutbits.springboot.emailservice.lib.model.EmailAttachment; +import it.aboutbits.springboot.emailservice.support.database.WithPostgres; +import it.aboutbits.springboot.emailservice.support.database.factory.EmailFactory; +import org.jspecify.annotations.NullMarked; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.test.context.bean.override.mockito.MockitoBean; +import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; +import org.springframework.mail.javamail.JavaMailSender; + +import java.util.HashSet; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.doAnswer; + +@SpringBootTest(properties = "aboutbits.emailservice.scheduling.batch-size=5") +@WithPostgres +@NullMarked +class CleanupAttachmentFilesTest { + private static final int BATCH_SIZE = 5; + + @MockitoBean + AttachmentDataSource attachmentDataSource; + + @MockitoSpyBean + JavaMailSender javaMailSender; + + @Autowired + EmailRepository emailRepository; + + @Autowired + CleanupAttachmentFiles cleanupAttachmentFiles; + + private ConcurrentHashMap releaseCallsByFileReference = new ConcurrentHashMap<>(); + + @BeforeEach + void setup() throws AttachmentException { + releaseCallsByFileReference = new ConcurrentHashMap<>(); + doAnswer(inv -> { + releaseCallsByFileReference.computeIfAbsent(inv.getArgument(0), k -> new AtomicInteger()).incrementAndGet(); + return null; + }).when(attachmentDataSource).releaseAttachment(anyLong()); + } + + @Test + void givenSentEmailsWithAttachments_claimAndCleanupAttachments_shouldReleasePayloadsAndMarkCleaned() { + seedSentEmailsWithAttachments(3); + + cleanupAttachmentFiles.claimAndCleanupAttachments(); + + assertThat(emailRepository.findAll()) + .hasSize(3) + .allMatch(Email::isAttachmentsCleaned); + assertThat(releaseCallsByFileReference).hasSize(3); + assertThat(releaseCallsByFileReference.values()) + .allSatisfy(count -> assertThat(count.get()).isEqualTo(1)); + } + + @Test + void givenMoreEmailsThanBatchSize_claimAndCleanupAttachments_shouldCleanOnlyBatchSizePerPass() { + var pending = 3 * BATCH_SIZE; + seedSentEmailsWithAttachments(pending); + + cleanupAttachmentFiles.claimAndCleanupAttachments(); + + var cleaned = emailRepository.findAll().stream() + .filter(Email::isAttachmentsCleaned) + .count(); + assertThat(cleaned).isEqualTo(BATCH_SIZE); + assertThat(releaseCallsByFileReference).hasSize(BATCH_SIZE); + } + + @Test + void givenConcurrentPasses_claimAndCleanupAttachments_shouldReleaseEachAttachmentExactlyOnce() throws Exception { + var totalEmails = 25; + var workerThreads = 4; + var passesPerWorker = 10; + seedSentEmailsWithAttachments(totalEmails); + + var startGate = new CountDownLatch(1); + var executor = Executors.newFixedThreadPool(workerThreads); + try { + for (var t = 0; t < workerThreads; t++) { + executor.submit(() -> { + startGate.await(); + for (var i = 0; i < passesPerWorker; i++) { + cleanupAttachmentFiles.claimAndCleanupAttachments(); + } + return null; + }); + } + startGate.countDown(); + executor.shutdown(); + assertThat(executor.awaitTermination(60, TimeUnit.SECONDS)).isTrue(); + } finally { + if (!executor.isTerminated()) { + executor.shutdownNow(); + } + } + + assertThat(emailRepository.findAll()) + .hasSize(totalEmails) + .allMatch(Email::isAttachmentsCleaned); + assertThat(releaseCallsByFileReference).hasSize(totalEmails); + assertThat(releaseCallsByFileReference.values()) + .allSatisfy(count -> assertThat(count.get()).isEqualTo(1)); + } + + private void seedSentEmailsWithAttachments(int count) { + var fileReferenceCounter = new AtomicInteger(1000); + var emails = EmailFactory.many(count).stream() + .map(builder -> { + var email = builder.state(EmailState.SENT).attachmentsCleaned(false).build(); + var attachment = new EmailAttachment(); + attachment.setEmail(email); + attachment.setFileName("payload.bin"); + attachment.setContentType("application/octet-stream"); + attachment.setFileReference(fileReferenceCounter.getAndIncrement()); + var attachments = new HashSet(); + attachments.add(attachment); + email.setAttachments(attachments); + return email; + }) + .toList(); + emailRepository.saveAll(emails); + } +} diff --git a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java new file mode 100644 index 0000000..02128e4 --- /dev/null +++ b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java @@ -0,0 +1,138 @@ +package it.aboutbits.springboot.emailservice.lib.application; + +import it.aboutbits.springboot.emailservice.lib.EmailState; +import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; +import it.aboutbits.springboot.emailservice.lib.model.Email; +import it.aboutbits.springboot.emailservice.support.database.WithPostgres; +import it.aboutbits.springboot.emailservice.support.database.factory.EmailFactory; +import jakarta.mail.internet.MimeMessage; +import org.jspecify.annotations.NullMarked; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.mail.MailSendException; +import org.springframework.mail.javamail.JavaMailSender; +import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +@SpringBootTest(properties = "aboutbits.emailservice.scheduling.batch-size=5") +@WithPostgres +@NullMarked +class SendScheduledEmailsTest { + private static final int BATCH_SIZE = 5; + + @MockitoSpyBean + JavaMailSender javaMailSender; + + @Autowired + EmailRepository emailRepository; + + @Autowired + SendScheduledEmails sendScheduledEmails; + + @BeforeEach + void setup() { + doNothing().when(javaMailSender).send(any(MimeMessage.class)); + } + + @Test + void givenPendingEmails_claimAndSendEmails_shouldMarkAllAsSent() { + emailRepository.saveAll(EmailFactory.many(3).stream().map(Email.EmailBuilder::build).toList()); + + sendScheduledEmails.claimAndSendEmails(); + + assertThat(emailRepository.findAll()) + .hasSize(3) + .allMatch(email -> email.getState() == EmailState.SENT); + verify(javaMailSender, times(3)).send(any(MimeMessage.class)); + } + + @Test + void givenMoreEmailsThanBatchSize_claimAndSendEmails_shouldSendOnlyBatchSizePerPass() { + var pending = 3 * BATCH_SIZE; + emailRepository.saveAll(EmailFactory.many(pending).stream().map(Email.EmailBuilder::build).toList()); + + sendScheduledEmails.claimAndSendEmails(); + + assertThat(emailRepository.findAll()) + .filteredOn(email -> email.getState() == EmailState.SENT) + .hasSize(BATCH_SIZE); + verify(javaMailSender, times(BATCH_SIZE)).send(any(MimeMessage.class)); + } + + @Test + void givenPreviouslyFailedEmail_claimAndSendEmails_shouldRetryAndMarkSent() { + var errored = EmailFactory.once().state(EmailState.ERROR).errorMessage("some SMTP hiccup").build(); + emailRepository.save(errored); + + sendScheduledEmails.claimAndSendEmails(); + + assertThat(emailRepository.findAll()) + .singleElement() + .satisfies(email -> { + assertThat(email.getState()).isEqualTo(EmailState.SENT); + assertThat(email.getSentAt()).isNotNull(); + assertThat(email.getErrorMessage()).isEmpty(); + }); + } + + @Test + void givenMailSenderThrows_claimAndSendEmails_shouldMarkEmailAsErrorAndPersistMessage() { + doThrow(new MailSendException("smtp down")).when(javaMailSender).send(any(MimeMessage.class)); + emailRepository.save(EmailFactory.once().build()); + + sendScheduledEmails.claimAndSendEmails(); + + assertThat(emailRepository.findAll()) + .singleElement() + .satisfies(email -> { + assertThat(email.getState()).isEqualTo(EmailState.ERROR); + assertThat(email.getErrorMessage()).contains("smtp down"); + }); + } + + @Test + void givenConcurrentPasses_claimAndSendEmails_shouldSendEachEmailExactlyOnce() throws Exception { + var totalEmails = 25; + var workerThreads = 4; + var passesPerWorker = 10; + emailRepository.saveAll(EmailFactory.many(totalEmails).stream().map(Email.EmailBuilder::build).toList()); + + var startGate = new CountDownLatch(1); + var executor = Executors.newFixedThreadPool(workerThreads); + try { + for (var t = 0; t < workerThreads; t++) { + executor.submit(() -> { + startGate.await(); + for (var i = 0; i < passesPerWorker; i++) { + sendScheduledEmails.claimAndSendEmails(); + } + return null; + }); + } + startGate.countDown(); + executor.shutdown(); + assertThat(executor.awaitTermination(60, TimeUnit.SECONDS)).isTrue(); + } finally { + if (!executor.isTerminated()) { + executor.shutdownNow(); + } + } + + assertThat(emailRepository.findAll()) + .hasSize(totalEmails) + .allMatch(email -> email.getState() == EmailState.SENT); + verify(javaMailSender, times(totalEmails)).send(any(MimeMessage.class)); + } +} diff --git a/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java b/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java index 1c37b19..f5f5992 100644 --- a/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java +++ b/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java @@ -6,6 +6,7 @@ import org.jspecify.annotations.NullMarked; import java.time.OffsetDateTime; +import java.util.ArrayList; import java.util.List; @NullMarked @@ -30,4 +31,25 @@ public static Email.EmailBuilder once() { .replyToName(FAKER.name().fullName()) .recipients(List.of(FAKER.internet().emailAddress(), FAKER.internet().emailAddress())); } + + public static List many(int number) { + var result = new ArrayList(); + for (int i = 0; i < number; i++) { + var body = FAKER.lorem().paragraph(); + result.add( + Email.builder() + .state(EmailState.PENDING) + .subject("Email subject") + .textBody(body) + .htmlBody("

" + body + "

") + .scheduledAt(OffsetDateTime.now()) + .fromAddress(FAKER.internet().emailAddress()) + .fromName(FAKER.name().fullName()) + .replyToAddress(FAKER.internet().emailAddress()) + .replyToName(FAKER.name().fullName()) + .recipients(List.of(FAKER.internet().emailAddress(), FAKER.internet().emailAddress())) + ); + } + return result; + } } From cc401e8df78afccc1b9caa883be334baba3b1e16 Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Fri, 14 Aug 2026 10:12:19 +0200 Subject: [PATCH 02/13] Revert "first draft for concurrency using an approach at the database level" This reverts commit 15c22d4bf65266c997963a8f127245322bf51dbb. --- readme.md | 24 +-- .../EmailServiceConfiguration.java | 32 +--- .../application/CleanupAttachmentFiles.java | 42 +++--- .../lib/application/ManageEmail.java | 63 +------- .../lib/application/QueryEmail.java | 30 ++-- .../lib/application/SendScheduledEmails.java | 44 +++--- ...itional-spring-configuration-metadata.json | 6 - .../CleanupAttachmentFilesTest.java | 142 ------------------ .../application/SendScheduledEmailsTest.java | 138 ----------------- .../database/factory/EmailFactory.java | 22 --- 10 files changed, 65 insertions(+), 478 deletions(-) delete mode 100644 src/test/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFilesTest.java delete mode 100644 src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java diff --git a/readme.md b/readme.md index 2dffdd3..19c5434 100644 --- a/readme.md +++ b/readme.md @@ -74,25 +74,11 @@ public class App { The following configuration options are available: -| Name | Default | Description | -|-------------------------------------------------|---------|----------------------------------------------------------------------------------------------------------------| -| `aboutbits.emailservice.migrations.enabled` | true | Enables database migrations. | -| `aboutbits.emailservice.scheduling.enabled` | true | Enables the scheduler sending the emails. | -| `aboutbits.emailservice.scheduling.cleanup.enabled` | true | Enables cleanup of attachment files after sending. | -| `aboutbits.emailservice.scheduling.interval` | 30000 | Specifies the milliseconds delay between runs of the scheduler. | -| `aboutbits.emailservice.scheduling.batch-size` | 50 | Maximum number of emails a single pod claims per scheduler pass (see [Multi-pod deployments](#multi-pod-deployments)). | - -## Multi-pod deployments - -The scheduler is safe to run in every pod concurrently. Each ready email is -claimed by exactly one pod using Postgres row-level `SELECT ... FOR UPDATE SKIP LOCKED`, -so pods work on disjoint rows in parallel and no email is ever sent twice by -different pods. The same guarantee applies to the attachment cleanup scheduler. - -`aboutbits.emailservice.scheduling.batch-size` caps the number of emails one -pod processes per pass. With the default 30s interval and 50 emails per pass, -a single pod can take up to 100 emails per minute; For higher-throughput deployments increase the batch size or -lower the interval. +| Name | Default | Description | +|----------------------------------------|-------------|-----------------------------------------------------------------------| +| `lib.emailservice.migrations.enabled` | true | Enables database migrations. | +| `lib.emailservice.scheduling.enabled` | true | Enables the scheduler sending the emails. | +| `lib.emailservice.scheduling.interval` | 30000 | Specifies the milliseconds delay between runs of the scheduler. | ## Local development: diff --git a/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java b/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java index 27ed84e..535054e 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java @@ -16,7 +16,6 @@ import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; import jakarta.persistence.EntityManager; import org.jspecify.annotations.NullMarked; -import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.AutoConfigurationPackage; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -46,44 +45,25 @@ public EmailMapper emailMapper() { } @Bean - public QueryEmail queryEmail( - EmailRepository emailRepository, - EmailMapper emailMapper, - EntityManager entityManager - ) { + public QueryEmail queryEmail(EmailRepository emailRepository, EmailMapper emailMapper, EntityManager entityManager) { return new QueryEmail(emailRepository, emailMapper, entityManager); } @Bean - public ManageEmail manageEmail( - EmailRepository emailRepository, - JavaMailSender javaMailSender, - AttachmentDataSource attachmentDataSource, - EmailMapper emailMapper - ) { + public ManageEmail manageEmail(EmailRepository emailRepository, JavaMailSender javaMailSender, AttachmentDataSource attachmentDataSource, EmailMapper emailMapper) { return new ManageEmail(emailRepository, javaMailSender, attachmentDataSource, emailMapper); } @Bean @ConditionalOnProperty(value = "aboutbits.emailservice.scheduling.enabled", matchIfMissing = true) - public SendScheduledEmails sendScheduledEmails( - QueryEmail queryEmail, - ManageEmail manageEmail, - List callbacks, - @Value("${aboutbits.emailservice.scheduling.batch-size:50}") int batchSize - ) { - return new SendScheduledEmails(queryEmail, manageEmail, callbacks, batchSize); + public SendScheduledEmails sendScheduledEmails(QueryEmail queryEmail, ManageEmail manageEmail, List callbacks) { + return new SendScheduledEmails(queryEmail, manageEmail, callbacks); } @Bean @ConditionalOnProperty(value = "aboutbits.emailservice.scheduling.cleanup.enabled", matchIfMissing = true) - public CleanupAttachmentFiles cleanupAttachments( - QueryEmail queryEmail, - ManageEmail manageEmail, - List callbacks, - @Value("${aboutbits.emailservice.scheduling.batch-size:50}") int batchSize - ) { - return new CleanupAttachmentFiles(queryEmail, manageEmail, callbacks, batchSize); + public CleanupAttachmentFiles cleanupAttachments(QueryEmail queryEmail, ManageEmail manageEmail, List callbacks) { + return new CleanupAttachmentFiles(queryEmail, manageEmail, callbacks); } @Bean diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFiles.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFiles.java index 057e4e9..df7209d 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFiles.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFiles.java @@ -2,14 +2,16 @@ import it.aboutbits.springboot.emailservice.lib.AttachmentCleanerCallback; +import it.aboutbits.springboot.emailservice.lib.exception.AttachmentException; +import lombok.RequiredArgsConstructor; import lombok.extern.log4j.Log4j2; import org.jspecify.annotations.NullMarked; import org.springframework.scheduling.annotation.Scheduled; -import org.springframework.transaction.annotation.Transactional; import java.time.Duration; import java.util.List; +@RequiredArgsConstructor @Log4j2 @NullMarked public class CleanupAttachmentFiles { @@ -18,39 +20,35 @@ public class CleanupAttachmentFiles { private final QueryEmail queryEmail; private final ManageEmail manageEmail; private final List callbacks; - private final int batchSize; private long lastInfoLogMillis = System.currentTimeMillis(); private long silentRuns = 0; private boolean firstRun = true; - public CleanupAttachmentFiles( - QueryEmail queryEmail, - ManageEmail manageEmail, - List callbacks, - int batchSize - ) { - this.queryEmail = queryEmail; - this.manageEmail = manageEmail; - this.callbacks = callbacks; - this.batchSize = batchSize; - } - @Scheduled(initialDelayString = "${aboutbits.emailservice.scheduling.interval:30000}", fixedDelayString = "${aboutbits.emailservice.scheduling.interval:30000}") - @Transactional - void claimAndCleanupAttachments() { + void cleanupAttachments() { logStartOfPass(); - var ids = queryEmail.claimReadyToCleanupIds(batchSize); - var outcome = manageEmail.cleanupBatch(ids); + var emailsToCleanup = queryEmail.readyToCleanup(); + + var countCleaned = 0; + var countError = 0; + for (var email : emailsToCleanup) { + try { + manageEmail.cleanupAttachments(email); + countCleaned++; + } catch (AttachmentException e) { + countError++; + } + } - logEndOfPass(outcome.cleaned(), outcome.errors()); + logEndOfPass(countCleaned, countError); for (var callback : callbacks) { callback.report(new AttachmentCleanerCallback.Report( - outcome.total(), - outcome.cleaned(), - outcome.errors() + emailsToCleanup.size(), + countCleaned, + countError )); } } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java index de43b83..3346aac 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java @@ -18,7 +18,6 @@ import org.springframework.mail.MailException; import org.springframework.mail.javamail.JavaMailSender; import org.springframework.mail.javamail.MimeMessageHelper; -import org.springframework.transaction.annotation.Transactional; import org.springframework.validation.annotation.Validated; import java.io.IOException; @@ -32,18 +31,6 @@ @Slf4j @NullMarked public class ManageEmail { - record SendBatchOutcome(int sent, int errors) { - int total() { - return sent + errors; - } - } - - record CleanupBatchOutcome(int cleaned, int errors) { - int total() { - return cleaned + errors; - } - } - private final EmailRepository emailRepository; private final JavaMailSender mailSender; private final AttachmentDataSource attachmentDataSource; @@ -53,7 +40,7 @@ public ManageEmail( EmailRepository emailRepository, JavaMailSender mailSender, AttachmentDataSource attachmentDataSource, - EmailMapper emailMapper + final EmailMapper emailMapper ) { this.emailRepository = emailRepository; this.mailSender = mailSender; @@ -92,54 +79,6 @@ public EmailDto sendOrFail(@Valid EmailParameter parameter) throws EmailExceptio return emailMapper.toDto(savedEmail); } - // Sends the emails identified by the given ids and persists SENT/ERROR state for each - @Transactional - SendBatchOutcome sendBatch(List ids) { - if (ids.isEmpty()) { - return new SendBatchOutcome(0, 0); - } - var emails = emailRepository.findByIdIn(ids); - var sent = 0; - var errors = 0; - for (var email : emails) { - try { - var updated = send(email); - if (updated.hasFailed()) { - errors++; - } else { - sent++; - } - } catch (RuntimeException e) { - // A single misbehaving email should not roll back the whole batch's committed state - // (which would risk duplicate delivery for siblings that already left the SMTP relay). - log.error("Unexpected failure while sending email: {}", email.getId(), e); - errors++; - } - } - return new SendBatchOutcome(sent, errors); - } - - // Releases attachment payloads for the given emails and marks them cleaned - @Transactional - CleanupBatchOutcome cleanupBatch(List ids) { - if (ids.isEmpty()) { - return new CleanupBatchOutcome(0, 0); - } - var emails = emailRepository.findByIdIn(ids); - var cleaned = 0; - var errors = 0; - for (var email : emails) { - try { - cleanupAttachments(email); - cleaned++; - } catch (AttachmentException | RuntimeException e) { - log.warn("Failed to cleanup attachments for email: {}", email.getId(), e); - errors++; - } - } - return new CleanupBatchOutcome(cleaned, errors); - } - Email send(Email email) { if (EmailState.SENT.equals(email.getState())) { return email; diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java index 268e279..7c43588 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java @@ -4,8 +4,8 @@ import it.aboutbits.springboot.emailservice.lib.EmailDto; import it.aboutbits.springboot.emailservice.lib.EmailState; import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; +import it.aboutbits.springboot.emailservice.lib.model.Email; import jakarta.persistence.EntityManager; -import jakarta.persistence.LockModeType; import lombok.RequiredArgsConstructor; import org.jspecify.annotations.NullMarked; import org.springframework.data.domain.Page; @@ -21,10 +21,6 @@ @RequiredArgsConstructor @NullMarked public class QueryEmail { - // Hibernate convention for "jakarta.persistence.lock.timeout": - // -2 translates to "SKIP LOCKED" at the database layer (matches "org.hibernate.Timeouts.SKIP_LOCKED_MILLI") - private static final int SKIP_LOCKED_TIMEOUT = -2; - private static final String LOCK_TIMEOUT_HINT = "jakarta.persistence.lock.timeout"; private final EmailRepository emailRepository; private final EmailMapper emailMapper; private final EntityManager entityManager; @@ -46,33 +42,29 @@ public List byIds(Collection ids) { return emailMapper.toDto(emailRepository.findByIdIn(ids)); } - List claimReadyToSendIds(int limit) { + List readyToSend() { + var entityGraph = entityManager.getEntityGraph("email_service_emails-entity-graph"); return entityManager.createQuery( """ - select e.id from Email e where e.scheduledAt < :scheduledBefore and e.state in ( + SELECT e from Email e WHERE e.scheduledAt < :scheduledBefore AND e.state IN ( it.aboutbits.springboot.emailservice.lib.EmailState.PENDING, it.aboutbits.springboot.emailservice.lib.EmailState.ERROR ) - order by e.scheduledAt - """, Long.class + """, Email.class ) .setParameter("scheduledBefore", OffsetDateTime.now()) - .setLockMode(LockModeType.PESSIMISTIC_WRITE) - .setHint(LOCK_TIMEOUT_HINT, SKIP_LOCKED_TIMEOUT) - .setMaxResults(limit) + .setHint("jakarta.persistence.fetchgraph", entityGraph) .getResultList(); } - List claimReadyToCleanupIds(int limit) { + List readyToCleanup() { + var entityGraph = entityManager.getEntityGraph("email_service_emails-entity-graph"); return entityManager.createQuery( """ - select e.id from Email e where e.attachmentsCleaned=false and e.state=it.aboutbits.springboot.emailservice.lib.EmailState.SENT - order by e.updatedAt - """, Long.class + SELECT e from Email e WHERE e.attachmentsCleaned=false AND e.state=it.aboutbits.springboot.emailservice.lib.EmailState.SENT + """, Email.class ) - .setLockMode(LockModeType.PESSIMISTIC_WRITE) - .setHint(LOCK_TIMEOUT_HINT, SKIP_LOCKED_TIMEOUT) - .setMaxResults(limit) + .setHint("jakarta.persistence.fetchgraph", entityGraph) .getResultList(); } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java index 988a5de..f1789d0 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java @@ -2,14 +2,15 @@ import it.aboutbits.springboot.emailservice.lib.EmailSchedulerCallback; +import lombok.RequiredArgsConstructor; import lombok.extern.log4j.Log4j2; import org.jspecify.annotations.NullMarked; import org.springframework.scheduling.annotation.Scheduled; -import org.springframework.transaction.annotation.Transactional; import java.time.Duration; import java.util.List; +@RequiredArgsConstructor @Log4j2 @NullMarked public class SendScheduledEmails { @@ -18,39 +19,38 @@ public class SendScheduledEmails { private final QueryEmail queryEmail; private final ManageEmail manageEmail; private final List callbacks; - private final int batchSize; private long lastInfoLogMillis = System.currentTimeMillis(); private long silentRuns = 0; private boolean firstRun = true; - public SendScheduledEmails( - QueryEmail queryEmail, - ManageEmail manageEmail, - List callbacks, - int batchSize - ) { - this.queryEmail = queryEmail; - this.manageEmail = manageEmail; - this.callbacks = callbacks; - this.batchSize = batchSize; - } - @Scheduled(initialDelayString = "${aboutbits.emailservice.scheduling.interval:30000}", fixedDelayString = "${aboutbits.emailservice.scheduling.interval:30000}") - @Transactional - void claimAndSendEmails() { + void sendEmails() { logStartOfPass(); - var ids = queryEmail.claimReadyToSendIds(batchSize); - var outcome = manageEmail.sendBatch(ids); + var emailsToSend = queryEmail.readyToSend(); + + var countSent = 0; + var countError = 0; + for (var email : emailsToSend) { + var updatedEmail = manageEmail.send(email); + switch (updatedEmail.getState()) { + case ERROR -> countError++; + case SENT -> countSent++; + default -> log.warn( + JOB_DESCRIPTION + " | Job produced an invalid notification result state: {}.", + updatedEmail.getState().name() + ); + } + } - logEndOfPass(outcome.sent(), outcome.errors()); + logEndOfPass(countSent, countError); for (var callback : callbacks) { callback.report(new EmailSchedulerCallback.Report( - outcome.total(), - outcome.sent(), - outcome.errors() + emailsToSend.size(), + countSent, + countError )); } } diff --git a/src/main/resources/META-INF/additional-spring-configuration-metadata.json b/src/main/resources/META-INF/additional-spring-configuration-metadata.json index 3e5a51b..6653b27 100644 --- a/src/main/resources/META-INF/additional-spring-configuration-metadata.json +++ b/src/main/resources/META-INF/additional-spring-configuration-metadata.json @@ -23,12 +23,6 @@ "type": "java.lang.Long", "description": "Specifies the milliseconds delay between runs of the scheduler.", "defaultValue": 30000 - }, - { - "name": "aboutbits.emailservice.scheduling.batch-size", - "type": "java.lang.Integer", - "description": "Maximum number of emails one pod claims per scheduler pass.", - "defaultValue": 50 } ] } diff --git a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFilesTest.java b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFilesTest.java deleted file mode 100644 index d30e4e1..0000000 --- a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/CleanupAttachmentFilesTest.java +++ /dev/null @@ -1,142 +0,0 @@ -package it.aboutbits.springboot.emailservice.lib.application; - -import it.aboutbits.springboot.emailservice.lib.AttachmentDataSource; -import it.aboutbits.springboot.emailservice.lib.EmailState; -import it.aboutbits.springboot.emailservice.lib.exception.AttachmentException; -import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; -import it.aboutbits.springboot.emailservice.lib.model.Email; -import it.aboutbits.springboot.emailservice.lib.model.EmailAttachment; -import it.aboutbits.springboot.emailservice.support.database.WithPostgres; -import it.aboutbits.springboot.emailservice.support.database.factory.EmailFactory; -import org.jspecify.annotations.NullMarked; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Test; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.test.context.bean.override.mockito.MockitoBean; -import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; -import org.springframework.mail.javamail.JavaMailSender; - -import java.util.HashSet; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicInteger; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.ArgumentMatchers.anyLong; -import static org.mockito.Mockito.doAnswer; - -@SpringBootTest(properties = "aboutbits.emailservice.scheduling.batch-size=5") -@WithPostgres -@NullMarked -class CleanupAttachmentFilesTest { - private static final int BATCH_SIZE = 5; - - @MockitoBean - AttachmentDataSource attachmentDataSource; - - @MockitoSpyBean - JavaMailSender javaMailSender; - - @Autowired - EmailRepository emailRepository; - - @Autowired - CleanupAttachmentFiles cleanupAttachmentFiles; - - private ConcurrentHashMap releaseCallsByFileReference = new ConcurrentHashMap<>(); - - @BeforeEach - void setup() throws AttachmentException { - releaseCallsByFileReference = new ConcurrentHashMap<>(); - doAnswer(inv -> { - releaseCallsByFileReference.computeIfAbsent(inv.getArgument(0), k -> new AtomicInteger()).incrementAndGet(); - return null; - }).when(attachmentDataSource).releaseAttachment(anyLong()); - } - - @Test - void givenSentEmailsWithAttachments_claimAndCleanupAttachments_shouldReleasePayloadsAndMarkCleaned() { - seedSentEmailsWithAttachments(3); - - cleanupAttachmentFiles.claimAndCleanupAttachments(); - - assertThat(emailRepository.findAll()) - .hasSize(3) - .allMatch(Email::isAttachmentsCleaned); - assertThat(releaseCallsByFileReference).hasSize(3); - assertThat(releaseCallsByFileReference.values()) - .allSatisfy(count -> assertThat(count.get()).isEqualTo(1)); - } - - @Test - void givenMoreEmailsThanBatchSize_claimAndCleanupAttachments_shouldCleanOnlyBatchSizePerPass() { - var pending = 3 * BATCH_SIZE; - seedSentEmailsWithAttachments(pending); - - cleanupAttachmentFiles.claimAndCleanupAttachments(); - - var cleaned = emailRepository.findAll().stream() - .filter(Email::isAttachmentsCleaned) - .count(); - assertThat(cleaned).isEqualTo(BATCH_SIZE); - assertThat(releaseCallsByFileReference).hasSize(BATCH_SIZE); - } - - @Test - void givenConcurrentPasses_claimAndCleanupAttachments_shouldReleaseEachAttachmentExactlyOnce() throws Exception { - var totalEmails = 25; - var workerThreads = 4; - var passesPerWorker = 10; - seedSentEmailsWithAttachments(totalEmails); - - var startGate = new CountDownLatch(1); - var executor = Executors.newFixedThreadPool(workerThreads); - try { - for (var t = 0; t < workerThreads; t++) { - executor.submit(() -> { - startGate.await(); - for (var i = 0; i < passesPerWorker; i++) { - cleanupAttachmentFiles.claimAndCleanupAttachments(); - } - return null; - }); - } - startGate.countDown(); - executor.shutdown(); - assertThat(executor.awaitTermination(60, TimeUnit.SECONDS)).isTrue(); - } finally { - if (!executor.isTerminated()) { - executor.shutdownNow(); - } - } - - assertThat(emailRepository.findAll()) - .hasSize(totalEmails) - .allMatch(Email::isAttachmentsCleaned); - assertThat(releaseCallsByFileReference).hasSize(totalEmails); - assertThat(releaseCallsByFileReference.values()) - .allSatisfy(count -> assertThat(count.get()).isEqualTo(1)); - } - - private void seedSentEmailsWithAttachments(int count) { - var fileReferenceCounter = new AtomicInteger(1000); - var emails = EmailFactory.many(count).stream() - .map(builder -> { - var email = builder.state(EmailState.SENT).attachmentsCleaned(false).build(); - var attachment = new EmailAttachment(); - attachment.setEmail(email); - attachment.setFileName("payload.bin"); - attachment.setContentType("application/octet-stream"); - attachment.setFileReference(fileReferenceCounter.getAndIncrement()); - var attachments = new HashSet(); - attachments.add(attachment); - email.setAttachments(attachments); - return email; - }) - .toList(); - emailRepository.saveAll(emails); - } -} diff --git a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java deleted file mode 100644 index 02128e4..0000000 --- a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java +++ /dev/null @@ -1,138 +0,0 @@ -package it.aboutbits.springboot.emailservice.lib.application; - -import it.aboutbits.springboot.emailservice.lib.EmailState; -import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; -import it.aboutbits.springboot.emailservice.lib.model.Email; -import it.aboutbits.springboot.emailservice.support.database.WithPostgres; -import it.aboutbits.springboot.emailservice.support.database.factory.EmailFactory; -import jakarta.mail.internet.MimeMessage; -import org.jspecify.annotations.NullMarked; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Test; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.mail.MailSendException; -import org.springframework.mail.javamail.JavaMailSender; -import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; - -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.Mockito.doNothing; -import static org.mockito.Mockito.doThrow; -import static org.mockito.Mockito.times; -import static org.mockito.Mockito.verify; - -@SpringBootTest(properties = "aboutbits.emailservice.scheduling.batch-size=5") -@WithPostgres -@NullMarked -class SendScheduledEmailsTest { - private static final int BATCH_SIZE = 5; - - @MockitoSpyBean - JavaMailSender javaMailSender; - - @Autowired - EmailRepository emailRepository; - - @Autowired - SendScheduledEmails sendScheduledEmails; - - @BeforeEach - void setup() { - doNothing().when(javaMailSender).send(any(MimeMessage.class)); - } - - @Test - void givenPendingEmails_claimAndSendEmails_shouldMarkAllAsSent() { - emailRepository.saveAll(EmailFactory.many(3).stream().map(Email.EmailBuilder::build).toList()); - - sendScheduledEmails.claimAndSendEmails(); - - assertThat(emailRepository.findAll()) - .hasSize(3) - .allMatch(email -> email.getState() == EmailState.SENT); - verify(javaMailSender, times(3)).send(any(MimeMessage.class)); - } - - @Test - void givenMoreEmailsThanBatchSize_claimAndSendEmails_shouldSendOnlyBatchSizePerPass() { - var pending = 3 * BATCH_SIZE; - emailRepository.saveAll(EmailFactory.many(pending).stream().map(Email.EmailBuilder::build).toList()); - - sendScheduledEmails.claimAndSendEmails(); - - assertThat(emailRepository.findAll()) - .filteredOn(email -> email.getState() == EmailState.SENT) - .hasSize(BATCH_SIZE); - verify(javaMailSender, times(BATCH_SIZE)).send(any(MimeMessage.class)); - } - - @Test - void givenPreviouslyFailedEmail_claimAndSendEmails_shouldRetryAndMarkSent() { - var errored = EmailFactory.once().state(EmailState.ERROR).errorMessage("some SMTP hiccup").build(); - emailRepository.save(errored); - - sendScheduledEmails.claimAndSendEmails(); - - assertThat(emailRepository.findAll()) - .singleElement() - .satisfies(email -> { - assertThat(email.getState()).isEqualTo(EmailState.SENT); - assertThat(email.getSentAt()).isNotNull(); - assertThat(email.getErrorMessage()).isEmpty(); - }); - } - - @Test - void givenMailSenderThrows_claimAndSendEmails_shouldMarkEmailAsErrorAndPersistMessage() { - doThrow(new MailSendException("smtp down")).when(javaMailSender).send(any(MimeMessage.class)); - emailRepository.save(EmailFactory.once().build()); - - sendScheduledEmails.claimAndSendEmails(); - - assertThat(emailRepository.findAll()) - .singleElement() - .satisfies(email -> { - assertThat(email.getState()).isEqualTo(EmailState.ERROR); - assertThat(email.getErrorMessage()).contains("smtp down"); - }); - } - - @Test - void givenConcurrentPasses_claimAndSendEmails_shouldSendEachEmailExactlyOnce() throws Exception { - var totalEmails = 25; - var workerThreads = 4; - var passesPerWorker = 10; - emailRepository.saveAll(EmailFactory.many(totalEmails).stream().map(Email.EmailBuilder::build).toList()); - - var startGate = new CountDownLatch(1); - var executor = Executors.newFixedThreadPool(workerThreads); - try { - for (var t = 0; t < workerThreads; t++) { - executor.submit(() -> { - startGate.await(); - for (var i = 0; i < passesPerWorker; i++) { - sendScheduledEmails.claimAndSendEmails(); - } - return null; - }); - } - startGate.countDown(); - executor.shutdown(); - assertThat(executor.awaitTermination(60, TimeUnit.SECONDS)).isTrue(); - } finally { - if (!executor.isTerminated()) { - executor.shutdownNow(); - } - } - - assertThat(emailRepository.findAll()) - .hasSize(totalEmails) - .allMatch(email -> email.getState() == EmailState.SENT); - verify(javaMailSender, times(totalEmails)).send(any(MimeMessage.class)); - } -} diff --git a/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java b/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java index f5f5992..1c37b19 100644 --- a/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java +++ b/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java @@ -6,7 +6,6 @@ import org.jspecify.annotations.NullMarked; import java.time.OffsetDateTime; -import java.util.ArrayList; import java.util.List; @NullMarked @@ -31,25 +30,4 @@ public static Email.EmailBuilder once() { .replyToName(FAKER.name().fullName()) .recipients(List.of(FAKER.internet().emailAddress(), FAKER.internet().emailAddress())); } - - public static List many(int number) { - var result = new ArrayList(); - for (int i = 0; i < number; i++) { - var body = FAKER.lorem().paragraph(); - result.add( - Email.builder() - .state(EmailState.PENDING) - .subject("Email subject") - .textBody(body) - .htmlBody("

" + body + "

") - .scheduledAt(OffsetDateTime.now()) - .fromAddress(FAKER.internet().emailAddress()) - .fromName(FAKER.name().fullName()) - .replyToAddress(FAKER.internet().emailAddress()) - .replyToName(FAKER.name().fullName()) - .recipients(List.of(FAKER.internet().emailAddress(), FAKER.internet().emailAddress())) - ); - } - return result; - } } From b25377f32de04dd6f21cd5163231de3f02d23b4c Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Fri, 14 Aug 2026 16:16:57 +0200 Subject: [PATCH 03/13] refactor Email entity by embedding EmailContent and adding new fields --- .../EmailServiceConfiguration.java | 40 ++++++++++++---- .../springboot/emailservice/lib/EmailDto.java | 8 ++-- .../emailservice/lib/EmailState.java | 1 + .../lib/application/EmailMapper.java | 9 ++++ .../lib/application/EmailServiceMigrator.java | 9 ++++ .../emailservice/lib/model/Email.java | 47 +++++-------------- .../emailservice/lib/model/EmailContent.java | 30 ++++++++++++ .../database/factory/EmailFactory.java | 19 ++++---- 8 files changed, 110 insertions(+), 53 deletions(-) create mode 100644 src/main/java/it/aboutbits/springboot/emailservice/lib/model/EmailContent.java diff --git a/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java b/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java index 535054e..8d526a3 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java @@ -14,8 +14,8 @@ import it.aboutbits.springboot.emailservice.lib.application.SendScheduledEmails; import it.aboutbits.springboot.emailservice.lib.application.UnavailableAttachmentDataSource; import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; -import jakarta.persistence.EntityManager; import org.jspecify.annotations.NullMarked; +import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.AutoConfigurationPackage; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -23,6 +23,7 @@ import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.mail.javamail.JavaMailSender; +import java.time.Duration; import java.util.List; @AutoConfigurationPackage @@ -45,24 +46,47 @@ public EmailMapper emailMapper() { } @Bean - public QueryEmail queryEmail(EmailRepository emailRepository, EmailMapper emailMapper, EntityManager entityManager) { - return new QueryEmail(emailRepository, emailMapper, entityManager); + public QueryEmail queryEmail(EmailRepository emailRepository, EmailMapper emailMapper) { + return new QueryEmail(emailRepository, emailMapper); } @Bean - public ManageEmail manageEmail(EmailRepository emailRepository, JavaMailSender javaMailSender, AttachmentDataSource attachmentDataSource, EmailMapper emailMapper) { - return new ManageEmail(emailRepository, javaMailSender, attachmentDataSource, emailMapper); + public ManageEmail manageEmail( + EmailRepository emailRepository, + JavaMailSender javaMailSender, + AttachmentDataSource attachmentDataSource, + EmailMapper emailMapper, + @Value("${aboutbits.emailservice.scheduling.max-attempts:3}") int maxAttempts, + @Value("${aboutbits.emailservice.scheduling.interval:30000}") long schedulerIntervalMillis + ) { + return new ManageEmail( + emailRepository, + javaMailSender, + attachmentDataSource, + emailMapper, + maxAttempts, + Duration.ofMillis(schedulerIntervalMillis) + ); } @Bean @ConditionalOnProperty(value = "aboutbits.emailservice.scheduling.enabled", matchIfMissing = true) - public SendScheduledEmails sendScheduledEmails(QueryEmail queryEmail, ManageEmail manageEmail, List callbacks) { - return new SendScheduledEmails(queryEmail, manageEmail, callbacks); + public SendScheduledEmails sendScheduledEmails( + QueryEmail queryEmail, + ManageEmail manageEmail, + List callbacks, + @Value("${aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold:PT5M}") Duration stuckSendingRecoveryThreshold + ) { + return new SendScheduledEmails(queryEmail, manageEmail, callbacks, stuckSendingRecoveryThreshold); } @Bean @ConditionalOnProperty(value = "aboutbits.emailservice.scheduling.cleanup.enabled", matchIfMissing = true) - public CleanupAttachmentFiles cleanupAttachments(QueryEmail queryEmail, ManageEmail manageEmail, List callbacks) { + public CleanupAttachmentFiles cleanupAttachments( + QueryEmail queryEmail, + ManageEmail manageEmail, + List callbacks + ) { return new CleanupAttachmentFiles(queryEmail, manageEmail, callbacks); } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/EmailDto.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/EmailDto.java index 11ff3ac..eefe73c 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/EmailDto.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/EmailDto.java @@ -32,10 +32,12 @@ public record EmailDto( OffsetDateTime scheduledAt, @Nullable - OffsetDateTime sentAt, - + OffsetDateTime executionStartTime, @Nullable - OffsetDateTime errorAt, + OffsetDateTime executionEndTime, + + int attempts, + @Nullable String errorMessage, diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/EmailState.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/EmailState.java index 4d2ed5d..5298f69 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/EmailState.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/EmailState.java @@ -5,6 +5,7 @@ @NullMarked public enum EmailState { PENDING, + SENDING, SENT, ERROR } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/EmailMapper.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/EmailMapper.java index 48aa8c9..ecf528c 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/EmailMapper.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/EmailMapper.java @@ -6,6 +6,7 @@ import org.jspecify.annotations.NullUnmarked; import org.mapstruct.AnnotateWith; import org.mapstruct.Mapper; +import org.mapstruct.Mapping; import org.springframework.data.domain.Page; import org.springframework.data.domain.PageImpl; @@ -15,6 +16,14 @@ @AnnotateWith(NullUnmarked.class) @NullMarked public interface EmailMapper { + @Mapping(source = "content.subject", target = "subject") + @Mapping(source = "content.fromAddress", target = "fromAddress") + @Mapping(source = "content.fromName", target = "fromName") + @Mapping(source = "content.replyToAddress", target = "replyToAddress") + @Mapping(source = "content.replyToName", target = "replyToName") + @Mapping(source = "content.recipients", target = "recipients") + @Mapping(source = "content.textBody", target = "textBody") + @Mapping(source = "content.htmlBody", target = "htmlBody") EmailDto toDto(Email model); List toDto(List model); diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/EmailServiceMigrator.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/EmailServiceMigrator.java index 617732a..919411a 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/EmailServiceMigrator.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/EmailServiceMigrator.java @@ -69,6 +69,15 @@ updated_at timestamp with time zone default now() not null, alter table email_service_emails add column if not exists reply_to_address text; alter table email_service_emails add column if not exists reply_to_name text; + + alter table email_service_emails add column if not exists attempts int default 0 not null; + alter table email_service_emails add column if not exists execution_start_time timestamp with time zone; + alter table email_service_emails add column if not exists execution_end_time timestamp with time zone; + + update email_service_emails + set execution_end_time = sent_at + where sent_at is not null + and execution_end_time is null; """ //@formatter:on ); diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/model/Email.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/model/Email.java index 88a990e..2c09720 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/model/Email.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/model/Email.java @@ -2,6 +2,7 @@ import it.aboutbits.springboot.emailservice.lib.EmailState; import jakarta.persistence.CascadeType; +import jakarta.persistence.Embedded; import jakarta.persistence.Entity; import jakarta.persistence.EnumType; import jakarta.persistence.Enumerated; @@ -18,24 +19,16 @@ import lombok.NoArgsConstructor; import lombok.Setter; import org.hibernate.annotations.CreationTimestamp; -import org.hibernate.annotations.JdbcTypeCode; import org.hibernate.annotations.UpdateTimestamp; -import org.hibernate.type.SqlTypes; import org.jspecify.annotations.NullUnmarked; import org.jspecify.annotations.Nullable; import java.time.OffsetDateTime; -import java.util.List; +import java.util.HashSet; import java.util.Set; import static it.aboutbits.springboot.emailservice.lib.model.Email.DEFAULT_ENTITY_GRAPH; -@NamedEntityGraph( - name = "email_service_emails-entity-graph", - attributeNodes = { - @NamedAttributeNode("attachments") - } -) @Entity @Getter @Setter @@ -55,34 +48,24 @@ public class Email { @Enumerated(EnumType.STRING) private EmailState state; - private String subject; - - private String fromAddress; - private String fromName; - - @Nullable - private String replyToAddress; - @Nullable - private String replyToName; - - @JdbcTypeCode(SqlTypes.JSON) - private List recipients; - - private String textBody; - private String htmlBody; + @Embedded + private EmailContent content; + @Builder.Default @OneToMany(cascade = CascadeType.PERSIST, mappedBy = "email", orphanRemoval = true) - private Set attachments; + private Set attachments = new HashSet<>(); @Builder.Default private boolean attachmentsCleaned = false; private OffsetDateTime scheduledAt; @Nullable - private OffsetDateTime sentAt; - + private OffsetDateTime executionStartTime; @Nullable - private OffsetDateTime errorAt; + private OffsetDateTime executionEndTime; + + private int attempts = 0; + @Nullable private String errorMessage; @@ -92,11 +75,7 @@ public class Email { @UpdateTimestamp private OffsetDateTime updatedAt; - public boolean isSent() { - return EmailState.SENT.equals(state); - } - - public boolean hasFailed() { - return EmailState.ERROR.equals(state); + public void incrementAttempts() { + this.attempts++; } } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/model/EmailContent.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/model/EmailContent.java new file mode 100644 index 0000000..dc68116 --- /dev/null +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/model/EmailContent.java @@ -0,0 +1,30 @@ +package it.aboutbits.springboot.emailservice.lib.model; + +import jakarta.persistence.Embeddable; +import org.hibernate.annotations.JdbcTypeCode; +import org.hibernate.type.SqlTypes; +import org.jspecify.annotations.NullMarked; +import org.jspecify.annotations.Nullable; + +import java.util.List; + +@Embeddable +@NullMarked +public record EmailContent( + String subject, + + String fromAddress, + String fromName, + + @Nullable + String replyToAddress, + @Nullable + String replyToName, + + @JdbcTypeCode(SqlTypes.JSON) + List recipients, + + String textBody, + String htmlBody +) { +} diff --git a/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java b/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java index 1c37b19..e8632cf 100644 --- a/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java +++ b/src/test/java/it/aboutbits/springboot/emailservice/support/database/factory/EmailFactory.java @@ -2,6 +2,7 @@ import it.aboutbits.springboot.emailservice.lib.EmailState; import it.aboutbits.springboot.emailservice.lib.model.Email; +import it.aboutbits.springboot.emailservice.lib.model.EmailContent; import it.aboutbits.springboot.testing.testdata.FakerExtended; import org.jspecify.annotations.NullMarked; @@ -20,14 +21,16 @@ public static Email.EmailBuilder once() { return Email.builder() .state(EmailState.PENDING) - .subject("Email subject") - .textBody(body) - .htmlBody("

" + body + "

") .scheduledAt(OffsetDateTime.now()) - .fromAddress(FAKER.internet().emailAddress()) - .fromName(FAKER.name().fullName()) - .replyToAddress(FAKER.internet().emailAddress()) - .replyToName(FAKER.name().fullName()) - .recipients(List.of(FAKER.internet().emailAddress(), FAKER.internet().emailAddress())); + .content(new EmailContent( + "Email subject", + FAKER.internet().emailAddress(), + FAKER.name().fullName(), + FAKER.internet().emailAddress(), + FAKER.name().fullName(), + List.of(FAKER.internet().emailAddress(), FAKER.internet().emailAddress()), + body, + "

" + body + "

" + )); } } From 6fa769db04f2a65c43adeedc39f0b4bb1585e7bc Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Fri, 14 Aug 2026 16:18:24 +0200 Subject: [PATCH 04/13] solve concurrency issue between pods using atomic compare-and-set; Add new configurables and make sendOrFail not retry --- .../lib/application/ManageEmail.java | 152 +++++++++++------- .../lib/application/QueryEmail.java | 29 +--- .../lib/application/SendScheduledEmails.java | 36 +++-- .../emailservice/lib/jpa/EmailRepository.java | 60 +++++++ ...itional-spring-configuration-metadata.json | 12 ++ 5 files changed, 202 insertions(+), 87 deletions(-) diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java index 3346aac..454d7c1 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java @@ -9,21 +9,24 @@ import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; import it.aboutbits.springboot.emailservice.lib.model.Email; import it.aboutbits.springboot.emailservice.lib.model.EmailAttachment; +import it.aboutbits.springboot.emailservice.lib.model.EmailContent; import jakarta.mail.MessagingException; import jakarta.validation.Valid; import lombok.extern.slf4j.Slf4j; import org.jspecify.annotations.NullMarked; import org.jspecify.annotations.Nullable; import org.springframework.core.io.ByteArrayResource; -import org.springframework.mail.MailException; import org.springframework.mail.javamail.JavaMailSender; import org.springframework.mail.javamail.MimeMessageHelper; +import org.springframework.transaction.annotation.Transactional; import org.springframework.validation.annotation.Validated; import java.io.IOException; +import java.time.Duration; import java.time.OffsetDateTime; import java.util.HashSet; import java.util.List; +import java.util.Optional; import java.util.Set; @@ -35,17 +38,23 @@ public class ManageEmail { private final JavaMailSender mailSender; private final AttachmentDataSource attachmentDataSource; private final EmailMapper emailMapper; + private final int maxAttempts; + private final Duration schedulerInterval; public ManageEmail( EmailRepository emailRepository, JavaMailSender mailSender, AttachmentDataSource attachmentDataSource, - final EmailMapper emailMapper + EmailMapper emailMapper, + int maxAttempts, + Duration schedulerInterval ) { this.emailRepository = emailRepository; this.mailSender = mailSender; this.attachmentDataSource = attachmentDataSource; this.emailMapper = emailMapper; + this.maxAttempts = maxAttempts; + this.schedulerInterval = schedulerInterval; } public EmailDto schedule(@Valid EmailParameter parameter) throws EmailException { @@ -69,9 +78,27 @@ public EmailDto sendOrFail(@Valid EmailParameter parameter) throws EmailExceptio throw new EmailException(e); } - var savedEmail = send(email); + email.setState(EmailState.SENDING); + email.setExecutionStartTime(OffsetDateTime.now()); + email.setExecutionEndTime(null); + email.setErrorMessage(null); + email.incrementAttempts(); + emailRepository.save(email); - if (savedEmail.hasFailed()) { + try { + sendMail(email); + email.setState(EmailState.SENT); + email.setExecutionEndTime(OffsetDateTime.now()); + } catch (MessagingException | AttachmentException | IOException | RuntimeException e) { + log.error("Failed to send email: {}", email.getId(), e); + email.setState(EmailState.ERROR); + email.setExecutionEndTime(OffsetDateTime.now()); + email.setErrorMessage(e.getMessage()); + } + + var savedEmail = emailRepository.save(email); + + if (savedEmail.getState() == EmailState.ERROR) { throw new EmailException("Failed to send email [id=%s, providerMessage=%s]" .formatted(savedEmail.getId(), savedEmail.getErrorMessage())); } @@ -79,22 +106,34 @@ public EmailDto sendOrFail(@Valid EmailParameter parameter) throws EmailExceptio return emailMapper.toDto(savedEmail); } - Email send(Email email) { - if (EmailState.SENT.equals(email.getState())) { - return email; - } + // Atomically transitions the row from PENDING/stale-SENDING into SENDING + @Transactional + Optional tryClaimForSend(long id, OffsetDateTime staleSendingBefore) { + var claimed = emailRepository.claimForSend(id, OffsetDateTime.now(), staleSendingBefore); + return claimed == 0 ? Optional.empty() : emailRepository.findById(id); + } + // Actually try sending the claimed email + Email completeClaimedSend(Email email) { try { sendMail(email); email.setState(EmailState.SENT); - email.setErrorMessage(""); - email.setSentAt(OffsetDateTime.now()); - } catch (MailException | MessagingException | AttachmentException | IOException e) { - log.error("Failed to send email: " + email.getId(), e); + email.setExecutionEndTime(OffsetDateTime.now()); + email.setErrorMessage(null); + } catch (Exception e) { + email.setExecutionEndTime(OffsetDateTime.now()); email.setErrorMessage(e.getMessage()); - email.setState(EmailState.ERROR); + if (email.getAttempts() > maxAttempts) { + log.error("Failed to send email: {}", email.getId(), e); + email.setState(EmailState.ERROR); + } else { + log.warn("Failed to send email: {}; Will be tried again", email.getId(), e); + email.setState(EmailState.PENDING); + email.setScheduledAt( + OffsetDateTime.now().plus(schedulerInterval.multipliedBy(email.getAttempts() * 2L)) + ); + } } - return emailRepository.save(email); } @@ -106,49 +145,17 @@ void cleanupAttachments(final Email email) throws AttachmentException { emailRepository.save(email); } - private Email fromParameter(EmailParameter parameter) throws AttachmentException { - var emailData = parameter.email(); - - final var email = new Email(); - email.setState(EmailState.PENDING); - email.setScheduledAt(parameter.scheduledAt()); - email.setSubject(emailData.subject()); - email.setTextBody(emailData.textBody()); - email.setHtmlBody(emailData.htmlBody()); - email.setRecipients(emailData.recipients()); - email.setFromAddress(emailData.fromAddress()); - email.setFromName(emailData.fromName()); - email.setReplyToAddress(emailData.replyToAddress()); - email.setReplyToName(emailData.replyToName()); - - var attachments = new HashSet(); - for (var attachment : parameter.email().attachments()) { - var reference = attachmentDataSource.storeAttachmentPayload(attachment.payload()); - - var emailAttachment = new EmailAttachment(); - emailAttachment.setEmail(email); - emailAttachment.setContentType(attachment.contentType()); - emailAttachment.setFileName(attachment.fileName()); - emailAttachment.setFileReference(reference); - - attachments.add(emailAttachment); - } - - email.setAttachments(attachments); - - return email; - } - private void sendMail(Email email) throws MessagingException, IOException, AttachmentException { + var content = email.getContent(); sendMail( - email.getFromAddress(), - email.getFromName(), - email.getReplyToAddress(), - email.getReplyToName(), - email.getRecipients(), - email.getSubject(), - email.getHtmlBody(), - email.getTextBody(), + content.fromAddress(), + content.fromName(), + content.replyToAddress(), + content.replyToName(), + content.recipients(), + content.subject(), + content.htmlBody(), + content.textBody(), email.getAttachments() ); } @@ -197,4 +204,39 @@ private void sendMail( mailSender.send(message); } + + private Email fromParameter(EmailParameter parameter) throws AttachmentException { + var emailData = parameter.email(); + + final var email = new Email(); + email.setState(EmailState.PENDING); + email.setScheduledAt(parameter.scheduledAt()); + email.setContent(new EmailContent( + emailData.subject(), + emailData.fromAddress(), + emailData.fromName(), + emailData.replyToAddress(), + emailData.replyToName(), + emailData.recipients(), + emailData.textBody(), + emailData.htmlBody() + )); + + var attachments = new HashSet(); + for (var attachment : parameter.email().attachments()) { + var reference = attachmentDataSource.storeAttachmentPayload(attachment.payload()); + + var emailAttachment = new EmailAttachment(); + emailAttachment.setEmail(email); + emailAttachment.setContentType(attachment.contentType()); + emailAttachment.setFileName(attachment.fileName()); + emailAttachment.setFileReference(reference); + + attachments.add(emailAttachment); + } + + email.setAttachments(attachments); + + return email; + } } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java index 7c43588..36da5bf 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/QueryEmail.java @@ -5,7 +5,6 @@ import it.aboutbits.springboot.emailservice.lib.EmailState; import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; import it.aboutbits.springboot.emailservice.lib.model.Email; -import jakarta.persistence.EntityManager; import lombok.RequiredArgsConstructor; import org.jspecify.annotations.NullMarked; import org.springframework.data.domain.Page; @@ -23,7 +22,6 @@ public class QueryEmail { private final EmailRepository emailRepository; private final EmailMapper emailMapper; - private final EntityManager entityManager; public Page paginatedByState(EmailState state, PageRequest pageParameter) { var pageRequest = PageRequest.of( @@ -42,30 +40,15 @@ public List byIds(Collection ids) { return emailMapper.toDto(emailRepository.findByIdIn(ids)); } - List readyToSend() { - var entityGraph = entityManager.getEntityGraph("email_service_emails-entity-graph"); - return entityManager.createQuery( - """ - SELECT e from Email e WHERE e.scheduledAt < :scheduledBefore AND e.state IN ( - it.aboutbits.springboot.emailservice.lib.EmailState.PENDING, - it.aboutbits.springboot.emailservice.lib.EmailState.ERROR - ) - """, Email.class - ) - .setParameter("scheduledBefore", OffsetDateTime.now()) - .setHint("jakarta.persistence.fetchgraph", entityGraph) - .getResultList(); + List candidateIdsToSend(OffsetDateTime staleSendingBefore) { + return emailRepository.findCandidateIdsToSend( + OffsetDateTime.now(), + staleSendingBefore + ); } List readyToCleanup() { - var entityGraph = entityManager.getEntityGraph("email_service_emails-entity-graph"); - return entityManager.createQuery( - """ - SELECT e from Email e WHERE e.attachmentsCleaned=false AND e.state=it.aboutbits.springboot.emailservice.lib.EmailState.SENT - """, Email.class - ) - .setHint("jakarta.persistence.fetchgraph", entityGraph) - .getResultList(); + return emailRepository.findReadyToCleanup(); } public Optional byId(long id) { diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java index f1789d0..4a293f0 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java @@ -2,15 +2,14 @@ import it.aboutbits.springboot.emailservice.lib.EmailSchedulerCallback; -import lombok.RequiredArgsConstructor; import lombok.extern.log4j.Log4j2; import org.jspecify.annotations.NullMarked; import org.springframework.scheduling.annotation.Scheduled; import java.time.Duration; +import java.time.OffsetDateTime; import java.util.List; -@RequiredArgsConstructor @Log4j2 @NullMarked public class SendScheduledEmails { @@ -19,27 +18,46 @@ public class SendScheduledEmails { private final QueryEmail queryEmail; private final ManageEmail manageEmail; private final List callbacks; + private final Duration stuckSendingRecoveryThreshold; private long lastInfoLogMillis = System.currentTimeMillis(); private long silentRuns = 0; private boolean firstRun = true; + public SendScheduledEmails( + QueryEmail queryEmail, + ManageEmail manageEmail, + List callbacks, + Duration stuckSendingRecoveryThreshold + ) { + this.queryEmail = queryEmail; + this.manageEmail = manageEmail; + this.callbacks = callbacks; + this.stuckSendingRecoveryThreshold = stuckSendingRecoveryThreshold; + } + @Scheduled(initialDelayString = "${aboutbits.emailservice.scheduling.interval:30000}", fixedDelayString = "${aboutbits.emailservice.scheduling.interval:30000}") void sendEmails() { logStartOfPass(); - var emailsToSend = queryEmail.readyToSend(); + var staleSendingBefore = OffsetDateTime.now().minus(stuckSendingRecoveryThreshold); + var candidateIds = queryEmail.candidateIdsToSend(staleSendingBefore); var countSent = 0; var countError = 0; - for (var email : emailsToSend) { - var updatedEmail = manageEmail.send(email); - switch (updatedEmail.getState()) { - case ERROR -> countError++; + for (var id : candidateIds) { + var claimed = manageEmail.tryClaimForSend(id, staleSendingBefore); + if (claimed.isEmpty()) { + // Lost race to another pod; Skip + continue; + } + var updated = manageEmail.completeClaimedSend(claimed.get()); + switch (updated.getState()) { + case ERROR, PENDING -> countError++; case SENT -> countSent++; default -> log.warn( JOB_DESCRIPTION + " | Job produced an invalid notification result state: {}.", - updatedEmail.getState().name() + updated.getState().name() ); } } @@ -48,7 +66,7 @@ void sendEmails() { for (var callback : callbacks) { callback.report(new EmailSchedulerCallback.Report( - emailsToSend.size(), + candidateIds.size(), countSent, countError )); diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/jpa/EmailRepository.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/jpa/EmailRepository.java index ba3100c..4cc1430 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/jpa/EmailRepository.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/jpa/EmailRepository.java @@ -8,7 +8,11 @@ import org.springframework.data.domain.PageRequest; import org.springframework.data.jpa.repository.EntityGraph; import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Modifying; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; +import java.time.OffsetDateTime; import java.util.Collection; import java.util.List; import java.util.Optional; @@ -32,4 +36,60 @@ public interface EmailRepository extends JpaRepository { @EntityGraph(value = Email.DEFAULT_ENTITY_GRAPH) List findByIdIn(Collection ids); + + // Plain read, no locking -> two pods may see overlapping candidate sets. + // The atomic UPDATE in claimForSend arbitrates the actual claim. + // Includes SENDING rows abandoned by crashed pods (past the stale threshold). + @Query(""" + select e.id from Email e + where e.scheduledAt < :now + and ( + e.state = it.aboutbits.springboot.emailservice.lib.EmailState.PENDING + or ( + e.state = it.aboutbits.springboot.emailservice.lib.EmailState.SENDING + and e.executionStartTime < :staleSendingBefore + ) + ) + order by e.scheduledAt + limit :limit + """) + List findCandidateIdsToSend( + @Param("now") OffsetDateTime now, + @Param("staleSendingBefore") OffsetDateTime staleSendingBefore + ); + + // Atomic compare-and-set claim: transitions a single row into SENDING if its + // current state is claimable (PENDING, or SENDING abandoned by a crashed pod past the stale threshold). + // Concurrent updates are serialized against the same row, so EXACTLY ONE caller gets returned 1. + @Modifying(clearAutomatically = true, flushAutomatically = true) + @Query(""" + update Email e + set e.state = it.aboutbits.springboot.emailservice.lib.EmailState.SENDING, + e.executionStartTime = :now, + e.executionEndTime = null, + e.errorMessage = null, + e.attempts = e.attempts + 1 + where e.id = :id + and e.scheduledAt < :now + and ( + e.state = it.aboutbits.springboot.emailservice.lib.EmailState.PENDING + or ( + e.state = it.aboutbits.springboot.emailservice.lib.EmailState.SENDING + and e.executionStartTime < :staleSendingBefore + ) + ) + """) + int claimForSend( + @Param("id") long id, + @Param("now") OffsetDateTime now, + @Param("staleSendingBefore") OffsetDateTime staleSendingBefore + ); + + @EntityGraph(value = Email.DEFAULT_ENTITY_GRAPH) + @Query(""" + select e from Email e + where e.attachmentsCleaned = false + and e.state = it.aboutbits.springboot.emailservice.lib.EmailState.SENT + """) + List findReadyToCleanup(); } diff --git a/src/main/resources/META-INF/additional-spring-configuration-metadata.json b/src/main/resources/META-INF/additional-spring-configuration-metadata.json index 6653b27..eccf79c 100644 --- a/src/main/resources/META-INF/additional-spring-configuration-metadata.json +++ b/src/main/resources/META-INF/additional-spring-configuration-metadata.json @@ -23,6 +23,18 @@ "type": "java.lang.Long", "description": "Specifies the milliseconds delay between runs of the scheduler.", "defaultValue": 30000 + }, + { + "name": "aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold", + "type": "java.time.Duration", + "description": "How long an email may stay in the SENDING state before being considered abandoned (crashed pod) and eligible to be re-claimed by another pod.", + "defaultValue": "PT5M" + }, + { + "name": "aboutbits.emailservice.scheduling.max-attempts", + "type": "java.lang.Integer", + "description": "Maximum number of send attempts before an email is marked as ERROR. Applies to the scheduled retry loop; the synchronous sendOrFail path is fail-fast and does not retry.", + "defaultValue": 3 } ] } From 7d6526f63e72b36bdcc962b90884eded2e7e56af Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Fri, 14 Aug 2026 16:18:50 +0200 Subject: [PATCH 05/13] add test for concurrency and schedule functionality --- .../application/SendScheduledEmailsTest.java | 235 ++++++++++++++++++ 1 file changed, 235 insertions(+) create mode 100644 src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java diff --git a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java new file mode 100644 index 0000000..c947f70 --- /dev/null +++ b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java @@ -0,0 +1,235 @@ +package it.aboutbits.springboot.emailservice.lib.application; + +import it.aboutbits.springboot.emailservice.lib.EmailState; +import it.aboutbits.springboot.emailservice.lib.jpa.EmailRepository; +import it.aboutbits.springboot.emailservice.support.database.WithPostgres; +import it.aboutbits.springboot.emailservice.support.database.factory.EmailFactory; +import jakarta.mail.internet.MimeMessage; +import org.jspecify.annotations.NullMarked; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.mail.MailSendException; +import org.springframework.mail.javamail.JavaMailSender; +import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; + +import java.time.OffsetDateTime; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.stream.IntStream; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +@SpringBootTest(properties = { + "aboutbits.emailservice.scheduling.batch-size=5", + "aboutbits.emailservice.scheduling.max-attempts=3", + "aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold=PT5M", + "aboutbits.emailservice.scheduling.interval=30000" +}) +@WithPostgres +@NullMarked +class SendScheduledEmailsTest { + private static final int BATCH_SIZE = 5; + private static final int MAX_ATTEMPTS = 3; + private static final int SCHEDULER_INTERVAL_SECONDS = 30; + + @MockitoSpyBean + JavaMailSender javaMailSender; + + @Autowired + EmailRepository emailRepository; + + @Autowired + SendScheduledEmails sendScheduledEmails; + + @BeforeEach + void setup() { + doNothing().when(javaMailSender).send(any(MimeMessage.class)); + } + + @Test + void givenPendingEmails_sendEmails_shouldMarkAllAsSent() { + emailRepository.saveAll(IntStream.range(0, 3).mapToObj(_ -> EmailFactory.once().build()).toList()); + + sendScheduledEmails.sendEmails(); + + assertThat(emailRepository.findAll()) + .hasSize(3) + .allMatch(email -> email.getState() == EmailState.SENT) + .allMatch(email -> email.getExecutionStartTime() != null) + .allMatch(email -> email.getExecutionEndTime() != null) + .allMatch(email -> email.getAttempts() == 1); + verify(javaMailSender, times(3)).send(any(MimeMessage.class)); + } + + @Test + void givenMoreEmailsThanBatchSize_sendEmails_shouldProcessOnlyBatchSizePerPass() { + var pending = 3 * BATCH_SIZE; + emailRepository.saveAll(IntStream.range(0, pending).mapToObj(_ -> EmailFactory.once().build()).toList()); + + sendScheduledEmails.sendEmails(); + + assertThat(emailRepository.findAll()) + .filteredOn(email -> email.getState() == EmailState.SENT) + .hasSize(BATCH_SIZE); + verify(javaMailSender, times(BATCH_SIZE)).send(any(MimeMessage.class)); + } + + @Test + void givenRetryableFailure_sendEmails_shouldRescheduleWithBackoffAndIncrementAttempts() { + doThrow(new MailSendException("smtp blip")).when(javaMailSender).send(any(MimeMessage.class)); + emailRepository.save(EmailFactory.once().build()); + + var beforePass = OffsetDateTime.now(); + sendScheduledEmails.sendEmails(); + + assertThat(emailRepository.findAll()) + .singleElement() + .satisfies(email -> { + assertThat(email.getState()).isEqualTo(EmailState.PENDING); + assertThat(email.getAttempts()).isEqualTo(1); + // First retry backoff: attempts (1) * 2 * scheduler interval. + assertThat(email.getScheduledAt()) + .isAfterOrEqualTo(beforePass.plusSeconds(2L * SCHEDULER_INTERVAL_SECONDS)); + assertThat(email.getErrorMessage()).contains("smtp blip"); + assertThat(email.getExecutionEndTime()).isNotNull(); + }); + } + + @Test + void givenAttemptsAlreadyAtBudget_sendEmails_shouldEscalateToError() { + doThrow(new MailSendException("smtp down")).when(javaMailSender).send(any(MimeMessage.class)); + emailRepository.save( + EmailFactory.once() + .attempts(MAX_ATTEMPTS) + .scheduledAt(OffsetDateTime.now().minusSeconds(30)) + .build() + ); + + sendScheduledEmails.sendEmails(); + + assertThat(emailRepository.findAll()) + .singleElement() + .satisfies(email -> { + // The atomic claim UPDATE incremented attempts from MAX_ATTEMPTS to MAX_ATTEMPTS+1, + // which is > threshold, so the failure escalates to ERROR. + assertThat(email.getState()).isEqualTo(EmailState.ERROR); + assertThat(email.getAttempts()).isEqualTo(MAX_ATTEMPTS + 1); + assertThat(email.getErrorMessage()).contains("smtp down"); + }); + } + + @Test + void givenErrorRow_sendEmails_shouldNotRetry() { + emailRepository.save( + EmailFactory.once() + .state(EmailState.ERROR) + .attempts(MAX_ATTEMPTS + 1) + .scheduledAt(OffsetDateTime.now().minusMinutes(1)) + .errorMessage("previous permanent failure") + .build() + ); + + sendScheduledEmails.sendEmails(); + + assertThat(emailRepository.findAll()) + .singleElement() + .satisfies(email -> { + assertThat(email.getState()).isEqualTo(EmailState.ERROR); + assertThat(email.getAttempts()).isEqualTo(MAX_ATTEMPTS + 1); + }); + verify(javaMailSender, times(0)).send(any(MimeMessage.class)); + } + + @Test + void givenStuckSendingRow_sendEmails_shouldRecoverAndResend() { + var staleStart = OffsetDateTime.now().minusMinutes(10); + emailRepository.save( + EmailFactory.once() + .state(EmailState.SENDING) + .attempts(1) + .executionStartTime(staleStart) + .scheduledAt(OffsetDateTime.now().minusSeconds(60)) + .build() + ); + + sendScheduledEmails.sendEmails(); + + assertThat(emailRepository.findAll()) + .singleElement() + .satisfies(email -> { + assertThat(email.getState()).isEqualTo(EmailState.SENT); + assertThat(email.getAttempts()).isEqualTo(2); + assertThat(email.getExecutionStartTime()).isAfter(staleStart); + }); + verify(javaMailSender, times(1)).send(any(MimeMessage.class)); + } + + @Test + void givenFreshSendingRowWithinThreshold_sendEmails_shouldNotStealFromOtherPod() { + // executionStartTime within the 5-minute threshold means another pod is legitimately + // sending this right now; we must not re-claim it. + var recentStart = OffsetDateTime.now().minusSeconds(30); + emailRepository.save( + EmailFactory.once() + .state(EmailState.SENDING) + .attempts(1) + .executionStartTime(recentStart) + .scheduledAt(OffsetDateTime.now().minusSeconds(60)) + .build() + ); + + sendScheduledEmails.sendEmails(); + + assertThat(emailRepository.findAll()) + .singleElement() + .satisfies(email -> { + assertThat(email.getState()).isEqualTo(EmailState.SENDING); + assertThat(email.getAttempts()).isEqualTo(1); + assertThat(email.getExecutionStartTime()).isEqualTo(recentStart); + }); + verify(javaMailSender, times(0)).send(any(MimeMessage.class)); + } + + @Test + void givenConcurrentPasses_sendEmails_shouldSendEachEmailExactlyOnce() throws Exception { + var totalEmails = 25; + var workerThreads = 4; + var passesPerWorker = 10; + emailRepository.saveAll(IntStream.range(0, totalEmails).mapToObj(_ -> EmailFactory.once().build()).toList()); + + var startGate = new CountDownLatch(1); + var executor = Executors.newFixedThreadPool(workerThreads); + try { + for (var t = 0; t < workerThreads; t++) { + executor.submit(() -> { + startGate.await(); + for (var i = 0; i < passesPerWorker; i++) { + sendScheduledEmails.sendEmails(); + } + return null; + }); + } + startGate.countDown(); + executor.shutdown(); + assertThat(executor.awaitTermination(60, TimeUnit.SECONDS)).isTrue(); + } finally { + if (!executor.isTerminated()) { + executor.shutdownNow(); + } + } + + assertThat(emailRepository.findAll()) + .hasSize(totalEmails) + .allMatch(email -> email.getState() == EmailState.SENT) + .allMatch(email -> email.getAttempts() == 1); + verify(javaMailSender, times(totalEmails)).send(any(MimeMessage.class)); + } +} From e429e2069710f0dde8d3280c746fb883563a29fa Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Fri, 14 Aug 2026 16:18:58 +0200 Subject: [PATCH 06/13] update readme.md --- readme.md | 39 ++++++++++++++++++++++++++++++++++----- 1 file changed, 34 insertions(+), 5 deletions(-) diff --git a/readme.md b/readme.md index 19c5434..85a0dd4 100644 --- a/readme.md +++ b/readme.md @@ -74,11 +74,40 @@ public class App { The following configuration options are available: -| Name | Default | Description | -|----------------------------------------|-------------|-----------------------------------------------------------------------| -| `lib.emailservice.migrations.enabled` | true | Enables database migrations. | -| `lib.emailservice.scheduling.enabled` | true | Enables the scheduler sending the emails. | -| `lib.emailservice.scheduling.interval` | 30000 | Specifies the milliseconds delay between runs of the scheduler. | +| Name | Default | Description | +|-------------------------------------------------------------------------|---------|------------------------------------------------------------------------------------------------------------------------| +| `aboutbits.emailservice.migrations.enabled` | true | Enables database migrations. | +| `aboutbits.emailservice.scheduling.enabled` | true | Enables the scheduler sending the emails. | +| `aboutbits.emailservice.scheduling.cleanup.enabled` | true | Enables cleanup of attachment files after sending. | +| `aboutbits.emailservice.scheduling.interval` | 30000 | Milliseconds delay between runs of the scheduler. | +| `aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold` | PT5M | How long an email may stay in `SENDING` before being considered abandoned (crashed pod) and eligible to be re-claimed. | +| `aboutbits.emailservice.scheduling.max-attempts` | 3 | Maximum number of send attempts before an email is marked as `ERROR`. Applies only to the scheduled retry loop. | + +## Multi-pod deployments + +The scheduler is safe to run on every pod concurrently. Each pass performs +three independent steps: + +1. **Candidate scan** — a plain `SELECT` returns up to `batch-size` ids of + emails whose `scheduled_at` has passed and whose state is `PENDING` or `SENDING` + older than `stuck-sending-recovery-threshold` (crashed-pod recovery). `ERROR` + is terminal: rows that exhausted `max-attempts` are never re-picked + automatically — an operator can reset them to `PENDING` if a retry is desired. + Two pods may see overlapping ids at this step, and that's fine — the atomic + claim below arbitrates. +2. **Atomic claim** — for each candidate id, a single compare-and-set + `UPDATE … SET state = SENDING … WHERE id = ? AND state IN (PENDING, ERROR, staleSENDING)` + is issued. The database serializes concurrent updates against the same row so + exactly one pod sees `rowsAffected = 1` and owns that email; the other pods + see `0` and move on. +3. **Send and persist** — the winning pod calls SMTP outside any database + transaction and then writes the final `SENT` / `ERROR` / `PENDING` state through a small `save()`. + +Crash recovery: if a pod dies between "claim" and "persist result", the row +stays in `SENDING` with its `execution_start_time` frozen. After +`stuck-sending-recovery-threshold` any pod's next pass picks it up as a +candidate again and retries. Repeated crashes therefore count against +`max-attempts` and eventually escalate the row to `ERROR` rather than looping forever. ## Local development: From be9d60e55925b4947a3bd173d748bb9d7a26849b Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Fri, 14 Aug 2026 16:24:33 +0200 Subject: [PATCH 07/13] test fixes --- .../emailservice/lib/jpa/EmailRepository.java | 1 - .../lib/application/SendScheduledEmailsTest.java | 15 --------------- 2 files changed, 16 deletions(-) diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/jpa/EmailRepository.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/jpa/EmailRepository.java index 4cc1430..4c8c150 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/jpa/EmailRepository.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/jpa/EmailRepository.java @@ -51,7 +51,6 @@ public interface EmailRepository extends JpaRepository { ) ) order by e.scheduledAt - limit :limit """) List findCandidateIdsToSend( @Param("now") OffsetDateTime now, diff --git a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java index c947f70..2775abd 100644 --- a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java +++ b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java @@ -28,7 +28,6 @@ import static org.mockito.Mockito.verify; @SpringBootTest(properties = { - "aboutbits.emailservice.scheduling.batch-size=5", "aboutbits.emailservice.scheduling.max-attempts=3", "aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold=PT5M", "aboutbits.emailservice.scheduling.interval=30000" @@ -36,7 +35,6 @@ @WithPostgres @NullMarked class SendScheduledEmailsTest { - private static final int BATCH_SIZE = 5; private static final int MAX_ATTEMPTS = 3; private static final int SCHEDULER_INTERVAL_SECONDS = 30; @@ -69,19 +67,6 @@ void givenPendingEmails_sendEmails_shouldMarkAllAsSent() { verify(javaMailSender, times(3)).send(any(MimeMessage.class)); } - @Test - void givenMoreEmailsThanBatchSize_sendEmails_shouldProcessOnlyBatchSizePerPass() { - var pending = 3 * BATCH_SIZE; - emailRepository.saveAll(IntStream.range(0, pending).mapToObj(_ -> EmailFactory.once().build()).toList()); - - sendScheduledEmails.sendEmails(); - - assertThat(emailRepository.findAll()) - .filteredOn(email -> email.getState() == EmailState.SENT) - .hasSize(BATCH_SIZE); - verify(javaMailSender, times(BATCH_SIZE)).send(any(MimeMessage.class)); - } - @Test void givenRetryableFailure_sendEmails_shouldRescheduleWithBackoffAndIncrementAttempts() { doThrow(new MailSendException("smtp blip")).when(javaMailSender).send(any(MimeMessage.class)); From 38db610a8b53a96ea7d21261c2c1f99446c979e5 Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Fri, 14 Aug 2026 16:42:37 +0200 Subject: [PATCH 08/13] truncate OffsetDateTime to MICROS for pipeline to run successfully --- .../springboot/emailservice/lib/application/ManageEmail.java | 1 - .../emailservice/lib/application/SendScheduledEmailsTest.java | 3 ++- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java index 454d7c1..0ba2b4b 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java @@ -80,7 +80,6 @@ public EmailDto sendOrFail(@Valid EmailParameter parameter) throws EmailExceptio email.setState(EmailState.SENDING); email.setExecutionStartTime(OffsetDateTime.now()); - email.setExecutionEndTime(null); email.setErrorMessage(null); email.incrementAttempts(); emailRepository.save(email); diff --git a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java index 2775abd..e0efdfb 100644 --- a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java +++ b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java @@ -15,6 +15,7 @@ import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; import java.time.OffsetDateTime; +import java.time.temporal.ChronoUnit; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -161,7 +162,7 @@ void givenStuckSendingRow_sendEmails_shouldRecoverAndResend() { void givenFreshSendingRowWithinThreshold_sendEmails_shouldNotStealFromOtherPod() { // executionStartTime within the 5-minute threshold means another pod is legitimately // sending this right now; we must not re-claim it. - var recentStart = OffsetDateTime.now().minusSeconds(30); + var recentStart = OffsetDateTime.now().minusSeconds(30).truncatedTo(ChronoUnit.MICROS); emailRepository.save( EmailFactory.once() .state(EmailState.SENDING) From f60de88d378f0dfd486c2e269b0e1f760e67cf61 Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Fri, 14 Aug 2026 16:49:24 +0200 Subject: [PATCH 09/13] remove unnecessary errorMessage set --- .../springboot/emailservice/lib/application/ManageEmail.java | 1 - 1 file changed, 1 deletion(-) diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java index 0ba2b4b..060803d 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java @@ -80,7 +80,6 @@ public EmailDto sendOrFail(@Valid EmailParameter parameter) throws EmailExceptio email.setState(EmailState.SENDING); email.setExecutionStartTime(OffsetDateTime.now()); - email.setErrorMessage(null); email.incrementAttempts(); emailRepository.save(email); From 1c03524b528bf6c5c63298e2af50221fdda305a9 Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Wed, 19 Aug 2026 16:54:24 +0200 Subject: [PATCH 10/13] requested changes --- readme.md | 29 +++---------------- .../EmailServiceConfiguration.java | 9 ++++-- .../lib/application/ManageEmail.java | 22 +++++++++----- ...itional-spring-configuration-metadata.json | 6 ++-- .../application/SendScheduledEmailsTest.java | 14 ++++----- 5 files changed, 34 insertions(+), 46 deletions(-) diff --git a/readme.md b/readme.md index 85a0dd4..cf86fff 100644 --- a/readme.md +++ b/readme.md @@ -80,34 +80,13 @@ The following configuration options are available: | `aboutbits.emailservice.scheduling.enabled` | true | Enables the scheduler sending the emails. | | `aboutbits.emailservice.scheduling.cleanup.enabled` | true | Enables cleanup of attachment files after sending. | | `aboutbits.emailservice.scheduling.interval` | 30000 | Milliseconds delay between runs of the scheduler. | -| `aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold` | PT5M | How long an email may stay in `SENDING` before being considered abandoned (crashed pod) and eligible to be re-claimed. | -| `aboutbits.emailservice.scheduling.max-attempts` | 3 | Maximum number of send attempts before an email is marked as `ERROR`. Applies only to the scheduled retry loop. | +| `aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold` | PT30M | How long an email may stay in `SENDING` before being considered abandoned (crashed pod) and eligible to be re-claimed. Must comfortably exceed the worst-case SMTP send duration: JavaMail's default connect/read/write timeouts are infinite, so configure `spring.mail.properties.mail.smtp.connectiontimeout`, `spring.mail.properties.mail.smtp.timeout` and `spring.mail.properties.mail.smtp.writetimeout` well below this threshold, otherwise a slow in-flight send can be re-claimed by another pod and delivered twice. | +| `aboutbits.emailservice.scheduling.max-attempts` | 3 | Maximum number of send attempts before an email is marked as `ERROR`. Applies only to the scheduled retry loop. Failed attempts are retried with exponential backoff (`scheduling.interval` × 2^attempts); with the defaults a persistently failing email runs attempt 1 → +60s → attempt 2 → +120s → attempt 3 → `ERROR`. | ## Multi-pod deployments -The scheduler is safe to run on every pod concurrently. Each pass performs -three independent steps: - -1. **Candidate scan** — a plain `SELECT` returns up to `batch-size` ids of - emails whose `scheduled_at` has passed and whose state is `PENDING` or `SENDING` - older than `stuck-sending-recovery-threshold` (crashed-pod recovery). `ERROR` - is terminal: rows that exhausted `max-attempts` are never re-picked - automatically — an operator can reset them to `PENDING` if a retry is desired. - Two pods may see overlapping ids at this step, and that's fine — the atomic - claim below arbitrates. -2. **Atomic claim** — for each candidate id, a single compare-and-set - `UPDATE … SET state = SENDING … WHERE id = ? AND state IN (PENDING, ERROR, staleSENDING)` - is issued. The database serializes concurrent updates against the same row so - exactly one pod sees `rowsAffected = 1` and owns that email; the other pods - see `0` and move on. -3. **Send and persist** — the winning pod calls SMTP outside any database - transaction and then writes the final `SENT` / `ERROR` / `PENDING` state through a small `save()`. - -Crash recovery: if a pod dies between "claim" and "persist result", the row -stays in `SENDING` with its `execution_start_time` frozen. After -`stuck-sending-recovery-threshold` any pod's next pass picks it up as a -candidate again and retries. Repeated crashes therefore count against -`max-attempts` and eventually escalate the row to `ERROR` rather than looping forever. +The scheduler is safe to run on every pod concurrently: the database arbitrates which pod sends each email. +Delivery is at-least-once - if a pod crashes after the SMTP server accepted the message but before the result was persisted, the email may be sent again on recovery. ## Local development: diff --git a/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java b/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java index 8d526a3..c882776 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/EmailServiceConfiguration.java @@ -22,6 +22,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.mail.javamail.JavaMailSender; +import org.springframework.transaction.PlatformTransactionManager; import java.time.Duration; import java.util.List; @@ -57,7 +58,8 @@ public ManageEmail manageEmail( AttachmentDataSource attachmentDataSource, EmailMapper emailMapper, @Value("${aboutbits.emailservice.scheduling.max-attempts:3}") int maxAttempts, - @Value("${aboutbits.emailservice.scheduling.interval:30000}") long schedulerIntervalMillis + @Value("${aboutbits.emailservice.scheduling.interval:30000}") long schedulerIntervalMillis, + PlatformTransactionManager transactionManager ) { return new ManageEmail( emailRepository, @@ -65,7 +67,8 @@ public ManageEmail manageEmail( attachmentDataSource, emailMapper, maxAttempts, - Duration.ofMillis(schedulerIntervalMillis) + Duration.ofMillis(schedulerIntervalMillis), + transactionManager ); } @@ -75,7 +78,7 @@ public SendScheduledEmails sendScheduledEmails( QueryEmail queryEmail, ManageEmail manageEmail, List callbacks, - @Value("${aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold:PT5M}") Duration stuckSendingRecoveryThreshold + @Value("${aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold:PT30M}") Duration stuckSendingRecoveryThreshold ) { return new SendScheduledEmails(queryEmail, manageEmail, callbacks, stuckSendingRecoveryThreshold); } diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java index 060803d..786f957 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java @@ -18,7 +18,10 @@ import org.springframework.core.io.ByteArrayResource; import org.springframework.mail.javamail.JavaMailSender; import org.springframework.mail.javamail.MimeMessageHelper; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.TransactionDefinition; import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionTemplate; import org.springframework.validation.annotation.Validated; import java.io.IOException; @@ -40,6 +43,7 @@ public class ManageEmail { private final EmailMapper emailMapper; private final int maxAttempts; private final Duration schedulerInterval; + private final TransactionTemplate transactionTemplate; public ManageEmail( EmailRepository emailRepository, @@ -47,7 +51,8 @@ public ManageEmail( AttachmentDataSource attachmentDataSource, EmailMapper emailMapper, int maxAttempts, - Duration schedulerInterval + Duration schedulerInterval, + PlatformTransactionManager transactionManager ) { this.emailRepository = emailRepository; this.mailSender = mailSender; @@ -55,6 +60,9 @@ public ManageEmail( this.emailMapper = emailMapper; this.maxAttempts = maxAttempts; this.schedulerInterval = schedulerInterval; + // Persist the final email state in its own, independent transaction for sendOrFail to never make it roll back + this.transactionTemplate = new TransactionTemplate(transactionManager); + this.transactionTemplate.setPropagationBehavior(TransactionDefinition.PROPAGATION_REQUIRES_NEW); } public EmailDto schedule(@Valid EmailParameter parameter) throws EmailException { @@ -78,10 +86,8 @@ public EmailDto sendOrFail(@Valid EmailParameter parameter) throws EmailExceptio throw new EmailException(e); } - email.setState(EmailState.SENDING); email.setExecutionStartTime(OffsetDateTime.now()); email.incrementAttempts(); - emailRepository.save(email); try { sendMail(email); @@ -94,7 +100,7 @@ public EmailDto sendOrFail(@Valid EmailParameter parameter) throws EmailExceptio email.setErrorMessage(e.getMessage()); } - var savedEmail = emailRepository.save(email); + var savedEmail = transactionTemplate.execute(_ -> emailRepository.save(email)); if (savedEmail.getState() == EmailState.ERROR) { throw new EmailException("Failed to send email [id=%s, providerMessage=%s]" @@ -121,14 +127,15 @@ Email completeClaimedSend(Email email) { } catch (Exception e) { email.setExecutionEndTime(OffsetDateTime.now()); email.setErrorMessage(e.getMessage()); - if (email.getAttempts() > maxAttempts) { + if (email.getAttempts() >= maxAttempts) { log.error("Failed to send email: {}", email.getId(), e); email.setState(EmailState.ERROR); } else { log.warn("Failed to send email: {}; Will be tried again", email.getId(), e); email.setState(EmailState.PENDING); email.setScheduledAt( - OffsetDateTime.now().plus(schedulerInterval.multipliedBy(email.getAttempts() * 2L)) + OffsetDateTime.now() + .plus(schedulerInterval.multipliedBy((long) Math.pow(2, email.getAttempts()))) ); } } @@ -199,14 +206,13 @@ private void sendMail( payload.close(); } - mailSender.send(message); } private Email fromParameter(EmailParameter parameter) throws AttachmentException { var emailData = parameter.email(); - final var email = new Email(); + var email = new Email(); email.setState(EmailState.PENDING); email.setScheduledAt(parameter.scheduledAt()); email.setContent(new EmailContent( diff --git a/src/main/resources/META-INF/additional-spring-configuration-metadata.json b/src/main/resources/META-INF/additional-spring-configuration-metadata.json index eccf79c..7b42404 100644 --- a/src/main/resources/META-INF/additional-spring-configuration-metadata.json +++ b/src/main/resources/META-INF/additional-spring-configuration-metadata.json @@ -27,13 +27,13 @@ { "name": "aboutbits.emailservice.scheduling.stuck-sending-recovery-threshold", "type": "java.time.Duration", - "description": "How long an email may stay in the SENDING state before being considered abandoned (crashed pod) and eligible to be re-claimed by another pod.", - "defaultValue": "PT5M" + "description": "How long an email may stay in the SENDING state before being considered abandoned (crashed pod) and eligible to be re-claimed by another pod. Must comfortably exceed the worst-case SMTP send duration: JavaMail's default connect/read/write timeouts are infinite, so configure spring.mail.properties.mail.smtp.connectiontimeout, spring.mail.properties.mail.smtp.timeout and spring.mail.properties.mail.smtp.writetimeout well below this threshold, otherwise a slow in-flight send can be re-claimed by another pod and delivered twice.", + "defaultValue": "PT30M" }, { "name": "aboutbits.emailservice.scheduling.max-attempts", "type": "java.lang.Integer", - "description": "Maximum number of send attempts before an email is marked as ERROR. Applies to the scheduled retry loop; the synchronous sendOrFail path is fail-fast and does not retry.", + "description": "Maximum number of send attempts before an email is marked as ERROR. Applies to the scheduled retry loop; the synchronous sendOrFail path is fail-fast and does not retry. Failed attempts are retried with exponential backoff (scheduling.interval x 2^attempts); with the defaults a persistently failing email runs attempt 1 -> +60s -> attempt 2 -> +120s -> attempt 3 -> ERROR.", "defaultValue": 3 } ] diff --git a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java index e0efdfb..1af6242 100644 --- a/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java +++ b/src/test/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmailsTest.java @@ -81,7 +81,7 @@ void givenRetryableFailure_sendEmails_shouldRescheduleWithBackoffAndIncrementAtt .satisfies(email -> { assertThat(email.getState()).isEqualTo(EmailState.PENDING); assertThat(email.getAttempts()).isEqualTo(1); - // First retry backoff: attempts (1) * 2 * scheduler interval. + // First retry backoff: 2^attempts * scheduler interval = 2 * scheduler interval assertThat(email.getScheduledAt()) .isAfterOrEqualTo(beforePass.plusSeconds(2L * SCHEDULER_INTERVAL_SECONDS)); assertThat(email.getErrorMessage()).contains("smtp blip"); @@ -94,7 +94,7 @@ void givenAttemptsAlreadyAtBudget_sendEmails_shouldEscalateToError() { doThrow(new MailSendException("smtp down")).when(javaMailSender).send(any(MimeMessage.class)); emailRepository.save( EmailFactory.once() - .attempts(MAX_ATTEMPTS) + .attempts(MAX_ATTEMPTS - 1) .scheduledAt(OffsetDateTime.now().minusSeconds(30)) .build() ); @@ -104,10 +104,10 @@ void givenAttemptsAlreadyAtBudget_sendEmails_shouldEscalateToError() { assertThat(emailRepository.findAll()) .singleElement() .satisfies(email -> { - // The atomic claim UPDATE incremented attempts from MAX_ATTEMPTS to MAX_ATTEMPTS+1, - // which is > threshold, so the failure escalates to ERROR. + // The atomic claim UPDATE incremented attempts from MAX_ATTEMPTS-1 to MAX_ATTEMPTS, + // which is equal to the threshold, so the failure escalates to ERROR. assertThat(email.getState()).isEqualTo(EmailState.ERROR); - assertThat(email.getAttempts()).isEqualTo(MAX_ATTEMPTS + 1); + assertThat(email.getAttempts()).isEqualTo(MAX_ATTEMPTS); assertThat(email.getErrorMessage()).contains("smtp down"); }); } @@ -117,7 +117,7 @@ void givenErrorRow_sendEmails_shouldNotRetry() { emailRepository.save( EmailFactory.once() .state(EmailState.ERROR) - .attempts(MAX_ATTEMPTS + 1) + .attempts(MAX_ATTEMPTS) .scheduledAt(OffsetDateTime.now().minusMinutes(1)) .errorMessage("previous permanent failure") .build() @@ -129,7 +129,7 @@ void givenErrorRow_sendEmails_shouldNotRetry() { .singleElement() .satisfies(email -> { assertThat(email.getState()).isEqualTo(EmailState.ERROR); - assertThat(email.getAttempts()).isEqualTo(MAX_ATTEMPTS + 1); + assertThat(email.getAttempts()).isEqualTo(MAX_ATTEMPTS); }); verify(javaMailSender, times(0)).send(any(MimeMessage.class)); } From 259f4e56e163f6d21f2f1fff65ff7cbefd637352 Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Wed, 19 Aug 2026 16:54:49 +0200 Subject: [PATCH 11/13] consider all error cases that can happen on a schedule pass --- .../lib/application/SendScheduledEmails.java | 54 +++++++++++++++---- 1 file changed, 45 insertions(+), 9 deletions(-) diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java index 4a293f0..a904205 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java @@ -2,6 +2,7 @@ import it.aboutbits.springboot.emailservice.lib.EmailSchedulerCallback; +import it.aboutbits.springboot.emailservice.lib.model.Email; import lombok.extern.log4j.Log4j2; import org.jspecify.annotations.NullMarked; import org.springframework.scheduling.annotation.Scheduled; @@ -9,6 +10,7 @@ import java.time.Duration; import java.time.OffsetDateTime; import java.util.List; +import java.util.Optional; @Log4j2 @NullMarked @@ -43,21 +45,55 @@ void sendEmails() { var staleSendingBefore = OffsetDateTime.now().minus(stuckSendingRecoveryThreshold); var candidateIds = queryEmail.candidateIdsToSend(staleSendingBefore); + var countClaimed = 0; var countSent = 0; var countError = 0; for (var id : candidateIds) { - var claimed = manageEmail.tryClaimForSend(id, staleSendingBefore); + Optional claimed; + try { + claimed = manageEmail.tryClaimForSend(id, staleSendingBefore); + } catch (Exception e) { + // Claim failed and rolled back: nothing was sent and the row is unchanged, so it stays claimable + // Another pod may pick it up in this same pass, or it is retried on a later pass. + log.error( + JOB_DESCRIPTION + " | Failed to claim email for sending: {}. " + + "Nothing was sent and the row is unchanged; it stays claimable and may be picked up " + + "by another pod in this cycle or retried on a later pass.", + id, + e + ); + continue; + } + if (claimed.isEmpty()) { // Lost race to another pod; Skip continue; } - var updated = manageEmail.completeClaimedSend(claimed.get()); - switch (updated.getState()) { - case ERROR, PENDING -> countError++; - case SENT -> countSent++; - default -> log.warn( - JOB_DESCRIPTION + " | Job produced an invalid notification result state: {}.", - updated.getState().name() + countClaimed++; + + try { + var updated = manageEmail.completeClaimedSend(claimed.get()); + switch (updated.getState()) { + // We do not count PENDING emails as errors since they will + // be retried and eventually end up in one of the two buckets. + case ERROR -> countError++; + case SENT -> countSent++; + default -> log.warn( + JOB_DESCRIPTION + " | Job produced an invalid notification result state: {}.", + updated.getState().name() + ); + } + } catch (Exception e) { + // Best-effort visibility only: we cannot reliably recover here. The email may already have been + // sent but persisting the result failed, so the row stays in SENDING and will be re-claimed and + // re-sent after the recovery threshold (possible duplicate delivery). + countError++; + log.error( + JOB_DESCRIPTION + " | Failed to complete send for email: {}. It may have been sent already; " + + "if the result could not be persisted the row stays in SENDING and will be re-sent " + + "after the recovery threshold (possible duplicate delivery).", + id, + e ); } } @@ -66,7 +102,7 @@ void sendEmails() { for (var callback : callbacks) { callback.report(new EmailSchedulerCallback.Report( - candidateIds.size(), + countClaimed, countSent, countError )); From 63ff6585899e521d62375ca5a2fa35e10d2448d4 Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Wed, 19 Aug 2026 17:01:09 +0200 Subject: [PATCH 12/13] remove the first try catch since it is overkill --- .../lib/application/SendScheduledEmails.java | 15 +-------------- 1 file changed, 1 insertion(+), 14 deletions(-) diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java index a904205..4926336 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/SendScheduledEmails.java @@ -50,20 +50,7 @@ void sendEmails() { var countError = 0; for (var id : candidateIds) { Optional claimed; - try { - claimed = manageEmail.tryClaimForSend(id, staleSendingBefore); - } catch (Exception e) { - // Claim failed and rolled back: nothing was sent and the row is unchanged, so it stays claimable - // Another pod may pick it up in this same pass, or it is retried on a later pass. - log.error( - JOB_DESCRIPTION + " | Failed to claim email for sending: {}. " - + "Nothing was sent and the row is unchanged; it stays claimable and may be picked up " - + "by another pod in this cycle or retried on a later pass.", - id, - e - ); - continue; - } + claimed = manageEmail.tryClaimForSend(id, staleSendingBefore); if (claimed.isEmpty()) { // Lost race to another pod; Skip From 3e89517b7feb4b212fb8dc686c5ca01836aee063 Mon Sep 17 00:00:00 2001 From: Jonas Mayr Date: Thu, 20 Aug 2026 10:15:10 +0200 Subject: [PATCH 13/13] make last save be in its on transaction so it can never be rolled back --- .../springboot/emailservice/lib/application/ManageEmail.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java index 786f957..7563888 100644 --- a/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java +++ b/src/main/java/it/aboutbits/springboot/emailservice/lib/application/ManageEmail.java @@ -139,7 +139,7 @@ Email completeClaimedSend(Email email) { ); } } - return emailRepository.save(email); + return transactionTemplate.execute(_ -> emailRepository.save(email)); } void cleanupAttachments(final Email email) throws AttachmentException {