diff --git a/.github/scripts/maven_publish.sh b/.github/scripts/maven_publish.sh index 655ad5c69..3c8c29096 100644 --- a/.github/scripts/maven_publish.sh +++ b/.github/scripts/maven_publish.sh @@ -43,6 +43,7 @@ echo "settings.xml written." echo "=== Step 3: Upload to Sonatype Central Portal ===" mvn clean deploy -s "${SETTINGS_FILE}" -pl sdk -P publishing -DskipTests --no-transfer-progress +mvn clean deploy -s "${SETTINGS_FILE}" -pl extra-serdes -P publishing -DskipTests --no-transfer-progress mvn clean deploy -s "${SETTINGS_FILE}" -pl sdk-testing -P publishing -DskipTests --no-transfer-progress mvn clean deploy -s "${SETTINGS_FILE}" -pl otel-plugin -P publishing -DskipTests --no-transfer-progress diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index 49ad9f01e..c155f2585 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -26,6 +26,7 @@ on: - '.github/workflows/ai-pr-review.yml' - '.github/prompts/ai-pr-review.md' - 'sdk/**' + - 'extra-serdes/**' - 'sdk-testing/**' - 'sdk-integration-tests/**' - 'examples/**' @@ -38,6 +39,7 @@ on: - '.github/workflows/ai-pr-review.yml' - '.github/prompts/ai-pr-review.md' - 'sdk/**' + - 'extra-serdes/**' - 'sdk-testing/**' - 'sdk-integration-tests/**' - 'examples/**' diff --git a/.github/workflows/e2e-tests.yml b/.github/workflows/e2e-tests.yml index ae25b371a..5d3ceb149 100644 --- a/.github/workflows/e2e-tests.yml +++ b/.github/workflows/e2e-tests.yml @@ -8,6 +8,7 @@ on: paths: - '.github/**' # for testing Github Actions - 'sdk/**' + - 'extra-serdes/**' - 'sdk-testing/**' - 'sdk-integration-tests/**' - 'examples/**' @@ -18,6 +19,7 @@ on: paths: - '.github/**' - 'sdk/**' + - 'extra-serdes/**' - 'sdk-testing/**' - 'sdk-integration-tests/**' - 'examples/**' diff --git a/.github/workflows/publish_maven.yml b/.github/workflows/publish_maven.yml index e9163696d..50c85c3a9 100644 --- a/.github/workflows/publish_maven.yml +++ b/.github/workflows/publish_maven.yml @@ -90,6 +90,7 @@ jobs: run: | gh release upload "$RELEASE_TAG" \ "sdk/target/aws-durable-execution-sdk-java-${RELEASE_VERSION}.jar" \ + "extra-serdes/target/aws-durable-execution-sdk-java-extra-serdes-${RELEASE_VERSION}.jar" \ "sdk-testing/target/aws-durable-execution-sdk-java-testing-${RELEASE_VERSION}.jar" \ "otel-plugin/target/aws-durable-execution-sdk-java-plugin-otel-${RELEASE_VERSION}.jar" \ --clobber diff --git a/README.md b/README.md index 766a71b02..d5d1c17c9 100644 --- a/README.md +++ b/README.md @@ -23,6 +23,7 @@ Build resilient, long-running AWS Lambda functions that automatically checkpoint - **Replay Safety** – Functions deterministically resume from checkpoints after interruptions - **Type Safety** – Full generic type support for step results - **Data-Driven Concurrency** – Apply a function across a collection with `map()`, with per-item error isolation and configurable completion criteria +- **Optional Lambda Event Models** – Parse SQS, SNS, S3, and other Lambda trigger events with the Java Lambda runtime mappings ## How It Works @@ -50,6 +51,20 @@ Your durable function extends `DurableHandler` and implements `handleReque ``` +For handlers that receive models from `aws-lambda-java-events`, also add the +optional event serialization module: + +```xml + + software.amazon.lambda.durable + aws-durable-execution-sdk-java-extra-serdes + VERSION + +``` + +Configure the module with +`DurableConfig.builder().withSerDes(new LambdaEventSerDes()).build()`. + ### Your First Durable Function ```java diff --git a/RELEASE.md b/RELEASE.md index 90c347ff1..2289bb177 100644 --- a/RELEASE.md +++ b/RELEASE.md @@ -43,9 +43,9 @@ The publication workflow: 1. Verifies that the tag is a semantic version, points to a commit on the default branch, and matches the Maven version in the tagged POM. -2. Builds, signs, and uploads the SDK, testing library, and OpenTelemetry plugin - to Sonatype Central Portal. -3. Uploads the three JARs to the existing GitHub release. +2. Builds, signs, and uploads the SDK, extra SerDes module, testing library, and + OpenTelemetry plugin to Sonatype Central Portal. +3. Uploads the four JARs to the existing GitHub release. 4. Opens a pull request for the next development version. A final release increments the patch version, so `2.1.1` produces `2.1.2-SNAPSHOT`. A prerelease keeps the same base version, so `2.1.1-rc1` produces @@ -56,7 +56,8 @@ After **Publish Maven Release** succeeds: 1. Open [Publishing Deployments](https://central.sonatype.com/publishing/deployments) in Sonatype Central Portal. 2. Find the deployments for the release version and verify that they contain - the expected SDK, testing library, and OpenTelemetry plugin artifacts. + the expected SDK, extra SerDes module, testing library, and OpenTelemetry + plugin artifacts. 3. Click **Publish** for each deployment and wait for publication to complete. The workflow uses `autoPublish=false`, so this manual action is required. 4. Confirm that the GitHub release contains the expected JARs and that the diff --git a/coverage-report/pom.xml b/coverage-report/pom.xml index 48300f209..0593b177b 100644 --- a/coverage-report/pom.xml +++ b/coverage-report/pom.xml @@ -27,6 +27,11 @@ aws-durable-execution-sdk-java-testing ${project.version} + + software.amazon.lambda.durable + aws-durable-execution-sdk-java-extra-serdes + ${project.version} + software.amazon.lambda.durable aws-durable-execution-sdk-java-integration-tests diff --git a/docs/advanced/configuration.md b/docs/advanced/configuration.md index bc8bf89d0..4e4db0a7e 100644 --- a/docs/advanced/configuration.md +++ b/docs/advanced/configuration.md @@ -40,6 +40,24 @@ public class OrderProcessor extends DurableHandler { The `withExecutorService()` option configures the thread pool used for running user-defined operations. Internal SDK coordination (checkpoint batching, polling) runs on an SDK-managed thread pool. +### Lambda trigger event inputs + +The optional `aws-durable-execution-sdk-java-extra-serdes` module provides +`LambdaEventSerDes`, which applies the official Java Lambda runtime mappings +for `SQSEvent`, `SNSEvent`, `S3Event`, and other supported event models: + +```java +@Override +protected DurableConfig createConfiguration() { + return DurableConfig.builder() + .withSerDes(new LambdaEventSerDes()) + .build(); +} +``` + +The module delegates non-event values, including generic types, to +`JacksonSerDes`. See the module README for installation details. + ### Dynamic plugin loading Dynamic plugin loading is an opt-in alternative to registering plugins in application code. Put provider JARs on the application class path, then set `DURABLE_EXECUTION_PLUGINS` to an ordered, comma-separated list of provider names: diff --git a/examples/README.md b/examples/README.md index c17750835..f0443fcee 100644 --- a/examples/README.md +++ b/examples/README.md @@ -87,6 +87,7 @@ mvn test -Dtest=CloudBasedIntegrationTest \ | [ErrorHandlingExample](src/main/java/software/amazon/lambda/durable/examples/general/ErrorHandlingExample.java) | Handling `StepFailedException` and `StepInterruptedException` | | [GenericTypesExample](src/main/java/software/amazon/lambda/durable/examples/general/GenericTypesExample.java) | Working with `List` and `Map` | | [CustomConfigExample](src/main/java/software/amazon/lambda/durable/examples/general/CustomConfigExample.java) | Custom Lambda client and SerDes | +| [LambdaEventSerDesExample](src/main/java/software/amazon/lambda/durable/examples/general/LambdaEventSerDesExample.java) | Deserializing SQS events with the optional extra SerDes module | | [WaitAtLeastExample](src/main/java/software/amazon/lambda/durable/examples/wait/WaitAtLeastExample.java) | Concurrent `stepAsync()` with `wait()` | | [WaitAsyncExample](src/main/java/software/amazon/lambda/durable/examples/wait/WaitAsyncExample.java) | Non-blocking `waitAsync()` with concurrent step | | [RetryInProcessExample](src/main/java/software/amazon/lambda/durable/examples/step/RetryInProcessExample.java) | In-process retry with concurrent operations | diff --git a/examples/pom.xml b/examples/pom.xml index 18654cf8e..56c28b46e 100644 --- a/examples/pom.xml +++ b/examples/pom.xml @@ -30,6 +30,11 @@ aws-durable-execution-sdk-java ${project.version} + + software.amazon.lambda.durable + aws-durable-execution-sdk-java-extra-serdes + ${project.version} + diff --git a/examples/src/main/java/software/amazon/lambda/durable/examples/general/LambdaEventSerDesExample.java b/examples/src/main/java/software/amazon/lambda/durable/examples/general/LambdaEventSerDesExample.java new file mode 100644 index 000000000..6e3ad7561 --- /dev/null +++ b/examples/src/main/java/software/amazon/lambda/durable/examples/general/LambdaEventSerDesExample.java @@ -0,0 +1,32 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.examples.general; + +import com.amazonaws.services.lambda.runtime.events.SQSEvent; +import software.amazon.lambda.durable.DurableConfig; +import software.amazon.lambda.durable.DurableContext; +import software.amazon.lambda.durable.DurableHandler; +import software.amazon.lambda.durable.events.LambdaEventSerDes; + +/** + * Example demonstrating Lambda runtime serialization for an SQS event. + * + *

