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/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 = ( 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; 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()