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