diff --git a/pom.xml b/pom.xml index 44e55a3f..9351f00d 100644 --- a/pom.xml +++ b/pom.xml @@ -14,7 +14,7 @@ UTF-8 UTF-8 uk.gov.pay.webhooks.app.WebhooksApp - 1.0.20260825154634 + 1.0.20260901094800 2.2.54 4.7.5 0.16.0 @@ -180,6 +180,11 @@ logging-dropwizard-5 ${pay-java-commons.version} + + uk.gov.service.payments + queue-dropwizard-5 + ${pay-java-commons.version} + diff --git a/src/main/java/uk/gov/pay/webhooks/app/WebhooksModule.java b/src/main/java/uk/gov/pay/webhooks/app/WebhooksModule.java index 93f0a593..57839932 100644 --- a/src/main/java/uk/gov/pay/webhooks/app/WebhooksModule.java +++ b/src/main/java/uk/gov/pay/webhooks/app/WebhooksModule.java @@ -27,6 +27,7 @@ import uk.gov.pay.webhooks.message.HttpPostFactory; import uk.gov.pay.webhooks.message.WebhookMessageSignatureGenerator; import uk.gov.pay.webhooks.util.IdGenerator; +import uk.gov.service.payments.commons.queue.sqs.SqsQueueService; import javax.net.ssl.SSLContext; import java.net.URI; @@ -166,4 +167,9 @@ public SqsClient sqsClient(WebhooksConfig webhooksConfig) { return clientBuilder.build(); } + + @Provides + public SqsQueueService sqsQueueService(SqsClient sqsClient, WebhooksConfig webhooksConfig) { + return new SqsQueueService(sqsClient, webhooksConfig.getSqsConfig().getMessageMaximumWaitTimeInSeconds(), webhooksConfig.getSqsConfig().getMessageMaximumBatchSize()); + } } diff --git a/src/main/java/uk/gov/pay/webhooks/message/WebhookMessageService.java b/src/main/java/uk/gov/pay/webhooks/message/WebhookMessageService.java index 34498a60..b52cbb24 100644 --- a/src/main/java/uk/gov/pay/webhooks/message/WebhookMessageService.java +++ b/src/main/java/uk/gov/pay/webhooks/message/WebhookMessageService.java @@ -11,7 +11,6 @@ import uk.gov.pay.webhooks.ledger.LedgerService; import uk.gov.pay.webhooks.ledger.model.LedgerTransaction; import uk.gov.pay.webhooks.message.apirepresentation.PaymentApiRepresentation; -import uk.gov.pay.webhooks.message.apirepresentation.RefundApiRepresentation; import uk.gov.pay.webhooks.message.dao.WebhookMessageDao; import uk.gov.pay.webhooks.message.dao.entity.WebhookMessageEntity; import uk.gov.pay.webhooks.queue.InternalEvent; diff --git a/src/main/java/uk/gov/pay/webhooks/queue/EventMessage.java b/src/main/java/uk/gov/pay/webhooks/queue/EventMessage.java index b955c281..7185acbf 100644 --- a/src/main/java/uk/gov/pay/webhooks/queue/EventMessage.java +++ b/src/main/java/uk/gov/pay/webhooks/queue/EventMessage.java @@ -1,6 +1,6 @@ package uk.gov.pay.webhooks.queue; -import uk.gov.pay.webhooks.queue.sqs.QueueMessage; +import uk.gov.service.payments.commons.queue.model.QueueMessage; public record EventMessage(EventMessageDto eventMessageDto, QueueMessage queueMessage) { @@ -20,5 +20,4 @@ public InternalEvent toInternalEvent() { eventMessageDto.resourceType() ); } - } diff --git a/src/main/java/uk/gov/pay/webhooks/queue/EventMessageHandler.java b/src/main/java/uk/gov/pay/webhooks/queue/EventMessageHandler.java index 9abaef44..fd68de92 100644 --- a/src/main/java/uk/gov/pay/webhooks/queue/EventMessageHandler.java +++ b/src/main/java/uk/gov/pay/webhooks/queue/EventMessageHandler.java @@ -7,7 +7,7 @@ import org.slf4j.LoggerFactory; import org.slf4j.MDC; import uk.gov.pay.webhooks.message.WebhookMessageService; -import uk.gov.pay.webhooks.queue.sqs.QueueException; +import uk.gov.service.payments.commons.queue.exception.QueueException; import java.util.UUID; @@ -38,7 +38,7 @@ public void handle() throws QueueException { for (EventMessage message : eventQueue.retrieveEvents()) { try { MDC.put(MDC_REQUEST_ID_KEY, UUID.randomUUID().toString()); - MDC.put(SQS_MESSAGE_ID, message.queueMessage().messageId()); + MDC.put(SQS_MESSAGE_ID, message.queueMessage().getMessageId()); MDC.put(SERVICE_EXTERNAL_ID, message.eventMessageDto().serviceId()); MDC.put(GATEWAY_ACCOUNT_ID, message.eventMessageDto().gatewayAccountId()); MDC.put(RESOURCE_IS_LIVE, String.valueOf(message.eventMessageDto().live())); diff --git a/src/main/java/uk/gov/pay/webhooks/queue/EventQueue.java b/src/main/java/uk/gov/pay/webhooks/queue/EventQueue.java index b5196296..a8559621 100644 --- a/src/main/java/uk/gov/pay/webhooks/queue/EventQueue.java +++ b/src/main/java/uk/gov/pay/webhooks/queue/EventQueue.java @@ -1,14 +1,14 @@ package uk.gov.pay.webhooks.queue; import com.fasterxml.jackson.databind.ObjectMapper; +import jakarta.inject.Inject; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import uk.gov.pay.webhooks.app.WebhooksConfig; -import uk.gov.pay.webhooks.queue.sqs.QueueException; -import uk.gov.pay.webhooks.queue.sqs.QueueMessage; -import uk.gov.pay.webhooks.queue.sqs.SqsQueueService; +import uk.gov.service.payments.commons.queue.exception.QueueException; +import uk.gov.service.payments.commons.queue.model.QueueMessage; +import uk.gov.service.payments.commons.queue.sqs.SqsQueueService; -import jakarta.inject.Inject; import java.io.IOException; import java.util.List; import java.util.Objects; @@ -43,22 +43,22 @@ public List retrieveEvents() throws QueueException { } public void markMessageAsProcessed(EventMessage message) throws QueueException { - sqsQueueService.deleteMessage(this.eventQueueUrl, message.queueMessage().receiptHandle()); + sqsQueueService.deleteMessage(this.eventQueueUrl, message.queueMessage().getReceiptHandle()); } public void scheduleMessageForRetry(EventMessage message) throws QueueException { - sqsQueueService.deferMessage(this.eventQueueUrl, message.queueMessage().receiptHandle(), retryDelayInSeconds); + sqsQueueService.deferMessage(this.eventQueueUrl, message.queueMessage().getReceiptHandle(), retryDelayInSeconds); } private EventMessage getMessage(QueueMessage queueMessage) { try { - SNSMessageDto snsMessageDto = objectMapper.readValue(queueMessage.messageBody(), SNSMessageDto.class); + SNSMessageDto snsMessageDto = objectMapper.readValue(queueMessage.getMessageBody(), SNSMessageDto.class); EventMessageDto eventMessageDto = objectMapper.readValue(snsMessageDto.Message(), EventMessageDto.class); return EventMessage.of(eventMessageDto, queueMessage); } catch (IOException e) { LOGGER.warn( "There was an exception parsing message [messageId={}] into an [{}] {}", - queueMessage.messageId(), + queueMessage.getMessageId(), EventMessage.class, e.getMessage()); return null; diff --git a/src/main/java/uk/gov/pay/webhooks/queue/sqs/QueueException.java b/src/main/java/uk/gov/pay/webhooks/queue/sqs/QueueException.java deleted file mode 100644 index 2067d33b..00000000 --- a/src/main/java/uk/gov/pay/webhooks/queue/sqs/QueueException.java +++ /dev/null @@ -1,9 +0,0 @@ -package uk.gov.pay.webhooks.queue.sqs; - -public class QueueException extends Exception { - - public QueueException(String message, Exception e) { - super(message, e); - } - -} diff --git a/src/main/java/uk/gov/pay/webhooks/queue/sqs/QueueMessage.java b/src/main/java/uk/gov/pay/webhooks/queue/sqs/QueueMessage.java deleted file mode 100644 index a59d2585..00000000 --- a/src/main/java/uk/gov/pay/webhooks/queue/sqs/QueueMessage.java +++ /dev/null @@ -1,21 +0,0 @@ -package uk.gov.pay.webhooks.queue.sqs; - -import software.amazon.awssdk.services.sqs.model.ReceiveMessageResponse; -import software.amazon.awssdk.services.sqs.model.SendMessageResponse; - -import java.util.List; - -public record QueueMessage(String messageId, String receiptHandle, String messageBody) { - - public static List of(ReceiveMessageResponse receiveMessageResult) { - return receiveMessageResult.messages() - .stream() - .map(message -> new QueueMessage(message.messageId(), message.receiptHandle(), message.body())) - .toList(); - } - - public static QueueMessage of(SendMessageResponse messageResult, String validJsonMessage) { - return new QueueMessage(messageResult.messageId(), null, validJsonMessage); - } - -} diff --git a/src/main/java/uk/gov/pay/webhooks/queue/sqs/SqsQueueService.java b/src/main/java/uk/gov/pay/webhooks/queue/sqs/SqsQueueService.java deleted file mode 100644 index 148cbf03..00000000 --- a/src/main/java/uk/gov/pay/webhooks/queue/sqs/SqsQueueService.java +++ /dev/null @@ -1,82 +0,0 @@ -package uk.gov.pay.webhooks.queue.sqs; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import software.amazon.awssdk.awscore.exception.AwsServiceException; -import software.amazon.awssdk.services.sqs.SqsClient; -import software.amazon.awssdk.services.sqs.model.ChangeMessageVisibilityRequest; -import software.amazon.awssdk.services.sqs.model.DeleteMessageRequest; -import software.amazon.awssdk.services.sqs.model.ReceiveMessageRequest; -import software.amazon.awssdk.services.sqs.model.ReceiveMessageResponse; -import software.amazon.awssdk.services.sqs.model.SqsException; -import uk.gov.pay.webhooks.app.WebhooksConfig; - -import jakarta.inject.Inject; - -import java.util.List; - -public class SqsQueueService { - private final Logger logger = LoggerFactory.getLogger(SqsQueueService.class); - - private final SqsClient sqsClient; - - private final int messageMaximumWaitTimeInSeconds; - private final int messageMaximumBatchSize; - - @Inject - public SqsQueueService(SqsClient sqsClient, WebhooksConfig webhooksConfig) { - this.sqsClient = sqsClient; - this.messageMaximumBatchSize = webhooksConfig.getSqsConfig().getMessageMaximumBatchSize(); - this.messageMaximumWaitTimeInSeconds = webhooksConfig.getSqsConfig().getMessageMaximumWaitTimeInSeconds(); - } - - public List receiveMessages(String queueUrl, String messageAttributeName) throws QueueException { - try { - ReceiveMessageRequest receiveMessageRequest = ReceiveMessageRequest.builder() - .queueUrl(queueUrl) - .messageAttributeNames(messageAttributeName) - .waitTimeSeconds(messageMaximumWaitTimeInSeconds) - .maxNumberOfMessages(messageMaximumBatchSize) - .build(); - - ReceiveMessageResponse receiveMessageResult = sqsClient.receiveMessage(receiveMessageRequest); - - return QueueMessage.of(receiveMessageResult); - } catch (Exception e) { - logger.error("Failed to receive messages from SQS queue - [{}] {}", e.getClass().getCanonicalName(), e.getMessage()); - throw new QueueException("Failed to receive messages from SQS queue", e); - } - } - - public void deleteMessage(String queueUrl, String messageReceiptHandle) throws QueueException { - try { - DeleteMessageRequest deleteMessageRequest = DeleteMessageRequest.builder() - .queueUrl(queueUrl) - .receiptHandle(messageReceiptHandle) - .build(); - sqsClient.deleteMessage(deleteMessageRequest); - } catch (SqsException | UnsupportedOperationException e) { - logger.error("Failed to delete message from SQS queue - {}", e.getMessage()); - throw new QueueException("Failed to delete message from SQS queue", e); - } catch (AwsServiceException e) { - logger.error("Failed to delete message from SQS queue - [errorMessage={}] [awsErrorCode={}]", e.getMessage(), e.awsErrorDetails().errorCode()); - throw new QueueException("Failed to delete message from SQS queue", e); - } - } - - public void deferMessage(String queueUrl, String messageReceiptHandle, int retryDelayInSeconds) throws QueueException { - try { - ChangeMessageVisibilityRequest changeMessageVisibilityRequest = ChangeMessageVisibilityRequest.builder() - .queueUrl(queueUrl) - .receiptHandle(messageReceiptHandle) - .visibilityTimeout(retryDelayInSeconds) - .build(); - - - sqsClient.changeMessageVisibility(changeMessageVisibilityRequest); - } catch (SqsException | UnsupportedOperationException e) { - logger.error("Failed to defer message from SQS queue - {}", e.getMessage()); - throw new QueueException("Failed to defer message from SQS queue", e); - } - } -} diff --git a/src/test/java/uk/gov/pay/webhooks/queue/EventQueueIT.java b/src/test/java/uk/gov/pay/webhooks/queue/EventQueueIT.java index b2e8405e..c78bc053 100644 --- a/src/test/java/uk/gov/pay/webhooks/queue/EventQueueIT.java +++ b/src/test/java/uk/gov/pay/webhooks/queue/EventQueueIT.java @@ -11,9 +11,9 @@ import uk.gov.pay.webhooks.app.QueueMessageReceiverConfig; import uk.gov.pay.webhooks.app.SqsConfig; import uk.gov.pay.webhooks.app.WebhooksConfig; -import uk.gov.pay.webhooks.queue.sqs.QueueException; -import uk.gov.pay.webhooks.queue.sqs.QueueMessage; -import uk.gov.pay.webhooks.queue.sqs.SqsQueueService; +import uk.gov.service.payments.commons.queue.exception.QueueException; +import uk.gov.service.payments.commons.queue.model.QueueMessage; +import uk.gov.service.payments.commons.queue.sqs.SqsQueueService; import java.io.IOException; import java.util.List; @@ -51,7 +51,7 @@ void shouldReceiveMessageFromTheQueue() throws QueueException { WebhooksConfig mockConfig = mock(WebhooksConfig.class); when(mockConfig.getSqsConfig()).thenReturn(sqsConfig); - SqsQueueService sqsQueueService = new SqsQueueService(client, mockConfig); + SqsQueueService sqsQueueService = new SqsQueueService(client, sqsConfig.getMessageMaximumWaitTimeInSeconds(), sqsConfig.getMessageMaximumBatchSize()); List result = sqsQueueService.receiveMessages(SqsTestDocker.getQueueUrl("event-queue"), "All"); assertFalse(result.isEmpty()); @@ -91,12 +91,12 @@ void shouldConvertValidMessageFromQueueToEventMessage() throws QueueException, I when(mockConfig.getSqsConfig()).thenReturn(sqsConfig); when(mockConfig.getQueueMessageReceiverConfig()).thenReturn(queueReceiverConfig); - SqsQueueService sqsQueueService = new SqsQueueService(client, mockConfig); + SqsQueueService sqsQueueService = new SqsQueueService(client, sqsConfig.getMessageMaximumWaitTimeInSeconds(), sqsConfig.getMessageMaximumBatchSize()); EventQueue eventQueue = new EventQueue(sqsQueueService, mockConfig, new ObjectMapper()); List result = eventQueue.retrieveEvents(); assertFalse(result.isEmpty()); - assertThat(result.get(0).eventMessageDto().resourceType(), is("payment")); - assertThat(result.get(0).eventMessageDto().gatewayAccountId(), is("100")); + assertThat(result.getFirst().eventMessageDto().resourceType(), is("payment")); + assertThat(result.getFirst().eventMessageDto().gatewayAccountId(), is("100")); } }