factory: state-machine workflow executor (factory.run/status/stop/resume, guarded transitions, re-entry) - #2401
sethkarten wants to merge 13 commits into
Conversation
Prime Agent performance — completedPR Overall: 0 regressed · 0 improved · 42 no clear change.
Python runtime
Session transport
UI interactions
Sandbox cost: ~$0.1107 — no inference calls. Methodology and samplesMain resolved at 2026-09-28T23:14:46.641848+00:00. Harness
|
…p race) Address the PR #2401 review findings: - foreach mixed-instance outcomes (finding 1): a node now fails at its FIRST permanently failed instance instead of waiting for all instances to settle. Previously a failure that settled before its siblings left the node stuck in running with every instance terminal and the failure policy never applied (the run ended via the stall guard); fail_fast could also never cancel in-flight siblings. The policy guard keeps the second and later permanent failures no-ops. - stop() race (finding 2): stop() sets a transitional "stopping" state before its first await; the control loop re-checks state before the completion check, _finalize only runs while the run is running, the node-failure policy records node errors but never overwrites a state another transition owns, and the executor-error path no longer overwrites a concurrent stop with failed. A stopped run now stays stopped: no new children admitted during the stop window, no spurious done + finished notice. Repeated stop() is idempotent (finding 7) and appends only one run_stopped ledger event. - resume() defers rate-limited admissions to the control loop with backoff exactly like run() (finding 3), so it never sleeps in the calling model turn. - fail_fast milestone detail now counts cancelled in-flight children. - tests: the host fake is now an async callable patched in directly (AsyncMock does not await async side effects) with per-child-id scripted outcomes and collect/delete suspension gates (finding 6) so races and mixed instances reach real paths; new tests cover mixed- instance foreach for every policy, the gated stop race, repeated stop, two-parent fan-in in one collect batch (finding 8), the ANSWER_CAPTURE_CAP slice (finding 9), and the resume deferral path.
|
Review round complete (for the morning review): the review found 2 must-fix bugs - foreach nodes with mixed instance outcomes never reached the failure policy, and stop() raced the control loop (run could finalize as done after an explicit stop). Both are fixed in commit 8344cdb with 7 regression tests (suite now 32, all green; gates clean): foreach now fails at the first permanently failed instance (fail_fast cancels in-flight siblings immediately), stop() uses a transitional stopping state with finalize guards, and resume() no longer sleeps in the model turn. Both fixes were verified by re-running the original reproductions. Remaining nits are documented V1 limitations in the PR body (settlement-only budgets, cancellation shortcut reporting). |
|
Integration note (verified by merging all open swarm branches together): main now enforces a deferred-boot contract - repl tests assert Version: ImageMagick 7.1.2-31 Q16-HDRI aarch64 8309dc92a:20260903 https://imagemagick.org Image Settings: Image Operators: Miscellaneous Options: By default, 'file' is written in the MIFF image format. To |
…p race) Address the PR #2401 review findings: - foreach mixed-instance outcomes (finding 1): a node now fails at its FIRST permanently failed instance instead of waiting for all instances to settle. Previously a failure that settled before its siblings left the node stuck in running with every instance terminal and the failure policy never applied (the run ended via the stall guard); fail_fast could also never cancel in-flight siblings. The policy guard keeps the second and later permanent failures no-ops. - stop() race (finding 2): stop() sets a transitional "stopping" state before its first await; the control loop re-checks state before the completion check, _finalize only runs while the run is running, the node-failure policy records node errors but never overwrites a state another transition owns, and the executor-error path no longer overwrites a concurrent stop with failed. A stopped run now stays stopped: no new children admitted during the stop window, no spurious done + finished notice. Repeated stop() is idempotent (finding 7) and appends only one run_stopped ledger event. - resume() defers rate-limited admissions to the control loop with backoff exactly like run() (finding 3), so it never sleeps in the calling model turn. - fail_fast milestone detail now counts cancelled in-flight children. - tests: the host fake is now an async callable patched in directly (AsyncMock does not await async side effects) with per-child-id scripted outcomes and collect/delete suspension gates (finding 6) so races and mixed instances reach real paths; new tests cover mixed- instance foreach for every policy, the gated stop race, repeated stop, two-parent fan-in in one collect batch (finding 8), the ANSWER_CAPTURE_CAP slice (finding 9), and the resume deferral path.
8344cdb to
39dac8d
Compare
|
Consolidated into #2452 (DAG swarm workflows as a single review/merge unit) per the DAG-first review plan. All commits, tests, and review history from this PR are included there; this branch is unchanged and can be deleted. |
…p race) Address the PR #2401 review findings: - foreach mixed-instance outcomes (finding 1): a node now fails at its FIRST permanently failed instance instead of waiting for all instances to settle. Previously a failure that settled before its siblings left the node stuck in running with every instance terminal and the failure policy never applied (the run ended via the stall guard); fail_fast could also never cancel in-flight siblings. The policy guard keeps the second and later permanent failures no-ops. - stop() race (finding 2): stop() sets a transitional "stopping" state before its first await; the control loop re-checks state before the completion check, _finalize only runs while the run is running, the node-failure policy records node errors but never overwrites a state another transition owns, and the executor-error path no longer overwrites a concurrent stop with failed. A stopped run now stays stopped: no new children admitted during the stop window, no spurious done + finished notice. Repeated stop() is idempotent (finding 7) and appends only one run_stopped ledger event. - resume() defers rate-limited admissions to the control loop with backoff exactly like run() (finding 3), so it never sleeps in the calling model turn. - fail_fast milestone detail now counts cancelled in-flight children. - tests: the host fake is now an async callable patched in directly (AsyncMock does not await async side effects) with per-child-id scripted outcomes and collect/delete suspension gates (finding 6) so races and mixed instances reach real paths; new tests cover mixed- instance foreach for every policy, the gated stop race, repeated stop, two-parent fan-in in one collect batch (finding 8), the ANSWER_CAPTURE_CAP slice (finding 9), and the resume deferral path.
8f57a69 to
b037304
Compare
…p race) Address the PR #2401 review findings: - foreach mixed-instance outcomes (finding 1): a node now fails at its FIRST permanently failed instance instead of waiting for all instances to settle. Previously a failure that settled before its siblings left the node stuck in running with every instance terminal and the failure policy never applied (the run ended via the stall guard); fail_fast could also never cancel in-flight siblings. The policy guard keeps the second and later permanent failures no-ops. - stop() race (finding 2): stop() sets a transitional "stopping" state before its first await; the control loop re-checks state before the completion check, _finalize only runs while the run is running, the node-failure policy records node errors but never overwrites a state another transition owns, and the executor-error path no longer overwrites a concurrent stop with failed. A stopped run now stays stopped: no new children admitted during the stop window, no spurious done + finished notice. Repeated stop() is idempotent (finding 7) and appends only one run_stopped ledger event. - resume() defers rate-limited admissions to the control loop with backoff exactly like run() (finding 3), so it never sleeps in the calling model turn. - fail_fast milestone detail now counts cancelled in-flight children. - tests: the host fake is now an async callable patched in directly (AsyncMock does not await async side effects) with per-child-id scripted outcomes and collect/delete suspension gates (finding 6) so races and mixed instances reach real paths; new tests cover mixed- instance foreach for every policy, the gated stop race, repeated stop, two-parent fan-in in one collect batch (finding 8), the ANSWER_CAPTURE_CAP slice (finding 9), and the resume deferral path.
b037304 to
706fb73
Compare
…p race) Address the PR #2401 review findings: - foreach mixed-instance outcomes (finding 1): a node now fails at its FIRST permanently failed instance instead of waiting for all instances to settle. Previously a failure that settled before its siblings left the node stuck in running with every instance terminal and the failure policy never applied (the run ended via the stall guard); fail_fast could also never cancel in-flight siblings. The policy guard keeps the second and later permanent failures no-ops. - stop() race (finding 2): stop() sets a transitional "stopping" state before its first await; the control loop re-checks state before the completion check, _finalize only runs while the run is running, the node-failure policy records node errors but never overwrites a state another transition owns, and the executor-error path no longer overwrites a concurrent stop with failed. A stopped run now stays stopped: no new children admitted during the stop window, no spurious done + finished notice. Repeated stop() is idempotent (finding 7) and appends only one run_stopped ledger event. - resume() defers rate-limited admissions to the control loop with backoff exactly like run() (finding 3), so it never sleeps in the calling model turn. - fail_fast milestone detail now counts cancelled in-flight children. - tests: the host fake is now an async callable patched in directly (AsyncMock does not await async side effects) with per-child-id scripted outcomes and collect/delete suspension gates (finding 6) so races and mixed instances reach real paths; new tests cover mixed- instance foreach for every policy, the gated stop race, repeated stop, two-parent fan-in in one collect batch (finding 8), the ANSWER_CAPTURE_CAP slice (finding 9), and the resume deferral path.
706fb73 to
fb55a82
Compare
Add the executor for stored swarm DAG entries on top of PR #2397's validator (rlm/swarm.py) and expose it as the rlm.swarm namespace (run/status/stop/resume) in rlm/__init__.py. The SwarmExecutor runs a canonicalized DAG through the existing RLM supervisor: nodes are admitted with rlm.spawn, settled through rlm.collect, and cancelled with rlm.delete_subagent. The supervisor owns the children; the executor owns run state in kernel memory. run() is a nonblocking admission phase: it re-validates the stored dag, resolves every subagent reference (reporting all failures before starting anything), reports the resolved node count and max_parallel, starts every ready node up to max_parallel, and returns while a background asyncio task continues the run. The control loop is iterative (tested at 100-node chain and 1000-node fan scale), polls collect with a 2s timeout, binds inputs ({name} placeholder render plus an appended Inputs section), expands foreach nodes with clamped instances, retries child failures, enforces node and run wall-clock budgets via an injectable clock, applies fail_fast/continue/escalate failure policies, and backs off rate-limited spawn admissions (doubling, capped at 60s, max 5 attempts) before failing the node. Progress reaches the parent through one quiet notice per run milestone (finished/failed/paused/budget-exceeded) via a new "swarm.progress" host handler in agent-session.ts (steer lane, mirroring bash.completed), and the run's event ledger records spawned/settled/answer_captured/retry/ error events with arrived/shown/delivered stages; status() returns node states plus the trailing 50 events and marks the ledger delivered. Runs do not survive a kernel restart (children are supervisor-owned and keep running); captured answers are collect previews capped by the host at 160 characters (compactRlmText).
…p race) Address the PR #2401 review findings: - foreach mixed-instance outcomes (finding 1): a node now fails at its FIRST permanently failed instance instead of waiting for all instances to settle. Previously a failure that settled before its siblings left the node stuck in running with every instance terminal and the failure policy never applied (the run ended via the stall guard); fail_fast could also never cancel in-flight siblings. The policy guard keeps the second and later permanent failures no-ops. - stop() race (finding 2): stop() sets a transitional "stopping" state before its first await; the control loop re-checks state before the completion check, _finalize only runs while the run is running, the node-failure policy records node errors but never overwrites a state another transition owns, and the executor-error path no longer overwrites a concurrent stop with failed. A stopped run now stays stopped: no new children admitted during the stop window, no spurious done + finished notice. Repeated stop() is idempotent (finding 7) and appends only one run_stopped ledger event. - resume() defers rate-limited admissions to the control loop with backoff exactly like run() (finding 3), so it never sleeps in the calling model turn. - fail_fast milestone detail now counts cancelled in-flight children. - tests: the host fake is now an async callable patched in directly (AsyncMock does not await async side effects) with per-child-id scripted outcomes and collect/delete suspension gates (finding 6) so races and mixed instances reach real paths; new tests cover mixed- instance foreach for every policy, the gated stop race, repeated stop, two-parent fan-in in one collect batch (finding 8), the ANSWER_CAPTURE_CAP slice (finding 9), and the resume deferral path.
…, re-entry, wait states) run() canonicalizes the stored spec to machine form (dag sugar compiles), enters the entry states, and drives the machine to quiescence: every settle is evaluated once and ALL guard-passing transitions fire (fan-out legal); max_entries blocking is recorded as transition_blocked events; self-loops and back-edges re-enter states with freshly re-bound inputs; foreach expands per entry; budgets/retries/failure policies stay V1 (per instance, per entry); wait states register rlm.watch.path/agent watches and settle the implicit event output (detail, timed_out) through the injectable clock; max_transitions pauses once with a budget_exceeded milestone. status() adds per-state entries_used/max_entries, per-entry breakdowns, and a transitions_fired count. Every V1 dag executor behavior is preserved through the compiler path: all 32 existing executor tests run unchanged.
Cover the state-machine semantics end to end with the fake host: a three-round review loop (guarded switch approved=false twice then true; 3 reviewing entries, 2 fixing entries, done), a mutually exclusive guard switch where only the matching branch enters, fan-out from one settle, max_entries blocking recorded as transition_blocked events with a quiescent done, a self-loop re-entering until its guard fails, wait-state settles from a completed path watch and from the injectable-clock timeout (event object bound into the dependent prompt), a max_transitions pause reported once and resumed, and stop() cancelling an active watch.
…ume, JSON-strict guards) Wait states are gated at the validator (the rlm.watch.* host handlers arrive with the communication series #2351/#2356), so the executor's wait registration/timeout/settle/cancel paths and their tests are removed entirely; quiescence, re-entry, and fan-out are unchanged. Review fixes: a max_transitions pause mid-settle now records the transition index and a resume continues after it (previously the settle re-fired already-fired transitions, duplicating spawned work); eq/ne guards compare JSON-strictly (a bool never equals a number, numbers compare numerically); contains requires a non-empty list value at validation and is defensively false at runtime; the max_transitions pause fires its own max_transitions_exceeded milestone kind (no collision with the run-budget budget_exceeded, both notice-able); the control-loop stall reports the pending entry whose input source never settled, with a test; EVENT_WINDOW grows 50 -> 200 (the pr-manager happy path is ~43 events before any retry); refinement validateEdit accepts machine-form swarm edits (either dag or machine object, never both). Also adds optional inputs: an input flagged optional binds a null sentinel instead of waiting when its source state never settled, which lets loop states (review with the previous fix report) run their first round before the fixer exists; compiled dags never emit it.
fb55a82 to
2c6e883
Compare
…ests check:test-policy flags every new sleep()/delay() call in test files, and asyncio.sleep(0) in the new factory executor tests tripped it on all three stack branches. Replace those cooperative yields with a yield_loop_turn() helper (call_soon + Event) that yields the loop with no wall-clock wait and no flagged token.
…observable stages Reviewer batch 1 findings on the factory executor: F1 (high): a child admitted after stop() was never cancelled. _admit now re-checks the run state after every await: a spawn that lands in a non-running run is retracted (delete_subagent + cancelled instance, never registered running) and a backoff that wakes in a non-running run never retries; both return stopped and _spawn_ready stops admitting. Pinned by test_stop_during_in_flight_admission_deletes_the_child and test_stop_during_admission_backoff_cancels_the_instance (FakeHost can now gate rlm.run admissions; GatedSleep suspends the backoff window). F2 (medium): _milestone returned before _event, so a repeated milestone (escalate pause -> resume -> fail again) left no ledger record. Every milestone is now appended; only the parent notice stays deduped per kind. Pinned by test_repeated_pause_records_every_milestone_in_the_ledger. F3 (question): documented that _NodeInstance.attempt counts admission attempts (429 deferrals included) and retries compares against it. F4 (question): _parse_json_output tried only the trailing fence; it now scans the fences trailing-first for the one carrying the output port, then falls back to the whole answer. Pinned by test_json_output_binds_from_the_fence_that_carries_the_port. F5 (question): status() collapsed every stage to delivered; it now advances only recorded/arrived events, so shown stays observable and the stage taxonomy is readable through status().
_halt_nonterminal cancelled the entry but left its prepared instances status pending, so a stopped/failed run reported pending instances the replay checker (factory-eval checkReplayLedger) rejects as impossible. The instance is now cancelled with its own ledger event (no child: it was cancelled before admission); _admit's stop-window retraction covers the in-flight case and deletes the child it returned.
…ations after failed deletes Reviewer round-2 findings on the factory executor: N1 (low): _halt_nonterminal built its cancellation list from state.status, which reports the LATEST entry only, so a re-entered state whose earlier entry was still in flight (later entry settled) was never reported or cancelled: stop() returned cancelled: [] and the stopped ledger still showed a running entry. The list is now built from entries (any non-terminal entry) plus never-entered states, so stop()/fail_fast cancel and report every in-flight entry. Pinned by test_stop_cancels_an_in_flight_earlier_entry_of_a_reentered_state (stop() -> cancelled ['x'], entries read [(0, cancelled), (1, done)]). N2 (low): a delete_subagent failure recorded only cancel_failed while the instance still reads cancelled, so the eval replay checker rejected the ledger (cancelled without a cancel event; spawned never cancelled). The ledger now records the cancellation alongside the failure (the slot is released either way), in both _cancel_running and _retract_admission. Pinned by test_failed_delete_records_cancel_failed_and_cancelled.
| "(the original DAG sugar in arguments['dag'] compiles to machine form): manage them with " | ||
| "create_factory/update_factory/delete_factory (create_factory validates either form at write time); run " | ||
| "them with rlm.factory.run(\"<id>\") once the executor lands in a follow-up PR.", | ||
| "them with await rlm.factory.run(\"<id>\"), watch with rlm.factory.status(run_id), stop with " |
There was a problem hiding this comment.
🟡 Medium rlm/harness.py:1054
The overview instructs callers to invoke rlm.factory.status, rlm.factory.stop, and rlm.factory.resume without await, so the REPL returns unexecuted coroutine objects instead of observing, stopping, or resuming the run. Add await to each async namespace call; otherwise a resident factory run can continue consuming children.
Also found in 2 other location(s)
packages/coding-agent/src/core/prompts/rlm.ts:49
The new prompt omits
awaitforrlm.factory.status(run_id),rlm.factory.stop(run_id), andrlm.factory.resume(run_id), although all three namespace methods are async. Following this instruction in the Python REPL merely creates unexecuted coroutine objects, so the model cannot actually inspect, stop, or resume a factory run.
packages/coding-agent/src/core/refinement/refinement.ts:721
The newly added factory usage hint omits
awaitfor bothrlm.factory.status(run_id)andrlm.factory.stop(run_id), although those methods are declaredasyncin the runtime. A model following this prompt merely creates unexecuted coroutine objects, so it neither observes a run nor stops it; in particular, a resident or stuck factory run continues consuming children.
🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @prime-agent-runtime/src/rlm/harness.py around line 1054:
The overview instructs callers to invoke `rlm.factory.status`, `rlm.factory.stop`, and `rlm.factory.resume` without `await`, so the REPL returns unexecuted coroutine objects instead of observing, stopping, or resuming the run. Add `await` to each async namespace call; otherwise a resident factory run can continue consuming children.
Evidence trail:
Reviewed commit c544820f. prime-agent-runtime/src/rlm/harness.py:1041-1055; prime-agent-runtime/src/rlm/__init__.py:531-551; prime-agent-runtime/src/rlm/repl.py:552-574; prime-agent-runtime/src/rlm/factory.py:1191-1206.
Also found in 2 other location(s):
- packages/coding-agent/src/core/prompts/rlm.ts:49 -- The new prompt omits `await` for `rlm.factory.status(run_id)`, `rlm.factory.stop(run_id)`, and `rlm.factory.resume(run_id)`, although all three namespace methods are async. Following this instruction in the Python REPL merely creates unexecuted coroutine objects, so the model cannot actually inspect, stop, or resume a factory run.
- packages/coding-agent/src/core/refinement/refinement.ts:721 -- The newly added factory usage hint omits `await` for both `rlm.factory.status(run_id)` and `rlm.factory.stop(run_id)`, although those methods are declared `async` in the runtime. A model following this prompt merely creates unexecuted coroutine objects, so it neither observes a run nor stops it; in particular, a resident or stuck factory run continues consuming children.
| if state.lifecycle == "resident" and entry.status == "running": | ||
| continue |
There was a problem hiding this comment.
🟠 High rlm/factory.py:1976
_run_complete reports quiescence while a resident running entry still has pending instances, so the run finalizes and exits before those children are admitted. This occurs when max_parallel leaves a resident transition queued; exclude resident entries with pending instances from the quiescent case.
| if state.lifecycle == "resident" and entry.status == "running": | |
| continue | |
| if state.lifecycle == "resident" and entry.status == "running" and not any( | |
| instance.status == "pending" for instance in entry.instances | |
| ): | |
| continue |
🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @prime-agent-runtime/src/rlm/factory.py around lines 1976-1977:
`_run_complete` reports quiescence while a resident `running` entry still has `pending` instances, so the run finalizes and exits before those children are admitted. This occurs when `max_parallel` leaves a resident transition queued; exclude resident entries with pending instances from the quiescent case.
Evidence trail:
prime-agent-runtime/src/rlm/factory.py:1465-1490, 1508-1512, 1966-1979, 2068-2073 at commit 3daf330
| for entry in state.entries: | ||
| for instance in entry.instances: | ||
| if instance.status == "pending": |
There was a problem hiding this comment.
🟠 High rlm/factory.py:1592
_next_pending_instance admits queued siblings even after a foreach entry has become terminal error, so failure_policy: "continue" still runs additional child work after the node failed. Skip pending instances belonging to non-running entries.
for entry in state.entries:
+ if entry.status != "running":
+ continue
for instance in entry.instances:🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @prime-agent-runtime/src/rlm/factory.py around lines 1592-1594:
`_next_pending_instance` admits queued siblings even after a `foreach` entry has become terminal `error`, so `failure_policy: "continue"` still runs additional child work after the node failed. Skip pending instances belonging to non-running entries.
Evidence trail:
c544820
prime-agent-runtime/src/rlm/factory.py:1465-1489
prime-agent-runtime/src/rlm/factory.py:1589-1596
prime-agent-runtime/src/rlm/factory.py:1785-1858
prime-agent-runtime/src/rlm/factory.py:1966-1979
| if self._run_complete(run): | ||
| await self._finalize(run) | ||
| return | ||
| if run.run_budget_ms is not None and not run.budget_reported: |
There was a problem hiding this comment.
🟠 High rlm/factory.py:2074
The initial admission phase can launch instances after run_budget_ms has expired, so a slow spawn fills max_parallel instead of stopping new work at the budget boundary. The budget check at line 2074 only runs in the background control loop, after run() has already called _spawn_ready; enforce the budget before each admission (including that initial path).
🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @prime-agent-runtime/src/rlm/factory.py around line 2074:
The initial admission phase can launch instances after `run_budget_ms` has expired, so a slow `spawn` fills `max_parallel` instead of stopping new work at the budget boundary. The budget check at line 2074 only runs in the background control loop, after `run()` has already called `_spawn_ready`; enforce the budget before each admission (including that initial path).
Evidence trail:
prime-agent-runtime/src/rlm/factory.py:1108-1112, 1304-1308, 1471-1481, 1624-1628, 2074-2089 at commit c544820f. Repository: https://github.com/PrimeIntellect-ai/prime-agent. Verification command: git show c544820f:prime-agent-runtime/src/rlm/factory.py
| f"resume with await rlm.factory.resume('{run.run_id}')", | ||
| ) | ||
| return | ||
| started = await self._spawn_ready(run, allow_backoff=True) |
There was a problem hiding this comment.
🟠 High rlm/factory.py:2089
_spawn_ready(..., allow_backoff=True) blocks the sole control loop during its rate-limit retry sleeps, so already-running children are not collected and run-budget checks are delayed for the entire backoff (up to 15 seconds). A child that completes during this interval is therefore observed late and can be classified as exceeding its per-state budget. Make admission backoff non-blocking to the control loop, or move retries into a separate task so collection and budget processing continue.
🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @prime-agent-runtime/src/rlm/factory.py around line 2089:
`_spawn_ready(..., allow_backoff=True)` blocks the sole control loop during its rate-limit retry sleeps, so already-running children are not collected and run-budget checks are delayed for the entire backoff (up to 15 seconds). A child that completes during this interval is therefore observed late and can be classified as exceeding its per-state budget. Make admission backoff non-blocking to the control loop, or move retries into a separate task so collection and budget processing continue.
Evidence trail:
prime-agent-runtime/src/rlm/factory.py:2046-2089, 1621-1643, 764-771, 1689-1724 at commit c544820f
| if run.state == "stopped": | ||
| return {"run_id": run.run_id, "state": "stopped", "cancelled": []} |
There was a problem hiding this comment.
🟡 Medium rlm/factory.py:1200
Concurrent stop() calls both enter _halt_nonterminal while the first call is awaiting child deletion, so the same children receive duplicate delete_subagent requests and a second run_stopped event is recorded. Guard the transitional "stopping" state as well as "stopped" before starting another cancellation pass.
| if run.state == "stopped": | |
| return {"run_id": run.run_id, "state": "stopped", "cancelled": []} | |
| if run.state in ("stopping", "stopped"): | |
| return {"run_id": run.run_id, "state": run.state, "cancelled": []} |
🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @prime-agent-runtime/src/rlm/factory.py around lines 1200-1201:
Concurrent `stop()` calls both enter `_halt_nonterminal` while the first call is awaiting child deletion, so the same children receive duplicate `delete_subagent` requests and a second `run_stopped` event is recorded. Guard the transitional `"stopping"` state as well as `"stopped"` before starting another cancellation pass.
Evidence trail:
Reviewed commit c544820. prime-agent-runtime/src/rlm/factory.py:1191-1206, 1868-1909, 1939-1962; prime-agent-runtime/src/rlm/__init__.py:460-479; prime-agent-runtime/test/test_factory_executor.py:964-968. Verify with `git show c544820 -- prime-agent-runtime/src/rlm/factory.py`.
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 2 potential issues.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 6dfd619. Configure here.
| parts.append(f"i{instance_index}") | ||
| if attempt > 1: | ||
| parts.append(f"a{attempt}") | ||
| return "-".join(parts) |
There was a problem hiding this comment.
Truncated child names can collide
Medium Severity
_child_name keeps only the first 20 characters of the state id, so two valid slugs that share a prefix (for example collect-findings-pass-1 and collect-findings-pass-2) produce the same sibling name. rlm.spawn requires unique sibling names, so the second admission fails and the state's failure_policy applies.
Reviewed by Cursor Bugbot for commit 6dfd619. Configure here.
| for instance in entry.instances: | ||
| if instance.status == "pending": | ||
| return state, entry, instance | ||
| return None |
There was a problem hiding this comment.
Failed foreach still admits siblings
Medium Severity
_next_pending_instance returns any pending instance and does not check whether its entry already failed. After a foreach entry hits a permanent failure under continue or escalate, later items that were never admitted are still spawned. Those children consume max_parallel slots and do extra work for an entry that is already terminal.
Additional Locations (1)
Reviewed by Cursor Bugbot for commit 6dfd619. Configure here.


