From 1fa68199bb02d110ca87e26fe550bfa17c0e8efa Mon Sep 17 00:00:00 2001 From: Alex Wang Date: Fri, 31 Jul 2026 23:57:25 +0000 Subject: [PATCH 1/5] feat: add custom map iteration naming --- .../src/main/java/map/MapItemNamer.java | 24 ++++ conformance-tests/template_map.yaml | 18 ++- .../lambda/durable/MapIntegrationTest.java | 134 ++++++++++++++++++ .../lambda/durable/config/MapConfig.java | 37 ++++- .../durable/context/DurableContextImpl.java | 18 +++ .../durable/operation/MapOperation.java | 34 ++++- .../lambda/durable/config/MapConfigTest.java | 22 +++ 7 files changed, 282 insertions(+), 5 deletions(-) create mode 100644 conformance-tests/src/main/java/map/MapItemNamer.java diff --git a/conformance-tests/src/main/java/map/MapItemNamer.java b/conformance-tests/src/main/java/map/MapItemNamer.java new file mode 100644 index 000000000..a8b2bf67a --- /dev/null +++ b/conformance-tests/src/main/java/map/MapItemNamer.java @@ -0,0 +1,24 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package map; + +import java.util.List; +import software.amazon.lambda.durable.DurableContext; +import software.amazon.lambda.durable.DurableHandler; +import software.amazon.lambda.durable.config.MapConfig; +import software.amazon.lambda.durable.model.MapResult; + +/** 9-13: Map with a custom item namer. */ +public class MapItemNamer extends DurableHandler, List> { + + @Override + public List handleRequest(List input, DurableContext context) { + var config = MapConfig.builder() + .maxConcurrency(1) + .itemNamer((item, index) -> "item-" + item) + .build(); + MapResult result = + context.map("named-items", input, Integer.class, (item, index, ctx) -> item * 10, config); + return result.results(); + } +} diff --git a/conformance-tests/template_map.yaml b/conformance-tests/template_map.yaml index 9819e4ba0..a9dbea6c1 100644 --- a/conformance-tests/template_map.yaml +++ b/conformance-tests/template_map.yaml @@ -59,8 +59,6 @@ Resources: reason: "items-only form (no name): every Java context.map overload requires a name argument" - id: "9-6" reason: "throw-if-error rethrow: MapResult has no throw-if-error and map exposes no per-item futures" - - id: "9-13" - reason: "custom item namer: Java MapConfig has no item-namer field" - id: "9-14" reason: "custom per-item serdes: Java MapConfig has a single serDes (no separate item-level serdes distinct from the result serde)" - id: "9-19" @@ -95,6 +93,22 @@ Resources: RetentionPeriodInDays: 7 ExecutionTimeout: 300 + MapItemNamer: + Type: AWS::Serverless::Function + TestingMetadata: + TestDescription: ["9-13"] + Properties: + CodeUri: . + Handler: map.MapItemNamer + Description: Map with custom iteration names + Role: + Fn::GetAtt: + - DurableFunctionRole + - Arn + DurableConfig: + RetentionPeriodInDays: 7 + ExecutionTimeout: 300 + MapEmpty: Type: AWS::Serverless::Function TestingMetadata: diff --git a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java index fef289c9b..93202c557 100644 --- a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java +++ b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java @@ -19,6 +19,7 @@ import software.amazon.lambda.durable.model.ConcurrencyCompletionStatus; import software.amazon.lambda.durable.model.ExecutionStatus; import software.amazon.lambda.durable.model.MapResult; +import software.amazon.lambda.durable.model.OperationSubType; import software.amazon.lambda.durable.model.WaitForConditionResult; import software.amazon.lambda.durable.retry.WaitStrategies; import software.amazon.lambda.durable.serde.JacksonSerDes; @@ -1888,4 +1889,137 @@ void testEmptyMapReplayUsesCheckpoint(NestingType nestingType, int events) { assertEquals(firstRunCount, executionCount.get(), "Map functions should not re-execute on replay"); assertEquals(events, result2.getHistoryEvents().size()); } + + @Test + void testItemNamerUsesCustomAndNullIterationNames() { + var runner = LocalDurableTestRunner.create(String.class, (input, context) -> { + var result = context.map( + "named-map", + List.of("a", "b", "c"), + String.class, + (item, index, ctx) -> item.toUpperCase(), + MapConfig.builder() + .itemNamer((item, index) -> index == 1 ? null : item + "-" + index) + .build()); + return String.join(",", result.results()); + }); + + var first = runner.runUntilComplete("test"); + assertEquals(ExecutionStatus.SUCCEEDED, first.getStatus()); + assertEquals("A,B,C", first.getResult(String.class)); + var iterationNames = first.getOperations().stream() + .filter(operation -> OperationSubType.MAP_ITERATION.getValue().equals(operation.getSubtype())) + .map(operation -> operation.getName()) + .toList(); + assertEquals(3, iterationNames.size()); + assertTrue(iterationNames.contains("a-0")); + assertTrue(iterationNames.contains("c-2")); + assertEquals( + 1, iterationNames.stream().filter(name -> name == null).count(), "iteration names: " + iterationNames); + assertFalse(iterationNames.contains("named-map-iteration-1")); + + var replay = runner.run("test"); + assertEquals(ExecutionStatus.SUCCEEDED, replay.getStatus()); + } + + @Test + void testInvalidItemNameDoesNotConsumeOperationId() { + var withRejectedMap = LocalDurableTestRunner.create(String.class, (input, context) -> { + try { + context.map( + "invalid-map", + List.of("a"), + String.class, + (item, index, ctx) -> item, + MapConfig.builder().itemNamer((item, index) -> "").build()); + } catch (IllegalArgumentException expected) { + // Continue so the next operation exposes whether the rejected map consumed an ID. + } + return context.step("after", String.class, stepContext -> "done"); + }); + var control = LocalDurableTestRunner.create( + String.class, (input, context) -> context.step("after", String.class, stepContext -> "done")); + + var attempted = withRejectedMap.runUntilComplete("test"); + var baseline = control.runUntilComplete("test"); + + assertEquals(ExecutionStatus.SUCCEEDED, attempted.getStatus()); + assertNull(attempted.getOperation("invalid-map")); + assertEquals( + baseline.getOperation("after").getId(), + attempted.getOperation("after").getId()); + } + + @Test + void testChangedItemNameFailsCachedReplay() { + var suffix = new AtomicReference<>("first"); + var runner = LocalDurableTestRunner.create(String.class, (input, context) -> { + var result = context.map( + "replay-map", + List.of("a", "b"), + String.class, + (item, index, ctx) -> item.toUpperCase(), + MapConfig.builder() + .itemNamer((item, index) -> item + "-" + suffix.get()) + .build()); + return String.join(",", result.results()); + }); + + assertEquals(ExecutionStatus.SUCCEEDED, runner.runUntilComplete("test").getStatus()); + suffix.set("second"); + + var replay = runner.run("test"); + + assertEquals(ExecutionStatus.FAILED, replay.getStatus()); + assertTrue(replay.getError().orElseThrow().errorType().contains("NonDeterministicExecutionException")); + } + + @Test + void testEmptyMapDoesNotInvokeItemNamer() { + var namerCalls = new AtomicInteger(); + var runner = LocalDurableTestRunner.create(String.class, (input, context) -> { + context.map( + "empty-map", + List.of(), + String.class, + (item, index, ctx) -> item, + MapConfig.builder() + .itemNamer((item, index) -> { + namerCalls.incrementAndGet(); + return "unused"; + }) + .build()); + return "done"; + }); + + var result = runner.runUntilComplete("test"); + + assertEquals(ExecutionStatus.SUCCEEDED, result.getStatus()); + assertEquals(0, namerCalls.get()); + } + + @Test + void testChangedItemNameFailsStartedMapReplay() { + var suffix = new AtomicReference<>("first"); + var runner = LocalDurableTestRunner.create(String.class, (input, context) -> { + var result = context.map( + "started-replay-map", + List.of("a"), + String.class, + (item, index, ctx) -> + ctx.waitForCallback("approval", String.class, (callbackId, stepContext) -> {}), + MapConfig.builder() + .itemNamer((item, index) -> item + "-" + suffix.get()) + .build()); + return result.getResult(0); + }); + + assertEquals(ExecutionStatus.PENDING, runner.run("test").getStatus()); + suffix.set("second"); + + var replay = runner.run("test"); + + assertEquals(ExecutionStatus.FAILED, replay.getStatus()); + assertTrue(replay.getError().orElseThrow().errorType().contains("NonDeterministicExecutionException")); + } } diff --git a/sdk/src/main/java/software/amazon/lambda/durable/config/MapConfig.java b/sdk/src/main/java/software/amazon/lambda/durable/config/MapConfig.java index d92572617..7462e6dde 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/config/MapConfig.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/config/MapConfig.java @@ -3,6 +3,7 @@ package software.amazon.lambda.durable.config; import java.util.Objects; +import java.util.function.BiFunction; import software.amazon.lambda.durable.serde.SerDes; /** @@ -15,12 +16,17 @@ public class MapConfig { private final CompletionConfig completionConfig; private final SerDes serDes; private final NestingType nestingType; + private final BiFunction itemNamer; private MapConfig(Builder builder) { this.maxConcurrency = Objects.requireNonNullElse(builder.maxConcurrency, Integer.MAX_VALUE); this.completionConfig = Objects.requireNonNullElse(builder.completionConfig, CompletionConfig.allCompleted()); this.nestingType = Objects.requireNonNullElse(builder.nestingType, NestingType.NESTED); this.serDes = builder.serDes; + this.itemNamer = builder.itemNamer; + if (itemNamer != null && nestingType == NestingType.FLAT) { + throw new IllegalArgumentException("itemNamer is not supported with FLAT map nesting"); + } } /** @return max concurrent items, or null for unlimited */ @@ -43,6 +49,18 @@ public NestingType nestingType() { return nestingType; } + /** + * Returns the function used to name map iterations. + * + *

The function receives the item and its zero-based index. A non-null result must satisfy the normal operation + * name constraints. A null result is preserved as an unnamed iteration. + * + * @return the item namer, or null when default iteration naming is used + */ + public BiFunction itemNamer() { + return itemNamer; + } + public static Builder builder() { return new Builder(); } @@ -52,7 +70,8 @@ public Builder toBuilder() { .maxConcurrency(maxConcurrency) .completionConfig(completionConfig) .serDes(serDes) - .nestingType(nestingType); + .nestingType(nestingType) + .itemNamer(itemNamer); } /** Builder for creating MapConfig instances. */ @@ -61,6 +80,7 @@ public static class Builder { private Integer maxConcurrency; private CompletionConfig completionConfig; private SerDes serDes; + private BiFunction itemNamer; private Builder() {} @@ -105,6 +125,21 @@ public Builder nestingType(NestingType nestingType) { return this; } + /** + * Sets a function that names each nested map iteration. + * + *

The function receives the item and its zero-based index. Returning null creates an unnamed iteration. + * Non-null names are validated before the map allocates an operation ID or emits a checkpoint. Item naming is + * not supported with {@link NestingType#FLAT} because flat iterations do not have context operations. + * + * @param itemNamer the item namer, or null to use default iteration naming + * @return this builder for method chaining + */ + public Builder itemNamer(BiFunction itemNamer) { + this.itemNamer = itemNamer; + return this; + } + public MapConfig build() { return new MapConfig(this); } diff --git a/sdk/src/main/java/software/amazon/lambda/durable/context/DurableContextImpl.java b/sdk/src/main/java/software/amazon/lambda/durable/context/DurableContextImpl.java index d75e5b8bd..21f64dd93 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/context/DurableContextImpl.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/context/DurableContextImpl.java @@ -268,6 +268,7 @@ public DurableFuture> mapAsync( // Convert to List for deterministic index-based access var itemList = List.copyOf(items); + var iterationNames = resolveMapIterationNames(name, itemList, config); var operationId = nextOperationId(); var operation = new MapOperation<>( @@ -276,11 +277,28 @@ public DurableFuture> mapAsync( function, resultType, config, + iterationNames, this); operation.execute(); return operation; } + private static List resolveMapIterationNames(String mapName, List items, MapConfig config) { + var namer = config.itemNamer(); + var branchPrefix = mapName == null ? "map-iteration-" : mapName + "-iteration-"; + var names = new java.util.ArrayList(items.size()); + for (int i = 0; i < items.size(); i++) { + if (namer == null) { + names.add(branchPrefix + i); + } else { + var iterationName = namer.apply(items.get(i), i); + ParameterValidator.validateOperationName(iterationName); + names.add(iterationName); + } + } + return names; + } + @Override public ParallelDurableFuture parallel(String name, ParallelConfig config) { Objects.requireNonNull(config, "config cannot be null"); diff --git a/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java b/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java index 8caf2a647..0b9e3e329 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java @@ -5,7 +5,9 @@ import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Objects; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import software.amazon.awssdk.services.lambda.model.ContextOptions; @@ -17,6 +19,7 @@ import software.amazon.lambda.durable.config.CompletionConfig; import software.amazon.lambda.durable.config.MapConfig; import software.amazon.lambda.durable.context.DurableContextImpl; +import software.amazon.lambda.durable.exception.NonDeterministicExecutionException; import software.amazon.lambda.durable.exception.UnrecoverableDurableExecutionException; import software.amazon.lambda.durable.execution.SuspendExecutionException; import software.amazon.lambda.durable.model.ConcurrencyCompletionStatus; @@ -46,6 +49,8 @@ public class MapOperation extends ConcurrencyOperation> { private final DurableContext.MapFunction function; private final TypeToken itemResultType; private final SerDes serDes; + private final List iterationNames; + private final boolean validateIterationNamesOnReplay; private volatile MapResult cachedResult; public MapOperation( @@ -54,6 +59,7 @@ public MapOperation( DurableContext.MapFunction function, TypeToken itemResultType, MapConfig config, + List iterationNames, DurableContextImpl durableContext) { super( operationIdentifier, @@ -73,6 +79,11 @@ public MapOperation( this.function = function; this.itemResultType = itemResultType; this.serDes = config.serDes(); + this.iterationNames = Collections.unmodifiableList(new ArrayList<>(iterationNames)); + if (this.iterationNames.size() != this.items.size()) { + throw new IllegalArgumentException("iterationNames must have one entry per item"); + } + this.validateIterationNamesOnReplay = config.itemNamer() != null; } private void addAllItems() { @@ -83,7 +94,6 @@ private void addUnskippedItems(List resultItems) // Enqueue all items first. // If the map is completed when replaying, mapResult != null and the items that have been skipped // will be skipped during replay. - var branchPrefix = getName() == null ? "map-iteration-" : getName() + "-iteration-"; for (int i = 0; i < items.size(); i++) { var index = i; var item = items.get(i); @@ -92,7 +102,7 @@ private void addUnskippedItems(List resultItems) var skip = status == MapResult.MapResultItem.Status.SKIPPED; enqueueItem( - branchPrefix + i, + iterationNames.get(i), childCtx -> function.apply(item, index, childCtx), itemResultType, serDes, @@ -101,6 +111,24 @@ private void addUnskippedItems(List resultItems) } } + private void validateIterationNamesAgainstCheckpoint() { + if (!validateIterationNamesOnReplay) { + return; + } + var checkpointedById = new HashMap(); + for (var child : getChildOperations()) { + checkpointedById.put(child.id(), child); + } + for (var branch : getBranches()) { + var checkpointed = checkpointedById.get(branch.getOperationId()); + if (checkpointed != null && !Objects.equals(checkpointed.name(), branch.getName())) { + throw terminateExecution(new NonDeterministicExecutionException(String.format( + "Map iteration name mismatch for \"%s\". Expected \"%s\", got \"%s\"", + branch.getOperationId(), checkpointed.name(), branch.getName()))); + } + } + } + @Override protected void start() { if (items.isEmpty()) { @@ -150,6 +178,7 @@ protected void replay(Operation existing) { throw terminateExecutionWithIllegalDurableOperationException( "Missing result in completed Map operation"); } + validateIterationNamesAgainstCheckpoint(); if (Boolean.TRUE.equals(existing.contextDetails().replayChildren())) { // Large result: re-execute children to reconstruct MapResult var expected = new ExpectedCompletionStatus( @@ -167,6 +196,7 @@ protected void replay(Operation existing) { // Map was in progress when interrupted — re-create children without sending // another START (the backend rejects duplicate START for existing operations) addAllItems(); + validateIterationNamesAgainstCheckpoint(); executeItems(); } default -> diff --git a/sdk/src/test/java/software/amazon/lambda/durable/config/MapConfigTest.java b/sdk/src/test/java/software/amazon/lambda/durable/config/MapConfigTest.java index fc8962de3..c21243d61 100644 --- a/sdk/src/test/java/software/amazon/lambda/durable/config/MapConfigTest.java +++ b/sdk/src/test/java/software/amazon/lambda/durable/config/MapConfigTest.java @@ -4,6 +4,7 @@ import static org.junit.jupiter.api.Assertions.*; +import java.util.function.BiFunction; import org.junit.jupiter.api.Test; import software.amazon.lambda.durable.serde.JacksonSerDes; @@ -59,6 +60,24 @@ void builderWithSerDes() { assertSame(serDes, config.serDes()); } + @Test + void builderWithItemNamer() { + BiFunction namer = (item, index) -> item + "-" + index; + + var config = MapConfig.builder().itemNamer(namer).build(); + + assertSame(namer, config.itemNamer()); + } + + @Test + void builderRejectsItemNamerWithFlatNesting() { + var builder = MapConfig.builder().nestingType(NestingType.FLAT).itemNamer((item, index) -> "item-" + index); + + var exception = assertThrows(IllegalArgumentException.class, builder::build); + + assertEquals("itemNamer is not supported with FLAT map nesting", exception.getMessage()); + } + @Test void builderChaining() { var completion = CompletionConfig.firstSuccessful(); @@ -79,10 +98,12 @@ void builderChaining() { void toBuilder_preservesValues() { var completion = CompletionConfig.minSuccessful(2); var serDes = new JacksonSerDes(); + BiFunction namer = (item, index) -> "item-" + index; var original = MapConfig.builder() .maxConcurrency(4) .completionConfig(completion) .serDes(serDes) + .itemNamer(namer) .build(); var copy = original.toBuilder().build(); @@ -90,6 +111,7 @@ void toBuilder_preservesValues() { assertEquals(4, copy.maxConcurrency()); assertSame(completion, copy.completionConfig()); assertSame(serDes, copy.serDes()); + assertSame(namer, copy.itemNamer()); } @Test From 87171bedc1ed1f80a1289ee6421eb227c3a7594d Mon Sep 17 00:00:00 2001 From: Alex Wang Date: Sat, 1 Aug 2026 00:26:31 +0000 Subject: [PATCH 2/5] fix: preserve map naming replay compatibility --- .../lambda/durable/MapIntegrationTest.java | 27 +++++++++++++ .../durable/operation/MapOperation.java | 39 ++++++++++++++++--- .../MapOperationCompatibilityTest.java | 27 +++++++++++++ 3 files changed, 88 insertions(+), 5 deletions(-) create mode 100644 sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java diff --git a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java index 93202c557..439bef3ad 100644 --- a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java +++ b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java @@ -7,6 +7,7 @@ import java.time.Duration; import java.util.ArrayList; import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import org.junit.jupiter.api.Test; @@ -2022,4 +2023,30 @@ void testChangedItemNameFailsStartedMapReplay() { assertEquals(ExecutionStatus.FAILED, replay.getStatus()); assertTrue(replay.getError().orElseThrow().errorType().contains("NonDeterministicExecutionException")); } + + @Test + void testRemovingItemNamerFailsCachedReplay() { + var useItemNamer = new AtomicBoolean(true); + var runner = LocalDurableTestRunner.create(String.class, (input, context) -> { + var configBuilder = MapConfig.builder(); + if (useItemNamer.get()) { + configBuilder.itemNamer((item, index) -> "custom-" + item); + } + var result = context.map( + "removed-namer-map", + List.of("a", "b"), + String.class, + (item, index, ctx) -> item.toUpperCase(), + configBuilder.build()); + return String.join(",", result.results()); + }); + + assertEquals(ExecutionStatus.SUCCEEDED, runner.runUntilComplete("test").getStatus()); + useItemNamer.set(false); + + var replay = runner.run("test"); + + assertEquals(ExecutionStatus.FAILED, replay.getStatus()); + assertTrue(replay.getError().orElseThrow().errorType().contains("NonDeterministicExecutionException")); + } } diff --git a/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java b/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java index 0b9e3e329..6722c42f1 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java @@ -28,6 +28,7 @@ import software.amazon.lambda.durable.model.OperationSubType; import software.amazon.lambda.durable.serde.SerDes; import software.amazon.lambda.durable.util.ExceptionHelper; +import software.amazon.lambda.durable.util.ParameterValidator; /** * Executes a map operation: applies a function to each item in a collection concurrently, with each item running in its @@ -50,9 +51,25 @@ public class MapOperation extends ConcurrencyOperation> { private final TypeToken itemResultType; private final SerDes serDes; private final List iterationNames; - private final boolean validateIterationNamesOnReplay; private volatile MapResult cachedResult; + public MapOperation( + OperationIdentifier operationIdentifier, + List items, + DurableContext.MapFunction function, + TypeToken itemResultType, + MapConfig config, + DurableContextImpl durableContext) { + this( + operationIdentifier, + items, + function, + itemResultType, + config, + resolveIterationNames(operationIdentifier.name(), items, config), + durableContext); + } + public MapOperation( OperationIdentifier operationIdentifier, List items, @@ -83,7 +100,22 @@ public MapOperation( if (this.iterationNames.size() != this.items.size()) { throw new IllegalArgumentException("iterationNames must have one entry per item"); } - this.validateIterationNamesOnReplay = config.itemNamer() != null; + } + + private static List resolveIterationNames(String mapName, List items, MapConfig config) { + var namer = config.itemNamer(); + var branchPrefix = mapName == null ? "map-iteration-" : mapName + "-iteration-"; + var names = new ArrayList(items.size()); + for (int i = 0; i < items.size(); i++) { + if (namer == null) { + names.add(branchPrefix + i); + } else { + var iterationName = namer.apply(items.get(i), i); + ParameterValidator.validateOperationName(iterationName); + names.add(iterationName); + } + } + return names; } private void addAllItems() { @@ -112,9 +144,6 @@ private void addUnskippedItems(List resultItems) } private void validateIterationNamesAgainstCheckpoint() { - if (!validateIterationNamesOnReplay) { - return; - } var checkpointedById = new HashMap(); for (var child : getChildOperations()) { checkpointedById.put(child.id(), child); diff --git a/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java b/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java new file mode 100644 index 000000000..dad78b009 --- /dev/null +++ b/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java @@ -0,0 +1,27 @@ +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.lambda.durable.operation; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; + +import java.util.List; +import org.junit.jupiter.api.Test; +import software.amazon.lambda.durable.DurableContext; +import software.amazon.lambda.durable.TypeToken; +import software.amazon.lambda.durable.config.MapConfig; +import software.amazon.lambda.durable.context.DurableContextImpl; +import software.amazon.lambda.durable.model.OperationIdentifier; + +class MapOperationCompatibilityTest { + + @Test + void retainsLegacyPublicConstructor() { + assertDoesNotThrow(() -> MapOperation.class.getConstructor( + OperationIdentifier.class, + List.class, + DurableContext.MapFunction.class, + TypeToken.class, + MapConfig.class, + DurableContextImpl.class)); + } +} From 33b5d7e9e264d87430f33918d44b23e413fb94a7 Mon Sep 17 00:00:00 2001 From: Alex Wang Date: Tue, 4 Aug 2026 20:12:08 +0000 Subject: [PATCH 3/5] feat: add typed map item namer overload --- .../lambda/durable/MapIntegrationTest.java | 25 +++++++ .../lambda/durable/config/MapConfig.java | 24 +++++++ .../lambda/durable/config/MapConfigTest.java | 65 +++++++++++++++++++ 3 files changed, 114 insertions(+) diff --git a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java index 439bef3ad..c7a5154ba 100644 --- a/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java +++ b/sdk-integration-tests/src/test/java/software/amazon/lambda/durable/MapIntegrationTest.java @@ -2049,4 +2049,29 @@ void testRemovingItemNamerFailsCachedReplay() { assertEquals(ExecutionStatus.FAILED, replay.getStatus()); assertTrue(replay.getError().orElseThrow().errorType().contains("NonDeterministicExecutionException")); } + + @Test + void testTypedItemNamerNamesIterationsEndToEnd() { + var runner = LocalDurableTestRunner.create(String.class, (input, context) -> { + var orders = List.of(new Order("a1"), new Order("b2")); + var result = context.map( + "typed-namer-map", + orders, + String.class, + (order, index, ctx) -> order.id().toUpperCase(), + MapConfig.builder() + .itemNamer(Order.class, (order, index) -> "order-" + order.id()) + .build()); + return String.join(",", result.results()); + }); + + var result = runner.runUntilComplete("test"); + + assertEquals(ExecutionStatus.SUCCEEDED, result.getStatus()); + assertEquals("A1,B2", result.getResult(String.class)); + assertNotNull(result.getOperation("order-a1")); + assertNotNull(result.getOperation("order-b2")); + } + + public record Order(String id) {} } diff --git a/sdk/src/main/java/software/amazon/lambda/durable/config/MapConfig.java b/sdk/src/main/java/software/amazon/lambda/durable/config/MapConfig.java index 7462e6dde..78bdcc293 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/config/MapConfig.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/config/MapConfig.java @@ -140,6 +140,30 @@ public Builder itemNamer(BiFunction itemNamer) { return this; } + /** + * Sets a function that names each nested map iteration, typed to the map's item class. + * + *

Equivalent to {@link #itemNamer(BiFunction)}, but the item type is declared explicitly so the namer can + * accept the item directly instead of {@link Object}: + * + *

{@code
+         * MapConfig.builder().itemNamer(Order.class, (order, index) -> order.id()).build();
+         * }
+ * + *

Each item is passed through {@link Class#cast}, so supplying a class that does not match the map's items + * fails with a {@link ClassCastException} naming the offending type. + * + * @param itemType the class of the map's items + * @param itemNamer the item namer, or null to use default iteration naming + * @param the map item type accepted by the namer + * @return this builder for method chaining + */ + public Builder itemNamer(Class itemType, BiFunction itemNamer) { + Objects.requireNonNull(itemType, "itemType cannot be null"); + this.itemNamer = itemNamer == null ? null : (item, index) -> itemNamer.apply(itemType.cast(item), index); + return this; + } + public MapConfig build() { return new MapConfig(this); } diff --git a/sdk/src/test/java/software/amazon/lambda/durable/config/MapConfigTest.java b/sdk/src/test/java/software/amazon/lambda/durable/config/MapConfigTest.java index c21243d61..77670529c 100644 --- a/sdk/src/test/java/software/amazon/lambda/durable/config/MapConfigTest.java +++ b/sdk/src/test/java/software/amazon/lambda/durable/config/MapConfigTest.java @@ -69,6 +69,69 @@ void builderWithItemNamer() { assertSame(namer, config.itemNamer()); } + @Test + void builderWithTypedItemNamer_receivesItemWithoutCast() { + // The lambda parameter is the domain type, so no cast is needed to reach its members. + var config = MapConfig.builder() + .itemNamer(Order.class, (order, index) -> order.id() + "-" + index) + .build(); + + assertEquals("a1-0", config.itemNamer().apply(new Order("a1"), 0)); + } + + @Test + void builderWithTypedItemNamer_acceptsAlreadyTypedFunction() { + // A BiFunction is rejected by itemNamer(BiFunction) but accepted here. + BiFunction namer = (order, index) -> order.id(); + + var config = MapConfig.builder().itemNamer(Order.class, namer).build(); + + assertEquals("a1", config.itemNamer().apply(new Order("a1"), 0)); + } + + @Test + void builderWithTypedItemNamer_preservesNullResult() { + var config = MapConfig.builder() + .itemNamer(Order.class, (order, index) -> null) + .build(); + + assertNull(config.itemNamer().apply(new Order("a1"), 0)); + } + + @Test + void builderWithTypedItemNamer_nullNamerLeavesDefaultNaming() { + var config = MapConfig.builder().itemNamer(Order.class, null).build(); + + assertNull(config.itemNamer()); + } + + @Test + void builderWithTypedItemNamer_nullItemTypeThrows() { + var exception = assertThrows( + NullPointerException.class, () -> MapConfig.builder().itemNamer(null, (order, index) -> "x")); + + assertEquals("itemType cannot be null", exception.getMessage()); + } + + @Test + void builderWithTypedItemNamer_mismatchedItemTypeThrows() { + var config = MapConfig.builder() + .itemNamer(Order.class, (order, index) -> order.id()) + .build(); + + assertThrows(ClassCastException.class, () -> config.itemNamer().apply("not-an-order", 0)); + } + + @Test + void builderRejectsTypedItemNamerWithFlatNesting() { + var builder = + MapConfig.builder().nestingType(NestingType.FLAT).itemNamer(Order.class, (order, index) -> order.id()); + + var exception = assertThrows(IllegalArgumentException.class, builder::build); + + assertEquals("itemNamer is not supported with FLAT map nesting", exception.getMessage()); + } + @Test void builderRejectsItemNamerWithFlatNesting() { var builder = MapConfig.builder().nestingType(NestingType.FLAT).itemNamer((item, index) -> "item-" + index); @@ -145,4 +208,6 @@ void builderWithNullMaxConcurrency_shouldPass() { var config = MapConfig.builder().maxConcurrency(null).build(); assertEquals(Integer.MAX_VALUE, config.maxConcurrency()); } + + private record Order(String id) {} } From 56a8fa05173193a81aade54690d492a1407109e9 Mon Sep 17 00:00:00 2001 From: Alex Wang Date: Tue, 4 Aug 2026 21:31:40 +0000 Subject: [PATCH 4/5] refactor: share map iteration name resolution --- .../durable/context/DurableContextImpl.java | 18 +---- .../durable/operation/MapOperation.java | 16 ++++- .../MapOperationCompatibilityTest.java | 69 +++++++++++++++++++ 3 files changed, 85 insertions(+), 18 deletions(-) diff --git a/sdk/src/main/java/software/amazon/lambda/durable/context/DurableContextImpl.java b/sdk/src/main/java/software/amazon/lambda/durable/context/DurableContextImpl.java index 21f64dd93..0c79165ec 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/context/DurableContextImpl.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/context/DurableContextImpl.java @@ -268,7 +268,7 @@ public DurableFuture> mapAsync( // Convert to List for deterministic index-based access var itemList = List.copyOf(items); - var iterationNames = resolveMapIterationNames(name, itemList, config); + var iterationNames = MapOperation.resolveIterationNames(name, itemList, config); var operationId = nextOperationId(); var operation = new MapOperation<>( @@ -283,22 +283,6 @@ public DurableFuture> mapAsync( return operation; } - private static List resolveMapIterationNames(String mapName, List items, MapConfig config) { - var namer = config.itemNamer(); - var branchPrefix = mapName == null ? "map-iteration-" : mapName + "-iteration-"; - var names = new java.util.ArrayList(items.size()); - for (int i = 0; i < items.size(); i++) { - if (namer == null) { - names.add(branchPrefix + i); - } else { - var iterationName = namer.apply(items.get(i), i); - ParameterValidator.validateOperationName(iterationName); - names.add(iterationName); - } - } - return names; - } - @Override public ParallelDurableFuture parallel(String name, ParallelConfig config) { Objects.requireNonNull(config, "config cannot be null"); diff --git a/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java b/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java index 6722c42f1..2f665547a 100644 --- a/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java +++ b/sdk/src/main/java/software/amazon/lambda/durable/operation/MapOperation.java @@ -102,7 +102,21 @@ public MapOperation( } } - private static List resolveIterationNames(String mapName, List items, MapConfig config) { + /** + * Resolves the operation name for every iteration of a map, applying the config's item namer when present and the + * default {@code "-iteration-N"} naming otherwise. A namer that returns null yields an unnamed iteration; + * any non-null name is validated here. + * + *

SDK-internal. This is the single source of iteration naming for both construction paths: the caller resolves + * names before an operation ID is allocated, and the legacy constructor resolves them on behalf of callers that do + * not. + * + * @param mapName the map operation's name, or null + * @param items the map's items, in iteration order + * @param config the map configuration supplying the optional item namer + * @return one name per item, in iteration order + */ + public static List resolveIterationNames(String mapName, List items, MapConfig config) { var namer = config.itemNamer(); var branchPrefix = mapName == null ? "map-iteration-" : mapName + "-iteration-"; var names = new ArrayList(items.size()); diff --git a/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java b/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java index dad78b009..9974fd9a6 100644 --- a/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java +++ b/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java @@ -3,8 +3,12 @@ package software.amazon.lambda.durable.operation; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.Test; import software.amazon.lambda.durable.DurableContext; import software.amazon.lambda.durable.TypeToken; @@ -24,4 +28,69 @@ void retainsLegacyPublicConstructor() { MapConfig.class, DurableContextImpl.class)); } + + // The resolver below is the single source of iteration naming for both the legacy constructor and + // DurableContextImpl.mapAsync, so these cases pin the contract that both paths share. + + @Test + void resolveIterationNames_withoutNamer_usesDefaultNaming() { + var names = MapOperation.resolveIterationNames( + "orders", List.of("a", "b"), MapConfig.builder().build()); + + assertEquals(List.of("orders-iteration-0", "orders-iteration-1"), names); + } + + @Test + void resolveIterationNames_withoutMapName_usesUnprefixedDefault() { + var names = MapOperation.resolveIterationNames( + null, List.of("a"), MapConfig.builder().build()); + + assertEquals(List.of("map-iteration-0"), names); + } + + @Test + void resolveIterationNames_withNamer_usesCustomNames() { + var config = MapConfig.builder() + .itemNamer((item, index) -> item + "-" + index) + .build(); + + var names = MapOperation.resolveIterationNames("orders", List.of("a", "b"), config); + + assertEquals(List.of("a-0", "b-1"), names); + } + + @Test + void resolveIterationNames_preservesNullNamerResult() { + var config = MapConfig.builder().itemNamer((item, index) -> null).build(); + + var names = MapOperation.resolveIterationNames("orders", List.of("a"), config); + + assertEquals(1, names.size()); + assertNull(names.get(0)); + } + + @Test + void resolveIterationNames_validatesCustomNames() { + var config = MapConfig.builder().itemNamer((item, index) -> "").build(); + + assertThrows( + IllegalArgumentException.class, + () -> MapOperation.resolveIterationNames("orders", List.of("a"), config)); + } + + @Test + void resolveIterationNames_emptyItemsDoesNotInvokeNamer() { + var calls = new AtomicInteger(); + var config = MapConfig.builder() + .itemNamer((item, index) -> { + calls.incrementAndGet(); + return "unused"; + }) + .build(); + + var names = MapOperation.resolveIterationNames("orders", List.of(), config); + + assertEquals(List.of(), names); + assertEquals(0, calls.get()); + } } From b18619902ea672933e3fb2280f95f2e983daab43 Mon Sep 17 00:00:00 2001 From: Alex Wang Date: Tue, 4 Aug 2026 21:34:58 +0000 Subject: [PATCH 5/5] test: drop map operation compatibility test --- .../MapOperationCompatibilityTest.java | 96 ------------------- 1 file changed, 96 deletions(-) delete mode 100644 sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java diff --git a/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java b/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java deleted file mode 100644 index 9974fd9a6..000000000 --- a/sdk/src/test/java/software/amazon/lambda/durable/operation/MapOperationCompatibilityTest.java +++ /dev/null @@ -1,96 +0,0 @@ -// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. -// SPDX-License-Identifier: Apache-2.0 -package software.amazon.lambda.durable.operation; - -import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertNull; -import static org.junit.jupiter.api.Assertions.assertThrows; - -import java.util.List; -import java.util.concurrent.atomic.AtomicInteger; -import org.junit.jupiter.api.Test; -import software.amazon.lambda.durable.DurableContext; -import software.amazon.lambda.durable.TypeToken; -import software.amazon.lambda.durable.config.MapConfig; -import software.amazon.lambda.durable.context.DurableContextImpl; -import software.amazon.lambda.durable.model.OperationIdentifier; - -class MapOperationCompatibilityTest { - - @Test - void retainsLegacyPublicConstructor() { - assertDoesNotThrow(() -> MapOperation.class.getConstructor( - OperationIdentifier.class, - List.class, - DurableContext.MapFunction.class, - TypeToken.class, - MapConfig.class, - DurableContextImpl.class)); - } - - // The resolver below is the single source of iteration naming for both the legacy constructor and - // DurableContextImpl.mapAsync, so these cases pin the contract that both paths share. - - @Test - void resolveIterationNames_withoutNamer_usesDefaultNaming() { - var names = MapOperation.resolveIterationNames( - "orders", List.of("a", "b"), MapConfig.builder().build()); - - assertEquals(List.of("orders-iteration-0", "orders-iteration-1"), names); - } - - @Test - void resolveIterationNames_withoutMapName_usesUnprefixedDefault() { - var names = MapOperation.resolveIterationNames( - null, List.of("a"), MapConfig.builder().build()); - - assertEquals(List.of("map-iteration-0"), names); - } - - @Test - void resolveIterationNames_withNamer_usesCustomNames() { - var config = MapConfig.builder() - .itemNamer((item, index) -> item + "-" + index) - .build(); - - var names = MapOperation.resolveIterationNames("orders", List.of("a", "b"), config); - - assertEquals(List.of("a-0", "b-1"), names); - } - - @Test - void resolveIterationNames_preservesNullNamerResult() { - var config = MapConfig.builder().itemNamer((item, index) -> null).build(); - - var names = MapOperation.resolveIterationNames("orders", List.of("a"), config); - - assertEquals(1, names.size()); - assertNull(names.get(0)); - } - - @Test - void resolveIterationNames_validatesCustomNames() { - var config = MapConfig.builder().itemNamer((item, index) -> "").build(); - - assertThrows( - IllegalArgumentException.class, - () -> MapOperation.resolveIterationNames("orders", List.of("a"), config)); - } - - @Test - void resolveIterationNames_emptyItemsDoesNotInvokeNamer() { - var calls = new AtomicInteger(); - var config = MapConfig.builder() - .itemNamer((item, index) -> { - calls.incrementAndGet(); - return "unused"; - }) - .build(); - - var names = MapOperation.resolveIterationNames("orders", List.of(), config); - - assertEquals(List.of(), names); - assertEquals(0, calls.get()); - } -}