sangkyoonnam opened a new issue, #1181: URL: https://github.com/apache/flink-agents/issues/1181
### Search before asking - [x] I searched in the [issues](https://github.com/apache/flink-agents/issues) and found nothing similar. ### Description Python durable calls serialize a raised exception with cloudpickle, and neither a dump failure nor a load failure is handled. Load fails for these SDK exceptions with required keyword-only constructor arguments: `httpx.HTTPStatusError`, `openai.APIStatusError` and subclasses such as `RateLimitError`, `anthropic.APIStatusError`. They pickle, but unpickling calls `cls(*args)`. Replay then raises `TypeError: HTTPStatusError.__init__() missing 2 required keyword-only arguments: 'request' and 'response'` from `_try_get_cached_result` (`flink_runner_context.py:907`); `gather` gets the same `TypeError` as the failed outcome (line 1214). The call is not run again, so code that catches the SDK error takes a different path on replay. Dump fails for exceptions that hold an unpicklable object, such as a lock or a `pemja.PyJObject`. For a fresh call without a reconciler, executed synchronously or awaited individually, `_record_call_completion` logs a warning and records nothing (line 960), so the call runs again if recovery re-enters it. For a call with a PENDING slot, which includes any call with a reconciler, the `TypeError: cannot pickle ...` raised under `_finalize_current_call` (line 968) replaces the original exception and the slot stays PENDING. Java stores the class name and message, tries a public `(String)` constructor, and falls back to `RuntimeException("<class>: <message>")`. I propose prefixing new exception payloads with a format marker and storing a versioned pickled dict with the class name, the message and optional pickle bytes of the exception, with `RuntimeError("<class>: <message>")` as the replay fallback unless you'd rather have a dedicated type. A payload without the marker is a legacy one: it is raised as before when it unpickles, and has no stored class or message to fall back on when it doesn't. Exceptions that round-trip keep their type. The fallback keeps the diagnostic but changes the type, so SDK-specific handlers are still bypassed on replay for the exceptions above. A dump fallback has one side effect to handle in the same change. In a job where a synchronous durable call was inside Java `Thread.sleep` when the task was cancelled, pemja raised `RuntimeError` whose only argument is the Java `InterruptedException`. Pickling failed, nothing was recorded, and recovery ran the call again. In a separate probe with a Java `RuntimeException` wrapping `InterruptedException`, a prototype dump fallback recorded FAILED; recovery replayed the saved error on all three restarts, and the probe action re-raised it each time, which used up the restart limit. With a check in front of the fallback, recovery ran the call again and the action completed. #1071 made Java skip recording interrupted calls, and I'd do the same here. The check walks the class and cause chain of that Java throwable for `InterruptedException` or `ClosedByInterruptException`, and leaves out `InterruptedIOException` for the reason `ModelRoutingResolver.isCancellation` gives. It misses a bar e `InterruptedIOException` and an interruption whose throwable survives in neither the Python cause or context chain nor the Java cause chain. ### How to reproduce ```python import threading import httpx def fetch(order_id: str) -> None: request = httpx.Request("GET", f"https://orders.example/{order_id}") httpx.Response(429, request=request).raise_for_status() ctx.durable_execute(fetch, "o-1") # httpx.HTTPStatusError: Client error '429 Too Many Requests' for url ... # After a failover the action replays from the saved action state: ctx.durable_execute(fetch, "o-1") # TypeError: HTTPStatusError.__init__() missing 2 required keyword-only # arguments: 'request' and 'response' class Declined(Exception): def __init__(self, message: str) -> None: super().__init__(message) self.lock = threading.Lock() def charge(order_id: str) -> None: raise Declined("payment declined") ctx.durable_execute(charge, "o-1") # Declined; warning logged, nothing recorded ctx.durable_execute(charge, "o-1") # on replay charge runs again ``` Both snippets are unit probes against `FlinkRunnerContext` with a fake Java store, sync and async. The job runs used PyFlink's MiniCluster with the Kafka action state store and a failover: a function raising an explicitly constructed `HTTPStatusError` replayed as the `TypeError` without being called again, for `durable_execute` and `durable_execute_async`. ### Version and environment main (`1618728b`). Python 3.12, cloudpickle 2.2.1, JDK 21. Python runtime only. ### 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]