The extra SerDes module preserves Lambda event property mappings such as {@code eventSourceARN} when converting + * the durable execution input into {@link SQSEvent}. + */ +public class LambdaEventSerDesExample extends DurableHandler { + + @Override + protected DurableConfig createConfiguration() { + return DurableConfig.builder().withSerDes(new LambdaEventSerDes()).build(); + } + + @Override + public String handleRequest(SQSEvent input, DurableContext context) { + var message = input.getRecords().get(0); + return context.step( + "read-sqs-message", + String.class, + stepContext -> message.getMessageId() + "|" + message.getBody() + "|" + message.getEventSourceArn()); + } +} diff --git a/examples/src/test/java/software/amazon/lambda/durable/examples/CloudBasedIntegrationTest.java b/examples/src/test/java/software/amazon/lambda/durable/examples/CloudBasedIntegrationTest.java index 7490961a5..72eea2ac2 100644 --- a/examples/src/test/java/software/amazon/lambda/durable/examples/CloudBasedIntegrationTest.java +++ b/examples/src/test/java/software/amazon/lambda/durable/examples/CloudBasedIntegrationTest.java @@ -5,6 +5,7 @@ import static org.junit.jupiter.api.Assertions.*; import static software.amazon.lambda.durable.TypeToken.get; +import com.amazonaws.services.lambda.runtime.events.SQSEvent; import java.time.Duration; import java.util.HashMap; import java.util.List; @@ -24,6 +25,7 @@ import software.amazon.awssdk.services.lambda.model.OperationStatus; import software.amazon.awssdk.services.sts.StsClient; import software.amazon.lambda.durable.TypeToken; +import software.amazon.lambda.durable.events.LambdaEventSerDes; import software.amazon.lambda.durable.examples.general.GenericTypesExample; import software.amazon.lambda.durable.examples.types.ApprovalRequest; import software.amazon.lambda.durable.examples.types.GreetingRequest; @@ -363,6 +365,32 @@ void testCustomConfigExample() { assertTrue(stepResult.contains("email_address")); } + @Test + void testLambdaEventSerDesExample() { + var message = new SQSEvent.SQSMessage(); + message.setMessageId("cloud-message-1"); + message.setBody("hello from cloud sqs"); + message.setEventSourceArn("arn:aws:sqs:us-west-2:123456789012:orders"); + + var event = new SQSEvent(); + event.setRecords(List.of(message)); + + var runner = CloudDurableTestRunner.create( + arn("lambda-event-ser-des-example"), SQSEvent.class, String.class, lambdaClient) + .withSerDes(new LambdaEventSerDes()); + var result = runner.run(event); + + assertEquals(ExecutionStatus.SUCCEEDED, result.getStatus()); + assertEquals( + "cloud-message-1|hello from cloud sqs|arn:aws:sqs:us-west-2:123456789012:orders", result.getResult()); + + var operation = runner.getOperation("read-sqs-message"); + assertNotNull(operation); + assertEquals( + "cloud-message-1|hello from cloud sqs|arn:aws:sqs:us-west-2:123456789012:orders", + operation.getStepResult(String.class)); + } + @Test void testErrorHandlingExample() { var runner = diff --git a/examples/src/test/java/software/amazon/lambda/durable/examples/general/LambdaEventSerDesExampleTest.java b/examples/src/test/java/software/amazon/lambda/durable/examples/general/LambdaEventSerDesExampleTest.java new file mode 100644 index 000000000..9ee5e1345 --- /dev/null +++ b/examples/src/test/java/software/amazon/lambda/durable/examples/general/LambdaEventSerDesExampleTest.java @@ -0,0 +1,33 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.examples.general; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import com.amazonaws.services.lambda.runtime.events.SQSEvent; +import java.util.List; +import org.junit.jupiter.api.Test; +import software.amazon.lambda.durable.model.ExecutionStatus; +import software.amazon.lambda.durable.testing.LocalDurableTestRunner; + +class LambdaEventSerDesExampleTest { + + @Test + void deserializesSqsEventWithLambdaRuntimeMappings() { + var message = new SQSEvent.SQSMessage(); + message.setMessageId("message-1"); + message.setBody("hello from sqs"); + message.setEventSourceArn("arn:aws:sqs:us-east-1:123456789012:orders"); + + var event = new SQSEvent(); + event.setRecords(List.of(message)); + + var result = LocalDurableTestRunner.create(SQSEvent.class, new LambdaEventSerDesExample()) + .run(event); + + assertEquals(ExecutionStatus.SUCCEEDED, result.getStatus()); + assertTrue(result.getResult(String.class).contains("message-1|hello from sqs|")); + assertTrue(result.getResult(String.class).contains("arn:aws:sqs:us-east-1:123456789012:orders")); + } +} diff --git a/extra-serdes/README.md b/extra-serdes/README.md new file mode 100644 index 000000000..e6b251f7d --- /dev/null +++ b/extra-serdes/README.md @@ -0,0 +1,41 @@ +# AWS Lambda Durable Execution Extra SerDes + +The `aws-durable-execution-sdk-java-extra-serdes` module provides +`LambdaEventSerDes`, which uses the Java Lambda runtime serializers for event +models from `aws-lambda-java-events`. This is useful for durable handlers that +receive SQS, SNS, S3, or other supported Lambda trigger events. + +The module is optional so applications that do not use Lambda trigger models do +not need the event model and runtime serialization dependencies. + +## Installation + +```xml + + software.amazon.lambda.durable + aws-durable-execution-sdk-java-extra-serdes + VERSION + +``` + +## Configuration + +```java +import software.amazon.lambda.durable.DurableConfig; +import software.amazon.lambda.durable.events.LambdaEventSerDes; + +@Override +protected DurableConfig createConfiguration() { + return DurableConfig.builder() + .withSerDes(new LambdaEventSerDes()) + .build(); +} +``` + +`LambdaEventSerDes` uses the official `aws-lambda-java-serialization` mappings +for supported Lambda event classes and delegates all other values to +`JacksonSerDes`. A custom delegate can be supplied for non-event values: + +```java +new LambdaEventSerDes(customSerDes) +``` diff --git a/extra-serdes/pom.xml b/extra-serdes/pom.xml new file mode 100644 index 000000000..b81a6a521 --- /dev/null +++ b/extra-serdes/pom.xml @@ -0,0 +1,84 @@ + + + 4.0.0 + + + software.amazon.lambda.durable + aws-durable-execution-sdk-java-parent + 2.2.1-SNAPSHOT + + + aws-durable-execution-sdk-java-extra-serdes + jar + + AWS Lambda Durable Execution SDK Extra SerDes + Additional serialization support for AWS Lambda event models + https://github.com/aws/aws-durable-execution-sdk-java + + + scm:git:https://github.com/aws/aws-durable-execution-sdk-java.git + scm:git:https://github.com/aws/aws-durable-execution-sdk-java.git + https://github.com/aws/aws-durable-execution-sdk-java + + + + + software.amazon.lambda.durable + aws-durable-execution-sdk-java + ${project.version} + + + com.amazonaws + aws-lambda-java-events + + + com.amazonaws + aws-lambda-java-serialization + + + + org.junit.jupiter + junit-jupiter + test + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + + + org.apache.maven.plugins + maven-surefire-plugin + + + org.apache.maven.plugins + maven-source-plugin + + + attach-sources + + jar-no-fork + + + + + + org.apache.maven.plugins + maven-javadoc-plugin + + + attach-javadocs + + jar + + + + + + + diff --git a/extra-serdes/src/main/java/software/amazon/lambda/durable/events/LambdaEventSerDes.java b/extra-serdes/src/main/java/software/amazon/lambda/durable/events/LambdaEventSerDes.java new file mode 100644 index 000000000..074a54c47 --- /dev/null +++ b/extra-serdes/src/main/java/software/amazon/lambda/durable/events/LambdaEventSerDes.java @@ -0,0 +1,96 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.events; + +import com.amazonaws.services.lambda.runtime.serialization.PojoSerializer; +import com.amazonaws.services.lambda.runtime.serialization.events.LambdaEventSerializers; +import java.io.ByteArrayOutputStream; +import java.nio.charset.StandardCharsets; +import java.util.Map; +import java.util.Objects; +import java.util.concurrent.ConcurrentHashMap; +import software.amazon.lambda.durable.TypeToken; +import software.amazon.lambda.durable.exception.SerDesException; +import software.amazon.lambda.durable.serde.JacksonSerDes; +import software.amazon.lambda.durable.serde.SerDes; + +/** + * A {@link SerDes} that uses the AWS Lambda Java runtime mappings for supported Lambda event models. + * + *

Lambda event payloads do not always follow the Java bean property names in {@code aws-lambda-java-events}. For + * example, an SQS payload uses {@code Records} and {@code eventSourceARN}, while {@code SQSEvent} exposes + * {@code records} and {@code eventSourceArn}. This implementation uses {@code aws-lambda-java-serialization} for those + * event classes so durable handlers receive the same model shape as standard Java Lambda handlers. + * + *

All other values, including generic types, are handled by a delegate {@link SerDes}. The default delegate is + * {@link JacksonSerDes}. + */ +public final class LambdaEventSerDes implements SerDes { + private final SerDes delegate; + private final Map, PojoSerializer> eventSerializers = new ConcurrentHashMap<>(); + + /** Creates a LambdaEventSerDes that delegates non-event values to {@link JacksonSerDes}. */ + public LambdaEventSerDes() { + this(new JacksonSerDes()); + } + + /** + * Creates a LambdaEventSerDes with a custom delegate for non-event values. + * + * @param delegate serializer used for values that are not supported Lambda event model classes + */ + public LambdaEventSerDes(SerDes delegate) { + this.delegate = Objects.requireNonNull(delegate, "Delegate SerDes cannot be null"); + } + + @Override + public String serialize(Object value) { + if (value == null) { + return null; + } + + var valueClass = value.getClass(); + if (!LambdaEventSerializers.isLambdaSupportedEvent(valueClass.getName())) { + return delegate.serialize(value); + } + + try { + var output = new ByteArrayOutputStream(); + PojoSerializer serializer = getEventSerializer(valueClass); + serializer.toJson(value, output); + return output.toString(StandardCharsets.UTF_8); + } catch (Exception e) { + throw new SerDesException("Serialization failed for Lambda event type: " + valueClass.getName(), e); + } + } + + @Override + public T deserialize(String data, TypeToken typeToken) { + if (data == null) { + return null; + } + + if (!(typeToken.getType() instanceof Class targetClass) + || !LambdaEventSerializers.isLambdaSupportedEvent(targetClass.getName())) { + return delegate.deserialize(data, typeToken); + } + + try { + PojoSerializer serializer = getEventSerializer(targetClass); + return serializer.fromJson(data); + } catch (Exception e) { + throw new SerDesException("Deserialization failed for Lambda event type: " + targetClass.getName(), e); + } + } + + @SuppressWarnings("unchecked") + private PojoSerializer getEventSerializer(Class eventClass) { + return (PojoSerializer) + eventSerializers.computeIfAbsent(eventClass, LambdaEventSerDes::createEventSerializer); + } + + @SuppressWarnings({"rawtypes", "unchecked"}) + private static synchronized PojoSerializer createEventSerializer(Class eventClass) { + return LambdaEventSerializers.serializerFor((Class) eventClass, eventClass.getClassLoader()); + } +} diff --git a/extra-serdes/src/test/java/software/amazon/lambda/durable/events/LambdaEventSerDesTest.java b/extra-serdes/src/test/java/software/amazon/lambda/durable/events/LambdaEventSerDesTest.java new file mode 100644 index 000000000..91a4626ec --- /dev/null +++ b/extra-serdes/src/test/java/software/amazon/lambda/durable/events/LambdaEventSerDesTest.java @@ -0,0 +1,207 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.events; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import com.amazonaws.services.lambda.runtime.events.S3Event; +import com.amazonaws.services.lambda.runtime.events.SNSEvent; +import com.amazonaws.services.lambda.runtime.events.SQSEvent; +import java.util.List; +import org.junit.jupiter.api.Test; +import software.amazon.lambda.durable.TypeToken; +import software.amazon.lambda.durable.exception.SerDesException; +import software.amazon.lambda.durable.serde.JacksonSerDes; +import software.amazon.lambda.durable.serde.SerDes; + +class LambdaEventSerDesTest { + private static final String SQS_EVENT_JSON = """ + { + "Records": [{ + "messageId": "message-1", + "receiptHandle": "receipt-1", + "body": "hello from sqs", + "attributes": {"ApproximateReceiveCount": "1"}, + "messageAttributes": {}, + "md5OfBody": "098f6bcd4621d373cade4e832627b4f6", + "eventSource": "aws:sqs", + "eventSourceARN": "arn:aws:sqs:us-east-1:123456789012:orders", + "awsRegion": "us-east-1" + }] + } + """; + + private static final String S3_EVENT_JSON = """ + { + "Records": [{ + "awsRegion": "us-east-1", + "eventName": "ObjectCreated:Put", + "eventSource": "aws:s3", + "eventTime": "2026-06-18T16:53:11.000Z", + "eventVersion": "2.1", + "requestParameters": {"sourceIPAddress": "127.0.0.1"}, + "responseElements": { + "x-amz-id-2": "extended-request-id", + "x-amz-request-id": "request-id" + }, + "s3": { + "configurationId": "configuration-id", + "bucket": { + "name": "orders", + "ownerIdentity": {"principalId": "principal-id"}, + "arn": "arn:aws:s3:::orders" + }, + "object": { + "key": "order.json", + "size": 42, + "eTag": "etag", + "sequencer": "sequencer" + }, + "s3SchemaVersion": "1.0" + }, + "userIdentity": {"principalId": "principal-id"} + }] + } + """; + + private static final String SNS_EVENT_JSON = """ + { + "Records": [{ + "EventSource": "aws:sns", + "EventVersion": "1.0", + "EventSubscriptionArn": "arn:aws:sns:us-east-1:123456789012:orders:subscription", + "Sns": { + "Type": "Notification", + "MessageId": "message-1", + "TopicArn": "arn:aws:sns:us-east-1:123456789012:orders", + "Subject": "order", + "Message": "hello from sns", + "Timestamp": "2026-06-18T16:53:11.000Z", + "SignatureVersion": "1", + "Signature": "signature", + "SigningCertUrl": "https://sns.us-east-1.amazonaws.com/cert.pem", + "UnsubscribeUrl": "https://sns.us-east-1.amazonaws.com/unsubscribe", + "MessageAttributes": {} + } + }] + } + """; + + private final LambdaEventSerDes serDes = new LambdaEventSerDes(); + + @Test + void deserializesSqsEventUsingLambdaRuntimePropertyNames() { + var event = serDes.deserialize(SQS_EVENT_JSON, TypeToken.get(SQSEvent.class)); + + assertNotNull(event.getRecords()); + assertEquals(1, event.getRecords().size()); + var message = event.getRecords().get(0); + assertEquals("message-1", message.getMessageId()); + assertEquals("hello from sqs", message.getBody()); + assertEquals("arn:aws:sqs:us-east-1:123456789012:orders", message.getEventSourceArn()); + } + + @Test + void deserializesSnsEventUsingLambdaRuntimePropertyNames() { + var event = serDes.deserialize(SNS_EVENT_JSON, TypeToken.get(SNSEvent.class)); + + assertNotNull(event.getRecords()); + assertEquals(1, event.getRecords().size()); + var record = event.getRecords().get(0); + assertNotNull(record.getSNS()); + assertEquals("hello from sns", record.getSNS().getMessage()); + assertEquals( + "arn:aws:sns:us-east-1:123456789012:orders", record.getSNS().getTopicArn()); + } + + @Test + void serializesSqsEventUsingLambdaRuntimePropertyNames() { + var event = serDes.deserialize(SQS_EVENT_JSON, TypeToken.get(SQSEvent.class)); + + var json = serDes.serialize(event); + + assertTrue(json.contains("\"Records\"")); + assertTrue(json.contains("\"eventSourceARN\"")); + assertFalse(json.contains("\"records\"")); + assertFalse(json.contains("\"eventSourceArn\"")); + + var roundTripped = serDes.deserialize(json, TypeToken.get(SQSEvent.class)); + assertEquals("hello from sqs", roundTripped.getRecords().get(0).getBody()); + } + + @Test + void roundTripsS3EventUsingLambdaRuntimeSerializer() { + var event = serDes.deserialize(S3_EVENT_JSON, TypeToken.get(S3Event.class)); + + assertNotNull(event.getRecords()); + assertEquals(1, event.getRecords().size()); + assertEquals("orders", event.getRecords().get(0).getS3().getBucket().getName()); + assertEquals("order.json", event.getRecords().get(0).getS3().getObject().getKey()); + + var json = serDes.serialize(event); + var roundTripped = serDes.deserialize(json, TypeToken.get(S3Event.class)); + + assertTrue(json.contains("\"Records\"")); + assertEquals( + "orders", roundTripped.getRecords().get(0).getS3().getBucket().getName()); + } + + @Test + void delegatesNonEventAndGenericTypes() { + var delegate = new TrackingSerDes(); + var eventSerDes = new LambdaEventSerDes(delegate); + + var json = eventSerDes.serialize(List.of("one", "two")); + var result = eventSerDes.deserialize(json, new TypeToken>() {}); + + assertEquals(List.of("one", "two"), result); + assertEquals(1, delegate.serializeCount); + assertEquals(1, delegate.deserializeCount); + } + + @Test + void handlesNullValues() { + assertNull(serDes.serialize(null)); + assertNull(serDes.deserialize(null, TypeToken.get(SQSEvent.class))); + } + + @Test + void wrapsEventDeserializationFailures() { + var exception = assertThrows( + SerDesException.class, () -> serDes.deserialize("{\"Records\":[", TypeToken.get(SQSEvent.class))); + + assertTrue(exception.getMessage().contains("Deserialization failed for Lambda event type")); + assertTrue(exception.getMessage().contains(SQSEvent.class.getName())); + assertNotNull(exception.getCause()); + } + + @Test + void rejectsNullDelegate() { + var exception = assertThrows(NullPointerException.class, () -> new LambdaEventSerDes(null)); + + assertEquals("Delegate SerDes cannot be null", exception.getMessage()); + } + + private static final class TrackingSerDes implements SerDes { + private final JacksonSerDes delegate = new JacksonSerDes(); + private int serializeCount; + private int deserializeCount; + + @Override + public String serialize(Object value) { + serializeCount++; + return delegate.serialize(value); + } + + @Override + public T deserialize(String data, TypeToken typeToken) { + deserializeCount++; + return delegate.deserialize(data, typeToken); + } + } +} diff --git a/pom.xml b/pom.xml index f4f616361..db0f5467c 100644 --- a/pom.xml +++ b/pom.xml @@ -40,6 +40,7 @@ sdk + extra-serdes sdk-testing sdk-integration-tests otel-plugin @@ -53,6 +54,8 @@ 17 UTF-8 2.54.7 + 3.16.1 + 1.4.1 2.22.2 6.1.3 5.23.0 @@ -78,6 +81,16 @@ aws-lambda-java-core 1.4.0 + + com.amazonaws + aws-lambda-java-events + ${aws.lambda.java.events.version} + + + com.amazonaws + aws-lambda-java-serialization + ${aws.lambda.java.serialization.version} + diff --git a/sdk-integration-tests/pom.xml b/sdk-integration-tests/pom.xml index e213bac09..60be41f4a 100644 --- a/sdk-integration-tests/pom.xml +++ b/sdk-integration-tests/pom.xml @@ -36,6 +36,12 @@ ${project.version} test + + software.amazon.lambda.durable + aws-durable-execution-sdk-java-extra-serdes + ${project.version} + test + org.junit.jupiter junit-jupiter diff --git a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/LambdaEventSerDesIntegrationTest.java b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/LambdaEventSerDesIntegrationTest.java new file mode 100644 index 000000000..327599a98 --- /dev/null +++ b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/LambdaEventSerDesIntegrationTest.java @@ -0,0 +1,52 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static software.amazon.lambda.durable.model.ExecutionStatus.SUCCEEDED; + +import com.amazonaws.services.lambda.runtime.events.SQSEvent; +import java.time.Duration; +import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.Test; +import software.amazon.lambda.durable.events.LambdaEventSerDes; +import software.amazon.lambda.durable.testing.LocalDurableTestRunner; + +class LambdaEventSerDesIntegrationTest { + + @Test + void parsesSqsEventInputDuringDurableExecutionAndReplay() { + var message = new SQSEvent.SQSMessage(); + message.setMessageId("message-1"); + message.setBody("hello from sqs"); + message.setEventSourceArn("arn:aws:sqs:us-east-1:123456789012:orders"); + + var event = new SQSEvent(); + event.setRecords(List.of(message)); + + var config = DurableConfig.builder().withSerDes(new LambdaEventSerDes()).build(); + var handlerRuns = new AtomicInteger(); + var stepRuns = new AtomicInteger(); + var runner = LocalDurableTestRunner.create( + SQSEvent.class, + (input, context) -> { + handlerRuns.incrementAndGet(); + var body = context.step("read-message", String.class, stepContext -> { + stepRuns.incrementAndGet(); + return input.getRecords().get(0).getBody(); + }); + context.wait("resume", Duration.ofSeconds(1)); + return body; + }, + config); + + var result = runner.runUntilComplete(event); + + assertEquals(SUCCEEDED, result.getStatus()); + assertEquals("hello from sqs", result.getResult(String.class)); + assertTrue(handlerRuns.get() >= 2); + assertEquals(1, stepRuns.get()); + } +}