Summary
PR G of the factory DAG feature: the executor. Builds on PR #2397 (factory entry kind,
validate_factory_spec/canonicalize_factory_spec/topological_order).prime-agent-runtime/src/rlm/factory.py+rlm/__init__.py): aFactoryExecutorexposed as therlm.factorynamespace —await rlm.factory.run(spec_id, name=?),rlm.factory.status(run_id),rlm.factory.stop(run_id),rlm.factory.resume(run_id)."factory.progress"host handler inagent-session.tsthat injects one quiet steer-lane notice per run milestone (mirroring thebash.completedpath), plus a hint sentence in the model prompt.State machine model
The factory specification is now a state machine stored in
arguments["machine"]. The original DAG form (arguments["dag"]) stays as sugar and compiles to machine form at canonicalization: a compiled dag is a restricted machine — each state entered at most once, forward-only, with guard-less transitions (one per effective dependency edge) andmax_entries: 1.{ "run": { "budget_ms": 600000, "failure_policy": "escalate", // fail_fast | continue | escalate (default) "max_parallel": 8, // 1..64 (default 8) "max_transitions": 100 // default 10 * states, hard cap 10000 }, "states": [ // 1..1024, unique slug ids, >= 1 entry state { "id": "reviewing", "entry": true, // admission enters every entry state "max_entries": 4, // bounded re-entry; further transitions are blocked "subagent": "reviewer", // required unless the state is a wait state (same forms as V1) "lifecycle": "task", // task | resident "inputs": [{ "name": "pr_url", "type": "text", "from": "entry.pr_url" }], // binds the LATEST settle of the source "outputs": [{ "name": "verdict", "type": "json" }], "budget_ms": 300000, "retries": 2, "failure_policy": "escalate", "foreach": { "over": "items", "max": 64 }, // expands per entry "wait": { "kind": "path", "target": "/tmp/x", "timeout_ms": 5000 } // wait states only } ], "transitions": [ { "from": "reviewing", "to": "fixing", "on": "settled", "when": { "output": "verdict", "path": "approved", "op": "eq", "value": false } } ] }max_entries; blocked transitions are recorded astransition_blockedledger events and the machine continues. Re-entry re-binds inputs from the current latest settles of the input sources. Self-loops and back-edge cycles are legal — the V1 acyclicity check is removed entirely. Completion is quiescence: no state in flight (running/pending/waiting) and no firable transition from any unconsumed settle.max_transitionsexceeded pauses the run once with abudget_exceededmilestone (resume-able).when, over the from state's latest settle output): opseq ne gt gte lt lte exists contains;pathresolves dotted paths inside json ports; numeric ops require numeric values,containsrequires a list,eq/nescalars.waitblock): no subagent, declared outputs, retries, or foreach, and never resident; entering registers anrlm.watch.path/rlm.watch.agenthost watch and settles the implicit json portevent: {"detail", "timed_out"}— from the watch outcome or its timeout through the injectable clock.stop(), notices, and the ledger's arrived/shown/read stages.This PR reworks the executor onto the machine form:
run()canonicalizes the stored spec (dag sugar compiles), admission enters the entry states, and the control loop evaluates transitions on every settle. All 32 V1 dag executor tests run unchanged through the compiler path (43 tests total with the machine suites). New ledger event kinds:state_entry(with entry index),transition_fired,transition_blocked(pluswait_settledaccepted by the replay checker for legacy ledgers);status()adds per-stateentries_used/max_entries, per-entry breakdowns, and atransitions_firedusage count over a 200-event window.Review-fix round: wait states are gated at the validator (the
rlm.watch.*host handlers arrive with the communication series #2351/#2356), so this PR removes the executor's wait registration/timeout/settle/cancel paths entirely — quiescence, re-entry, and fan-out are unchanged. Amax_transitionspause mid-settle now records the transition index where it landed and a resume continues after it (previously the settle re-fired already-fired transitions, duplicating spawned work);eq/neguards compare JSON-strictly (a bool never equals a number, numbers compare numerically);containsrequires a non-empty list value at validation and is defensively false at runtime; themax_transitionspause fires its ownmax_transitions_exceededmilestone kind (no collision with the run-budgetbudget_exceeded— both notice-able); the control-loop stall names the pending entry whose input source never settled (with a test); EVENT_WINDOW grows 50 → 200 (the pr-manager happy path is ~43 events before any retry); refinementvalidateEditaccepts machine-form factory edits (eitherarguments.dagorarguments.machineobject, never both). Machine inputs may declare"optional": trueto bind a null sentinel instead of waiting when their source state never settled — the enabler for the closed pr-manager review/fix loop in the eval PR (compiled dags never emit it, so dag semantics stay V1-exact).Spec: Swarm DAGs: declarative orchestration in Continual Harness (sections "Executor" and "Runtime contract").
Executor semantics
Ownership split. The RLM supervisor owns the children (admission via
rlm.spawn, settlement viarlm.collect, cancellation viarlm.delete_subagent); the executor owns run state in kernel memory. Every host call resolves through the module-levelrlmfunctions /host_requestat call time, so tests patchrlm.host_requestdirectly.Dry run = write time + run time.
create_factoryvalidates the graph at write time.run()re-validates and canonicalizes the stored dag, then resolves every node's subagent reference — a string is a harness subagent entry id or title (entry content becomes the prompt template; optionalmetadata.model/metadata.thinkingcarry spawn settings), an inline object uses its own fields. All reference failures are reported together in oneValueErrorand nothing starts. The resolved node count andmax_parallelare reported in the result. Actual admission limits (concurrency, recursion depth, provider rate limits) are enforced at spawn time through the backoff path; that split is the V1 "runtime tree limits" dry run.Nonblocking admission.
run()initializes per-state state, creates an entry for every entry state (astate_entryledger event each), then prepares and admits instances up tomax_parallel(respecting foreach expansion), records handles, and returns{"run_id", "spec_id", "name", "nodes", "max_parallel", "started", "pending"}. A kernel asyncio task (the sameloop.create_taskpattern as the job-watch poller, guarded so a dead bridge fails the run gracefully instead of wedging) continues the work; the calling model turn ends immediately.Control loop (iterative, no recursion). Polls in-flight instances with
await rlm.collect(ids, timeout_ms=2000). For each settled instance: capture the answer,duration_ms, mark done, append a ledger event; when an entry's last instance settles, the entry settles and captures the state's declared output ports (json ports parsed as in binding, text ports the captured string). Every settle is then evaluated once: ALL guard-passing outgoing transitions fire (fan-out is legal), each fire enters its target unless it is out ofmax_entries(recorded astransition_blockedand the machine continues), and re-entry re-binds inputs from the current latest settles. Completion is quiescence: no entry in flight (pending/running) and no unevaluated settle. (Wait states are specified but gated: therlm.watch.*host handlers arrive with the communication series #2351/#2356, so the validator rejectswaitblocks and this PR ships no wait execution paths.) Yields once per iteration so an instantly-settling host cannot hot-spin. Tested at 100-node chain and 1000-node wide-fan scale against a fake host.Input binding. For each input
{name, type, from: "<node>.<output>"}: text ports take the upstream captured answer; json ports parse the upstream text — preferring the trailing fenced```jsonblock whose object contains the output name, else the whole text, else the node errors. Rendering replaces{input_name}placeholders in a single pass; inputs without a placeholder are appended in a trailing## Inputssection, so no bound value is ever dropped. Binding failures never retry (a deterministic binding error would recur on every re-render); the node fails and its failure_policy applies.foreach. When the
overjson input resolves to a list of K items, the node spawns K instances (each item bound as that input's value), K clamped toforeach.max. A foreach node is done when all instances settle; its captured answer is its instances' answers joined with blank lines. A node fails at its FIRST permanently failed instance — it does not wait for its remaining instances — sofail_fastcancels in-flight siblings while they are still running, and a failure that settles before its successes can never strand the node in a mixed terminal state.Retries and failure policies. A child failure (collect
errorstatus/field) re-spawns the same rendered prompt whileattempts <= retries; then the nodefailure_policyapplies:fail_fast→ cancel every running child of the run viarlm.delete_subagentand mark the run failed;continue→ mark the node error and keep going (dependents without a data edge still run; data dependents fail at binding);escalate(default) → mark the node error and pause the run, spawn nothing further untilresume().Budgets (injectable clock). Wall-clock budgets measure admission to settlement. Per node: if the spawn-to-settle span exceeds
budget_ms, the attempt is treated as failed — the budget is spent, so no retry; the failure_policy applies. Run-level: elapsed time pastrun.budget_mspauses new spawns (children already in flight keep running) with abudget_exceededmilestone; the budget milestone reports once per run, and resuming after it is an explicit operator decision.Rate limits. If spawn admission fails with a rate-limit/429-style error, the executor retries with exponential backoff — doubling delays capped at 60s, at most 5 admissions per attempt — then fails the node through its failure_policy; the loop never crashes. In the
run()admission phase (and inresume(), exactly the same way) a rate limit does not sleep inside the calling turn: the node stays pending and the control loop retries it with backoff (keeps both calls nonblocking).stop() and races.
stop()sets a transitionalstoppingstate before its first await, so the control loop can neither admit new children nor finalize the run while the cancellations are in flight; the loop re-checks state before completion,_finalizeonly runs while the run isrunning, and the executor-error path never overwrites a concurrent stop. A stopped run therefore staysstopped(neverdone+ a spuriousfinishednotice). Repeatedstop()is idempotent and appends only onerun_stoppedledger event.Event ledger. Kernel memory, per run: spawned/attempts, settled (with
duration_ms), answer captured, retries, errors, and run milestones. Stages follow the spec:arrived(child answer settled and captured),shown(milestone notice injected),delivered(parent read viastatus()— marked on each call).status()returns state states (with per-stateentries_used/max_entries, per-entry breakdowns, and per-instance details), the trailing 50 events,elapsed_ms, and usage (spawns/settled/tool uses/running/transitions_fired).status/stop/resumeraiseValueErroron an unknown run id.Milestone notices. One quiet runtime notice per milestone kind per run (
finished/failed/paused/budget_exceeded), never per node, sent viahost_request("factory.progress", {run_id, kind, node?, detail}). The TS handler (createFactoryProgressHostHandler+createFactoryProgressMessage) validates the payload and injects a quiet custom message (customTypefactory_progress_notice, steer lane,queueIfBusy/resumeIfIdle,suppressAutonomousContinuation, mirroringbash.completed) that wakes the parent with[factory-progress run:<id>] finished|failed|paused|budget-exceeded: <detail>. A dead bridge cannot be told; the ledger keeps the milestone andstatus()still surfaces it.Resident nodes (V1 wake sources). A resident node is admitted like any other node and then stays alive under the parent session; its wake source is its own subagent prompt/tooling (watches, heartbeats, or continuous work inside the child) — no structural dry-run rule in V1. A run's declarative work completes when all non-resident nodes settle; resident children keep running (they cannot outlive the parent session) until
rlm.factory.stop(run_id)tears them down.Known limitations (documented deliberately)
rlm.list_subagents()can still see and stop them.rlm.collectpreviews, which the host caps at 160 characters (compactRlmText); the executor additionally caps at 200 defensively. Input binding and downstream prompts work on these previews; full child outputs stay in the child's own session.resume()(completed children keep their results for a latercollect).failed(continue policy finishes the graph, but the run is marked failed so the parent sees the error); stop() reportsstate: "stopped"and lists every non-terminal node as cancelled.delete_subagentfails, the instance is still reportedcancelledin the ledger and acancelledevent is recorded next tocancel_failed(the child keeps running supervisor-side; the executor treats its slot as released). Instances of cancelled entries are markedcancelledwith their own ledger event — including never-admitted (rate-limit-deferred) ones and admissions that land afterstop()(the returned child is deleted, not registered running).Tests
prime-agent-runtime/test/test_factory_executor.py(49 tests, deterministic, no live model; patchedrlm.host_requestasync fake routing by request type with per-child-id scripted outcomes, collect/delete suspension gates, injected fake clock and sleeps): dry-run rejection (invalid dag via a bypassed store entry, missing references listing all failures, unknown spec), admission starting exactly the in-degree-0 nodes withmax_parallelrespected, propagation with exact rendered prompts (text, fenced-json, whole-text json, bad json → node error without spawning, unplaced inputs appended, two-parent fan-in in one collect batch, answer-capture cap slice), foreach expansion + clamp + zero items, all three failure policies (fail_fast cascade cancels running children; continue finishes remaining nodes; escalate pauses + notices +resume()continues) plus mixed-instance foreach per policy (fail_fast cancels in-flight siblings; escalate pauses; continue finishes with node error), retries until attempts exhausted, node and run budgets via the fake clock, stop() cascade + states + idempotency, a gated stop-race test proving a stop mid-collect keepsstopped, finalizes nothing, and admits no new children, 429 admission deferral → backoff → success and backoff exhaustion → node error (including theresume()deferral path), 100-node chain and 1000-node fan scale under a generous wall bound, status delivered-marking + unknown-runValueError, resident node spawn/stay/stop, plus stop-window admission races (a spawn landing afterstop()is deleted and never registered; a backoff waking in a stopped run never retries), repeated-pause milestones in the ledger, multi-fence json binding,shownstages surviving astatus()read, re-entered-state stop cancelling an in-flight earlier entry, and thecancel_failed/cancelledpairing on failed deletes.packages/coding-agent/test/factory-executor.test.ts(13 tests, mirroring the async-bash completion test): bracket-grammar notice content,convertToLlmpass-through,startsAgentRunwake, payload validation and forwarding, invalid payload rejection.prime-agent-runtime:uv run python -m unittest discover -s test— 483 tests, only the pre-existingtest_bashfailure(s) reproduced on the clean base branch (test_windows_without_bash_raises_teaching_error;test_unconfigured_journal_stays_permissivealso fails on some machines).packages/coding-agent:npx vitest --run test/factory-executor.test.ts test/refinement.test.ts— 105 passed. Repo root:npx tsgo --noEmitandnpx biome checkon changed files — clean.Stacked: base branch is swarm/swarm-spec-harness-entry (#2397); merge that first.
Draft only - review withheld per workflow.
Note
High Risk
Introduces async orchestration over subagent spawn/collect/delete with in-memory run state, race-prone stop/resume paths, and parent-session message injection—core delegation and session continuity are affected.
Overview
Adds the state-machine factory executor behind
await rlm.factory.run/status/stop/resume: storeddagormachinespecs are canonicalized, entry states are admitted viarlm.spawn, and a background control loop settles children, evaluates guarded transitions (fan-out, re-entry,max_entriesblocks), and applies retries, budgets, and fail_fast/continue/escalate policies until quiescence or pause/stop.Run milestones reach the parent through a new
factory.progresshost path that injects quietfactory_progress_noticecustom messages (treated as agent-run boundaries). Harness/refinement docs and validation now accept eitherarguments.dagorarguments.machine(not both); machine validation gains optional inputs and strictercontainsguards. Wait-state execution remains rejected untilrlm.watch.*lands.Large test suites cover executor semantics (Python) and progress messaging (TypeScript).
Reviewed by Cursor Bugbot for commit 6dfd619. Bugbot is set up for automated code reviews on this repo. Configure here.
Note
Add
FactoryExecutorstate-machine workflow runner withrlm.factoryrun/status/stop/resume APIFactoryExecutorin factory.py. It runs canonical factory machines through RLM supervisor children, tracks runs in an in-memory registry, and evaluates guarded transitions, retries, failure policies, budgets, re-entry, and cancellation in an async control looprlm.factory.run,status,stop, andresumethrough the new_RLMFactoryNamespacein init.py, backed by module-level wrappers on a lazy default executorfactory.progresshost handling in messages.ts, agent-session.ts, and rlm-runtime.ts so run milestones appear as agent messages that start a new agent runmachineordagfactory arguments in refinement edits (refinement.ts) and rejects both, neither, or non-object forms; tightens validation of boolean inputoptionalflags and non-empty contains-guard values📊 Macroscope summarized 6dfd619. 10 files reviewed, 8 issues evaluated, 2 issues filtered, 6 comments posted
🗂️ Filtered Issues
packages/coding-agent/src/core/prompts/rlm.ts — 0 comments posted, 1 evaluated, 1 filtered
awaitforrlm.factory.status(run_id),rlm.factory.stop(run_id), andrlm.factory.resume(run_id), although all three namespace methods are async. Following this instruction in the Python REPL merely creates unexecuted coroutine objects, so the model cannot actually inspect, stop, or resume a factory run. [ Cross-file consolidated ]packages/coding-agent/src/core/refinement/refinement.ts — 0 comments posted, 1 evaluated, 1 filtered
awaitfor bothrlm.factory.status(run_id)andrlm.factory.stop(run_id), although those methods are declaredasyncin the runtime. A model following this prompt merely creates unexecuted coroutine objects, so it neither observes a run nor stops it; in particular, a resident or stuck factory run continues consuming children. [ Cross-file consolidated ]