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
93 changes: 51 additions & 42 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -6,34 +6,37 @@

<groupId>uk.gov.ons.census</groupId>
<artifactId>census-rm-caseprocessor</artifactId>
<version>1.0-SNAPSHOT</version>
<version>1.0.0-SNAPSHOT</version>

<properties>
<maven.compiler.release>21</maven.compiler.release>
<java.version>21</java.version>
<maven.compiler.release>${java.version}</maven.compiler.release>

<census-rm-common-entity-model.version>1.0.0</census-rm-common-entity-model.version>
<census-rm-shared-sample-validation.version>1.0.0</census-rm-shared-sample-validation.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<container.cli>docker</container.cli>

<!-- BOM / dependency management versions -->
<spring-cloud.version>2025.0.3</spring-cloud.version>
<spring-cloud-gcp.version>7.4.8</spring-cloud-gcp.version>
<spring-cloud.version>2025.1.3</spring-cloud.version>
<spring-cloud-gcp.version>8.1.0</spring-cloud-gcp.version>

<!-- Dependency versions -->
<lombok.version>1.18.30</lombok.version>
<jakarta-xml-bind-api.version>4.0.0</jakarta-xml-bind-api.version>
<javax-jaxb-api.version>2.3.0</javax-jaxb-api.version>
<hypersistence-utils.version>3.15.4</hypersistence-utils.version>
<commons-validator.version>1.10.1</commons-validator.version>
<logstash-logback-encoder.version>7.4</logstash-logback-encoder.version>
<opencsv.version>5.9</opencsv.version>
<aspectjweaver.version>1.9.24</aspectjweaver.version>
<hypersistence-utils.version>3.15.5</hypersistence-utils.version>
<commons-validator.version>1.11.0</commons-validator.version>
<logstash-logback-encoder.version>9.0</logstash-logback-encoder.version>
<opencsv.version>5.12.0</opencsv.version>
<spring-retry.version>2.0.13</spring-retry.version>

<!-- Plugin versions -->
<maven-pmd-plugin.version>3.24.0</maven-pmd-plugin.version>
<exec-maven-plugin.version>3.1.1</exec-maven-plugin.version>
<spotless-maven-plugin.version>2.43.0</spotless-maven-plugin.version>
<google-java-format.version>1.22.0</google-java-format.version>
<jacoco-maven-plugin.version>0.8.11</jacoco-maven-plugin.version>
<error-prone-core.version>2.23.0</error-prone-core.version>
<maven-pmd-plugin.version>3.28.0</maven-pmd-plugin.version>
<pmd.version>7.26.0</pmd.version>
<exec-maven-plugin.version>3.6.3</exec-maven-plugin.version>
<spotless-maven-plugin.version>3.10.0</spotless-maven-plugin.version>
<google-java-format.version>1.36.1</google-java-format.version>
<jacoco-maven-plugin.version>0.8.15</jacoco-maven-plugin.version>
<error-prone-core.version>2.50.0</error-prone-core.version>
</properties>

<profiles>
Expand All @@ -56,7 +59,7 @@
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.5.16</version>
<version>4.1.1</version>
</parent>

