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]

Reply via email to