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]

Reply via email to