sangkyoonnam opened a new pull request, #1182:
URL: https://github.com/apache/flink-agents/pull/1182

   Linked issue: #1181
   
   ### Purpose of change
   
   #### User-visible outcome
   
   A failed Python durable call recorded by this version replays as a failure 
that names the recorded error. An exception that pickles but can't be rebuilt, 
such as `httpx.HTTPStatusError`, replayed as the unpickling `TypeError`; it now 
replays as `RuntimeError("httpx.HTTPStatusError: <message>")`. An exception 
that can't be pickled went unrecorded, so the call ran again on recovery; it is 
now recorded and replays the same way. A detected Java interruption is not 
recorded, so recovery executes or reconciles the call again.
   
   #### Runtime flow
   
   `_serialize_call_payloads` writes a marker, then a pickled dict with a 
version, the class name, the message and the pickled exception, or `None` when 
pickling fails. `_try_get_cached_result` and `_read_terminal_outcome` return 
the rebuilt exception, or `RuntimeError("<class>: <message>")` when the pickle 
is missing or fails to load. A payload without the marker is read as the bare 
pickle earlier versions wrote.
   
   When `is_java_interruption` is true, `_record_call_completion` and 
`_finalize_current_call` return without writing and their callers raise the 
original exception; `_finalize_batch_execution` raises it, leaving that slot 
and later ones pending as `JavaRunnerContextImpl` does since #1111; 
`DurableFuture` leaves its handle unresolved.
   
   #### Key decisions
   
   - Keep pickle and add a fallback, so exceptions that round-trip keep their 
type and existing payloads stay readable. `RuntimeError` as the fallback 
matches Java, unless you'd rather have a dedicated type.
   - A Java interruption is detected from the Java throwable in `args`, by 
class and cause chain, never by message; `InterruptedIOException` is left out 
for the reason `ModelRoutingResolver.isCancellation` gives. Only the explicit 
Python cause chain is followed.
   - An interruption the check can't see is now recorded; before, it went 
unrecorded only because pickling failed. The job runs hit this with a Java 
`InterruptedIOException` and a wrapper that dropped the cause. Java behaves the 
same way, so I kept the parity. The alternative, leaving any unpicklable 
exception that carries a Java throwable unrecorded, also keeps re-executing 
ordinary Java failures; I'll switch to it if you'd rather.
   - The message copy is bounded at 16,384 characters. Stored in full it 
doubled the record, and a 450,000-character failure was lost without a logged 
error, Kafka's request size limit being the inferred cause.
   
   ### Behavioral Semantics
   
   #### Interaction decisions
   
   | Exception raised by the call | Recorded | Replayed as |
   |---|---|---|
   | Detected Java interruption, or caused by one through `from` (takes 
precedence) | No; a pending slot stays PENDING | The call executes or 
reconciles again |
   | Pickles and rebuilds | Yes (unchanged) | The same exception (unchanged) |
   | Pickles, can't be rebuilt | Yes (unchanged) | `RuntimeError("<class>: 
<message>")` (was the unpickling `TypeError`) |
   | Can't be pickled, no pending slot | Yes (was not recorded) | 
`RuntimeError("<class>: <message>")` (the call used to run again) |
   | Can't be pickled, pending slot | Yes (was left PENDING, pickling 
`TypeError` raised) | `RuntimeError("<class>: <message>")` |
   | Recorded by an earlier version | n/a | The same exception when it 
unpickles; the unpickling error otherwise (unchanged) |
   
   #### Behavioral contracts
   
   1. A failure recorded by this version whose exception can't be rebuilt 
replays as `RuntimeError` with its class name and message, without executing 
the call again.
   2. A failure whose exception can't be pickled is recorded, on a fresh call 
and on a pending slot; the first run raises the original exception.
   3. An exception that round-trips replays with its type; a payload written 
before this change that unpickles replays as before.
   4. A detected Java interruption is not recorded on any durable path, the 
original exception reaches the caller, and the call index does not advance.
   5. A failure raised while an interruption was being handled, or whose 
message only looks like one, is recorded.
   6. A record that can't be read raises `RuntimeError` naming the record; a 
message formatter that raises does not prevent recording.
   
   #### Failure behavior
   
   A record with an unknown version or missing or mistyped fields raises 
`RuntimeError` when read. Inspecting an exception for a Java interruption never 
raises; one that can't be inspected is recorded as an ordinary failure.
   
   ### Tests
   
   #### Contracts to tests
   
   All in `test_durable_exception.py`.
   
   | Contract | Tests |
   |---|---|
   | 1 | `test_replay_of_exception_that_cannot_be_rebuilt`, 
