wenjin272 opened a new pull request, #1220: URL: https://github.com/apache/flink-agents/pull/1220
Linked issue: #1086 ### Purpose of change Actions can use Java `ctx.chat(...).await()` or Python `await ctx.chat(...)` without splitting an invocation across request and response actions. The event API remains available. Scope: chat and its tool rounds, not standalone Tool/RAG calls. #### Runtime flow 1. `ctx.chat` snapshots the request and returns a lazy `DurableFuture<ChatMessage>`. 2. First await replays a recorded result or queues a private `ChatCallEvent` bootstrap. The caller suspends while existing Chat/Tool actions run. 3. Chat/Tool events stay inside the invocation. Its matching `ChatResponseEvent` completes it without reaching user response listeners. 4. The caller records the response when durability is enabled, then returns the message or raises its error. Java waits via `resolveAsync`; Python yields through the coroutine bridge. #### Key decisions - Reuse event orchestration, not direct resource calls, to retain routing, retries, tool rounds and structured output. - Reuse `DurableFuture`, without an extra durable-execute wrapper. `ChatContext` groups task-local references; an operator-wide manager coordinates Java/Python calls. - Release only internal bootstraps while suspended; flushing ordinary events would change action ordering. ### Behavioral Semantics #### Interaction decisions | Conditions | Result | |---|---| | Java caller + JDK 21 continuation context | Caller yields; an async worker waits for the response. | | Java caller + no continuation support/context | Await rejects before dispatch. | | Python async caller + JDK 11 or 21 | Coroutine yields without occupying a response-wait worker. | | Durable execution + recorded terminal response | Replay success or failure without bootstrapping another chat. | | No recorded response, or durability disabled | Run the event flow; incomplete external effects may repeat after recovery. | | Caller inside an internal subagent | Use the child plan's resources and the caller's memory scope. | | Java wait + asynchronous model/tool work | Both use the shared pool; insufficient capacity can deadlock (known limitation #1213). | #### Behavioral contracts 1. Creating an unawaited future dispatches no request. Repeated awaits in the creating action execution reuse the terminal outcome. 2. Await returns a full `ChatMessage`, not only text; a failed response raises a catchable chat exception. 3. Private chat responses do not trigger user `ChatResponseEvent` actions. Ordinary pending events remain buffered until the calling action finishes. 4. Chat uses the caller's resources and memory, including an internal subagent's scope. 5. Durable replay reuses recorded terminal outcomes. This is not exactly-once execution of provider or tool side effects. 6. `gather` rejects chat futures before starting them. Futures must stay within their creating action execution; cross-execution use is unsupported and not checked at runtime. 7. Existing model/tool timeout settings remain applicable; there is no new whole-chat timeout. #### Failure behavior Built-in retries/fallback run unchanged. A terminal FAILED response raises `ChatResponseEvent.ChatResponseException` (Java) or `ChatResponseError` (Python), and repeated awaits raise again without restarting the call. Unsupported Java runtimes, invalid requests/lineage and unsupported composition raise rather than silently falling back. Cancellation, fatal errors, bridge/serialization errors, event-delivery errors and persistence failures propagate instead of becoming FAILED chat responses. They are not cached as successful chat outcomes. ### Tests | Contract | Coverage | |---|---| | Lazy invocation and repeated success/failure awaits (1, 2) | `ChatCallTest`, `ChatDurableFutureTest`, Python `test_chat_call.py` | | Private routing and ordinary-event ordering (3) | `ChatCallTest.lazyRepeatedAwaitAndOrdinaryEventIsolation`, MiniCluster `chat_call_test.py` | | Caller memory and child resources (4) | `ChatCallTest.chatInsideInternalSubagentUsesChildResources`, MiniCluster chat test | | Recorded success/failure replay; interrupted recovery (5) | `ChatCallManagerTest.terminalResponseReplaysWithoutBootstrap`, `ChatCallTest.checkpointRecoveryBeforeDispatchAfterResponseAndAfterAwait` | | Reject gather; same-execution usage constraint (6) | Java/Python gather-rejection tests; cross-execution misuse intentionally has no runtime guard/test | | Existing timeout behavior (7) | No dedicated chat-call timeout test; no timeout policy changed | | Runtime failures remain runtime failures | Java interrupted-wait/persistence/fatal-provider tests; Python cancellation and bridge-failure tests | Passed locally on this commit's source: - JDK 21 API/Plan/Runtime/MCP: **2,395 passed, 44 skipped**. - JDK 11 focused chat/subagent/async/durable tests: **126 passed, 24 skipped**. - Python API/runtime: **860 passed, 14 skipped**. - Rebuilt Flink 2.3 JARs and wheel: **4 MiniCluster chat/subagent tests passed on JDK 21**, and **1 Python chat test passed on JDK 11**. - Spotless, Ruff and `git diff --check` passed. Not verified: live providers, chat-specific routing/structured-output/media E2E, Python chat checkpoint recovery E2E, low-capacity deadlock resolution, other Flink versions, or older saved-state migration. <details> <summary>Implementation invariants and verification details</summary> - `ChatCallOwner` allocates replay-stable invocation IDs from caller identity and call order. It is runtime-only, not checkpointed. `ChatContext` survives continuation transfer and is replaced on task switches; `ActionTaskContextManagerTest` covers both. - `ChatInvocation.complete` checks request correlation and rejects duplicate terminal responses. `ChatCallManagerTest` covers both checks and record-scoped cleanup. - Completed callers do not replay consumed bootstrap envelopes. Recovery restarts from root callers rather than restoring orphaned private chat tasks. Operator tests exercise checkpoints before dispatch, after response, after await and during a tool round. - Terminal business failures are persisted as response events, then exposed as exceptions on access. The wait itself adds no separate durable slot. - Java module command: `mvn --batch-mode --no-transfer-progress test -pl runtime -am` on JDK 21. Python command: `python -m pytest flink_agents/api/tests flink_agents/runtime/tests -m 'not integration' -q`. - MiniCluster tests use deterministic providers, not external services. Python tests use the active environment's site-packages in `PYTHONPATH`; E2E uses rebuilt/installed artifacts. </details> ### API Add Java `RunnerContext.chat(ChatRequestEvent)` / `chat(String, List<ChatMessage>)` and Python `RunnerContext.chat(model, messages, *, prompt_args=None, output_schema=None)`, returning `DurableFuture`. Default implementations reject runtimes without event-backed support. Existing event APIs and YAML declarations are unchanged; no YAML call syntax or saved-state migration is added. Java awaits require JDK 21 with the continuation export flag. ### Documentation - [ ] `doc-needed` - [ ] `doc-not-needed` - [x] `doc-included` ### Was this patch authored or co-authored using generative AI tooling? - [x] Yes - [ ] No Generated-by: Codex CLI 0.153.4 (GPT-5) -- 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]
