diff --git a/timeless-api/pom.xml b/timeless-api/pom.xml index 6e5c675..055fbaf 100644 --- a/timeless-api/pom.xml +++ b/timeless-api/pom.xml @@ -62,6 +62,10 @@ quarkus-quinoa ${quarkus-quinoa.version} + + io.quarkiverse.amazonservices + quarkus-messaging-amazon-sqs + io.quarkiverse.amazonservices quarkus-amazon-sqs diff --git a/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessage.java b/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessage.java new file mode 100644 index 0000000..310ea4a --- /dev/null +++ b/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessage.java @@ -0,0 +1,101 @@ +package dev.matheuscruz.infra.outbox; + +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.EnumType; +import jakarta.persistence.Enumerated; +import jakarta.persistence.GeneratedValue; +import jakarta.persistence.GenerationType; +import jakarta.persistence.Id; +import jakarta.persistence.Table; +import java.time.Instant; +import java.util.UUID; + +@Entity +@Table(name = "outbox_messages") +public class OutboxMessage { + + @Id + @GeneratedValue(strategy = GenerationType.UUID) + private UUID id; + + @Column(nullable = false, columnDefinition = "TEXT") + private String payload; + + @Column(name = "queue_url", nullable = false) + private String queueUrl; + + @Column(name = "message_group_id", nullable = false) + private String messageGroupId; + + @Enumerated(EnumType.STRING) + @Column(nullable = false) + private OutboxStatus status; + + @Column(name = "created_at", nullable = false) + private Instant createdAt; + + @Column(name = "processed_at") + private Instant processedAt; + + @Column(name = "retry_count", nullable = false) + private int retryCount; + + protected OutboxMessage() { + } + + public OutboxMessage(String payload, String queueUrl, String messageGroupId) { + this.payload = payload; + this.queueUrl = queueUrl; + this.messageGroupId = messageGroupId; + this.status = OutboxStatus.PENDING; + this.createdAt = Instant.now(); + this.retryCount = 0; + } + + public UUID getId() { + return id; + } + + public String getPayload() { + return payload; + } + + public String getQueueUrl() { + return queueUrl; + } + + public String getMessageGroupId() { + return messageGroupId; + } + + public OutboxStatus getStatus() { + return status; + } + + public Instant getCreatedAt() { + return createdAt; + } + + public Instant getProcessedAt() { + return processedAt; + } + + public int getRetryCount() { + return retryCount; + } + + public void markAsSent() { + this.status = OutboxStatus.SENT; + this.processedAt = Instant.now(); + } + + public void markAsFailed() { + this.status = OutboxStatus.FAILED; + this.processedAt = Instant.now(); + } + + public void incrementRetryCount() { + this.retryCount++; + } +} diff --git a/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessageRelay.java b/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessageRelay.java new file mode 100644 index 0000000..9f2d901 --- /dev/null +++ b/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessageRelay.java @@ -0,0 +1,70 @@ +package dev.matheuscruz.infra.outbox; + +import io.quarkus.narayana.jta.QuarkusTransaction; +import io.quarkus.scheduler.Scheduled; +import jakarta.enterprise.context.ApplicationScoped; +import java.util.List; +import org.jboss.logging.Logger; +import software.amazon.awssdk.services.sqs.SqsClient; +import software.amazon.awssdk.services.sqs.model.SendMessageRequest; + +@ApplicationScoped +public class OutboxMessageRelay { + + private static final int MAX_RETRIES = 10; + + private final OutboxMessageRepository outboxMessageRepository; + private final SqsClient sqsClient; + private final Logger logger = Logger.getLogger(OutboxMessageRelay.class); + + public OutboxMessageRelay(OutboxMessageRepository outboxMessageRepository, SqsClient sqsClient) { + this.outboxMessageRepository = outboxMessageRepository; + this.sqsClient = sqsClient; + } + + @Scheduled(every = "5s") + public void processOutbox() { + List pending = outboxMessageRepository.findPendingMessages(); + + for (OutboxMessage message : pending) { + if (message.getRetryCount() >= MAX_RETRIES) { + failMessage(message); + continue; + } + + try { + sqsClient.sendMessage(SendMessageRequest.builder().queueUrl(message.getQueueUrl()) + .messageBody(message.getPayload()).messageGroupId(message.getMessageGroupId()).build()); + + updateMessageStatus(message, true); + } catch (Exception e) { + logger.errorf(e, "Failed to send outbox message %s to queue %s", message.getId(), + message.getQueueUrl()); + updateMessageStatus(message, false); + } + } + } + + private void updateMessageStatus(OutboxMessage message, boolean success) { + QuarkusTransaction.requiringNew().run(() -> { + OutboxMessage managed = outboxMessageRepository.findById(message.getId()); + if (managed != null) { + if (success) { + managed.markAsSent(); + } else { + managed.incrementRetryCount(); + } + } + }); + } + + private void failMessage(OutboxMessage message) { + QuarkusTransaction.requiringNew().run(() -> { + OutboxMessage managed = outboxMessageRepository.findById(message.getId()); + if (managed != null) { + managed.markAsFailed(); + logger.errorf("Outbox message %s has exceeded max retries, marking as FAILED", message.getId()); + } + }); + } +} diff --git a/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessageRepository.java b/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessageRepository.java new file mode 100644 index 0000000..0a00257 --- /dev/null +++ b/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxMessageRepository.java @@ -0,0 +1,18 @@ +package dev.matheuscruz.infra.outbox; + +import io.quarkus.hibernate.orm.panache.PanacheRepositoryBase; +import jakarta.enterprise.context.ApplicationScoped; +import java.util.List; +import java.util.UUID; + +@ApplicationScoped +public class OutboxMessageRepository implements PanacheRepositoryBase { + + private static final int MAX_RETRIES = 10; + private static final int BATCH_SIZE = 20; + + public List findPendingMessages() { + return find("status = ?1 and retryCount < ?2 order by createdAt", OutboxStatus.PENDING, MAX_RETRIES) + .range(0, BATCH_SIZE - 1).list(); + } +} diff --git a/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxStatus.java b/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxStatus.java new file mode 100644 index 0000000..d3e8be8 --- /dev/null +++ b/timeless-api/src/main/java/dev/matheuscruz/infra/outbox/OutboxStatus.java @@ -0,0 +1,5 @@ +package dev.matheuscruz.infra.outbox; + +public enum OutboxStatus { + PENDING, SENT, FAILED +} diff --git a/timeless-api/src/main/java/dev/matheuscruz/infra/queue/SQS.java b/timeless-api/src/main/java/dev/matheuscruz/infra/queue/SQS.java index 6aae662..a815d9a 100644 --- a/timeless-api/src/main/java/dev/matheuscruz/infra/queue/SQS.java +++ b/timeless-api/src/main/java/dev/matheuscruz/infra/queue/SQS.java @@ -13,122 +13,126 @@ import dev.matheuscruz.infra.ai.data.RecognizedOperation; import dev.matheuscruz.infra.ai.data.RecognizedTransaction; import dev.matheuscruz.infra.ai.data.SimpleMessage; +import dev.matheuscruz.infra.outbox.OutboxMessage; +import dev.matheuscruz.infra.outbox.OutboxMessageRepository; import io.quarkus.narayana.jta.QuarkusTransaction; -import io.quarkus.scheduler.Scheduled; import jakarta.enterprise.context.ApplicationScoped; -import java.io.IOException; import java.util.Optional; -import java.util.UUID; +import java.util.concurrent.CompletionStage; import org.eclipse.microprofile.config.inject.ConfigProperty; +import org.eclipse.microprofile.reactive.messaging.Incoming; +import org.eclipse.microprofile.reactive.messaging.Message; import org.jboss.logging.Logger; -import software.amazon.awssdk.services.sqs.SqsClient; @ApplicationScoped public class SQS { - final String incomingMessagesUrl; - final String processedMessagesUrl; - final SqsClient sqs; final ObjectMapper objectMapper; final TextAiService aiService; final RecordRepository recordRepository; final UserRepository userRepository; + final OutboxMessageRepository outboxMessageRepository; + final String recognizedQueueUrl; final Logger logger = Logger.getLogger(SQS.class); private static final ObjectReader INCOMING_MESSAGE_READER = new ObjectMapper().readerFor(IncomingMessage.class); - private static final ObjectReader AI_RESPONSE_READER = new ObjectMapper().readerFor(RecognizedOperation.class); + public SQS(ObjectMapper objectMapper, TextAiService aiService, RecordRepository recordRepository, + UserRepository userRepository, OutboxMessageRepository outboxMessageRepository, + @ConfigProperty(name = "whatsapp.recognized-message.queue-url") String recognizedQueueUrl) { - public SQS(SqsClient sqs, @ConfigProperty(name = "whatsapp.incoming-message.queue-url") String incomingMessagesUrl, - @ConfigProperty(name = "whatsapp.recognized-message.queue-url") String messagesProcessedUrl, - ObjectMapper objectMapper, TextAiService aiService, RecordRepository recordRepository, - UserRepository userRepository) { - - this.sqs = sqs; - this.incomingMessagesUrl = incomingMessagesUrl; - this.processedMessagesUrl = messagesProcessedUrl; this.objectMapper = objectMapper; this.aiService = aiService; this.recordRepository = recordRepository; this.userRepository = userRepository; + this.outboxMessageRepository = outboxMessageRepository; + this.recognizedQueueUrl = recognizedQueueUrl; } - @Scheduled(every = "5s") - public void receiveMessages() { - sqs.receiveMessage(req -> req.maxNumberOfMessages(10).queueUrl(incomingMessagesUrl)).messages() - .forEach(message -> processMessage(message.body(), message.receiptHandle())); - } - - private void processMessage(String body, String receiptHandle) { + @Incoming("whatsapp-incoming") + public CompletionStage receiveMessages(Message message) { + String body = message.getPayload(); IncomingMessage incomingMessage = parseIncomingMessage(body); - if (!MessageKind.TEXT.equals(incomingMessage.kind())) - return; + + if (!MessageKind.TEXT.equals(incomingMessage.kind())) { + return message.ack(); + } Optional user = this.userRepository.findByPhoneNumber(incomingMessage.sender()); if (user.isEmpty()) { - logger.error("User not found. Deleting message from queue."); - deleteMessageUsing(receiptHandle); - return; + logger.error("User not found."); + return message.nack(new RuntimeException("User not found for phone: " + incomingMessage.sender())); } - handleUserMessage(user.get(), incomingMessage, receiptHandle); + try { + handleUserMessage(user.get(), incomingMessage); + return message.ack(); + } catch (Exception e) { + logger.error("Failed to process message: " + incomingMessage.messageId(), e); + return message.nack(e); + } } - private void handleUserMessage(User user, IncomingMessage message, String receiptHandle) { - try { - AllRecognizedOperations allRecognizedOperations = aiService.handleMessage(message.messageBody(), - user.getId()); - - for (RecognizedOperation recognizedOperation : allRecognizedOperations.all()) { - switch (recognizedOperation.operation()) { - case AiOperations.ADD_TRANSACTION -> - processAddTransactionMessage(user, message, receiptHandle, recognizedOperation); - case AiOperations.GET_BALANCE -> { - logger.info("Processing GET_BALANCE operation" + recognizedOperation.recognizedTransaction()); - processSimpleMessage(user, message, receiptHandle, recognizedOperation); - } - default -> logger.warnf("Unknown operation type: %s", recognizedOperation.operation()); + private void handleUserMessage(User user, IncomingMessage message) { + AllRecognizedOperations allRecognizedOperations = aiService.handleMessage(message.messageBody(), user.getId()); + + for (RecognizedOperation recognizedOperation : allRecognizedOperations.all()) { + switch (recognizedOperation.operation()) { + case AiOperations.ADD_TRANSACTION -> processAddTransactionMessage(user, message, recognizedOperation); + case AiOperations.GET_BALANCE -> { + logger.info("Processing GET_BALANCE operation" + recognizedOperation.recognizedTransaction()); + processSimpleMessage(user, message, recognizedOperation); } + default -> logger.warnf("Unknown operation type: %s", recognizedOperation.operation()); } - - } catch (Exception e) { - logger.error("Failed to process message: " + message.messageId(), e); } } - private void processAddTransactionMessage(User user, IncomingMessage message, String receiptHandle, - RecognizedOperation recognizedOperation) throws IOException { + private void processAddTransactionMessage(User user, IncomingMessage message, + RecognizedOperation recognizedOperation) { RecognizedTransaction recognizedTransaction = recognizedOperation.recognizedTransaction(); - sendProcessedMessage(new TransactionMessageProcessed(AiOperations.ADD_TRANSACTION.commandName(), - message.messageId(), MessageStatus.PROCESSED, user.getPhoneNumber(), recognizedTransaction.withError(), - recognizedTransaction)); Record record = new Record.Builder().userId(user.getId()).amount(recognizedTransaction.amount()) .description(recognizedTransaction.description()).transaction(recognizedTransaction.type()) .category(recognizedTransaction.category()).build(); - QuarkusTransaction.requiringNew().run(() -> recordRepository.persist(record)); + TransactionMessageProcessed processed = new TransactionMessageProcessed( + AiOperations.ADD_TRANSACTION.commandName(), message.messageId(), MessageStatus.PROCESSED, + user.getPhoneNumber(), recognizedTransaction.withError(), recognizedTransaction); + + OutboxMessage outboxMessage = new OutboxMessage(serialize(processed), recognizedQueueUrl, + user.getPhoneNumber()); - deleteMessageUsing(receiptHandle); + QuarkusTransaction.requiringNew().run(() -> { + recordRepository.persist(record); + outboxMessageRepository.persist(outboxMessage); + }); logger.infof("Message %s processed as ADD_TRANSACTION", message.messageId()); } - private void processSimpleMessage(User user, IncomingMessage message, String receiptHandle, - RecognizedOperation recognizedOperation) throws IOException { + private void processSimpleMessage(User user, IncomingMessage message, RecognizedOperation recognizedOperation) { logger.infof("Processing simple message for user %s", recognizedOperation.recognizedTransaction()); SimpleMessage response = new SimpleMessage(recognizedOperation.recognizedTransaction().description()); - sendProcessedMessage(new SimpleMessageProcessed(AiOperations.GET_BALANCE.commandName(), message.messageId(), - MessageStatus.PROCESSED, user.getPhoneNumber(), response)); - deleteMessageUsing(receiptHandle); + + SimpleMessageProcessed processed = new SimpleMessageProcessed(AiOperations.GET_BALANCE.commandName(), + message.messageId(), MessageStatus.PROCESSED, user.getPhoneNumber(), response); + + OutboxMessage outboxMessage = new OutboxMessage(serialize(processed), recognizedQueueUrl, + user.getPhoneNumber()); + + QuarkusTransaction.requiringNew().run(() -> outboxMessageRepository.persist(outboxMessage)); + logger.infof("Message %s processed as GET_BALANCE", message.messageId()); } - private void sendProcessedMessage(Object processedMessage) throws JsonProcessingException { - String messageBody = objectMapper.writeValueAsString(processedMessage); - sqs.sendMessage(req -> req.messageBody(messageBody).messageGroupId("ProcessedMessages") - .messageDeduplicationId(UUID.randomUUID().toString()).queueUrl(processedMessagesUrl)); + private String serialize(Object message) { + try { + return objectMapper.writeValueAsString(message); + } catch (JsonProcessingException e) { + throw new RuntimeException("Failed to serialize message", e); + } } private IncomingMessage parseIncomingMessage(String messageBody) { @@ -139,10 +143,6 @@ private IncomingMessage parseIncomingMessage(String messageBody) { } } - private void deleteMessageUsing(String receiptHandle) { - sqs.deleteMessage(req -> req.queueUrl(incomingMessagesUrl).receiptHandle(receiptHandle)); - } - public record TransactionMessageProcessed(String kind, String messageId, MessageStatus status, String user, Boolean withError, RecognizedTransaction content) { } diff --git a/timeless-api/src/main/resources/application.properties b/timeless-api/src/main/resources/application.properties index f6c765b..f52c612 100644 --- a/timeless-api/src/main/resources/application.properties +++ b/timeless-api/src/main/resources/application.properties @@ -5,6 +5,11 @@ security.sensible.secret=${SECURITY_KEY} whatsapp.incoming-message.queue-url=${INCOMING_MESSAGE_FIFO_URL} whatsapp.recognized-message.queue-url=${RECOGNIZED_MESSAGE_FIFO_URL} +# smallrye reactive messaging sqs +mp.messaging.incoming.whatsapp-incoming.connector=smallrye-sqs +mp.messaging.incoming.whatsapp-incoming.queue=${INCOMING_MESSAGE_FIFO_URL} +mp.messaging.incoming.whatsapp-incoming.visibility-timeout=30 + # aws sqs quarkus.sqs.devservices.enabled=false