`test_gather_replays_exception_that_cannot_be_rebuilt` |
   | 2 | `test_*exception_that_cannot_be_pickled*` (fresh, pending slot, 
gather) |
   | 3 | `test_picklable_exception_round_trips_with_its_type`, 
`test_payload_written_*` (legacy, every pickle protocol) |
   | 4 | `test_java_interruption_*` (six), 
`test_interrupted_*_is_executed_again_when_awaited_again` (future, gather) |
   | 5 | `test_failure_raised_while_handling_an_interruption_is_not_one`, 
`test_other_*` (Python and Java failures) |
   | 6 | `test_unreadable_record_is_reported`, 
`test_exception_whose_message_cannot_be_formatted_is_recorded*` |
   
   #### Coverage and what was not verified
   
   `pytest flink_agents/runtime flink_agents/plan flink_agents/api`: 1326 
passed, 14 skipped (1266 before the 60 new tests). `tools/lint.sh -c` passes on 
JDK 11. The change at `56ce602c`, before the message bound, was also run as 
PyFlink MiniCluster jobs with the Kafka action state store under failover and 
cancel with restore, main against the branch, 216 runs; cases in the details 
block. A state written by main replayed unchanged on the branch; one written by 
the branch made main fail on replay with `ValueError: could not convert string 
to float: 'LINK-AGENTS-EXC'`, as expected for a rollback.
   
   Not verified: Kafka recovery with the message bound in place; JDK 11 and 17; 
pemja other than 0.5.7; a TaskManager process kill; the Fluss store; a real 
HTTP client whose interrupted I/O surfaces as `InterruptedIOException`; a Java 
interruption on the async thread pool, which the job runs never produced.
   
   <details>
   <summary>Implementation invariants and supporting evidence</summary>
   
   - With the source changes reverted and the tests kept, 22 of the 60 new 
tests fail.
   - A `ValueError("boom")` record grows from 49 to 153 bytes; 
`UnpicklingError` is the other error main raised on the new format.
   - Job cases: `httpx.HTTPStatusError` from `raise_for_status()`, 
`openai.RateLimitError` and `anthropic.RateLimitError` built the way the SDKs 
build them, an exception holding a lock, a picklable exception with attributes, 
`gather` over a mix of them, and Java interruptions of six shapes including a 
Java chat model resource, with reconciler variants. The repository's 
`execute_test.py` passes against the branch.
   - In the batch, later slots whose call already ran execute or reconcile 
again on recovery.
   - The marker can't start a valid pickle of any protocol: protocols 2 and up 
start with `0x80`, and in protocols 0 and 1 `F` opens a float whose text 
`LINK-AGENTS-EXC` doesn't parse.
   - The exception payload is opaque bytes on the Java side (`CallResult`, 
`ActionStateSerde`).
   - The fake `PyJObject` in the tests refuses pickling and prints like a Java 
throwable, matching what pemja 0.5.7 produced in the job runs.
   - The check runs where a failure is recorded rather than where it is caught, 
so the sync, async, reconciler and batch paths share it.
   - Two existing tests in `test_flink_runner_context_reconcilable.py` read the 
stored payload with `cloudpickle.loads`; they now use 
`deserialize_durable_exception`, with their assertions unchanged.
   </details>
   
   ### API
   
   #### Compatibility impact
   
   No public signature changes. `DurableFuture` gains a private hook, 
`_is_cancellation`, returning `False` unless a runtime future overrides it. 
`_DurableExecutionException` now tolerates `str(exception)` raising.
   
   The recorded exception format changes. This version reads what earlier 
versions wrote; an earlier version can't read what this one writes, so a 
rollback with failed durable calls in the store fails when it replays them. 
Replay of an exception that can't be rebuilt changes type, from the unpickling 
`TypeError` to `RuntimeError`. A failure that couldn't be pickled is now 
recorded, so its call no longer runs again on recovery.
   
   ### Documentation
   
   <!-- Do not remove this section. Check the proper box only. -->
   
   - [ ] `doc-needed` <!-- Your PR changes impact docs -->
   - [ ] `doc-not-needed` <!-- Your PR changes do not impact docs -->
   - [x] `doc-included` <!-- Your PR already contains the necessary 
documentation updates -->
   
   ### Was this patch authored or co-authored using generative AI tooling?
   
   <!-- Do not remove this section. Check the proper box only. -->
   
   - [x] Yes
   - [ ] No
   
   Generated-by: Claude Code 2.1.284 (Claude Fable 5.1)
   


-- 
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