dangzitou opened a new issue, #1189:
URL: https://github.com/apache/flink-agents/issues/1189

   ### Search before asking
   
   - [x] I searched in the 
[issues](https://github.com/apache/flink-agents/issues) and found nothing 
similar.
   
   ### Description
   
   A Python durable call can return a successful result even though the Java 
runtime explicitly failed to persist its call result. This affects fresh, 
completion-only `durable_execute(...)` calls and individually awaited 
`durable_execute_async(...)` calls without a reconciler.
   
   `FlinkRunnerContext._record_call_completion` catches every `Exception` 
around serialization **and** 
`self._j_runner_context.recordCallCompletion(...)`, then returns normally. It 
only logs when the exception text does not contain `recordCallCompletion`. 
Consequently, a Java persistence error is swallowed after a callable 
successfully returns a trivially serializable string.
   
   Expected: a state-persistence failure must escape to the caller so the 
action cannot proceed as if the durable result had been recorded. Actual: the 
caller receives the successful result while the persisted journal still has no 
call result. Restoring the pending action from that journal executes the 
operation again.
   
   This is a persistence-boundary failure, not an unpicklable result/exception. 
Related #1181 concerns exception serialization; its reproducer is different. 
Related #1173 concerns Kafka errors not reaching the Java caller; this report 
starts with a Java persistence error that **does** reach the Python boundary. 
The Java `DurableExecutionManager.persist` already wraps store errors as 
`RuntimeException("Failed to persist ActionState", cause)`.
   
   Source at the tested revision:
   
   - [Python catch around the Java 
write](https://github.com/apache/flink-agents/blob/99103da672f897e9c241b51922f2c31a8a728e43/python/flink_agents/runtime/flink_runner_context.py#L925-L960).
   - [Java call-result 
persistence](https://github.com/apache/flink-agents/blob/99103da672f897e9c241b51922f2c31a8a728e43/runtime/src/main/java/org/apache/flink/agents/runtime/context/RunnerContextImpl.java#L1171-L1190).
   - [Store failures deliberately propagated by the Java 
manager](https://github.com/apache/flink-agents/blob/99103da672f897e9c241b51922f2c31a8a728e43/runtime/src/main/java/org/apache/flink/agents/runtime/operator/DurableExecutionManager.java#L335-L342).
   
   ### How to reproduce
   
   The local integration probe below uses the real `FlinkRunnerContext`, 
`RunnerContextImpl`, `DurableExecutionContext`, and `ActionStateSerde`. Only 
the persister is a small in-memory test component that can throw a write error. 
Snapshots are serialized bytes, never references to the live mutable 
`ActionState`.
   
   For each mode, the probe executes an operation returning `"receipt-1"`, 
restores a new Java context from persisted bytes, and repeats the same call. 
The local operation counter stands in for a non-idempotent external operation; 
no external payment or model service is called.
   
   Observed in **4 independent journal instances per case**:
   
   | Mode | Write outcome | Initial Python result | Persisted call results 
before recovery | Operation executions after recovery |
   | --- | --- | --- | --- | --- |
   | sync | success (control) | `receipt-1` | 1 | 1 |
   | sync | Java write throws | `receipt-1` | 0 | 2 |
   | individually awaited async | success (control) | `receipt-1` | 1 | 1 |
   | individually awaited async | Java write throws | `receipt-1` | 0 | 2 |
   
   Calling the same Java `recordCallCompletion` directly raises 
`java.lang.RuntimeException: Failed to persist ActionState`, caused by 
`java.io.IOException: local test write failed`. Passing through the Python 
durable API instead returns success. Both the arguments and result pickle 
successfully.
   
   The recovery step deliberately occurs **before any later successful final 
action-state write**. A later successful write may save the in-memory call 
result; this report does not claim that every failed intermediate write 
necessarily causes a duplicate.
   
   <details>
   <summary>Complete local Java journal fixture: JournalGateway.java</summary>
   
   ```java
   import java.net.InetAddress;
   import java.util.HashMap;
   import java.util.Map;
   
   import org.apache.flink.agents.api.InputEvent;
   import org.apache.flink.agents.plan.AgentPlan;
   import org.apache.flink.agents.runtime.actionstate.ActionState;
   import org.apache.flink.agents.runtime.actionstate.ActionStateSerde;
   import org.apache.flink.agents.runtime.context.RunnerContextImpl;
   import py4j.GatewayServer;
   
   /** Local test entry point: production Java durable context and serde, 
in-memory persistence. */
   public class JournalGateway {
       private final Map<String, Journal> journals = new HashMap<>();
   
       public Journal create(String name) {
           Journal journal = new Journal();
           journals.put(name, journal);
           return journal;
       }
   
       public Journal get(String name) { return journals.get(name); }
   
       public static final class Journal {
           private final InputEvent event = new InputEvent("local-repro-input");
           private byte[] persisted = ActionStateSerde.serialize(new 
ActionState(event));
           private RunnerContextImpl context;
           private boolean failPersist;
           private int persistCalls;
   
           public Journal() { recover(); }
   
           public RunnerContextImpl getContext() { return context; }
           public void setFailPersist(boolean value) { failPersist = value; }
           public void setPersistFailure(boolean value) { failPersist = value; }
           public void resetFromPersistedSnapshot() { recover(); }
           public int getPersistCalls() { return persistCalls; }
           public int getPersistedCount() {
               return 
ActionStateSerde.deserialize(persisted).getCallResultCount();
           }
           public String getPersistedJson() {
               return new String(persisted, 
java.nio.charset.StandardCharsets.UTF_8);
           }
   
           public void recover() {
               context = new RunnerContextImpl(null, () -> {}, new 
AgentPlan(new HashMap<>()), null, "local-repro");
               context.setDurableExecutionContext(new 
RunnerContextImpl.DurableExecutionContext(
                   "local-key", 1L, null, event, 
ActionStateSerde.deserialize(persisted),
                   (key, seq, action, input, state) -> {
                       persistCalls++;
                       if (failPersist) {
                           throw new RuntimeException("Failed to persist 
ActionState", new java.io.IOException("local test write failed"));
                       }
                       persisted = ActionStateSerde.serialize(state);
                   }));
           }
       }
   
       public static void main(String[] args) throws Exception {
           int port = args.length == 0 ? 0 : Integer.parseInt(args[0]);
           GatewayServer server = new GatewayServer.GatewayServerBuilder(new 
JournalGateway())
               
.javaPort(port).javaAddress(InetAddress.getLoopbackAddress()).build();
           server.start();
           System.out.println("PORT=" + server.getListeningPort());
       }
   }
   ```
   
   </details>
   
   <details>
   <summary>Complete Python probe: repro_persistence_failure.py</summary>
   
   ```python
   """Local Python -> production Java durable journal contract check.
   
   Start java-runtime/JournalGateway.java first and pass its localhost port.
   The journal persister is the only injected component; no external service 
runs.
   """
   import argparse
   import asyncio
   import json
   from concurrent.futures import ThreadPoolExecutor
   
   from py4j.java_gateway import GatewayParameters, JavaGateway
   from py4j.protocol import Py4JJavaError
   
   from flink_agents.runtime.flink_runner_context import FlinkRunnerContext
   
   PLAN = '{"actions":{},"resource_providers":{},"config":{"conf_data":{}}}'
   
   
   async def consume(future):
       return await future
   
   
   def run(gateway, mode, fail, repeat):
       journal = 
gateway.entry_point.create(f"persistence-{mode}-{fail}-{repeat}")
       journal.setFailPersist(fail)
       calls = []
   
       def operation(order_id):
           calls.append(order_id)
           return "receipt-1"
   
       with ThreadPoolExecutor(max_workers=1) as executor:
           def invoke():
               ctx = FlinkRunnerContext(journal.getContext(), PLAN, executor, 
None)
               try:
                   if mode == "sync":
                       return ctx.durable_execute(operation, "order-1")
                   return 
asyncio.run(consume(ctx.durable_execute_async(operation, "order-1")))
               finally:
                   ctx.close()
   
           first = invoke()
           writes = journal.getPersistCalls()
           persisted_before_recovery = journal.getPersistedCount()
           assert first == "receipt-1"
           assert writes == 1
           assert persisted_before_recovery == (0 if fail else 1)
           journal.setFailPersist(False)
           journal.recover()  # Production serde reconstructs only persisted 
bytes.
           recovered = invoke()
           assert recovered == "receipt-1"
           assert len(calls) == (2 if fail else 1)
   
       return dict(mode=mode, failed_write=fail, run=repeat,
                   returned_success=True, persist_attempts=writes,
                   persisted_before_recovery=persisted_before_recovery,
                   operation_calls_after_recovery=len(calls))
   
   
   if __name__ == "__main__":
       parser = argparse.ArgumentParser()
       parser.add_argument("port", type=int)
       args = parser.parse_args()
       gateway = 
JavaGateway(gateway_parameters=GatewayParameters(port=args.port))
       try:
           direct = gateway.entry_point.create("direct-java-control")
           direct.setFailPersist(True)
           try:
               direct.getContext().recordCallCompletion("control", "digest", 
b"ok", None)
           except Py4JJavaError as error:
               assert "Failed to persist ActionState" in str(error)
               print(json.dumps({"direct_java_write": "raised", "error": 
str(error.java_exception)}))
           else:
               raise AssertionError("Java control did not raise")
           for repeat in range(1, 5):
               for mode in ("sync", "async"):
                   for fail in (False, True):
                       print(json.dumps(run(gateway, mode, fail, repeat)), 
flush=True)
       finally:
           gateway.close()
   ```
   
   </details>
   
   Build the checkout and run the fixture with JDK 17 and Python 3.12. The test 
environment had the Python dependencies needed to import `FlinkRunnerContext`, 
including PyFlink 2.2.0, Pydantic 2.11.4, cloudpickle 2.2.1, and Py4J 0.10.9.7. 
From the repository root, with the two files above saved in a separate 
directory:
   
   ```bash
   # Activate the Python environment; set JAVA_HOME to JDK 17.
   mvn -B -ntp -pl runtime -am test-compile jar:test-jar 
dependency:build-classpath \
     -DskipTests -Dspotless.skip=true \
     -Dmdep.outputFile=target/repro-classpath.txt -Dmdep.includeScope=test
   
   # Set REPRO_DIR to the directory containing the two files above.
   PY4J_JAR="$(python -c 'import sys; 
print(sys.prefix)')/share/py4j/py4j0.10.9.7.jar"
   
CP="$REPRO_DIR:runtime/target/classes:plan/target/classes:api/target/classes:integrations/mcp/target/classes:$PY4J_JAR:$(cat
 runtime/target/repro-classpath.txt)"
   "$JAVA_HOME/bin/javac" -cp "$CP" "$REPRO_DIR/JournalGateway.java"
   "$JAVA_HOME/bin/java" -cp "$CP" JournalGateway 25339
   
   # In a second terminal, activate the same Python environment, use the 
repository
   # root as cwd, and set REPRO_DIR again to the directory containing the probe:
   export PYTHONPATH="$PWD/python:$(python -c 'import sysconfig; 
print(sysconfig.get_paths()["purelib"])')"
   python "$REPRO_DIR/repro_persistence_failure.py" 25339
   ```
   
   The fixture binds to loopback. The Python program asserts all table entries 
and the direct-Java error control, and prints one JSON record per case.
   
   ### Version and environment
   
   - Flink Agents `main`, commit `99103da672f897e9c241b51922f2c31a8a728e43` 
(`0.4-SNAPSHOT` / Python `0.4.dev0`), unmodified source.
   - macOS arm64; Python 3.12.14; Pydantic 2.11.4; cloudpickle 2.2.1; PyFlink 
2.2.0 used for Python imports; Py4J 0.10.9.7; OpenJDK 17.0.20.
   - Java runtime and dependency modules compiled from this revision. Python 
imports resolve to this checkout.
   - This is a local Python-to-Java integration/recovery test through Py4J with 
a controlled persister failure. It is not a deployed Flink/Kafka failover or a 
Pemja end-to-end test.
   - Existing focused Python durable/reconciler tests: 56 passed; Java 
durable-context tests: 34 passed.
   
   ### Are you willing to submit a PR?
   
   - [ ] I'm willing to submit a PR!
   
   AI disclosure: This report and reproducer were generated with OpenAI Codex. 
The observations above come from executing the included local probes.
   


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