diff --git a/pom.xml b/pom.xml index 2be9bbc..b178bfd 100644 --- a/pom.xml +++ b/pom.xml @@ -6,34 +6,37 @@ uk.gov.ons.census census-rm-caseprocessor - 1.0-SNAPSHOT + 1.0.0-SNAPSHOT 21 21 + ${java.version} + + 1.0.0 + 1.0.0 + UTF-8 docker - 2025.0.3 - 7.4.8 + 2025.1.3 + 8.1.0 - 1.18.30 - 4.0.0 - 2.3.0 - 3.15.4 - 1.10.1 - 7.4 - 5.9 - 1.9.24 + 3.15.5 + 1.11.0 + 9.0 + 5.12.0 + 2.0.13 - 3.24.0 - 3.1.1 - 2.43.0 - 1.22.0 - 0.8.11 - 2.23.0 + 3.28.0 + 7.26.0 + 3.6.3 + 3.10.0 + 1.36.1 + 0.8.15 + 2.50.0 @@ -56,7 +59,7 @@ org.springframework.boot spring-boot-starter-parent - 3.5.16 + 4.1.1 @@ -107,17 +110,21 @@ uk.gov.ons.census census-rm-common-entity-model - 0.0.3 + ${census-rm-common-entity-model.version} uk.gov.ons.census census-rm-shared-sample-validation - 0.1.0 + ${census-rm-shared-sample-validation.version} org.springframework.boot spring-boot-starter + + org.springframework.boot + spring-boot-starter-jackson + org.springframework.boot spring-boot-starter-integration @@ -137,21 +144,6 @@ jakarta.xml.bind jakarta.xml.bind-api - ${jakarta-xml-bind-api.version} - - - javax.xml.bind - jaxb-api - ${javax-jaxb-api.version} - - - - com.fasterxml.jackson.datatype - jackson-datatype-jsr310 - - - com.fasterxml.jackson.datatype - jackson-datatype-jdk8 org.postgresql @@ -160,12 +152,11 @@ org.projectlombok lombok - ${lombok.version} provided io.hypersistence - hypersistence-utils-hibernate-63 + hypersistence-utils-hibernate-73 ${hypersistence-utils.version} @@ -186,7 +177,6 @@ org.aspectj aspectjweaver - ${aspectjweaver.version} @@ -195,6 +185,12 @@ spring-boot-starter-test test + + org.springframework.retry + spring-retry + ${spring-retry.version} + + @@ -204,8 +200,20 @@ org.apache.maven.plugins maven-pmd-plugin ${maven-pmd-plugin.version} + + + net.sourceforge.pmd + pmd-core + ${pmd.version} + + + net.sourceforge.pmd + pmd-java + ${pmd.version} + + - 21 + ${java.version} exclude-pmd.properties 3 true @@ -270,7 +278,6 @@ org.springframework.boot spring-boot-maven-plugin - true uk.gov.ons.census.caseprocessor.Application @@ -345,11 +352,13 @@ maven-compiler-plugin - 21 - 21 + ${maven.compiler.release} UTF-8 -XDcompilePolicy=simple + --should-stop=ifError=FLOW + + -XDaddTypeAnnotationsToSymbol=true -Xplugin:ErrorProne diff --git a/src/main/java/uk/gov/ons/census/caseprocessor/Application.java b/src/main/java/uk/gov/ons/census/caseprocessor/Application.java index 7eef4b5..83b91b2 100644 --- a/src/main/java/uk/gov/ons/census/caseprocessor/Application.java +++ b/src/main/java/uk/gov/ons/census/caseprocessor/Application.java @@ -2,7 +2,7 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.boot.autoconfigure.domain.EntityScan; +import org.springframework.boot.persistence.autoconfigure.EntityScan; import org.springframework.integration.annotation.IntegrationComponentScan; @SpringBootApplication diff --git a/src/main/java/uk/gov/ons/census/caseprocessor/config/DefaultListenerSupport.java b/src/main/java/uk/gov/ons/census/caseprocessor/config/DefaultListenerSupport.java index fc6adc5..b76ba9d 100644 --- a/src/main/java/uk/gov/ons/census/caseprocessor/config/DefaultListenerSupport.java +++ b/src/main/java/uk/gov/ons/census/caseprocessor/config/DefaultListenerSupport.java @@ -4,7 +4,10 @@ import org.springframework.retry.RetryContext; import org.springframework.retry.RetryListener; -public class DefaultListenerSupport implements RetryListener { +/* Bridge both listener contracts so one bean supports migrated runtime wiring and legacy @Retryable +listeners. */ +public class DefaultListenerSupport + implements org.springframework.core.retry.RetryListener, RetryListener { @Override public void close( @@ -15,7 +18,6 @@ public void close( @Override public void onError( RetryContext context, RetryCallback callback, Throwable throwable) { - RetryListener.super.onError(context, callback, throwable); } diff --git a/src/main/java/uk/gov/ons/census/caseprocessor/config/MessageConsumerConfig.java b/src/main/java/uk/gov/ons/census/caseprocessor/config/MessageConsumerConfig.java index ac6c45a..596b29c 100644 --- a/src/main/java/uk/gov/ons/census/caseprocessor/config/MessageConsumerConfig.java +++ b/src/main/java/uk/gov/ons/census/caseprocessor/config/MessageConsumerConfig.java @@ -5,18 +5,25 @@ import com.google.cloud.spring.pubsub.core.PubSubTemplate; import com.google.cloud.spring.pubsub.integration.AckMode; import com.google.cloud.spring.pubsub.integration.inbound.PubSubInboundChannelAdapter; +import java.time.Duration; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.core.retry.RetryListener; +import org.springframework.core.retry.RetryPolicy; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.handler.advice.RequestHandlerRetryAdvice; import org.springframework.messaging.MessageChannel; -import org.springframework.retry.RetryListener; import uk.gov.ons.census.caseprocessor.messaging.ManagedMessageRecoverer; @Configuration public class MessageConsumerConfig { + // Spring core retry defaults to maxRetries = 3, i.e. 3 retries after the initial call + // (4 total invocations). Pin this to 3 total invocations to preserve the pre-migration + // behaviour. + private static final int MESSAGE_TOTAL_ATTEMPTS = 3; + private final ManagedMessageRecoverer managedMessageRecoverer; private final PubSubTemplate pubSubTemplate; @@ -209,6 +216,8 @@ private PubSubInboundChannelAdapter makeAdapter(MessageChannel channel, String s @Bean public RequestHandlerRetryAdvice retryAdvice() { RequestHandlerRetryAdvice requestHandlerRetryAdvice = new RequestHandlerRetryAdvice(); + requestHandlerRetryAdvice.setRetryPolicy( + RetryPolicy.builder().maxRetries(MESSAGE_TOTAL_ATTEMPTS - 1).delay(Duration.ZERO).build()); requestHandlerRetryAdvice.setRecoveryCallback(managedMessageRecoverer); return requestHandlerRetryAdvice; } diff --git a/src/main/java/uk/gov/ons/census/caseprocessor/messaging/ManagedMessageRecoverer.java b/src/main/java/uk/gov/ons/census/caseprocessor/messaging/ManagedMessageRecoverer.java index cc19ad6..15b5bc7 100644 --- a/src/main/java/uk/gov/ons/census/caseprocessor/messaging/ManagedMessageRecoverer.java +++ b/src/main/java/uk/gov/ons/census/caseprocessor/messaging/ManagedMessageRecoverer.java @@ -7,10 +7,11 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; +import org.springframework.core.AttributeAccessor; +import org.springframework.integration.core.RecoveryCallback; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.MessagingException; -import org.springframework.retry.RecoveryCallback; import org.springframework.retry.RetryContext; import org.springframework.stereotype.Component; import uk.gov.ons.census.caseprocessor.client.ExceptionManagerClient; @@ -19,7 +20,8 @@ import uk.gov.ons.census.caseprocessor.utils.HashHelper; @Component -public class ManagedMessageRecoverer implements RecoveryCallback { +public class ManagedMessageRecoverer + implements RecoveryCallback, org.springframework.retry.RecoveryCallback { private static final Logger log = LoggerFactory.getLogger(ManagedMessageRecoverer.class); private static final String SERVICE_NAME = "Case Processor"; @@ -32,16 +34,23 @@ public ManagedMessageRecoverer(ExceptionManagerClient exceptionManagerClient) { this.exceptionManagerClient = exceptionManagerClient; } + @Override + public Object recover(AttributeAccessor context, Throwable throwable) { + return recoverFromThrowable(throwable); + } + @Override public Object recover(RetryContext retryContext) { - if (!(retryContext.getLastThrowable() instanceof MessagingException)) { - log.error( - "Super duper unexpected kind of error, so going to fail very noisily", - retryContext.getLastThrowable()); - throw new RuntimeException(retryContext.getLastThrowable()); + return recoverFromThrowable(retryContext.getLastThrowable()); + } + + private Object recoverFromThrowable(Throwable throwable) { + MessagingException messagingException = findMessagingException(throwable); + if (messagingException == null) { + log.error("Super duper unexpected kind of error, so going to fail very noisily", throwable); + throw new RuntimeException(throwable); } - MessagingException messagingException = (MessagingException) retryContext.getLastThrowable(); Message message = messagingException.getFailedMessage(); BasicAcknowledgeablePubsubMessage originalMessage = (BasicAcknowledgeablePubsubMessage) @@ -53,14 +62,9 @@ public Object recover(RetryContext retryContext) { String messageHash = HashHelper.hash(rawMessageBody); - String stackTraceRootCause = findUsefulRootCauseInStackTrace(retryContext.getLastThrowable()); + String stackTraceRootCause = findUsefulRootCauseInStackTrace(throwable); - Throwable cause = retryContext.getLastThrowable(); - if (retryContext.getLastThrowable() != null - && retryContext.getLastThrowable().getCause() != null - && retryContext.getLastThrowable().getCause().getCause() != null) { - cause = retryContext.getLastThrowable().getCause().getCause(); - } + Throwable cause = findReportableCause(messagingException); ExceptionReportResponse reportResult = getExceptionReportResponse(cause, messageHash, stackTraceRootCause, subscriptionName); @@ -71,14 +75,38 @@ public Object recover(RetryContext retryContext) { peekMessage(reportResult, messageHash, rawMessageBody); - logMessage( - reportResult, retryContext.getLastThrowable().getCause(), messageHash, stackTraceRootCause); + logMessage(reportResult, cause, messageHash, stackTraceRootCause); // Reject the original message (auto nack'ed). It will be retried at some future point in time throw new MessageHandlingException( message, "Cannot process this message at this time, but it will be retried"); } + private MessagingException findMessagingException(Throwable throwable) { + Throwable current = throwable; + while (current != null) { + if (current instanceof MessagingException messagingException) { + return messagingException; + } + current = current.getCause(); + } + return null; + } + + private Throwable findReportableCause(MessagingException messagingException) { + Throwable cause = messagingException.getCause(); + + if (cause == null) { + return messagingException; + } + + if (cause.getCause() != null) { + return cause.getCause(); + } + + return cause; + } + private ExceptionReportResponse getExceptionReportResponse( Throwable cause, String messageHash, String stackTraceRootCause, String subscriptionName) { ExceptionReportResponse reportResult = null; diff --git a/src/main/java/uk/gov/ons/census/caseprocessor/utils/JsonHelper.java b/src/main/java/uk/gov/ons/census/caseprocessor/utils/JsonHelper.java index da7d825..d3a26f0 100644 --- a/src/main/java/uk/gov/ons/census/caseprocessor/utils/JsonHelper.java +++ b/src/main/java/uk/gov/ons/census/caseprocessor/utils/JsonHelper.java @@ -2,18 +2,19 @@ import static uk.gov.ons.census.caseprocessor.utils.Constants.ALLOWED_INBOUND_EVENT_SCHEMA_VERSIONS; -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; -import java.io.IOException; +import tools.jackson.core.JacksonException; +import tools.jackson.databind.ObjectMapper; import uk.gov.ons.census.caseprocessor.model.dto.EventDTO; public class JsonHelper { private static final ObjectMapper objectMapper = ObjectMapperFactory.objectMapper(); public static String convertObjectToJson(Object obj) { + // JacksonException is unchecked in Jackson 3. The catch is kept deliberately + // so the failure mode and message are unchanged. try { return objectMapper.writeValueAsString(obj); - } catch (JsonProcessingException e) { + } catch (JacksonException e) { throw new RuntimeException("Failed converting Object To Json", e); } } @@ -22,7 +23,7 @@ public static EventDTO convertJsonBytesToEvent(byte[] bytes) { EventDTO event; try { event = objectMapper.readValue(bytes, EventDTO.class); - } catch (IOException e) { + } catch (JacksonException e) { throw new RuntimeException(e); } diff --git a/src/main/java/uk/gov/ons/census/caseprocessor/utils/ObjectMapperFactory.java b/src/main/java/uk/gov/ons/census/caseprocessor/utils/ObjectMapperFactory.java index 94bb8d7..6838183 100644 --- a/src/main/java/uk/gov/ons/census/caseprocessor/utils/ObjectMapperFactory.java +++ b/src/main/java/uk/gov/ons/census/caseprocessor/utils/ObjectMapperFactory.java @@ -1,17 +1,18 @@ package uk.gov.ons.census.caseprocessor.utils; -import com.fasterxml.jackson.databind.DeserializationFeature; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.fasterxml.jackson.databind.SerializationFeature; -import com.fasterxml.jackson.datatype.jdk8.Jdk8Module; -import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; +import tools.jackson.databind.DeserializationFeature; +import tools.jackson.databind.MapperFeature; +import tools.jackson.databind.ObjectMapper; +import tools.jackson.databind.cfg.DateTimeFeature; +import tools.jackson.databind.json.JsonMapper; public class ObjectMapperFactory { public static ObjectMapper objectMapper() { - return new ObjectMapper() - .registerModule(new JavaTimeModule()) - .registerModule(new Jdk8Module()) - .disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS) - .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); + return JsonMapper.builder() + .disable(MapperFeature.SORT_PROPERTIES_ALPHABETICALLY) + .disable(DateTimeFeature.WRITE_DATES_AS_TIMESTAMPS) + .disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES) + .disable(DeserializationFeature.FAIL_ON_TRAILING_TOKENS) + .build(); } } diff --git a/src/main/java/uk/gov/ons/census/caseprocessor/utils/RedactHelper.java b/src/main/java/uk/gov/ons/census/caseprocessor/utils/RedactHelper.java index a355d13..c963aa2 100644 --- a/src/main/java/uk/gov/ons/census/caseprocessor/utils/RedactHelper.java +++ b/src/main/java/uk/gov/ons/census/caseprocessor/utils/RedactHelper.java @@ -1,7 +1,5 @@ package uk.gov.ons.census.caseprocessor.utils; -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; import java.lang.reflect.InvocationTargetException; import java.lang.reflect.Method; import java.lang.reflect.Modifier; @@ -10,6 +8,8 @@ import lombok.AllArgsConstructor; import lombok.Data; import org.springframework.util.StringUtils; +import tools.jackson.core.JacksonException; +import tools.jackson.databind.ObjectMapper; public class RedactHelper { @@ -40,7 +40,7 @@ public static Object redact(Object rootObjectToRedact) { recursivelyRedact( rootObjectToRedactDeepCopy, rootObjectToRedactDeepCopy.getClass().getPackageName()); return rootObjectToRedactDeepCopy; - } catch (JsonProcessingException e) { + } catch (JacksonException e) { throw new RuntimeException(REDACTION_FAILURE, e); } } diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 44dc284..0d91500 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -4,6 +4,14 @@ spring: pool: size: 30 + jackson: + # Boot 4 migration: freeze the wire format at Jackson 2 behaviour so this + # release is a pure infrastructure change with no downstream coordination. + # Jackson 3's defaults flip property ordering to alphabetical and serialise + # java.time as ISO-8601 strings instead of epoch numbers. + # REMOVE in H2 2027 as its own announced change (ticket). + use-jackson2-defaults: true + datasource: url: jdbc:postgresql://localhost:6432/rm username: appuser diff --git a/src/main/resources/logback.xml b/src/main/resources/logback-spring.xml similarity index 100% rename from src/main/resources/logback.xml rename to src/main/resources/logback-spring.xml diff --git a/src/test/java/uk/gov/ons/census/caseprocessor/config/DefaultListenerSupportTest.java b/src/test/java/uk/gov/ons/census/caseprocessor/config/DefaultListenerSupportTest.java new file mode 100644 index 0000000..87efb64 --- /dev/null +++ b/src/test/java/uk/gov/ons/census/caseprocessor/config/DefaultListenerSupportTest.java @@ -0,0 +1,55 @@ +package uk.gov.ons.census.caseprocessor.config; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.mockito.Mockito.mock; + +import org.junit.jupiter.api.Test; + +class DefaultListenerSupportTest { + private final DefaultListenerSupport underTest = new DefaultListenerSupport(); + + @Test + void shouldSupportLegacyRetryableCallbacks() { + org.springframework.retry.RetryContext retryContext = + mock(org.springframework.retry.RetryContext.class); + @SuppressWarnings("unchecked") + org.springframework.retry.RetryCallback retryCallback = + (org.springframework.retry.RetryCallback) + mock(org.springframework.retry.RetryCallback.class); + + assertThat(underTest.open(retryContext, retryCallback)).isTrue(); + assertThatCode( + () -> { + underTest.onError(retryContext, retryCallback, new RuntimeException("failure")); + underTest.close(retryContext, retryCallback, new RuntimeException("failure")); + }) + .doesNotThrowAnyException(); + } + + @Test + void shouldAllowCoreRetryListenerDefaultsToBeCalled() { + org.springframework.core.retry.RetryListener coreRetryListener = underTest; + org.springframework.core.retry.RetryPolicy retryPolicy = + mock(org.springframework.core.retry.RetryPolicy.class); + @SuppressWarnings("unchecked") + org.springframework.core.retry.Retryable retryable = + (org.springframework.core.retry.Retryable) + mock(org.springframework.core.retry.Retryable.class); + + assertThatCode( + () -> { + coreRetryListener.beforeRetry(retryPolicy, retryable); + coreRetryListener.onRetryFailure( + retryPolicy, retryable, new RuntimeException("failure")); + }) + .doesNotThrowAnyException(); + } + + @Test + void shouldImplementBothRetryListenerContracts() { + assertThat(underTest) + .isInstanceOf(org.springframework.core.retry.RetryListener.class) + .isInstanceOf(org.springframework.retry.RetryListener.class); + } +} diff --git a/src/test/java/uk/gov/ons/census/caseprocessor/config/MessageConsumerConfigTest.java b/src/test/java/uk/gov/ons/census/caseprocessor/config/MessageConsumerConfigTest.java new file mode 100644 index 0000000..e660751 --- /dev/null +++ b/src/test/java/uk/gov/ons/census/caseprocessor/config/MessageConsumerConfigTest.java @@ -0,0 +1,74 @@ +package uk.gov.ons.census.caseprocessor.config; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import com.google.cloud.spring.pubsub.core.PubSubTemplate; +import java.util.Objects; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.core.retry.RetryException; +import org.springframework.core.retry.RetryListener; +import org.springframework.core.retry.RetryTemplate; +import org.springframework.integration.handler.advice.RequestHandlerRetryAdvice; +import org.springframework.test.util.ReflectionTestUtils; +import uk.gov.ons.census.caseprocessor.messaging.ManagedMessageRecoverer; + +@ExtendWith(MockitoExtension.class) +class MessageConsumerConfigTest { + + @Mock private ManagedMessageRecoverer managedMessageRecoverer; + + @Mock private PubSubTemplate pubSubTemplate; + + @Test + void shouldCreateRetryAdviceUsingRecovererCallback() { + MessageConsumerConfig underTest = + new MessageConsumerConfig(managedMessageRecoverer, pubSubTemplate); + + RequestHandlerRetryAdvice retryAdvice = underTest.retryAdvice(); + + assertThat( + org.springframework.test.util.ReflectionTestUtils.getField( + retryAdvice, "recoveryCallback")) + .isEqualTo(managedMessageRecoverer); + } + + @Test + void shouldRetryThreeTotalInvocations() { + MessageConsumerConfig underTest = + new MessageConsumerConfig(managedMessageRecoverer, pubSubTemplate); + + RequestHandlerRetryAdvice retryAdvice = underTest.retryAdvice(); + RetryTemplate retryTemplate = + (RetryTemplate) + Objects.requireNonNull(ReflectionTestUtils.getField(retryAdvice, "retryTemplate")); + AtomicInteger attempts = new AtomicInteger(); + + assertThrows( + RetryException.class, + () -> + retryTemplate.execute( + () -> { + attempts.incrementAndGet(); + throw new IllegalStateException("test"); + })); + + assertThat(attempts).hasValue(3); + } + + @Test + void shouldExposeDefaultListenerSupportAsTheCoreRetryListener() { + MessageConsumerConfig underTest = + new MessageConsumerConfig(managedMessageRecoverer, pubSubTemplate); + + RetryListener retryListener = underTest.retryListener(); + + assertThat(retryListener).isInstanceOf(DefaultListenerSupport.class); + assertThat(retryListener).isInstanceOf(org.springframework.core.retry.RetryListener.class); + assertThat(retryListener).isInstanceOf(org.springframework.retry.RetryListener.class); + } +} diff --git a/src/test/java/uk/gov/ons/census/caseprocessor/messaging/FulfilmentRequestReceiverIT.java b/src/test/java/uk/gov/ons/census/caseprocessor/messaging/FulfilmentRequestReceiverIT.java index 5895f81..9faf0b4 100644 --- a/src/test/java/uk/gov/ons/census/caseprocessor/messaging/FulfilmentRequestReceiverIT.java +++ b/src/test/java/uk/gov/ons/census/caseprocessor/messaging/FulfilmentRequestReceiverIT.java @@ -8,8 +8,6 @@ import ch.qos.logback.classic.Logger; import ch.qos.logback.classic.spi.ILoggingEvent; import ch.qos.logback.core.read.ListAppender; -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; import java.time.OffsetDateTime; import java.util.List; import java.util.Map; @@ -28,6 +26,8 @@ import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit.jupiter.SpringExtension; import org.springframework.transaction.annotation.Transactional; +import tools.jackson.core.JacksonException; +import tools.jackson.databind.ObjectMapper; import uk.gov.ons.census.caseprocessor.model.dto.*; import uk.gov.ons.census.caseprocessor.model.repository.*; import uk.gov.ons.census.caseprocessor.service.FulfilmentRequestService; @@ -240,7 +240,7 @@ void testFulfilmentRequestForSms() throws InterruptedException { } @Test - void testSmsRequestEnrichedReceiver() throws InterruptedException, JsonProcessingException { + void testSmsRequestEnrichedReceiver() throws InterruptedException, JacksonException { // Given // Set up all the data required Survey survey = new Survey(); diff --git a/src/test/java/uk/gov/ons/census/caseprocessor/messaging/FulfilmentRequestReceiverTest.java b/src/test/java/uk/gov/ons/census/caseprocessor/messaging/FulfilmentRequestReceiverTest.java index 7aaf646..4f46ec1 100644 --- a/src/test/java/uk/gov/ons/census/caseprocessor/messaging/FulfilmentRequestReceiverTest.java +++ b/src/test/java/uk/gov/ons/census/caseprocessor/messaging/FulfilmentRequestReceiverTest.java @@ -1,11 +1,11 @@ package uk.gov.ons.census.caseprocessor.messaging; -import static org.junit.jupiter.api.Assertions.*; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.*; -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; import java.util.Optional; import java.util.UUID; import org.junit.jupiter.api.Test; @@ -15,6 +15,9 @@ import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; +import tools.jackson.core.JacksonException; +import tools.jackson.databind.ObjectMapper; +import tools.jackson.databind.json.JsonMapper; import uk.gov.ons.census.caseprocessor.logging.EventLogger; import uk.gov.ons.census.caseprocessor.model.dto.*; import uk.gov.ons.census.caseprocessor.model.repository.FulfilmentToProcessRepository; @@ -489,8 +492,8 @@ void testReceiveMessage_sms_fulfilment_success_case_not_HH() throws Exception { eq(msg)); // TODO: Check warning and fix it. } - private Message buildMessage(EventDTO event) throws JsonProcessingException { - ObjectMapper mapper = new ObjectMapper(); + private Message buildMessage(EventDTO event) throws JacksonException { + final ObjectMapper mapper = JsonMapper.builder().build(); byte[] payload = mapper.writeValueAsBytes(event); return MessageBuilder.withPayload(payload).build(); } diff --git a/src/test/java/uk/gov/ons/census/caseprocessor/messaging/ManagedMessageRecovererTest.java b/src/test/java/uk/gov/ons/census/caseprocessor/messaging/ManagedMessageRecovererTest.java index 781b554..58fb46c 100644 --- a/src/test/java/uk/gov/ons/census/caseprocessor/messaging/ManagedMessageRecovererTest.java +++ b/src/test/java/uk/gov/ons/census/caseprocessor/messaging/ManagedMessageRecovererTest.java @@ -22,10 +22,13 @@ import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.core.AttributeAccessor; +import org.springframework.core.retry.RetryException; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.MessagingException; +import org.springframework.messaging.support.MessageBuilder; import org.springframework.retry.RetryContext; import uk.gov.ons.census.caseprocessor.client.ExceptionManagerClient; import uk.gov.ons.census.caseprocessor.model.dto.ExceptionReportResponse; @@ -182,10 +185,103 @@ public void testRecoverPeek() { verify(exceptionManagerClient).respondToPeek(TEST_MESSAGE_HASH, "TEST PAYLOAD".getBytes()); } + @Test + void testRecoverUnwrapsRetryExceptionAndReportsUnderlyingCause() { + RetryContext retryContext = mock(RetryContext.class); + + ProjectSubscriptionName projectSubscriptionName = mock(ProjectSubscriptionName.class); + when(originalMessage.getProjectSubscriptionName()).thenReturn(projectSubscriptionName); + when(projectSubscriptionName.getSubscription()).thenReturn("TEST SUBSCRIPTION"); + + ByteString byteString = ByteString.copyFrom("TEST PAYLOAD".getBytes()); + PubsubMessage pubsubMessage = PubsubMessage.newBuilder().setData(byteString).build(); + when(originalMessage.getPubsubMessage()).thenReturn(pubsubMessage); + + Message message = + MessageBuilder.withPayload("TEST PAYLOAD".getBytes()) + .setHeader("gcp_pubsub_original_message", originalMessage) + .build(); + MessagingException messagingException = + new MessageHandlingException(message, new RuntimeException("qid '555555' not found!")); + + when(retryContext.getLastThrowable()) + .thenReturn(new RetryException("retry exhausted", messagingException)); + + when(exceptionManagerClient.reportException( + anyString(), anyString(), anyString(), any(Throwable.class), anyString())) + .thenReturn(new ExceptionReportResponse()); + + MessageHandlingException thrownException = + assertThrows(MessageHandlingException.class, () -> underTest.recover(retryContext)); + + ArgumentCaptor causeCaptor = ArgumentCaptor.forClass(Throwable.class); + verify(exceptionManagerClient) + .reportException( + eq(TEST_MESSAGE_HASH), + eq("Case Processor"), + eq("TEST SUBSCRIPTION"), + causeCaptor.capture(), + anyString()); + + assertThat(causeCaptor.getValue()).isInstanceOf(RuntimeException.class); + assertThat(causeCaptor.getValue().getMessage()).isEqualTo("qid '555555' not found!"); + assertThat(thrownException.getMessage()) + .isEqualTo("Cannot process this message at this time, but it will be retried"); + } + + @Test + void testRecoverWithAttributeAccessorUnwrapsRetryExceptionAndReportsUnderlyingCause() { + AttributeAccessor context = mock(AttributeAccessor.class); + + ProjectSubscriptionName projectSubscriptionName = mock(ProjectSubscriptionName.class); + when(originalMessage.getProjectSubscriptionName()).thenReturn(projectSubscriptionName); + when(projectSubscriptionName.getSubscription()).thenReturn("TEST SUBSCRIPTION"); + + ByteString byteString = ByteString.copyFrom("TEST PAYLOAD".getBytes()); + PubsubMessage pubsubMessage = PubsubMessage.newBuilder().setData(byteString).build(); + when(originalMessage.getPubsubMessage()).thenReturn(pubsubMessage); + + Message message = + MessageBuilder.withPayload("TEST PAYLOAD".getBytes()) + .setHeader("gcp_pubsub_original_message", originalMessage) + .build(); + MessagingException messagingException = + new MessageHandlingException(message, new RuntimeException("qid '555555' not found!")); + + Throwable wrappedFailure = new RetryException("retry exhausted", messagingException); + + when(exceptionManagerClient.reportException( + anyString(), anyString(), anyString(), any(Throwable.class), anyString())) + .thenReturn(new ExceptionReportResponse()); + + MessageHandlingException thrownException = + assertThrows( + MessageHandlingException.class, () -> underTest.recover(context, wrappedFailure)); + + ArgumentCaptor causeCaptor = ArgumentCaptor.forClass(Throwable.class); + verify(exceptionManagerClient) + .reportException( + eq(TEST_MESSAGE_HASH), + eq("Case Processor"), + eq("TEST SUBSCRIPTION"), + causeCaptor.capture(), + anyString()); + + assertThat(causeCaptor.getValue()).isInstanceOf(RuntimeException.class); + assertThat(causeCaptor.getValue().getMessage()).isEqualTo("qid '555555' not found!"); + assertThat(thrownException.getMessage()) + .isEqualTo("Cannot process this message at this time, but it will be retried"); + } + private RetryContext testSetupTestRecover(ExceptionReportResponse exceptionReportResponse) { + return testSetupTestRecover( + exceptionReportResponse, new RuntimeException(new RuntimeException("TEST EXCEPTION"))); + } + + private RetryContext testSetupTestRecover( + ExceptionReportResponse exceptionReportResponse, Throwable reportedCause) { MessagingException messagingException = mock(MessagingException.class); - when(messagingException.getCause()) - .thenReturn(new RuntimeException(new RuntimeException("TEST EXCEPTION"))); + when(messagingException.getCause()).thenReturn(reportedCause); RetryContext retryContext = mock(RetryContext.class); when(retryContext.getLastThrowable()).thenReturn(messagingException); diff --git a/src/test/java/uk/gov/ons/census/caseprocessor/testutils/JsonHelper.java b/src/test/java/uk/gov/ons/census/caseprocessor/testutils/JsonHelper.java index d2288cb..aa857d3 100644 --- a/src/test/java/uk/gov/ons/census/caseprocessor/testutils/JsonHelper.java +++ b/src/test/java/uk/gov/ons/census/caseprocessor/testutils/JsonHelper.java @@ -1,7 +1,7 @@ package uk.gov.ons.census.caseprocessor.testutils; -import com.fasterxml.jackson.databind.ObjectMapper; -import java.io.IOException; +import tools.jackson.core.JacksonException; +import tools.jackson.databind.ObjectMapper; import uk.gov.ons.census.caseprocessor.utils.ObjectMapperFactory; public class JsonHelper { @@ -10,7 +10,7 @@ public class JsonHelper { public static T convertJsonBytesToObject(byte[] bytes, Class clazz) { try { return objectMapper.readValue(bytes, clazz); - } catch (IOException e) { + } catch (JacksonException e) { throw new RuntimeException(e); } } diff --git a/src/test/java/uk/gov/ons/census/caseprocessor/testutils/PubsubHelper.java b/src/test/java/uk/gov/ons/census/caseprocessor/testutils/PubsubHelper.java index 4218b6f..397ae20 100644 --- a/src/test/java/uk/gov/ons/census/caseprocessor/testutils/PubsubHelper.java +++ b/src/test/java/uk/gov/ons/census/caseprocessor/testutils/PubsubHelper.java @@ -4,11 +4,9 @@ import static com.google.cloud.spring.pubsub.support.PubSubTopicUtils.toProjectTopicName; import static uk.gov.ons.census.caseprocessor.testutils.TestConstants.OUR_PUBSUB_PROJECT; -import com.fasterxml.jackson.databind.ObjectMapper; import com.google.cloud.pubsub.v1.Subscriber; import com.google.cloud.spring.autoconfigure.pubsub.GcpPubSubProperties; import com.google.cloud.spring.pubsub.core.PubSubTemplate; -import java.io.IOException; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.CompletableFuture; @@ -27,6 +25,8 @@ import org.springframework.test.context.ActiveProfiles; import org.springframework.web.client.HttpClientErrorException; import org.springframework.web.client.RestTemplate; +import tools.jackson.core.JacksonException; +import tools.jackson.databind.ObjectMapper; import uk.gov.ons.census.caseprocessor.utils.ObjectMapperFactory; @Component @@ -62,7 +62,7 @@ public QueueSpy listen(String subscription, Class contentClass) { message.getPubsubMessage().getData().toByteArray(), contentClass); queue.add(messageObject); message.ack(); - } catch (IOException e) { + } catch (JacksonException e) { System.out.println("ERROR: Cannot unmarshal bad data on PubSub subscription"); } finally { // Always want to ack, to get rid of dodgy messages @@ -116,7 +116,7 @@ private void purgeMessages(String subscription, String topic, String project) { // There's no concept of a 'purge' with pubsub. Crudely, we have to delete & recreate restTemplate.delete(subscriptionUrl); } catch (HttpClientErrorException exception) { - if (exception.getRawStatusCode() != 404) { + if (exception.getStatusCode().value() != 404) { throw exception; } } @@ -125,7 +125,7 @@ private void purgeMessages(String subscription, String topic, String project) { restTemplate.put( subscriptionUrl, new SubscriptionTopic("projects/" + project + "/topics/" + topic)); } catch (HttpClientErrorException exception) { - if (exception.getRawStatusCode() != 409) { + if (exception.getStatusCode().value() != 409) { throw exception; } } diff --git a/src/test/java/uk/gov/ons/census/caseprocessor/utils/ObjectMapperFactoryTest.java b/src/test/java/uk/gov/ons/census/caseprocessor/utils/ObjectMapperFactoryTest.java new file mode 100644 index 0000000..93c947f --- /dev/null +++ b/src/test/java/uk/gov/ons/census/caseprocessor/utils/ObjectMapperFactoryTest.java @@ -0,0 +1,36 @@ +package uk.gov.ons.census.caseprocessor.utils; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.time.Instant; +import org.junit.jupiter.api.Test; +import tools.jackson.databind.ObjectMapper; + +class ObjectMapperFactoryTest { + private static final ObjectMapper OBJECT_MAPPER = ObjectMapperFactory.objectMapper(); + private static final String CASE_ID = "10000000001"; + private static final Instant EVENT_TIME = Instant.parse("2024-06-01T10:15:30Z"); + private static final String ACTION = "CREATE"; + + @Test + void shouldPreserveFieldOrderAndWriteDatesAsIsoStrings() { + CaseEvent caseEvent = new CaseEvent(CASE_ID, EVENT_TIME, ACTION); + + assertThat(OBJECT_MAPPER.writeValueAsString(caseEvent)) + .isEqualTo( + "{\"caseId\":\"10000000001\",\"eventTime\":\"2024-06-01T10:15:30Z\",\"action\":\"CREATE\"}"); + } + + @Test + void shouldIgnoreUnknownPropertiesWhenReadingJson() { + CaseEvent caseEvent = + OBJECT_MAPPER.readValue( + "{\"caseId\":\"10000000001\",\"eventTime\":\"2024-06-01T10:15:30Z\"," + + "\"action\":\"CREATE\",\"extraField\":\"ignored\"}", + CaseEvent.class); + + assertThat(caseEvent).isEqualTo(new CaseEvent(CASE_ID, EVENT_TIME, ACTION)); + } + + private record CaseEvent(String caseId, Instant eventTime, String action) {} +}