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]