<dependencyManagement>
Expand Down Expand Up @@ -107,17 +110,21 @@
<dependency>
<groupId>uk.gov.ons.census</groupId>
<artifactId>census-rm-common-entity-model</artifactId>
<version>0.0.3</version>
<version>${census-rm-common-entity-model.version}</version>
</dependency>
<dependency>
<groupId>uk.gov.ons.census</groupId>
<artifactId>census-rm-shared-sample-validation</artifactId>
<version>0.1.0</version>
<version>${census-rm-shared-sample-validation.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jackson</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-integration</artifactId>
Expand All @@ -137,21 +144,6 @@
<dependency>
<groupId>jakarta.xml.bind</groupId>
<artifactId>jakarta.xml.bind-api</artifactId>
<version>${jakarta-xml-bind-api.version}</version>
</dependency>
<dependency>
<groupId>javax.xml.bind</groupId>
<artifactId>jaxb-api</artifactId>
<version>${javax-jaxb-api.version}</version>
</dependency>

<dependency>
<groupId>com.fasterxml.jackson.datatype</groupId>
<artifactId>jackson-datatype-jsr310</artifactId>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.datatype</groupId>
<artifactId>jackson-datatype-jdk8</artifactId>
</dependency>
<dependency>
<groupId>org.postgresql</groupId>
Expand All @@ -160,12 +152,11 @@
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>${lombok.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>io.hypersistence</groupId>
<artifactId>hypersistence-utils-hibernate-63</artifactId>
<artifactId>hypersistence-utils-hibernate-73</artifactId>
<version>${hypersistence-utils.version}</version>
</dependency>
<dependency>
Expand All @@ -186,7 +177,6 @@
<dependency>
<groupId>org.aspectj</groupId>
<artifactId>aspectjweaver</artifactId>
<version>${aspectjweaver.version}</version>
</dependency>

<!-- Test Dependencies below this point -->
Expand All @@ -195,6 +185,12 @@
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.retry</groupId>
<artifactId>spring-retry</artifactId>
<version>${spring-retry.version}</version>
</dependency>

</dependencies>

<build>
Expand All @@ -204,8 +200,20 @@
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-pmd-plugin</artifactId>
<version>${maven-pmd-plugin.version}</version>
<dependencies>
<dependency>
<groupId>net.sourceforge.pmd</groupId>
<artifactId>pmd-core</artifactId>
<version>${pmd.version}</version>
</dependency>
<dependency>
<groupId>net.sourceforge.pmd</groupId>
<artifactId>pmd-java</artifactId>
<version>${pmd.version}</version>
</dependency>
</dependencies>
<configuration>
<targetJdk>21</targetJdk>
<targetJdk>${java.version}</targetJdk>
<excludeFromFailureFile>exclude-pmd.properties</excludeFromFailureFile>
<failurePriority>3</failurePriority>
<failOnViolation>true</failOnViolation>
Expand Down Expand Up @@ -270,7 +278,6 @@
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<executable>true</executable>
<mainClass>uk.gov.ons.census.caseprocessor.Application</mainClass>
</configuration>
<executions>
Expand Down Expand Up @@ -345,11 +352,13 @@
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>21</source>
<target>21</target>
<release>${maven.compiler.release}</release>
<encoding>UTF-8</encoding>
<compilerArgs>
<arg>-XDcompilePolicy=simple</arg>
<arg>--should-stop=ifError=FLOW</arg>
<!-- Required by Error Prone 2.44+ on JDK 21 -->
<arg>-XDaddTypeAnnotationsToSymbol=true</arg>
<arg>-Xplugin:ErrorProne</arg>
</compilerArgs>
<annotationProcessorPaths>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 <T, E extends Throwable> void close(
Expand All @@ -15,7 +18,6 @@ public <T, E extends Throwable> void close(
@Override
public <T, E extends Throwable> void onError(
RetryContext context, RetryCallback<T, E> callback, Throwable throwable) {

RetryListener.super.onError(context, callback, throwable);
}

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

Expand Down Expand Up @@ -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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -19,7 +20,8 @@
import uk.gov.ons.census.caseprocessor.utils.HashHelper;

@Component
public class ManagedMessageRecoverer implements RecoveryCallback<Object> {
public class ManagedMessageRecoverer
implements RecoveryCallback<Object>, org.springframework.retry.RecoveryCallback<Object> {
private static final Logger log = LoggerFactory.getLogger(ManagedMessageRecoverer.class);
private static final String SERVICE_NAME = "Case Processor";

Expand All @@ -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)
Expand All @@ -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);
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
Expand All @@ -22,11 +23,11 @@
EventDTO event;
try {
event = objectMapper.readValue(bytes, EventDTO.class);
} catch (IOException e) {
} catch (JacksonException e) {
throw new RuntimeException(e);
}

if (!ALLOWED_INBOUND_EVENT_SCHEMA_VERSIONS.contains((event.getHeader().getVersion()))) {

Check warning on line 30 in src/main/java/uk/gov/ons/census/caseprocessor/utils/JsonHelper.java

View workflow job for this annotation

GitHub Actions / Java Checks and Tests

[UnnecessaryParentheses] These parentheses are unnecessary; it is unlikely the code will be misinterpreted without them
throw new RuntimeException(
String.format(
"Unsupported message version. Got %s but RM only supports %s",
Expand Down
Loading
Loading