From f9bc89f38933da7ad4a561673f7848073de3b6a4 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 21 Sep 2026 21:11:08 +0800 Subject: [PATCH 1/3] fix(collaboration): scope delegation readiness and preserve peer task candidates Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../codex-subagent-orchestration.md | 52 ++++-- .../collaboration/delegation_context.py | 5 +- .../control_plane/effect_runtime_handlers.ts | 5 + .../control_plane/quota/peer_orchestration.ts | 62 +++++++ .../control_plane/quota/task_orchestration.py | 170 +++++------------- loopx/control_plane/quota/turn_envelope.ts | 27 ++- loopx/control_plane/subagent_context.ts | 9 +- .../presentation/renderers/quota_markdown.py | 28 ++- .../control_plane/test_delegation_context.py | 39 +++- .../test_quota_cli_projection.py | 21 +++ .../test_task_orchestration_admission.py | 51 ++++++ tests/control_plane_ts/agent_context.test.ts | 6 +- .../peer_orchestration.test.ts | 43 +++++ tests/control_plane_ts/turn_envelope.test.ts | 35 ++++ 14 files changed, 399 insertions(+), 154 deletions(-) create mode 100644 loopx/control_plane/quota/peer_orchestration.ts create mode 100644 tests/control_plane_ts/peer_orchestration.test.ts diff --git a/docs/integrations/codex-subagent-orchestration.md b/docs/integrations/codex-subagent-orchestration.md index 9f51b18a24..f7ab71dd8a 100644 --- a/docs/integrations/codex-subagent-orchestration.md +++ b/docs/integrations/codex-subagent-orchestration.md @@ -295,12 +295,23 @@ loopx quota should-run \ --available-capability peer_agent_activation ``` -The contract includes only peer lanes that are currently actionable. Dormant -registered agents and closed, blocked, or deferred todos are not coordinator -candidates. A dormant or non-resumable lane is projected under -`blocked_peer_lanes`; if no peer lane can run, the bundle has +The peer contract is scoped to `execution_scope=peer_agent_activation`. +Its `task_selection=canonical_claimed_candidates` identifies open claimed +Todo candidates, not the task pinned to a running session. Canonical inventory +outranks display rows; all eligible tasks for the same peer remain visible. +Large inventories stay in the full decision. The thin TurnEnvelope preserves +scoped gates, counts and a signed content hash plus a required detail read; +counts alone never authorize a task selection. Native child lanes are preserved +when only their nested peer diagnostic needs compaction. +Only currently actionable candidates appear under `eligible_peer_lanes`. +Closed, blocked, or deferred Todos are excluded. An open candidate whose peer +is dormant or whose dependency is not ready appears under `blocked_peer_lanes`; +if no peer lane can run, the bundle has `execution_state=blocked`, `terminal_outcome=blocked`, and `retry_policy=material_peer_state_change_only`. +When native child admission succeeds, its adaptive contract remains actionable +and retains the blocked peer contract as `peer_activation_diagnostic`. Neither +path grants permission to use a different entrypoint. That blocked diagnostic does not replace the coordinator's own runnable lane or re-arm an activation obligation on every heartbeat. If the coordinator also has no in-scope runnable fallback, the final interaction mode is @@ -467,7 +478,7 @@ At `before_plan`, `loopx agent-context` considers at most six bindings authorize for the current requester and projects as many as fit the existing context byte budget; `authorized_count` and `routes_truncated` make omissions explicit. Each route carries only binding/Agent/Todo/runtime -identity, `ready|blocked|unknown`, an optional public-safe execution profile and +identity, separate `runtime_readiness` and `readiness` observations, an optional public-safe execution profile and one stable `loopx delegation` entrypoint. A separate, explicit `agent-context --phase after_delegate_result` read may include bounded operation-status and recovery-required counts. Automatic planning and managed @@ -475,15 +486,28 @@ return paths do not enumerate the operation journal. These reads do not launch, resume or accept a worker, and they never expose raw host arguments, workspaces, output references, credentials or child results. -Runtime availability and business adoption are separate. `ready` says that the -existing runtime owner did not find a launch blocker; it does not select the -route, establish task fit or prove execution. `blocked` or `unknown` never -prevents useful native work. Before dispatch, the coordinator must use the -authorized binding, recheck its runtime/model/budget and record a stable -operation id; it must not silently substitute another runtime or model. No -heartbeat is required to use every route. When Lark or another managed surface -exposes these facts, it must consume the same signed capability context rather -than own another route configuration. +Runtime availability, entrypoint admission and business adoption are separate. +`runtime_readiness=ready` only reports the existing runtime probe. Planning does +not run the execution preflight: `execution_scope=bound_delegation` and +`preflight=required` accompany `readiness=unknown` (or `blocked` for a known +runtime failure). Use `loopx delegation inspect` on the selected binding to +check canonical authority, task validation and Turn admission before dispatch. +A peer-activation blocker does not assess this route or native children; neither +does a successful runtime probe bypass authority promotion or another gate. +Before dispatch, recheck runtime/model/budget and record a stable operation id; +never silently substitute a runtime or model. User preferences guide independent +batch selection; no heartbeat must relaunch every route. Lark and other managed +surfaces consume this same capability context, not another route configuration. + +中文:peer 合约的阻塞范围仅为 `peer_agent_activation`,候选来自完整权威任务源, +不再以显示列表第一条任务冒充绑定;同一 Agent 的多条候选保留。 +候选过多时,全量决策保留全部条目,简版携带状态、数量、哈希和必读详情引用; +不能凭数量选择任务。若 native 子 Agent +通过独立准入,使用 adaptive 合约并保留 `peer_activation_diagnostic`,不因另一个 +入口受阻而提前返回。运行库/凭据可用只记为 `runtime_readiness`;委派入口尚未预检时 +`readiness=unknown`、`preflight=required`,通过现有 `loopx delegation inspect` +核验权威状态、任务验证与 Turn 准入。未知不等于禁用,ready 运行库也不等于可执行。 +用户的异构偏好影响批次选择,不要求每次心跳重启所有路线,不绕过任一真实门禁。 中文:可在现有 `multi_subagent` 能力中配置 `.loopx/config/delegations.json` 指针,让当前请求 Agent 在规划前看到自己已获授权的 diff --git a/loopx/control_plane/collaboration/delegation_context.py b/loopx/control_plane/collaboration/delegation_context.py index 44ea509794..4d8231ad9e 100644 --- a/loopx/control_plane/collaboration/delegation_context.py +++ b/loopx/control_plane/collaboration/delegation_context.py @@ -56,7 +56,8 @@ def _route(binding: dict[str, Any]) -> dict[str, Any]: "todo_id": binding["todo_id"], "runtime_id": executor.get("executor") or "unknown", "executor_kind": executor.get("executor_kind") or "generic", - "readiness": readiness, + "runtime_readiness": readiness, + "readiness": "blocked" if available is False else "unknown", } profile = str(executor.get("execution_profile") or "").strip() if profile: @@ -116,6 +117,8 @@ def project_delegation_context( result = { "schema_version": "loopx_delegation_context_v0", "configuration_state": "ready", + "execution_scope": "bound_delegation", + "preflight": "required", "observed_at": observed_at, "authorized_count": len(bindings), "projected_count": len(routes), diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 386fad6a8d..4d509c101e 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -1,5 +1,6 @@ import {selectPeriodicReportProgress, selectPeriodicReportApprovalRetry} from "./capabilities/periodic_report_progress.ts"; import {planIssueFixMonitorReconciliation} from "./capabilities/issue_fix_monitor_reconciliation.ts"; +import {projectPeerOrchestration} from "./quota/peer_orchestration.ts"; import {inspectTaskLease} from "./work_items/task_lease_inspection.ts"; import {evaluateTodoPriority} from "./todos/priority.ts"; import {evaluateUserCompletion} from "./todos/user_completion.ts"; @@ -623,6 +624,10 @@ export function createEffectRuntimeHandlers( "governed_capability.settlement_status", (params) => governedCapabilitySettlementStatus(params.failure), ], + [ + "quota.peer_orchestration.project", + (params) => projectPeerOrchestration(params), + ], [ "capability_hook.agent_context.describe", () => describeSubagentContext(), diff --git a/loopx/control_plane/quota/peer_orchestration.ts b/loopx/control_plane/quota/peer_orchestration.ts new file mode 100644 index 0000000000..073130c52e --- /dev/null +++ b/loopx/control_plane/quota/peer_orchestration.ts @@ -0,0 +1,62 @@ +/** Peer activation admission only; native children and bound delegation have separate gates. */ +import type { JsonObject } from "../effect_program.ts"; +import { jsonObject, requireJsonObject } from "../runtime_decode.ts"; + +type ActivationState = "ready" | "blocked"; +const activeStates = new Set(["running", "monitoring", "executing", "bound", "launchable"]); +const rows = (value: unknown): JsonObject[] => Array.isArray(value) + ? value.flatMap(item => { const row = jsonObject(item); return row ? [row] : []; }) : []; + +export function projectPeerOrchestration(value: unknown): JsonObject | null { + const input = requireJsonObject(value, "peer orchestration"); + const coordinator = String(input.agent_id ?? ""); + const registered = new Set(Array.isArray(input.registered_agents) ? input.registered_agents : []); + const activation = Array.isArray(input.available_capabilities) + && input.available_capabilities.includes("peer_agent_activation"); + const runtime = new Map(rows(input.agents).map(row => [row.agent_id, row])); + // A peer can own several open tasks. Never let the first display row hide + // another canonical task, or describe a claimed candidate as a session binding. + const candidates = rows(input.items).filter(row => row.done !== true + && ["", "open"].includes(String(row.status ?? "").trim().toLowerCase()) + && row.task_class === "advancement_task" && row.claimed_by + && row.claimed_by !== coordinator && registered.has(row.claimed_by) + && typeof row.todo_id === "string" && row.todo_id.length > 0); + const unique = new Map(candidates.map(row => [JSON.stringify([row.claimed_by, row.todo_id]), row])); + const eligible: JsonObject[] = [], blocked: JsonObject[] = []; + for (const row of [...unique.values()].sort((a, b) => + String(a.claimed_by).localeCompare(String(b.claimed_by)) + || String(a.todo_id).localeCompare(String(b.todo_id)))) { + const lane: JsonObject = { + agent_id: row.claimed_by, todo_id: row.todo_id, + priority: row.priority ?? null, task_class: row.task_class, + action_kind: row.action_kind ?? null, + title: String(row.title ?? row.text ?? "").trim(), + resume_when: row.resume_when ?? null, resume_ready: row.resume_ready ?? null, + }; + const reasons: string[] = []; + if (!activation) reasons.push("peer_agent_activation_unavailable"); + const peer = runtime.get(row.claimed_by); + if (!peer) reasons.push("peer_liveness_unavailable"); + else if (peer.stale_claim_hint) reasons.push("peer_runtime_stale"); + else if (!activeStates.has(String(peer.state))) reasons.push("peer_runtime_not_active"); + if (lane.resume_when && lane.resume_ready !== true) reasons.push("peer_lane_not_resume_ready"); + if (reasons.length) blocked.push({ ...lane, reason_codes: reasons }); + else eligible.push(lane); + } + if (!eligible.length && !blocked.length) return null; + const state: ActivationState = eligible.length ? "ready" : "blocked"; + return { + schema_version: "task_orchestration_contract_v1", mode: "task_scoped_peer", + coordinator_agent_id: coordinator, + execution_scope: "peer_agent_activation", task_selection: "canonical_claimed_candidates", + execution_state: state, activation_required: eligible.length > 0, + activation_allowed: activation, required_capability: "peer_agent_activation", + eligible_peer_lanes: eligible, blocked_peer_lanes: blocked, + retry_policy: "material_peer_state_change_only", + terminal_outcome: state === "blocked" ? "blocked" : null, + writeback_owner: "task_coordinator", + coordinator_obligation: eligible.length + ? "Activate or resume eligible peer tasks; multiple tasks may share a peer. Review returned evidence before writeback." + : "Peer activation is blocked for these tasks; retry this entrypoint only after its gates change. This does not assess native children or bound delegation; inspect their own admission before use.", + }; +} diff --git a/loopx/control_plane/quota/task_orchestration.py b/loopx/control_plane/quota/task_orchestration.py index d73b5499d3..84a982deb5 100644 --- a/loopx/control_plane/quota/task_orchestration.py +++ b/loopx/control_plane/quota/task_orchestration.py @@ -4,6 +4,7 @@ from ..agents.agent_scope_frontier import AgentScopeFrontierAction from ..agents.runtime_model import peer_work_key +from ..effect_runtime import effect_runtime_result from ..todos.contract import ( normalize_required_capabilities, normalize_todo_claimed_by, @@ -261,6 +262,7 @@ def _task_orchestration_contract( configured_coordinator = normalize_todo_claimed_by( peer_coordination.get("coordinator_agent_id") ) + peer_contract = None if ( peer_coordination.get("enabled") is True and configured_coordinator == agent_id @@ -269,20 +271,21 @@ def _task_orchestration_contract( agent_id=agent_id, agent_identity=agent_identity, raw_agent_todo_summary=raw_agent_todo_summary, + agent_todo_source_items=agent_todo_source_items, available_capabilities=available, agent_management_projection=agent_management_projection, ) - if peer_contract: + if task_orchestration_contract_is_actionable(peer_contract): return peer_contract if orchestration.get("mode") != "multi_subagent": - return None + return peer_contract if orchestration.get("spawn_allowed") is not True: - return None + return peer_contract max_children = orchestration.get("max_children") if not isinstance(max_children, int) or max_children <= 0: - return None + return peer_contract if SUBAGENT_SPAWN_CAPABILITY in available: - return build_adaptive_task_orchestration_contract( + native_contract = build_adaptive_task_orchestration_contract( agent_id=agent_id, agent_identity=agent_identity, goal_boundary=goal_boundary, @@ -295,7 +298,11 @@ def _task_orchestration_contract( parent_goal_id=parent_goal_id, max_children=max_children, ) - return None + if native_contract: + if peer_contract: + native_contract["peer_activation_diagnostic"] = peer_contract + return native_contract + return peer_contract def _registered_peer_task_orchestration_contract( @@ -303,129 +310,46 @@ def _registered_peer_task_orchestration_contract( agent_id: str, agent_identity: dict[str, Any], raw_agent_todo_summary: dict[str, Any] | None, + agent_todo_source_items: list[dict[str, Any]] | None, available_capabilities: list[str], agent_management_projection: dict[str, Any] | None, ) -> dict[str, Any] | None: - source_items = ( - raw_agent_todo_summary.get("items") - if isinstance(raw_agent_todo_summary, dict) - and isinstance(raw_agent_todo_summary.get("items"), list) - else [] - ) - candidate_lanes: list[dict[str, Any]] = [] - seen_agents: set[str] = set() - registered_agents = set(agent_identity.get("registered_agents") or []) - for item in source_items: - if not isinstance(item, dict): - continue - status = str(item.get("status") or "").strip().lower() - if item.get("done") is True or status not in {"", "open"}: - continue - if str(item.get("task_class") or "") != "advancement_task": - continue - peer_agent = normalize_todo_claimed_by(item.get("claimed_by")) - if ( - not peer_agent - or peer_agent in seen_agents - or peer_agent not in registered_agents - or peer_agent == agent_id - ): - continue - candidate_lanes.append( - { - "agent_id": peer_agent, - "todo_id": str(item.get("todo_id") or "").strip() or None, - "priority": item.get("priority"), - "task_class": item.get("task_class"), - "action_kind": item.get("action_kind"), - "title": str(item.get("title") or item.get("text") or "").strip(), - "resume_when": item.get("resume_when"), - "resume_ready": item.get("resume_ready"), - } - ) - seen_agents.add(peer_agent) - if not candidate_lanes: + # Canonical inventory, including an explicitly empty list, outranks display rows. + source_items = agent_todo_source_items + if source_items is None: + source_items = (raw_agent_todo_summary or {}).get("items", []) + fields = ("todo_id", "done", "status", "task_class", "priority", "action_kind", + "title", "text", "resume_when", "resume_ready") + items = [ + {**{key: item[key] for key in fields if key in item}, + "claimed_by": normalize_todo_claimed_by(item.get("claimed_by"))} + for item in (source_items or []) if isinstance(item, dict) + ] + management_rows = (agent_management_projection or {}).get("agents") + agents = [ + {"agent_id": normalize_todo_claimed_by(row.get("agent_id")), + "state": row.get("state"), "stale_claim_hint": row.get("stale_claim_hint")} + for row in (management_rows if isinstance(management_rows, list) else []) + if isinstance(row, dict) + ] + contract = effect_runtime_result("quota.peer_orchestration.project", { + "agent_id": agent_id, + "registered_agents": agent_identity.get("registered_agents") or [], + "items": items, + "available_capabilities": available_capabilities, + "agents": agents, + }) + if contract is None: return None - assignment_key = peer_work_key( - { - "mode": "task_scoped_peer", - "lanes": sorted( - [ - {"agent_id": lane["agent_id"], "todo_id": lane["todo_id"]} - for lane in candidate_lanes - ], - key=lambda lane: (lane["agent_id"], lane["todo_id"] or ""), - ), - }, - fallback="task_orchestration", - ) - activation_available = PEER_AGENT_ACTIVATION_CAPABILITY in available_capabilities - peer_runtime_state = _peer_runtime_state_by_agent(agent_management_projection) - eligible_peer_lanes: list[dict[str, Any]] = [] - blocked_peer_lanes: list[dict[str, Any]] = [] - for lane in candidate_lanes: - reason_codes: list[str] = [] - if not activation_available: - reason_codes.append("peer_agent_activation_unavailable") - peer_state = peer_runtime_state.get(str(lane["agent_id"])) - if not peer_state: - reason_codes.append("peer_liveness_unavailable") - elif peer_state.get("stale_claim_hint"): - reason_codes.append("peer_runtime_stale") - # New work-state refinements retain the old running admission policy. - # Activity/claim freshness, activation capability and dependency readiness - # remain independent gates; a binding alone never grants activation. - elif peer_state.get("state") not in { - "running", "monitoring", "executing", "bound", "launchable", - }: - reason_codes.append("peer_runtime_not_active") - if lane.get("resume_when") and lane.get("resume_ready") is not True: - reason_codes.append("peer_lane_not_resume_ready") - if reason_codes: - blocked_peer_lanes.append({**lane, "reason_codes": reason_codes}) - else: - eligible_peer_lanes.append(lane) - execution_state = "ready" if eligible_peer_lanes else "blocked" - return { - "schema_version": "task_orchestration_contract_v1", + lanes = contract["eligible_peer_lanes"] + contract["blocked_peer_lanes"] + contract["assignment_key"] = peer_work_key({ "mode": "task_scoped_peer", - "coordinator_agent_id": agent_id, - "assignment_key": assignment_key, - "execution_state": execution_state, - "activation_required": bool(eligible_peer_lanes), - "activation_allowed": activation_available, - "required_capability": PEER_AGENT_ACTIVATION_CAPABILITY, - "eligible_peer_lanes": eligible_peer_lanes, - "blocked_peer_lanes": blocked_peer_lanes, - "retry_policy": "material_peer_state_change_only", - "terminal_outcome": "blocked" if execution_state == "blocked" else None, - "writeback_owner": "task_coordinator", - "coordinator_obligation": ( - "activate or resume eligible peer lanes, review returned evidence, " - "then write accepted state/todos for this task bundle" - if eligible_peer_lanes - else "peer task bundle is blocked; do not retry until peer activation " - "capability or peer readiness changes" + "lanes": sorted( + [{"agent_id": lane["agent_id"], "todo_id": lane["todo_id"]} for lane in lanes], + key=lambda lane: (lane["agent_id"], lane["todo_id"]), ), - } - - -def _peer_runtime_state_by_agent( - projection: dict[str, Any] | None, -) -> dict[str, dict[str, Any]]: - agents = ( - projection.get("agents") - if isinstance(projection, dict) - and isinstance(projection.get("agents"), list) - else [] - ) - return { - agent_id: item - for item in agents - if isinstance(item, dict) - for agent_id in [normalize_todo_claimed_by(item.get("agent_id"))] - if agent_id - } + }, fallback="task_orchestration") + return contract def _task_orchestration_work_lane_contract( diff --git a/loopx/control_plane/quota/turn_envelope.ts b/loopx/control_plane/quota/turn_envelope.ts index 34ba28c525..fc2d1e64e2 100644 --- a/loopx/control_plane/quota/turn_envelope.ts +++ b/loopx/control_plane/quota/turn_envelope.ts @@ -9,7 +9,7 @@ import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; import { turnStartPromptBudgetBytes } from "../capability_hooks.ts"; import { requireJsonObject } from "../runtime_decode.ts"; import { projectPendingCapabilityIntent } from "../work_items/pending_capability_intent.ts"; -import { measureTurnEnvelope, turnEnvelopeBudgetBytes } from "./turn_envelope_budget.ts"; +import { measureTurnEnvelope, turnEnvelopeBudgetBytes, TURN_ENVELOPE_SECTION_TARGETS } from "./turn_envelope_budget.ts"; export { TURN_ENVELOPE_BUDGET_BYTES } from "./turn_envelope_budget.ts"; export const TURN_ENVELOPE_SCHEMA_VERSION = "loopx_turn_envelope_v0"; @@ -689,8 +689,33 @@ function actionProjection(payload: JsonObject, protocolActionFields: JsonObject) return projection; } +function compactPeerActivation(contract: JsonObject, detailRef: string): JsonObject { + if (contract.mode !== "task_scoped_peer" + || Buffer.byteLength(JSON.stringify(contract), "utf8") <= TURN_ENVELOPE_SECTION_TARGETS.context) return contract; + const summary: JsonObject = {}; + for (const field of ["schema_version", "mode", "coordinator_agent_id", "assignment_key", + "execution_scope", "task_selection", "execution_state", "activation_required", + "activation_allowed", "required_capability", "retry_policy", "terminal_outcome", "writeback_owner"]) { + if (field in contract) summary[field] = contract[field]; + } + return { ...summary, delivery: "projected", read_required: true, + eligible_peer_count: Array.isArray(contract.eligible_peer_lanes) ? contract.eligible_peer_lanes.length : 0, + blocked_peer_count: Array.isArray(contract.blocked_peer_lanes) ? contract.blocked_peer_lanes.length : 0, + content_hash: canonicalHash(contract), detail_ref: detailRef, + instruction: "Read the envelope full_decision detail command before selecting or activating a peer task; counts are not executable lanes.", + }; +} + function turnActionProjection(payload: JsonObject, protocolActionFields: JsonObject): JsonObject { const projection = actionProjection(payload, protocolActionFields); + const orchestration = object(projection.task_orchestration_contract); + if (Object.keys(orchestration).length > 0) { + const peer = object(orchestration.peer_activation_diagnostic); + projection.task_orchestration_contract = orchestration.mode === "adaptive" && Object.keys(peer).length > 0 + ? { ...orchestration, peer_activation_diagnostic: compactPeerActivation(peer, + "full_decision.task_orchestration_contract.peer_activation_diagnostic") } + : compactPeerActivation(orchestration, "full_decision.task_orchestration_contract"); + } const action = object(projection.action); const horizon = object(action.planning_horizon); if (Object.keys(object(horizon.detail_refs)).length > 0) { diff --git a/loopx/control_plane/subagent_context.ts b/loopx/control_plane/subagent_context.ts index 785249aa3d..31cc76b35e 100644 --- a/loopx/control_plane/subagent_context.ts +++ b/loopx/control_plane/subagent_context.ts @@ -4,14 +4,14 @@ import type { JsonObject } from "./effect_program.ts"; import { jsonObject, requireJsonObject } from "./runtime_decode.ts"; export const subagentContextProvider: AgentContextProvider = { - hookId: "multi_subagent.coordinator", capabilityId: "multi_subagent", revision: "v4", + hookId: "multi_subagent.coordinator", capabilityId: "multi_subagent", revision: "v5", phases: AGENT_CONTEXT_PHASES, produce(input, config) { const guidance = { before_plan: [ "Prefer bounded independent delegation; max_children is a configured ceiling, not live availability. Admit native children incrementally, avoid duplicate reads, and keep one parent question.", "Native child tools can read loopx agent-context at before_delegate and after_delegate_result; these calls do not start Turns or spend quota.", - "Authorized routes are observations, not obligations. Use a ready route only when it fits; blocked/unknown routes never block native work. No heartbeat must use every route.", + "Choose independent work from current routes and user preferences. Runtime availability is not entrypoint admission; inspect bound delegation. Peer-activation blocks do not assess native children or delegation. Do not relaunch every route on every heartbeat.", ], before_delegate: [ "Give each child a bounded question, sources, read/write limits, expected evidence and stopping condition; identify dependencies and the coordinator's concurrent question.", @@ -119,7 +119,9 @@ function boundedDelegationContext(value: unknown): JsonObject | null { || !["ready", "blocked", "unknown"].includes(readiness)) return []; const compact: JsonObject = { binding_id: bindingId, agent_id: agentId, todo_id: todoId, - runtime_id: runtimeId, readiness, + runtime_id: runtimeId, readiness: readiness === "blocked" ? "blocked" : "unknown", + runtime_readiness: ["ready", "blocked", "unknown"].includes(String(route.runtime_readiness)) + ? route.runtime_readiness! : readiness, }; const executorKind = identifier(route.executor_kind); const profile = typeof route.execution_profile === "string" @@ -149,6 +151,7 @@ function boundedDelegationContext(value: unknown): JsonObject | null { authorized_count: boundedCount(source.authorized_count), projected_count: 0, entrypoint: "loopx delegation", + execution_scope: "bound_delegation", preflight: "required", routes: [], }; if (rawReceipts) result.operation_receipts = operationReceipts; diff --git a/loopx/presentation/renderers/quota_markdown.py b/loopx/presentation/renderers/quota_markdown.py index 409b3fab1b..1e39c1cd81 100644 --- a/loopx/presentation/renderers/quota_markdown.py +++ b/loopx/presentation/renderers/quota_markdown.py @@ -485,8 +485,13 @@ def render_quota_should_run_markdown(payload: dict[str, Any]) -> str: f"resume_when={vision_wait_state.get('resume_when')} " f"automatic_resume={vision_wait_state.get('automatic_resume')}" ) - task_orchestration = as_dict(payload.get("task_orchestration_contract")) - if task_orchestration: + orchestration = as_dict(payload.get("task_orchestration_contract")) + for label, task_orchestration in ( + ("task_orchestration", orchestration), + ("peer_activation_diagnostic", as_dict(orchestration.get("peer_activation_diagnostic"))), + ): + if not task_orchestration: + continue lanes = as_list(task_orchestration.get("eligible_child_lanes")) if not lanes: lanes = as_list(task_orchestration.get("eligible_peer_lanes")) @@ -494,13 +499,26 @@ def render_quota_should_run_markdown(payload: dict[str, Any]) -> str: if not blocked_lanes: blocked_lanes = as_list(task_orchestration.get("blocked_peer_lanes")) lines.append( - "- task_orchestration: " + f"- {label}: " f"mode={task_orchestration.get('mode')} " f"activation_required={task_orchestration.get('activation_required')} " - f"lanes={len(lanes)} " - f"blocked_lanes={len(blocked_lanes)} " + f"lanes={task_orchestration.get('eligible_peer_count', len(lanes))} " + f"blocked_lanes={task_orchestration.get('blocked_peer_count', len(blocked_lanes))} " f"writeback_owner={task_orchestration.get('writeback_owner')}" ) + fields = " ".join( + f"{key}={markdown_scalar(task_orchestration[key])}" + for key in ("execution_scope", "execution_state", "task_selection", "activation_allowed") + if key in task_orchestration + ) + if fields: + lines.append(f"- {label}_admission: {fields}") + reasons = sorted({str(reason) for row in blocked_lanes if isinstance(row, dict) + for reason in as_list(row.get("reason_codes"))}) + if reasons: + lines.append(f"- {label}_blockers: {', '.join(reasons)}") + if task_orchestration.get("read_required") is True: + lines.append(f"- {label}_required_detail: {task_orchestration.get('detail_ref')}") replan_decision = as_dict(payload.get("autonomous_replan_decision")) if replan_decision: lines.append( diff --git a/tests/control_plane/test_delegation_context.py b/tests/control_plane/test_delegation_context.py index 00215e13d2..f91dbc7fc5 100644 --- a/tests/control_plane/test_delegation_context.py +++ b/tests/control_plane/test_delegation_context.py @@ -40,7 +40,7 @@ def _fixture(tmp_path: Path) -> tuple[Path, Path, Path]: { "id": "independent-review", "agent_id": "worker", - "todo_id": "todo-review", + "todo_id": "todo_review0002", "requesters": ["coordinator"], "workspace": str(tmp_path / "worker"), "host_args": ["--host", "generic-cli"], @@ -54,7 +54,7 @@ def _fixture(tmp_path: Path) -> tuple[Path, Path, Path]: return project, registry, tmp_path / "runtime" -def test_projects_authorized_route_without_private_binding_material( +def test_runtime_availability_does_not_admit_uninspected_delegation( tmp_path: Path, monkeypatch ) -> None: project, registry, runtime = _fixture(tmp_path) @@ -87,16 +87,19 @@ def test_projects_authorized_route_without_private_binding_material( execution_config=".loopx/config/delegations.json", ) assert packet["configuration_state"] == "ready" + assert packet["preflight"] == "required" + assert packet["execution_scope"] == "bound_delegation" assert packet["authorized_count"] == packet["projected_count"] == 1 assert "operation_receipts" not in packet assert packet["routes"] == [ { "binding_id": "independent-review", "agent_id": "worker", - "todo_id": "todo-review", + "todo_id": "todo_review0002", "runtime_id": "managed-runtime", "executor_kind": "managed", - "readiness": "ready", + "runtime_readiness": "ready", + "readiness": "unknown", "execution_profile": "model-a@high", } ] @@ -157,6 +160,9 @@ def test_repeated_live_quota_planning_never_reads_operation_inventory( "execution_config": ".loopx/config/delegations.json", } payload["goals"][0]["spawn_policy"] = policy + coordination = {"registered_agents": ["coordinator", "worker"], + "peer_task_coordination": {"coordinator_agent_id": "coordinator"}} + payload["goals"][0]["coordination"] = coordination registry.write_text(json.dumps(payload)) from loopx.collaboration_mcp import Delegations @@ -171,7 +177,7 @@ def test_repeated_live_quota_planning_never_reads_operation_inventory( goal_id="goal-a", status="active", recommended_action="Inspect delegated evidence", - coordination={"registered_agents": ["coordinator", "worker"]}, + coordination=coordination, agent_todo_items=[ { "todo_id": "todo_review0001", @@ -181,7 +187,11 @@ def test_repeated_live_quota_planning_never_reads_operation_inventory( "priority": "P1", "role": "agent", "task_class": "advancement_task", - } + }, + {"todo_id": "todo_old_peer", "status": "open", "priority": "P0", + "role": "agent", "task_class": "advancement_task", "claimed_by": "worker", "text": "Older still-open task"}, + {"todo_id": "todo_review0002", "status": "open", "priority": "P0", + "role": "agent", "task_class": "advancement_task", "claimed_by": "worker", "text": "Bound review task"}, ], goal_extra={"repo": str(project), "spawn_policy": policy}, ) @@ -207,6 +217,23 @@ def test_repeated_live_quota_planning_never_reads_operation_inventory( ]["facts"] assert facts["delegation_context"]["projected_count"] == 1 assert "operation_receipts" not in facts["delegation_context"] + route = facts["delegation_context"]["routes"][0] + assert route["todo_id"] == "todo_review0002" + assert route["readiness"] == "unknown" + assert facts["delegation_context"]["preflight"] == "required" + peer = packet["task_orchestration_contract"] + assert peer["execution_scope"] == "peer_agent_activation" + assert peer["execution_state"] == "blocked" + assert {row["todo_id"] for row in peer["blocked_peer_lanes"]} == {"todo_old_peer", "todo_review0002"} + assert packet["should_run"] is True + from loopx.control_plane.quota.turn_envelope import build_turn_envelope + from loopx.control_plane.turn_driver.driver import build_loopx_turn_plan + + envelope = build_turn_envelope(packet) + planned = build_loopx_turn_plan(envelope, host="codex-cli", execution_mode="interactive-visible") + # The parent Turn owns its selected Todo; peer candidates and the bound + # delegation target must not silently replace the coordinator's identity. + assert planned["route"]["selected_todo"]["todo_id"] == "todo_review0001" def test_missing_config_is_blocked_observation_not_empty_success(tmp_path: Path) -> None: diff --git a/tests/control_plane/test_quota_cli_projection.py b/tests/control_plane/test_quota_cli_projection.py index 6fd144375d..35101edf42 100644 --- a/tests/control_plane/test_quota_cli_projection.py +++ b/tests/control_plane/test_quota_cli_projection.py @@ -32,6 +32,27 @@ def _items(count: int, *, prefix: str) -> list[dict[str, object]]: ] +def test_markdown_distinguishes_peer_blockers_from_native_admission() -> None: + peer = {"mode": "task_scoped_peer", "execution_scope": "peer_agent_activation", + "execution_state": "blocked", "task_selection": "canonical_claimed_candidates", + "activation_allowed": False, "eligible_peer_lanes": [], + "blocked_peer_lanes": [{"todo_id": "todo-peer", "reason_codes": ["peer_agent_activation_unavailable"]}]} + for contract in [peer, {"mode": "adaptive", "execution_state": "ready", "peer_activation_diagnostic": peer}]: + rendered = render_quota_should_run_markdown({"task_orchestration_contract": contract}) + assert "execution_scope=peer_agent_activation" in rendered + assert "execution_state=blocked" in rendered + assert "task_selection=canonical_claimed_candidates" in rendered + assert "peer_agent_activation_unavailable" in rendered + if contract["mode"] == "adaptive": + assert "peer_activation_diagnostic_admission:" in rendered + assert "execution_state=ready" in rendered + projected = {**peer, "blocked_peer_lanes": [], "blocked_peer_count": 50, + "read_required": True, "detail_ref": "full_decision.task_orchestration_contract"} + rendered = render_quota_should_run_markdown({"task_orchestration_contract": projected}) + assert "blocked_lanes=50" in rendered + assert "task_orchestration_required_detail: full_decision.task_orchestration_contract" in rendered + + def test_compact_quota_should_run_cli_payload_keeps_decision_lanes_and_counts() -> None: backlog = _items(40, prefix="backlog") first_open = _items(5, prefix="open") diff --git a/tests/control_plane/test_task_orchestration_admission.py b/tests/control_plane/test_task_orchestration_admission.py index 394ae9d014..4ed92c4b44 100644 --- a/tests/control_plane/test_task_orchestration_admission.py +++ b/tests/control_plane/test_task_orchestration_admission.py @@ -824,3 +824,54 @@ def test_admission_requires_scope_for_unknown_child_action() -> None: "reason_codes": ["write_scope_missing"], } ] + + +def test_peer_inventory_does_not_confuse_first_display_row_with_binding() -> None: + old = {**_todo("todo_old"), "claimed_by": "worker"} + current = {**_todo("todo_current"), "claimed_by": "worker"} + kwargs = dict( + fallback_work_lane_contract={"lane": "advancement_task"}, + goal_boundary={"peer_task_coordination": {"enabled": True, "coordinator_agent_id": AGENT_ID}}, + agent_identity={"agent_id": AGENT_ID, "registered_agents": [AGENT_ID, "worker"]}, + agent_todo_summary={"items": [old]}, + raw_agent_todo_summary={"items": [old]}, + available_capabilities=[], + agent_management_projection=_peer_management("worker"), + ) + # The canonical source replaces the stale display, not merely its ordering. + contract, lane = apply_task_orchestration_contract(**kwargs, agent_todo_source_items=[current]) + assert [r["todo_id"] for r in contract["blocked_peer_lanes"]] == ["todo_current"] + assert contract["execution_scope"] == "peer_agent_activation" + assert contract["task_selection"] == "canonical_claimed_candidates" + assert lane == {"lane": "advancement_task"} + assert apply_task_orchestration_contract(**kwargs, agent_todo_source_items=[])[0] is None + # Multiple genuinely open tasks are retained; neither becomes an implicit binding. + a, _ = apply_task_orchestration_contract(**kwargs, agent_todo_source_items=[old, current]) + b, _ = apply_task_orchestration_contract(**kwargs, agent_todo_source_items=[current, old]) + assert a == b + assert len(a["blocked_peer_lanes"]) == 2 + + +def test_blocked_peer_activation_does_not_hide_admitted_native_children() -> None: + items = [{**_todo("todo_one"), "claimed_by": AGENT_ID}, + {**_todo("todo_two"), "claimed_by": AGENT_ID}, + {**_todo("todo_peer", resume_ready=False), "claimed_by": "worker"}] + summary = {"items": items} + contract, lane = apply_task_orchestration_contract( + fallback_work_lane_contract={"lane": "advancement_task"}, + goal_boundary={ + "write_scope": ["loopx/**"], + "peer_task_coordination": {"enabled": True, "coordinator_agent_id": AGENT_ID}, + "orchestration": {"mode": "multi_subagent", "spawn_allowed": True, "max_children": 2}, + }, + agent_identity={"agent_id": AGENT_ID, "registered_agents": [AGENT_ID, "worker"]}, + agent_todo_summary=summary, raw_agent_todo_summary=summary, + agent_todo_source_items=items, available_capabilities=["subagent_spawn"], + agent_management_projection=_peer_management("worker"), + ) + assert contract["mode"] == "adaptive" + assert lane["must_attempt_work"] is True + peer = contract["peer_activation_diagnostic"] + assert peer["execution_state"] == "blocked" + assert peer["activation_allowed"] is False + assert peer["blocked_peer_lanes"][0]["reason_codes"] == ["peer_agent_activation_unavailable", "peer_lane_not_resume_ready"] diff --git a/tests/control_plane_ts/agent_context.test.ts b/tests/control_plane_ts/agent_context.test.ts index db8ac54372..cf7147244d 100644 --- a/tests/control_plane_ts/agent_context.test.ts +++ b/tests/control_plane_ts/agent_context.test.ts @@ -129,6 +129,10 @@ test("delegation routes and explicit result receipts are bounded public-safe fac const beforeFacts = (before.contributions as any[])[0].facts; assert.equal(beforeFacts.delegation_context.projected_count, 1); assert.equal(beforeFacts.delegation_context.entrypoint, "loopx delegation"); + assert.equal(beforeFacts.delegation_context.execution_scope, "bound_delegation"); + assert.equal(beforeFacts.delegation_context.preflight, "required"); + assert.equal(beforeFacts.delegation_context.routes[0].readiness, "unknown"); + assert.equal(beforeFacts.delegation_context.routes[0].runtime_readiness, "ready"); assert.equal(beforeFacts.delegation_context.routes[0].execution_profile, "model-a@high"); assert.equal(beforeFacts.delegation_context.operation_receipts, undefined); assert.ok(!JSON.stringify(before).includes("credential")); @@ -188,7 +192,7 @@ test("coordinator participation guidance survives all bounded lifecycle projecti } })!; assert.deepEqual(packet.failures, []); const [contribution] = packet.contributions as Record[]; - assert.equal(contribution.revision, "v4"); + assert.equal(contribution.revision, "v5"); assert.equal(packet.authority, "guidance_only"); assert.ok(Buffer.byteLength(JSON.stringify(contribution)) <= 2048); assert.ok(Buffer.byteLength(JSON.stringify(packet)) <= 3072); diff --git a/tests/control_plane_ts/peer_orchestration.test.ts b/tests/control_plane_ts/peer_orchestration.test.ts new file mode 100644 index 0000000000..2d8ac84377 --- /dev/null +++ b/tests/control_plane_ts/peer_orchestration.test.ts @@ -0,0 +1,43 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { projectPeerOrchestration } from "../../loopx/control_plane/quota/peer_orchestration.ts"; + +const task = (todo_id: string, extra = {}) => ({ todo_id, claimed_by: "worker", + status: "open", task_class: "advancement_task", ...extra }); +const input = { agent_id: "parent", registered_agents: ["parent", "worker"], + available_capabilities: ["peer_agent_activation"], agents: [{ agent_id: "worker", state: "running" }] }; + +test("activation is scoped and preserves distinct tasks without inventing session binding", () => { + const a = task("old"), b = task("new"); + const result = projectPeerOrchestration({ ...input, items: [a, b, a] })!; + assert.equal(result.execution_scope, "peer_agent_activation"); + assert.equal(result.task_selection, "canonical_claimed_candidates"); + assert.equal(result.execution_state, "ready"); + assert.deepEqual((result.eligible_peer_lanes as any[]).map(row => row.todo_id), ["new", "old"]); + assert.deepEqual(result, projectPeerOrchestration({ ...input, items: [b, a] })); +}); + +test("runtime, activation and dependency gates remain independently enforced", () => { + const cases = [ + { available_capabilities: [], agents: input.agents, extra: {}, reasons: ["peer_agent_activation_unavailable"] }, + { available_capabilities: input.available_capabilities, agents: [], extra: {}, reasons: ["peer_liveness_unavailable"] }, + { available_capabilities: input.available_capabilities, agents: [{ agent_id: "worker", state: "running", stale_claim_hint: true }], extra: {}, reasons: ["peer_runtime_stale"] }, + { available_capabilities: input.available_capabilities, agents: [{ agent_id: "worker", state: "dormant" }], extra: {}, reasons: ["peer_runtime_not_active"] }, + { available_capabilities: input.available_capabilities, agents: input.agents, extra: { resume_when: "todo_done:dependency", resume_ready: false }, reasons: ["peer_lane_not_resume_ready"] }, + ]; + for (const row of cases) { + const result = projectPeerOrchestration({ ...input, ...row, items: [task("task", row.extra)] })!; + assert.equal(result.execution_state, "blocked"); + assert.deepEqual(result.eligible_peer_lanes, []); + assert.deepEqual((result.blocked_peer_lanes as any[])[0].reason_codes, row.reasons); + } +}); + +test("closed, deferred, unregistered, self and non-advancement rows never grant activation", () => { + const items = [task("done", { done: true }), task("blocked", { status: "blocked" }), + task("deferred", { status: "deferred" }), task("foreign", { claimed_by: "foreign" }), + task("self", { claimed_by: "parent" }), task("monitor", { task_class: "continuous_monitor" }), + task("unclaimed", { claimed_by: null }), task("")]; + assert.equal(projectPeerOrchestration({ ...input, items }), null); + assert.equal(projectPeerOrchestration({ ...input, items: [] }), null); +}); diff --git a/tests/control_plane_ts/turn_envelope.test.ts b/tests/control_plane_ts/turn_envelope.test.ts index 0cdaafd091..99eefb2d58 100644 --- a/tests/control_plane_ts/turn_envelope.test.ts +++ b/tests/control_plane_ts/turn_envelope.test.ts @@ -12,6 +12,7 @@ import { import { EffectRuntimeRequestError } from "../../loopx/control_plane/effect_runtime_errors.ts"; import type { JsonObject } from "../../loopx/control_plane/effect_program.ts"; import { TURN_ENVELOPE_SECTION_TARGETS } from "../../loopx/control_plane/quota/turn_envelope_budget.ts"; +import { projectPeerOrchestration } from "../../loopx/control_plane/quota/peer_orchestration.ts"; function payload(): Record { return { @@ -72,6 +73,40 @@ const protocolActionFields = { agent_action: "advance one bounded segment", }; +test("large peer inventories retain scoped gates and signed detail without hiding tasks", () => { + const source = payload(); + const items = Array.from({ length: 50 }, (_, i) => ({ todo_id: `todo_${i}`, + claimed_by: "worker", status: "open", task_class: "advancement_task", + title: "Inspect independent source and return verified evidence" })); + for (const available of [[], ["peer_agent_activation"]]) { + const peer = projectPeerOrchestration({ agent_id: "parent", registered_agents: ["worker"], + items, available_capabilities: available, agents: [{ agent_id: "worker", state: "running" }] })!; + for (const nested of [false, true]) { + source.task_orchestration_contract = nested + ? { mode: "adaptive", execution_state: "ready", eligible_child_lanes: [{ todo_id: "local-child" }], peer_activation_diagnostic: peer } + : peer; + const render = () => buildTurnEnvelope({ payload: source, + protocol_action_fields: protocolActionFields, scheduler_execution_args: " --available-capability shell" }); + const envelope = render(); + const contract = envelope.task_orchestration_contract as JsonObject; + const compact = (nested ? contract.peer_activation_diagnostic : contract) as JsonObject; + assert.equal(compact.execution_scope, "peer_agent_activation"); + assert.equal(compact.execution_state, available.length ? "ready" : "blocked"); + assert.equal(compact.activation_allowed, available.length > 0); + assert.equal(Number(compact.eligible_peer_count) + Number(compact.blocked_peer_count), 50); + assert.equal(compact.read_required, true); + assert.equal(compact.eligible_peer_lanes, undefined); + assert.equal((peer.eligible_peer_lanes as unknown[]).length + (peer.blocked_peer_lanes as unknown[]).length, 50); + assert.equal((envelope.compaction as JsonObject).within_budget, true); + assert.deepEqual(quotaActionSignatureDocument(source, protocolActionFields), turnEnvelopeActionSignatureDocument(envelope)); + const rows = (available.length ? peer.eligible_peer_lanes : peer.blocked_peer_lanes) as JsonObject[]; + rows[0].todo_id = `changed-identity-${nested}`; + assert.notEqual((render().action_signature as JsonObject).source_hash, (envelope.action_signature as JsonObject).source_hash); + if (nested) assert.deepEqual(contract.eligible_child_lanes, [{ todo_id: "local-child" }]); + } + } +}); + test("Turn preserves checkpointed scope approval without lifting other gates", () => { const source = payload(); const scope = source.goal_boundary as Record; From 19c1ff6c43d7428ea1f1a2f1d895ebeb2abe82e2 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 22 Sep 2026 06:20:07 +0800 Subject: [PATCH 2/3] test(collaboration): migrate delegation readiness consumer Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- tests/test_turn_machine_credential.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/tests/test_turn_machine_credential.py b/tests/test_turn_machine_credential.py index effa214b10..36722d9f55 100644 --- a/tests/test_turn_machine_credential.py +++ b/tests/test_turn_machine_credential.py @@ -113,6 +113,9 @@ def test_delegation_uses_machine_readiness_without_expanding_requesters(tmp_path ) assert packet["authorized_count"] == count if count: - assert packet["routes"][0]["readiness"] == "ready" + route = packet["routes"][0] + assert route["runtime_readiness"] == "ready" + assert route["readiness"] == "unknown" + assert packet["preflight"] == "required" assert KEY not in json.dumps(packet) assert not provider.operator_provider_store_path(runtime).exists() From e44a44b764a9b0af0031850dbc1d3c64148cfc3c Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 22 Sep 2026 13:24:48 +0800 Subject: [PATCH 3/3] test(docs): keep the sub-agent orchestration contract smoke in sync The peer-contract documentation now scopes the contract to `execution_scope=peer_agent_activation`, states that only currently actionable candidates appear under `eligible_peer_lanes`, and places a dormant or dependency-blocked candidate under `blocked_peer_lanes` instead of dropping it. That intentional rewrite removed the sentence this smoke asserted, so the smoke failed while the contract itself stayed true. Assert the replacement sentences so the durable invariant (closed, blocked, and deferred Todos never become activation candidates, and a dormant peer is surfaced as a blocked lane) is still pinned to the documentation. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- examples/codex-subagent-orchestration-contract-smoke.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/examples/codex-subagent-orchestration-contract-smoke.py b/examples/codex-subagent-orchestration-contract-smoke.py index 52fe86deaf..8198fdcc23 100644 --- a/examples/codex-subagent-orchestration-contract-smoke.py +++ b/examples/codex-subagent-orchestration-contract-smoke.py @@ -30,7 +30,9 @@ '"agent_model": "peer_v1"', "independent worktrees", "Review remains `action_kind=review`", - "Dormant registered agents and closed, blocked, or deferred todos are not coordinator candidates.", + "Only currently actionable candidates appear under `eligible_peer_lanes`.", + "Closed, blocked, or deferred Todos are excluded.", + "appears under `blocked_peer_lanes`", ) FORBIDDEN_PHRASES = (