Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<mainClass>uk.gov.pay.webhooks.app.WebhooksApp</mainClass>
<pay-java-commons.version>1.0.20260825154634</pay-java-commons.version>
<pay-java-commons.version>1.0.20260901094800</pay-java-commons.version>
<swagger-version>2.2.54</swagger-version>
<pact.version>4.7.5</pact.version>
<prometheus.version>0.16.0</prometheus.version>
Expand Down Expand Up @@ -180,6 +180,11 @@
<artifactId>logging-dropwizard-5</artifactId>
<version>${pay-java-commons.version}</version>
</dependency>
<dependency>
<groupId>uk.gov.service.payments</groupId>
<artifactId>queue-dropwizard-5</artifactId>
<version>${pay-java-commons.version}</version>
</dependency>

<!-- Test dependencies that are imported from one of the BOMs specified
in <dependencyManagement> so no explicit versions needed -->
Expand Down
6 changes: 6 additions & 0 deletions src/main/java/uk/gov/pay/webhooks/app/WebhooksModule.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
3 changes: 1 addition & 2 deletions src/main/java/uk/gov/pay/webhooks/queue/EventMessage.java
Original file line number Diff line number Diff line change
@@ -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) {

Expand All @@ -20,5 +20,4 @@ public InternalEvent toInternalEvent() {
eventMessageDto.resourceType()
);
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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()));
Expand Down
16 changes: 8 additions & 8 deletions src/main/java/uk/gov/pay/webhooks/queue/EventQueue.java
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -43,22 +43,22 @@ public List<EventMessage> 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;
Expand Down

This file was deleted.

21 changes: 0 additions & 21 deletions src/main/java/uk/gov/pay/webhooks/queue/sqs/QueueMessage.java

This file was deleted.

82 changes: 0 additions & 82 deletions src/main/java/uk/gov/pay/webhooks/queue/sqs/SqsQueueService.java

This file was deleted.

14 changes: 7 additions & 7 deletions src/test/java/uk/gov/pay/webhooks/queue/EventQueueIT.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<QueueMessage> result = sqsQueueService.receiveMessages(SqsTestDocker.getQueueUrl("event-queue"), "All");
assertFalse(result.isEmpty());
Expand Down Expand Up @@ -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<EventMessage> 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"));
}
}
Loading