sangkyoonnam opened a new issue, #1217: URL: https://github.com/apache/flink-agents/issues/1217
### Search before asking - [x] I searched in the [issues](https://github.com/apache/flink-agents/issues) and found nothing similar. ### Description When `gather` mixes `durableExecuteAsync` / `durable_execute_async` children with `executeAsync` / `execute_async` children, an interruption reported by an ordinary child is rethrown before the batch's durable outcomes are finalized. A newly executed durable child that returned successfully earlier in input order stays PENDING, and if recovery restores that state and the call has no reconciler, it runs again. In Java, `JavaRunnerContextImpl.resolveAsyncBatch` calls `rethrowCancellation` for ordinary children in the outcome loop (`JavaRunnerContextImpl.java:135-142`), before `finalizeExecutedOutcomes` (line 150). In Python, `_BatchAsyncExecutionResult.__await__` raises the ordinary child's interruption at `flink_runner_context.py:413-418`, before `_finalize_batch_execution` (line 420). Durable-only batches don't do this. `testGatherInterruptionLeavesRemainingSlotsPendingAndPropagates` keeps the outcomes before the interrupted slot and leaves the interrupted and later slots PENDING. The mixed case is pinned the other way: #1208's `test_ordinary_java_interruption_propagates_without_caching` expects the completed durable sibling to stay PENDING. Unless that was deliberate, I'd change that expectation to match the durable-only case, so a side-effecting durable call (a payment, an external write) gathered with an ordinary call isn't replayed after a task cancellation when its outcome is already known. In Java, a `CancellationException` thrown by the ordinary callable takes the same path. The change I have in mind: walk the outcomes in input order, finalize started durable children before the first cancellation signal, then propagate it. The interrupted slot and later new durable slots stay PENDING, and ordinary children take no durable index. Existing terminal slots stay as they are, and saving a durable outcome doesn't complete its handle on cancellation: re-awaiting the same gather rereads the saved prefix at the same base and retries the interrupted ordinary child. On cancellation the call index stays at the batch's base; on normal completion, including ordinary batch timeouts, it advances by the number of durable children as today. In Python, `_finalize_batch_execution` also advances the index (line 1427), so the fix has to separate the two. ### How to reproduce ```java TestDurableCallable<String> charge = new TestDurableCallable<>("charge", String.class, () -> "charged"); AsyncFuture<String> ordinary = context.executeAsync(() -> { throw new InterruptedException("cancelled"); }); context.gather(List.of(context.durableExecuteAsync(charge), ordinary)).await(); // throws InterruptedException // charge callCount=1, call result 0: pending=true, success=false // durable-only batch, same shape: call result 0 success=true ``` ```python await ctx.gather( ctx.durable_execute_async(charge, durable_id="charge"), ctx.execute_async(ordinary), # raises RuntimeError wrapping a Java InterruptedException ) # call result statuses: ['PENDING'] # re-entering the action with that state: charge runs a second time # durable-only batch, same shape: ['SUCCEEDED', 'PENDING'] ``` Both are unit probes. The Java one runs on JDK 21 without a continuation executor (the inline fallback), and the same result shows on JDK 17 through `AsyncBatchExecutor`. The Python one uses the fake Java context from `test_flink_runner_context_reconcilable.py`, injects a Java interruption, and simulates recovery by re-entering with the retained state. I haven't reproduced it with JDK 21 continuations or a real task cancellation and checkpoint restore. ### Version and environment main (`54db2844`, the #1208 merge). JDK 21 and JDK 17, Python 3.12. ### Are you willing to submit a PR? - [x] I'm willing to submit a PR! -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
