aws / aws/aws-durable-execution-sdk-java

[Feature]: Support CompletionStage-returning durable step functions

Abierto
#657 0 comentarios 0 reacciones 0 asignados Ver en GitHub
enhancement pkg:sdk require adr
Lenguaje dominante
Java
Estrellas
28
Forks
11
Merge medio
1 d 8 h
PR fusionados (30 d)
47

Descripción

### What would you like?

Add a supported API for durable step functions whose user function returns a `CompletionStage` (typically a `CompletableFuture`), so applications can use asynchronous Java clients without blocking an SDK user-executor thread.

An illustrative API is:

```java
DurableFuture response = context.stepStage(
"fetch",
Response.class,
stepContext -> asyncClient.sendAsync(request));
```

`CompletionStage` should be the public abstraction rather than requiring the concrete `CompletableFuture` type. The exact name is open to design; a distinct method such as `stepStage` avoids erasure and overload ambiguity with the existing `Function` methods.

This is intended as a targeted addition, not a proposal to make the entire SDK `CompletableFuture`-first. The existing synchronous APIs, `DurableFuture`, and virtual-thread executor support should remain.

#### Benefits

- **Direct async-client integration:** Applications can return stages from AWS SDK async clients, async HTTP clients, database drivers, and other non-blocking libraries without calling `.join()` inside a step.
- **Better Java 17 scalability:** The SDK targets Java 17, where virtual threads are unavailable. A stage-returning step avoids holding one platform thread for each in-flight asynchronous request.
- **Lower per-operation overhead at high concurrency:** Virtual threads are much cheaper than platform threads, but each pending virtual thread still carries stack and thread-local state. Tracking a stage directly can reduce that overhead for very large fan-out workloads.
- **Clear ownership of lifecycle semantics:** The SDK can track the stage as active in-process work, checkpoint its terminal result, normalize failures, and preserve retry/replay behavior instead of requiring every application to build a blocking adapter.
- **Complementary to virtual threads:** Virtual threads remain the simplest option for blocking libraries, while stage-returning steps provide a natural path for libraries that are already asynchronous.

This feature does not reduce Lambda compute time for an in-flight network request: the invocation must remain active until the stage completes because the request and continuation are in-memory. Durable waits already handle compute-free suspension. The benefit here is avoiding a blocked user thread while preserving the same invocation-lifecycle behavior.

#### Example use cases

1. Call `DynamoDbAsyncClient`, `S3AsyncClient`, `LambdaAsyncClient`, or another AWS async client from a durable step.
2. Use Java 11+ `HttpClient.sendAsync()` or a third-party async HTTP client for remote API calls.
3. Compose several asynchronous calls inside one checkpointed step with `thenCompose`/`thenCombine`, then checkpoint the final value.
4. Run high-fan-out I/O workloads on Java 17 without configuring a large platform-thread pool.
5. Integrate third-party libraries that expose `CompletionStage` as their primary API without converting them to blocking calls.

### Possible Implementation

Add `stepStage` overloads parallel to the existing `stepAsync` overloads:

```java
DurableFuture stepStage(
String name,
Class resultType,
Function> function);

DurableFuture stepStage(
String name,
TypeToken resultType,
Function> function,
StepConfig config);
```

The RFC should determine the final naming and whether a dedicated functional interface would produce a clearer API.

Suggested lifecycle:

1. Start the step using the same operation identity, replay lookup, retry policy, plugin hooks, and serialization rules as a synchronous step.
2. Invoke the user function with an attached `StepContext`.
3. Validate that the returned stage is non-null.
4. Treat the pending stage as active in-process work even though no user thread is blocked. The execution must not suspend while the stage is pending, because its request and continuation would be lost with the invocation.
5. On completion, run SDK-owned terminal processing with the required context attached, normalize the value or exception, and checkpoint it through the existing step machinery.
6. On replay of a completed step, return the checkpointed result without invoking the user function or creating a new stage.

Design considerations:

- The current execution manager tracks active threads. Supporting this safely may require an active-work token or equivalent mechanism for asynchronous work that is not represented by a live SDK user thread.
- Completion can occur on an arbitrary client/event-loop thread. SDK checkpointing, plugin hooks, logging metadata, OpenTelemetry context, and thread-local SDK context must be propagated or restored deliberately.
- Synchronously thrown exceptions and exceptional stage completion should have identical retry and step-failure semantics.
- Step timeout behavior should cover the lifetime of the returned stage.
- Cancellation semantics must be documented. Cancelling a local stage may be best-effort and is not equivalent to cancelling a checkpointed durable operation.
- User-defined completion callbacks are ordinary in-process step code and are not individually checkpointed. Only the final step result is durable.
- The implementation should not use the common fork-join pool implicitly for SDK terminal processing.

Acceptance criteria:

- A durable step can return a `CompletionStage` and still expose its operation result as `DurableFuture`.
- Already-completed and asynchronously-completed stages both checkpoint successfully.
- Exceptional and cancelled stages produce documented failure/retry behavior.
- The execution cannot suspend early while a stage-backed step has pending in-process work.
- Completed steps replay without invoking the asynchronous function again.
- Existing step retry, timeout, serialization, and per-retry/per-run semantics apply consistently.
- SDK context, logging metadata, and plugin/OpenTelemetry lifecycle behavior are correct when completion originates on a non-SDK thread.
- Unit tests cover success, failure, cancellation, timeout, retry, replay, null stages, and suspension races.
- An integration test uses `LocalDurableTestRunner` to demonstrate delayed asynchronous completion and replay.
- Documentation includes an async HTTP or AWS SDK async-client example and explains when to choose stages versus virtual threads.
- Existing synchronous and `DurableFuture` APIs remain source- and binary-compatible.

### Is this a breaking change?

No

### Does this require an RFC?

Yes

### Additional Context

The SDK already uses `CompletableFuture` internally for operation coordination, and `DurableFuture.get()` integrates blocking waits with active-thread deregistration and durable suspension. This proposal is specifically about accepting asynchronous user work as a step body while retaining SDK-controlled checkpoint/replay semantics.

A raw `CompletableFuture` for a durable operation should not be exposed as though it survives invocation suspension. The returned public operation handle should remain `DurableFuture`; the `CompletionStage` is only the in-process implementation of the step body.

Virtual threads and this feature solve related but different problems:

- Virtual threads make blocking code scalable on Java 21+.
- Stage-returning steps integrate naturally with already-asynchronous libraries and also improve scalability on Java 17.
- Durable suspension remains responsible for stopping an invocation when no in-process work can make progress.

Guía de contribución

Abrir la guía de contribución

Línea de trabajo

Empieza leyendo las sobrecargas existentes de stepAsync, el manejo de DurableFuture y el seguimiento de hilos activos del gestor de ejecución. Usa LocalDurableTestRunner para probar la finalización retardada, el replay, los fallos, el timeout, el retry, la cancelación, las etapas null y las carreras de suspensión. Se considera terminado cuando los pasos respaldados por CompletionStage conservan el checkpointing, el contexto y el comportamiento del ciclo de vida, mientras las APIs síncronas existentes siguen siendo compatibles.

Escrito por el modelo de indexación a partir del texto del issue.

Evaluación

Stack tecnológico
java
Área
api, backend, distributed-systems
Tipo de issue
Nueva funcionalidad
Dificultad
5/5
Tiempo estimado
Más de una semana
Estado de actividad
Activo
Claridad
Bastante claro
Aptitud para principiantes
35/100

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.