dangzitou opened a new pull request, #1199: URL: https://github.com/apache/flink-agents/pull/1199
Linked issue: #1140 This draft is stacked on #1138 and includes its internal sub-agent framework prerequisite. The batch-specific change is [6775b12b](https://github.com/dangzitou/flink-agents/commit/6775b12b8ede8d0fae8004e259a8f22a28530b2e). ### Purpose of change Parents can combine internal sub-agent handles into a durable batch, run a bounded fan-out across the same or different children, and receive ordered per-call results. This removes serial resolution and shared-worker starvation from the grouped invocation path. #### Runtime flow `submit` allocates deterministic identities without starting internal children. Group resolution validates admission, prepares deferred handles, and gathers their durable futures. The runtime reserves every outstanding slot, starts a bounded window on the mailbox, suspends the parent, persists each terminal result before admitting its replacement, and returns results in handle order. Recovery replays the parent first to reconstruct child scopes and reuses child action journals. #### Key decisions - Reuse the existing group API and durable batch machinery; internal completion uses futures/coroutine polling instead of blocking async workers. - Bound active invocations per gathered batch, including different resource names. Nested batches have independent capacity so waiting ancestors cannot starve descendants. - Bound queued state by rejecting an oversized group before reservation; queued calls use existing `PENDING` slots rather than a second journal format. ### Behavioral Semantics #### Interaction decisions | Conditions | Behavior | | --- | --- | | Internal or mixed internal/deferred group | One sub-agent window; tool-call parallelism and timeout do not apply. | | External deferred group only | Existing tool-batch scheduling options apply. | | Already submitted async handles | Await existing runs; grouping cannot retroactively control their submission. | | Completed and outstanding handles | Reuse completed results; schedule outstanding deferred calls in their original order. | | Nested internal group | Its own window; the ancestor occupies only its parent's batch slot. | | Child failure with unfinished actions/events | Retain the slot until the child scope quiesces, then return its failed result. | #### Behavioral contracts 1. A group returns one result per handle, in input order, after all child calls settle; successful siblings survive child failures. 2. Internal groups admit at most `subagent.parallelism` invocations at once and work with one shared async worker, including nested calls. 3. Oversized groups, duplicate handles, different owning contexts, and non-positive limits fail before internal preparation/reservation. 4. Recovery preserves terminal results and replays running/queued calls under their original identities, including intermediate child events and failure summaries. 5. Cancellation and infrastructure errors propagate from resolution; infrastructure errors leave unfinished durable slots pending. #### Failure behavior Admission errors raise. Java internal calls require JDK 21+ continuation support; Python callers must be async. Child action failures become `SubagentResult.error` after quiescence and their summaries survive replay. Bootstrap, polling, serialization, and journal errors propagate rather than becoming child results. Unfinished external effects may replay under the existing durability model, so callers still need idempotency or a reconciler. ### Tests | Contract | Tests | | --- | --- | | Ordered results and partial failures | `InternalSubagentBatchTest`, Python `test_internal_subagent_batch`, Java/Python batch E2E | | Bounded concurrency, one worker, nesting | Java/Python batch E2E; existing nested-call tests with parallelism/worker count 1 | | Admission without side effects | Invalid-group cases in Java/Python batch unit tests | | Terminal/running/queued recovery | `InternalSubagentRecoveryTest`, Java checkpoint/restart E2E, Python partial-result replay test | | Exceptions versus durable child results | Java bootstrap-failure/quiescence tests; Python bootstrap/polling-failure tests; existing cancellation tests | Coverage targets scheduling, worker starvation, durable slot indexing, checkpoint queue rotation, nested event replay, and the interval between child failure and parent result persistence. Not verified: a Python parent through a MiniCluster restart, batch recovery against a production Kafka/Fluss journal, rescaling across TaskManagers, or E2E on Flink versions other than 2.3. <details> <summary>Implementation invariants and verification evidence</summary> All slots are reserved before bootstrap. Admission counts submitted slots until their completion callback persists them. Recovery drops checkpointed child tasks and regenerates their graph from the waiting parent, avoiding dispatch before transient scopes exist. Completed actions replay same-call forwarded events and outputs, but not nested bootstrap commands. Failed actions persist their failure summary and release discarded event counts. Nested AgentPlan JSON now closes its root object explicitly so resource-provider type markers remain in the correct object. Verification uses JDK 21, Python 3.12 and Flink 2.3. Distribution artifacts are rebuilt and copied to the Python library directories as in `tools/build.sh`. - `mvn -B --no-transfer-progress -pl runtime -am -Dspotless.skip=true test`: 2,190 tests, 39 skipped, no failures. - `./tools/ut.sh -p`: 1,670 passed, 14 skipped; final core regression: 1,289 passed, 14 skipped. - Java `InternalSubagentBatchE2ETest`: bounded mixed-target/partial-failure batch and checkpoint-triggered restart, using a serialized test journal in the MiniCluster JVM. - Python `internal_subagent_test.py`: single call, failure, nested call and bounded 12-call mixed-target batch. - `./tools/lint.sh -c` with JDK 17; `./tools/check-license.sh`; `tools/check-agents-md.py`. </details> ### API No invocation signatures change. Combining deferred handles now resolves one durable batch; sequential waits remain sequential. New Java/Python/YAML configuration keys are `subagent.parallelism` (16) and `subagent.max-batch-size` (1024), both positive. These bound call counts per group, not payload bytes or total job concurrency. Single internal waits use the same scheduler. External async submission remains unchanged. ### 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.154.0 (GPT-6) -- 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]
