diff --git a/README.md b/README.md index 8c21799ea7..1e6bcf41b9 100644 --- a/README.md +++ b/README.md @@ -414,6 +414,7 @@ evidence → recovery; continuation → governance. | --- | --- | --- | | Goal state and status | Tracks active state, todos, claims, gates, evidence, run history, and first-screen attention. | `loopx status`, `loopx diagnose`, `loopx review-packet` | | Quota and interaction contract | Decides whether a turn should deliver, ask, wait, self-repair, or stay quiet. | `loopx quota should-run`, [quota allocation](docs/quota-allocation.md) | +| Company control loop | Routes company work to AI, human decisions, human execution, monitors, or blockers; reconciles evidence into the next planning cycle. | `loopx company-control-loop`, [company control loop](docs/reference/company-control-loop.md) | | Agent runtime bridges | Keeps Codex App, Codex CLI, Claude Code, and generic workers aligned with the same guard. | `loopx heartbeat-prompt`, `loopx codex-cli-bootstrap-message`, `loopx worker-bridge` | | Operator surfaces | Renders compact status without making the browser the state authority. | `loopx serve-status`, [dashboard](apps/presentation/dashboard/README.md) | | Session dash | Starts a live single-page panel that tracks fleet progress: sessions, their goals, and each goal's status/todo progress, with result statistics; auto-refreshes in place. | `loopx dash`, [session dash design](docs/product/surfaces/session-dash-panel-design.md) | diff --git a/README.zh-CN.md b/README.zh-CN.md index 254e5a0670..1445b3d076 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -370,6 +370,7 @@ Kernel 把控制面归结为五个用户可以直接行动的问题。每个问 | --- | --- | --- | | Goal state 与 status | 跟踪 active state、todo、claim、gate、evidence、run history 和首屏关注点。 | `loopx status`、`loopx diagnose`、`loopx review-packet` | | Quota 与 interaction contract | 决定一轮应该执行、提问、等待、自修复还是静默。 | `loopx quota should-run`、[Quota Allocation](docs/quota-allocation.md) | +| Company Control Loop | 将公司工作路由给 AI、人类决策、人类执行、监控或 blocker,并把证据汇入下一轮规划。 | `loopx company-control-loop`、[Company Control Loop](docs/reference/company-control-loop.md) | | Agent runtime bridge | 让 Codex App、Codex CLI、Claude Code 和 generic worker 服从同一 guard。 | `loopx heartbeat-prompt`、`loopx codex-cli-bootstrap-message`、`loopx worker-bridge` | | Operator surface | 呈现紧凑状态,但不让浏览器成为状态事实源。 | `loopx serve-status`、[Dashboard](apps/presentation/dashboard/README.md) | | External projection | 把 todo / gate 投影到协作表面,同时保持 LoopX 权威。 | `loopx lark-kanban`、[Lark Kanban adapter](docs/integrations/lark-kanban-control-plane-adapter.md) | diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md index 64d1086771..5461eead13 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -56,6 +56,7 @@ These directly determine whether a long-running team is usable. A directory or R | Area and existing entry | Streams | Next requirement in this roadmap | | --- | --- | --- | | [Goal Vision/replan](../../reference/protocols/goal-vision-replan-contract-v0.md), [work graph](../../reference/protocols/task-graph-projection-v0.md), [peer runtime](../../reference/protocols/peer-agent-runtime-v1.md), [supervisor](../../reference/protocols/peer-supervisor-v0.md) | S2/S3 | Exercise dependencies, replanning, acceptance and handoff in one real case; aggregate closeout consumes acceptance facts | +| [Company control-loop profile](../../reference/company-control-loop.md) | S1/S3 | Keep it a bounded caller of Goal/Todo/evidence owners; qualify the broader steward claim only through R2/R3 multi-Agent adoption, dependent artifacts, independent acceptance, recovery and result return | | [Quota](../../quota-allocation.md), [cadence](../../operations/long-task-cadence-policy.md), [attention](../../operations/attention-queue.md) | S5/S7 | Budget exhaustion, deferral and blocking expose next triggers/readback; scale without frequent full-state polling | | [Material lifecycle](../../reference/protocols/material-lifecycle-architecture-v0.md), [material frontier](../../reference/protocols/agent-material-frontier-v0.md), [authority registration](../../operations/authority-source-registration.md) | S6 | Agents discover roadmap/RFC revisions and record reads; reading grants neither agreement nor authority; archival preserves raw-source ownership | | [Decision Context](../../../loopx/capabilities/decision_context/README.md), [Reward Memory](../../../loopx/capabilities/reward_memory/README.md), [Semantic Preference](../../../loopx/capabilities/semantic_preference/README.md), [Turn Recall](../../../loopx/capabilities/agent_turn_recall/README.md) | S6/S11 | Distinguish facts/preferences/advice/attribution/authority; scoped recall, expiry and outcome-feedback counterexamples first | diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md index a9ea996afe..ab033f02b1 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md @@ -56,6 +56,7 @@ LoopX 的目标是让人用本地前端或 Lark 提出、修订和验收复杂 | 领域与已有入口 | 所属工作流 | 当前总纲要求的下一步 | | --- | --- | --- | | [Goal Vision/replan](../../reference/protocols/goal-vision-replan-contract-v0.md)、[work graph](../../reference/protocols/task-graph-projection-v0.md)、[peer runtime](../../reference/protocols/peer-agent-runtime-v1.md)、[监督](../../reference/protocols/peer-supervisor-v0.md) | S2/S3 | 将跨工作依赖、重规划、验收与 handoff 放进同一个真实案例;aggregate closeout 必须消费验收事实 | +| [公司控制闭环 profile](../../reference/company-control-loop.md) | S1/S3 | 保持为 Goal/Todo/证据 owner 的有界 caller;更广的 steward 主张必须通过 R2/R3 的多 Agent adoption、依赖产物、独立验收、恢复和结果返回来验收 | | [quota](../../quota-allocation.md)、[cadence](../../operations/long-task-cadence-policy.md)、[attention](../../operations/attention-queue.md) | S5/S7 | 预算耗尽/延期/被阻塞时有明确下一次触发及用户回读;百 Agent 不靠高频全文轮询 | | [材料生命周期](../../reference/protocols/material-lifecycle-architecture-v0.zh-CN.md)、[材料 frontier](../../reference/protocols/agent-material-frontier-v0.md)、[authority 注册](../../operations/authority-source-registration.md) | S6 | 路线/RFC 更新能由 Agent 按 revision 发现并登记阅读;read receipt 不表示同意或获得权限;归档不丢原始来源 | | [Decision Context](../../../loopx/capabilities/decision_context/README.md)、[Reward Memory](../../../loopx/capabilities/reward_memory/README.md)、[Semantic Preference](../../../loopx/capabilities/semantic_preference/README.md)、[Turn Recall](../../../loopx/capabilities/agent_turn_recall/README.md) | S6/S11 | 区分事实、偏好、建议、归因和权威;同一 scope 的回忆/失效/结果反馈负例先行 | diff --git a/docs/reference/company-control-loop.md b/docs/reference/company-control-loop.md new file mode 100644 index 0000000000..013467764a --- /dev/null +++ b/docs/reference/company-control-loop.md @@ -0,0 +1,121 @@ +# Company Control Loop + +The Company Control Loop is LoopX's provider-neutral planning layer for a +long-running company direction. It routes bounded work to AI or people, keeps +the result under one Goal, reconciles Todo evidence, and produces the next +planning cycle. + +## Roadmap placement + +This command is a bounded planning profile for the persistent steward path in +the [LoopX overall roadmap](../architecture/rfcs/loopx-overall-roadmap-v0.md), +primarily S1 and S3. It exercises existing Goal, Todo, evidence, quota, and +replan owners; it does not create a second steward, work ledger, scheduler, or +authority model. Its current acceptance boundary is the documented CLI and +packaged runtime lifecycle. The broader R2/R3 journey still requires real +multi-Agent adoption, dependent artifacts, independent acceptance, restart +recovery, and automatic result return through their existing owners. + +## Authority boundary + +The command does not grant execution authority. Existing LoopX Todo rules own +claims, user gates, leases, validation, and completion. The company layer owns +only these decisions: + +- `ai_execute` becomes an agent advancement Todo. +- `human_decide` becomes a blocking user gate. +- `human_execute` becomes a user action. +- `observe` becomes a bounded continuous monitor. +- prohibited work becomes a blocker. + +A Todo marked done is accepted only when reconciliation also receives an +evidence reference. Completion without evidence becomes `awaiting_evidence`; +a blocked Todo becomes `replanning`. + +## State lifecycle + +Start with a `outcome_routing_plan_request_v0` JSON object. It contains one +direction, a cycle number, outcomes, work items, and feedback. + +`company-control-loop` is the product profile. Its control-plane contracts, +effect IDs, state directory, and schemas use the domain-neutral +`outcome_routing_*` family so the shared work-item kernel does not acquire +company-specific vocabulary. + +```sh +loopx company-control-loop project --state-json company.json +loopx company-control-loop save \ + --goal-id company-goal \ + --state-json company.json +``` + +Both commands are read-only at this point. Add `--execute` to `save` after +review. Replacing existing state also requires the exact revision returned by +`show` or the prior write: + +```sh +loopx company-control-loop save \ + --goal-id company-goal \ + --state-json company.json \ + --expected-revision REVISION \ + --execute + +loopx company-control-loop show --goal-id company-goal +``` + +## Materialize work as Todos + +Preview the idempotent plan first, then execute it: + +```sh +loopx company-control-loop sync-todos \ + --goal-id company-goal \ + --agent-id company-ceo \ + --project /path/to/project + +loopx company-control-loop sync-todos \ + --goal-id company-goal \ + --agent-id company-ceo \ + --project /path/to/project \ + --execute +``` + +The profile persists a revisioned `work_item_id` to `todo_id` binding after +Todo readback. Agent Todos may still use the kernel's existing `target_key` +identity. Human Todos remain ordinary `user_gate` and `user_action` records; +their correlation identity stays inside profile state. Failed readback or a +stale state revision stops the command. + +## Reconcile and plan the next cycle + +Reconciliation is also dry-run by default: + +```sh +loopx company-control-loop reconcile-todos \ + --goal-id company-goal \ + --agent-id company-ceo \ + --project /path/to/project + +loopx company-control-loop reconcile-todos \ + --goal-id company-goal \ + --agent-id company-ceo \ + --project /path/to/project \ + --execute + +loopx company-control-loop next-cycle --goal-id company-goal +``` + +`next-cycle` returns a new request object. Evidence-backed completed work leaves +the active frontier. Failures and missing evidence become typed feedback and +set `replan_required`. When no work remains, `goal_converged` is true. + +Review the returned state before saving it as the next cycle. Revision checks +prevent an older planner or restarted worker from overwriting newer state. + +## Always-on operation + +An always-on host should run the ordinary LoopX heartbeat contract. Each wake +must enter through `quota should-run`, advance only the selected Todo, validate +the result, write state, and spend the matching slot. Scheduler cadence and +human notification remain host responsibilities; this command does not create +an independent hidden scheduler. diff --git a/loopx/cli.py b/loopx/cli.py index e50388ab13..0329476701 100644 --- a/loopx/cli.py +++ b/loopx/cli.py @@ -2,10 +2,15 @@ import argparse import sys +from pathlib import Path from .cli_commands.agent_capabilities import register_agent_capabilities, handle_agent_capabilities from .cli_commands.agent_directory import register_agent_directory, handle_agent_directory from .cli_commands.agent_context import register_agent_context, handle_agent_context +from .cli_commands.company_control_loop import ( + handle_company_control_loop_command, + register_company_control_loop_command, +) from .cli_commands.todo_continuation import register_todo_continuation, handle_todo_continuation from .cli_commands.manager_inbox import register_manager_inbox, handle_manager_inbox from .capabilities.content_ops.cli import ( @@ -330,6 +335,7 @@ def build_parser() -> LoopXArgumentParser: register_manager_inbox(sub, add_subcommand_format) register_agent_capabilities(sub, add_subcommand_format) register_agent_context(sub, add_subcommand_format) + register_company_control_loop_command(sub, add_subcommand_format) register_agent_directory(sub, add_subcommand_format) register_lark_inbox_commands(sub, add_subcommand_format) register_lark_kanban_commands(sub, add_subcommand_format) @@ -770,6 +776,26 @@ def main(argv: list[str] | None = None) -> int: if args.command == "agent-context": return handle_agent_context(args, registry_path, print_payload, output_format) + company_control_loop_result = handle_company_control_loop_command( + args, + output_format=output_format, + print_payload=print_payload, + runtime_root=( + ( + Path(args.runtime_root).expanduser().resolve() + if args.runtime_root + else effective_runtime_root(registry_path, None) + ) + if args.command == "company-control-loop" + and args.company_control_loop_command + in {"save", "show", "sync-todos", "reconcile-todos", "next-cycle"} + else None + ), + registry_path=registry_path, + ) + if company_control_loop_result is not None: + return company_control_loop_result + if args.command == "agent-directory": return handle_agent_directory( args, registry_path, effective_runtime_root(registry_path, args.runtime_root), diff --git a/loopx/cli_commands/company_control_loop.py b/loopx/cli_commands/company_control_loop.py new file mode 100644 index 0000000000..07982836a9 --- /dev/null +++ b/loopx/cli_commands/company_control_loop.py @@ -0,0 +1,492 @@ +from __future__ import annotations + +import argparse +import json +from collections.abc import Callable +from datetime import UTC, datetime +from pathlib import Path +from typing import Any + +from ..control_plane.effect_runtime import effect_runtime_result +from ..todos import add_goal_todo, list_goal_todos + + +def register_company_control_loop_command( + subparsers: argparse._SubParsersAction, + add_subcommand_format: Callable[[argparse.ArgumentParser], None], +) -> None: + parser = subparsers.add_parser( + "company-control-loop", + help="Validate and project a company planning snapshot into LoopX Todo lanes.", + ) + add_subcommand_format(parser) + actions = parser.add_subparsers( + dest="company_control_loop_command", + required=True, + ) + project = actions.add_parser( + "project", + help="Project a outcome_routing_plan_request_v0 JSON object without writing state.", + ) + add_subcommand_format(project) + project.add_argument( + "--state-json", + required=True, + help="Path to a outcome_routing_plan_request_v0 JSON object.", + ) + save = actions.add_parser( + "save", + help="Validate and persist company control state under one Goal runtime.", + ) + add_subcommand_format(save) + save.add_argument("--goal-id", required=True, help="Goal that owns the company state.") + save.add_argument( + "--state-json", + required=True, + help="Path to a outcome_routing_plan_request_v0 JSON object.", + ) + save.add_argument( + "--expected-revision", + help="Exact revision returned by show/save. Required to replace existing state.", + ) + save.add_argument( + "--execute", + action="store_true", + help="Persist the validated projection. Without this flag, return a preview.", + ) + show = actions.add_parser( + "show", + help="Read the persisted company control state for one Goal.", + ) + add_subcommand_format(show) + show.add_argument("--goal-id", required=True, help="Goal that owns the company state.") + sync = actions.add_parser( + "sync-todos", + help="Create missing LoopX Todos from persisted company work and verify readback.", + ) + add_subcommand_format(sync) + sync.add_argument( + "--goal-id", required=True, help="Goal that owns the company state and Todos." + ) + sync.add_argument( + "--agent-id", + required=True, + help="Registered agent that owns routed agent work.", + ) + sync.add_argument("--project", help="Project containing the Goal active state.") + sync.add_argument( + "--execute", + action="store_true", + help="Create missing Todos. Without this flag, return the idempotent plan.", + ) + reconcile = actions.add_parser( + "reconcile-todos", + help="Reconcile Todo status and evidence into persisted company state.", + ) + add_subcommand_format(reconcile) + reconcile.add_argument( + "--goal-id", required=True, help="Goal that owns the company state and Todos." + ) + reconcile.add_argument( + "--agent-id", + required=True, + help="Registered agent whose routed Todo lane is reconciled.", + ) + reconcile.add_argument("--project", help="Project containing the Goal active state.") + reconcile.add_argument( + "--execute", + action="store_true", + help="Persist reconciliation. Without this flag, return a preview.", + ) + next_cycle = actions.add_parser( + "next-cycle", + help="Plan the next company cycle from reconciled Todo outcomes.", + ) + add_subcommand_format(next_cycle) + next_cycle.add_argument( + "--goal-id", required=True, help="Goal that owns the reconciled company state." + ) + + +def render_company_control_loop_markdown(payload: dict[str, Any]) -> str: + lines = [ + "# LoopX Company Control Loop", + "", + f"- ok: `{payload.get('ok')}`", + ] + if payload.get("error"): + lines.append(f"- error: {payload['error']}") + return "\n".join(lines) + lines.extend([ + f"- schema_version: `{payload.get('schema_version')}`", + f"- direction: {payload.get('direction')}", + f"- cycle: {payload.get('cycle')}", + f"- replan_required: `{payload.get('replan_required')}`", + "", + "## Work routing", + "", + ]) + work_items = payload.get("work_items") + if not isinstance(work_items, list) or not work_items: + lines.append("- No work items.") + else: + for item in work_items: + if isinstance(item, dict): + lines.append( + f"- `{item.get('work_item_id')}` -> `{item.get('route')}` " + f"({item.get('status')}): {item.get('title')}" + ) + return "\n".join(lines) + + +def _read_json_object(path_text: str) -> dict[str, Any]: + payload = json.loads(Path(path_text).expanduser().read_text(encoding="utf-8")) + if not isinstance(payload, dict): + raise TypeError("company control state JSON must contain an object") + return payload + + +def handle_company_control_loop_command( + args: argparse.Namespace, + *, + output_format: Callable[..., str], + print_payload: Callable[[dict[str, Any], str, Callable[[dict[str, Any]], str]], None], + runtime_root: Path | None = None, + registry_path: Path | None = None, +) -> int | None: + if args.command != "company-control-loop": + return None + try: + command = args.company_control_loop_command + if command in {"show", "sync-todos", "reconcile-todos", "next-cycle"}: + if runtime_root is None: + raise ValueError("company control state requires a runtime root") + projection = effect_runtime_result( + "work_item.outcome_routing_state.load", + { + "schema_version": "outcome_routing_state_store_request_v0", + "runtime_root": str(runtime_root), + "goal_id": args.goal_id, + }, + ) + if command == "show": + payload = {"ok": True, **projection} + elif command == "next-cycle": + payload = { + "ok": True, + **effect_runtime_result( + "work_item.outcome_routing_state.next_cycle", + { + "schema_version": "outcome_routing_next_cycle_request_v0", + "goal_id": args.goal_id, + "state": projection.get("state"), + }, + ), + } + elif command == "sync-todos": + if registry_path is None: + raise ValueError("company Todo sync requires a registry") + payload = _sync_todos( + stored=projection, + goal_id=args.goal_id, + agent_id=args.agent_id, + project=Path(args.project).expanduser() if args.project else None, + registry_path=registry_path, + runtime_root=runtime_root, + execute=bool(args.execute), + ) + else: + if registry_path is None: + raise ValueError("company Todo reconciliation requires a registry") + payload = _reconcile_todos( + stored=projection, + goal_id=args.goal_id, + agent_id=args.agent_id, + project=Path(args.project).expanduser() if args.project else None, + registry_path=registry_path, + runtime_root=runtime_root, + execute=bool(args.execute), + ) + else: + request = _read_json_object(args.state_json) + projection = effect_runtime_result( + "work_item.outcome_routing_plan.project", + request, + ) + if command != "save": + payload = {"ok": True, **projection} + else: + if not args.execute: + payload = { + "ok": True, + "dry_run": True, + "goal_id": args.goal_id, + "projection": projection, + } + else: + if runtime_root is None: + raise ValueError("company control state requires a runtime root") + write_request: dict[str, Any] = { + "schema_version": "outcome_routing_state_store_request_v0", + "runtime_root": str(runtime_root), + "goal_id": args.goal_id, + "state": request, + "updated_at": datetime.now(UTC).isoformat(), + } + if args.expected_revision: + write_request["expected_revision"] = args.expected_revision + saved = effect_runtime_result( + "work_item.outcome_routing_state.write", + write_request, + ) + payload = {"ok": True, "dry_run": False, **saved} + exit_code = 0 + except Exception as exc: + payload = {"ok": False, "error": str(exc)} + exit_code = 1 + print_payload( + payload, + output_format(args), + render_company_control_loop_markdown, + ) + return exit_code + + +def _sync_todos( + *, + stored: dict[str, Any], + goal_id: str, + agent_id: str, + project: Path | None, + registry_path: Path, + runtime_root: Path, + execute: bool, +) -> dict[str, Any]: + state = stored.get("state") + if not isinstance(state, dict): + raise ValueError("persisted company control state does not exist") + company = state.get("projection") + if not isinstance(company, dict): + raise TypeError("persisted company control projection is invalid") + work_items = company.get("work_items") + if not isinstance(work_items, list): + raise TypeError("persisted company work_items must be an array") + listing = list_goal_todos( + registry_path=registry_path, + runtime_root_arg=str(runtime_root), + goal_id=goal_id, + agent_id=agent_id, + project=project, + limit=500, + ) + todos = [item for item in listing.get("todos", []) if isinstance(item, dict)] + by_id = { + str(todo.get("todo_id")): todo + for todo in todos + if todo.get("todo_id") + } + by_target: dict[str, dict[str, Any]] = {} + for todo in todos: + target = str(todo.get("target_key") or "").strip() + if not target: + continue + if target in by_target: + raise ValueError(f"multiple LoopX Todos use target_key {target!r}") + by_target[target] = todo + stored_bindings = { + str(binding.get("work_item_id")): binding + for binding in state.get("todo_bindings", []) + if isinstance(binding, dict) and binding.get("work_item_id") + } + actions: list[dict[str, Any]] = [] + seen_targets: set[str] = set() + for raw in work_items: + if not isinstance(raw, dict): + raise TypeError("persisted company work item is invalid") + todo_projection = raw.get("todo_projection") + if not isinstance(todo_projection, dict): + raise TypeError("persisted company Todo projection is invalid") + target = str(todo_projection.get("target_key") or "").strip() + if not target or target in seen_targets: + raise ValueError("company work target_key must be present and unique") + seen_targets.add(target) + binding = stored_bindings.get(str(raw.get("work_item_id"))) + matched = by_id.get(str(binding.get("todo_id"))) if binding else None + if matched is None and todo_projection.get("role") == "agent": + matched = by_target.get(target) + if matched: + actions.append({ + "work_item_id": raw.get("work_item_id"), + "target_key": target, + "action": "linked_existing", + "todo_id": matched.get("todo_id"), + }) + continue + action: dict[str, Any] = { + "work_item_id": raw.get("work_item_id"), + "target_key": target, + "action": "would_create", + "role": todo_projection.get("role"), + "task_class": todo_projection.get("task_class"), + } + if execute: + role = str(todo_projection.get("role") or "") + task_class = str(todo_projection.get("task_class") or "") + monitor_metadata: dict[str, Any] = {} + if role == "agent": + monitor_metadata["target_key"] = target + if task_class == "continuous_monitor": + monitor_metadata["watch_only"] = "true" + created = add_goal_todo( + registry_path=registry_path, + runtime_root_arg=str(runtime_root), + goal_id=goal_id, + project=project, + role=role, + text=f"[P1] {todo_projection.get('text')}", + status="open", + note=f"Acceptance: {todo_projection.get('acceptance')}", + task_class=task_class, + action_kind=str(todo_projection.get("action_kind") or ""), + claimed_by=agent_id if role == "agent" else None, + agent_id=agent_id, + blocks_agent=agent_id if task_class == "user_gate" else None, + bound_agent=agent_id if task_class == "user_action" else None, + decision_scope=( + f"direction:action:{target}" + if task_class == "user_gate" + else None + ), + monitor_metadata=monitor_metadata, + ) + action["action"] = "created" + action["todo_id"] = created.get("todo_id") + actions.append(action) + if execute: + readback = list_goal_todos( + registry_path=registry_path, + runtime_root_arg=str(runtime_root), + goal_id=goal_id, + agent_id=agent_id, + project=project, + limit=500, + ) + readback_by_id = { + str(item.get("todo_id")): item + for item in readback.get("todos", []) + if isinstance(item, dict) and item.get("todo_id") + } + missing = sorted( + str(action.get("todo_id")) + for action in actions + if str(action.get("todo_id")) not in readback_by_id + ) + if missing: + raise RuntimeError(f"LoopX Todo readback missing ids: {missing}") + binding_result = effect_runtime_result( + "work_item.outcome_routing_state.bind", + { + "schema_version": "outcome_routing_state_bind_request_v0", + "runtime_root": str(runtime_root), + "goal_id": goal_id, + "expected_revision": state.get("revision"), + "updated_at": datetime.now(UTC).isoformat(), + "todo_bindings": [ + { + "work_item_id": action["work_item_id"], + "target_key": action["target_key"], + "todo_id": action["todo_id"], + "role": action.get("role") or next( + str(item["todo_projection"].get("role")) + for item in work_items + if isinstance(item, dict) + and item.get("work_item_id") == action["work_item_id"] + and isinstance(item.get("todo_projection"), dict) + ), + } + for action in actions + ], + }, + ) + state = binding_result.get("state", state) + return { + "ok": True, + "dry_run": not execute, + "goal_id": goal_id, + "state_revision": state.get("revision"), + "actions": actions, + "readback_verified": execute, + } + + +def _reconcile_todos( + *, + stored: dict[str, Any], + goal_id: str, + agent_id: str, + project: Path | None, + registry_path: Path, + runtime_root: Path, + execute: bool, +) -> dict[str, Any]: + state = stored.get("state") + if not isinstance(state, dict): + raise ValueError("persisted company control state does not exist") + company = state.get("projection") + if not isinstance(company, dict): + raise TypeError("persisted company control projection is invalid") + work_items = company.get("work_items") + if not isinstance(work_items, list): + raise TypeError("persisted company work_items must be an array") + bindings = [item for item in state.get("todo_bindings", []) if isinstance(item, dict)] + if not bindings: + raise ValueError("persisted outcome routing state has no Todo bindings; run sync-todos first") + bindings_by_todo = { + str(item.get("todo_id")): item + for item in bindings + if item.get("todo_id") + } + listing = list_goal_todos( + registry_path=registry_path, + runtime_root_arg=str(runtime_root), + goal_id=goal_id, + agent_id=agent_id, + project=project, + limit=500, + ) + observations: list[dict[str, Any]] = [] + seen_targets: set[str] = set() + for raw in listing.get("todos", []): + if not isinstance(raw, dict): + continue + binding = bindings_by_todo.get(str(raw.get("todo_id") or "")) + if binding is None: + continue + target = str(binding.get("target_key") or "").strip() + if target in seen_targets: + raise ValueError(f"multiple LoopX Todos use target_key {target!r}") + seen_targets.add(target) + observation: dict[str, Any] = { + "target_key": target, + "todo_id": raw.get("todo_id"), + "status": raw.get("status"), + } + evidence = raw.get("evidence") + if isinstance(evidence, str) and evidence.strip(): + observation["evidence_ref"] = evidence.strip() + observations.append(observation) + if not observations: + raise ValueError("no LoopX Todos match persisted outcome routing bindings") + result = effect_runtime_result( + "work_item.outcome_routing_state.reconcile", + { + "schema_version": "outcome_routing_state_reconcile_request_v0", + "runtime_root": str(runtime_root), + "goal_id": goal_id, + "expected_revision": state.get("revision"), + "updated_at": datetime.now(UTC).isoformat(), + "execute": execute, + "observations": observations, + }, + ) + return {"ok": True, **result} diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index e926acbc42..373f1ed5a6 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -172,6 +172,16 @@ import { projectTodoPlanningInventoryDetail, } from "./work_items/planning_inventory.ts"; import { resolveRefreshRecommendation } from "./work_items/refresh_recommendation.ts"; +import { + projectOutcomeRoutingPlan, +} from "./work_items/outcome_routing_plan.ts"; +import { + bindOutcomeRoutingTodos, + loadOutcomeRoutingState, + planOutcomeRoutingNextCycle, + reconcileOutcomeRoutingState, + writeOutcomeRoutingState, +} from "./work_items/outcome_routing_state.ts"; import { validateInteractionProjectionHookInvocation, validateInteractionProjectionHookRegistration, @@ -451,6 +461,12 @@ export function createEffectRuntimeHandlers( ["work_item.planning_inventory.project", projectTodoPlanningInventory], ["work_item.planning_inventory.detail", projectTodoPlanningInventoryDetail], ["work_item.refresh_recommendation.resolve", resolveRefreshRecommendation], + ["work_item.outcome_routing_plan.project", projectOutcomeRoutingPlan], + ["work_item.outcome_routing_state.bind", bindOutcomeRoutingTodos], + ["work_item.outcome_routing_state.load", loadOutcomeRoutingState], + ["work_item.outcome_routing_state.next_cycle", planOutcomeRoutingNextCycle], + ["work_item.outcome_routing_state.reconcile", reconcileOutcomeRoutingState], + ["work_item.outcome_routing_state.write", writeOutcomeRoutingState], ["work_item.delivery_history.project", projectDeliveryHistory], ["work_item.delivery_response.project", projectDeliveryResponse], ["work_item.delivery_claim.validate", validateDeliveryClaim], diff --git a/loopx/control_plane/work_items/outcome_routing_plan.ts b/loopx/control_plane/work_items/outcome_routing_plan.ts new file mode 100644 index 0000000000..7e5ce15699 --- /dev/null +++ b/loopx/control_plane/work_items/outcome_routing_plan.ts @@ -0,0 +1,341 @@ +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { + optionalNonEmptyString, + requireBoolean, + requireInteger, + requireJsonObject, + requireNonEmptyString, + requireStringArray, + requireStringLiteral, +} from "../runtime_decode.ts"; + +import type { JsonObject } from "../effect_program.ts"; + +export const OUTCOME_ROUTING_PLAN_REQUEST_SCHEMA_VERSION = + "outcome_routing_plan_request_v0"; +export const OUTCOME_ROUTING_PLAN_SCHEMA_VERSION = "outcome_routing_plan_v0"; +const MAX_OUTCOMES = 128; +const MAX_WORK_ITEMS = 256; +const MAX_FEEDBACK_ITEMS = 256; +const PUBLIC_ID = /^[a-z][a-z0-9_-]{2,127}$/; + +export const OUTCOME_WORK_ROUTES = [ + "ai_execute", + "human_decide", + "human_execute", + "observe", + "reject", +] as const; + +export type OutcomeWorkRoute = (typeof OUTCOME_WORK_ROUTES)[number]; +export type OutcomeAuthorityTier = "A" | "B" | "C" | "D"; + +interface OutcomeWorkItem extends JsonObject { + work_item_id: string; + outcome_id: string; + title: string; + acceptance: string; + authority_tier: OutcomeAuthorityTier; + ai_capable: boolean; + prohibited: boolean; + material_decision: boolean; + human_identity_required: boolean; + wait_for?: string; + target_key: string; +} + +interface RoutedOutcomeWorkItem extends OutcomeWorkItem { + route: OutcomeWorkRoute; + route_reason: string; + status: + | "ready" + | "waiting_human_decision" + | "waiting_human_execution" + | "waiting_external_evidence" + | "cancelled"; + todo_projection: JsonObject; +} + +function boundedArray( + value: unknown, + label: string, + maximum: number, +): unknown[] { + if (!Array.isArray(value)) { + throw new EffectRuntimeRequestError(`${label} must be an array`); + } + if (value.length > maximum) { + throw new EffectRuntimeRequestError( + `${label} must contain at most ${maximum} items`, + ); + } + return value; +} + +function publicId(value: unknown, label: string): string { + const normalized = requireNonEmptyString(value, label); + if (!PUBLIC_ID.test(normalized)) { + throw new EffectRuntimeRequestError(`${label} must be a public-safe id`); + } + return normalized; +} + +function requireUniqueIds( + values: readonly JsonObject[], + field: string, + label: string, +): void { + const seen = new Set(); + for (const value of values) { + const identifier = String(value[field]); + if (seen.has(identifier)) { + throw new EffectRuntimeRequestError(`${label} must be unique`); + } + seen.add(identifier); + } +} + +function outcomeWorkItem(value: unknown, label: string): OutcomeWorkItem { + const raw = requireJsonObject(value, label); + const waitFor = optionalNonEmptyString(raw.wait_for, `${label}.wait_for`); + return { + work_item_id: publicId(raw.work_item_id, `${label}.work_item_id`), + outcome_id: publicId(raw.outcome_id, `${label}.outcome_id`), + title: requireNonEmptyString(raw.title, `${label}.title`), + acceptance: requireNonEmptyString(raw.acceptance, `${label}.acceptance`), + authority_tier: requireStringLiteral( + raw.authority_tier, + ["A", "B", "C", "D"] as const, + `${label}.authority_tier`, + ), + ai_capable: requireBoolean(raw.ai_capable, `${label}.ai_capable`), + prohibited: raw.prohibited === undefined + ? false + : requireBoolean(raw.prohibited, `${label}.prohibited`), + material_decision: raw.material_decision === undefined + ? false + : requireBoolean(raw.material_decision, `${label}.material_decision`), + human_identity_required: raw.human_identity_required === undefined + ? false + : requireBoolean( + raw.human_identity_required, + `${label}.human_identity_required`, + ), + ...(waitFor === null ? {} : { wait_for: waitFor }), + target_key: publicId(raw.target_key, `${label}.target_key`), + }; +} + +export function routeOutcomeWorkItem( + item: OutcomeWorkItem, +): { route: OutcomeWorkRoute; reason: string } { + if (item.prohibited || item.authority_tier === "D") { + return { + route: "reject", + reason: "policy or current authority prohibits execution", + }; + } + if (item.wait_for) { + return { + route: "observe", + reason: "work depends on a future external state", + }; + } + if (item.material_decision || item.authority_tier === "B") { + return { + route: "human_decide", + reason: "a material choice or authority grant is required", + }; + } + if (item.human_identity_required || item.authority_tier === "C") { + return { + route: "human_execute", + reason: "a human identity or physical action is required", + }; + } + if (item.ai_capable && item.authority_tier === "A") { + return { + route: "ai_execute", + reason: "AI capability, authority, and acceptance criteria are present", + }; + } + return { + route: "human_decide", + reason: "AI execution preconditions are incomplete", + }; +} + +function routeStatus(route: OutcomeWorkRoute): RoutedOutcomeWorkItem["status"] { + switch (route) { + case "ai_execute": return "ready"; + case "human_decide": return "waiting_human_decision"; + case "human_execute": return "waiting_human_execution"; + case "observe": return "waiting_external_evidence"; + case "reject": return "cancelled"; + } +} + +function todoProjection(item: OutcomeWorkItem, route: OutcomeWorkRoute): JsonObject { + const mapping: Record = { + ai_execute: ["agent", "advancement_task"], + human_decide: ["user", "user_gate"], + human_execute: ["user", "user_action"], + observe: ["agent", "continuous_monitor"], + reject: ["agent", "blocker"], + }; + const [role, taskClass] = mapping[route]; + return { + role, + task_class: taskClass, + action_kind: route, + target_key: item.target_key, + text: item.title, + acceptance: item.acceptance, + }; +} + +function projectWorkItem(value: unknown, label: string): RoutedOutcomeWorkItem { + const item = outcomeWorkItem(value, label); + const decision = routeOutcomeWorkItem(item); + return { + ...item, + route: decision.route, + route_reason: decision.reason, + status: routeStatus(decision.route), + todo_projection: todoProjection(item, decision.route), + }; +} + +function projectFeedback(value: unknown, label: string): JsonObject { + const raw = requireJsonObject(value, label); + const kind = requireStringLiteral( + raw.kind, + [ + "fact", + "decision", + "execution_result", + "risk", + "metric_change", + "comment", + ] as const, + `${label}.kind`, + ); + return { + feedback_id: publicId(raw.feedback_id, `${label}.feedback_id`), + source: requireNonEmptyString(raw.source, `${label}.source`), + subject: requireNonEmptyString(raw.subject, `${label}.subject`), + kind, + observed_at: requireNonEmptyString(raw.observed_at, `${label}.observed_at`), + evidence_ref: requireNonEmptyString(raw.evidence_ref, `${label}.evidence_ref`), + affected_outcome_ids: requireStringArray( + raw.affected_outcome_ids, + `${label}.affected_outcome_ids`, + ).map((item, index) => publicId( + item, + `${label}.affected_outcome_ids[${index}]`, + )), + disposition: kind === "comment" ? "recorded" : "replan", + }; +} + +/** + * Validate one outcome-level planning snapshot and project each open unit of + * work into LoopX's existing Todo lanes. This is a pure control-plane + * contract: providers own collection and execution, while LoopX owns routing + * precedence and the provider-neutral projection. + */ +export function projectOutcomeRoutingPlan(value: unknown): JsonObject { + const request = requireJsonObject(value, "outcome_routing_plan_request"); + if (request.schema_version !== OUTCOME_ROUTING_PLAN_REQUEST_SCHEMA_VERSION) { + throw new EffectRuntimeRequestError( + `outcome_routing_plan_request.schema_version must be ${OUTCOME_ROUTING_PLAN_REQUEST_SCHEMA_VERSION}`, + ); + } + const cycle = requireInteger(request.cycle, "outcome_routing_plan_request.cycle"); + if (cycle < 0 || !Number.isSafeInteger(cycle)) { + throw new EffectRuntimeRequestError( + "outcome_routing_plan_request.cycle must be a non-negative safe integer", + ); + } + const outcomes = boundedArray( + request.outcomes, + "outcome_routing_plan_request.outcomes", + MAX_OUTCOMES, + ).map((value, index) => { + const raw = requireJsonObject( + value, + `outcome_routing_plan_request.outcomes[${index}]`, + ); + return { + outcome_id: publicId( + raw.outcome_id, + `outcome_routing_plan_request.outcomes[${index}].outcome_id`, + ), + title: requireNonEmptyString( + raw.title, + `outcome_routing_plan_request.outcomes[${index}].title`, + ), + metric: requireNonEmptyString( + raw.metric, + `outcome_routing_plan_request.outcomes[${index}].metric`, + ), + target: requireNonEmptyString( + raw.target, + `outcome_routing_plan_request.outcomes[${index}].target`, + ), + evidence_source: requireNonEmptyString( + raw.evidence_source, + `outcome_routing_plan_request.outcomes[${index}].evidence_source`, + ), + }; + }); + requireUniqueIds(outcomes, "outcome_id", "outcome_id values"); + const outcomeIds = new Set(outcomes.map((outcome) => outcome.outcome_id)); + const workItems = boundedArray( + request.work_items, + "outcome_routing_plan_request.work_items", + MAX_WORK_ITEMS, + ).map((item, index) => projectWorkItem( + item, + `outcome_routing_plan_request.work_items[${index}]`, + )); + requireUniqueIds(workItems, "work_item_id", "outcome work_item_id values"); + requireUniqueIds(workItems, "target_key", "outcome work target_key values"); + for (const [index, item] of workItems.entries()) { + if (!outcomeIds.has(item.outcome_id)) { + throw new EffectRuntimeRequestError( + `outcome_routing_plan_request.work_items[${index}].outcome_id must reference an outcome`, + ); + } + } + const feedback = boundedArray( + request.feedback, + "outcome_routing_plan_request.feedback", + MAX_FEEDBACK_ITEMS, + ).map((item, index) => projectFeedback( + item, + `outcome_routing_plan_request.feedback[${index}]`, + )); + requireUniqueIds(feedback, "feedback_id", "feedback_id values"); + for (const [index, item] of feedback.entries()) { + for (const outcomeId of item.affected_outcome_ids as string[]) { + if (!outcomeIds.has(outcomeId)) { + throw new EffectRuntimeRequestError( + `outcome_routing_plan_request.feedback[${index}].affected_outcome_ids must reference outcomes`, + ); + } + } + } + return { + schema_version: OUTCOME_ROUTING_PLAN_SCHEMA_VERSION, + direction: requireNonEmptyString( + request.direction, + "outcome_routing_plan_request.direction", + ), + cycle, + outcomes, + work_items: workItems, + feedback, + replan_required: feedback.some((item) => item.disposition === "replan"), + }; +} diff --git a/loopx/control_plane/work_items/outcome_routing_state.ts b/loopx/control_plane/work_items/outcome_routing_state.ts new file mode 100644 index 0000000000..db73ba55e8 --- /dev/null +++ b/loopx/control_plane/work_items/outcome_routing_state.ts @@ -0,0 +1,543 @@ +import { createHash } from "node:crypto"; +import { readFile } from "node:fs/promises"; +import { isAbsolute, join } from "node:path"; + +import type { JsonObject } from "../effect_program.ts"; +import { + EffectRuntimeConflictError, + EffectRuntimeRequestError, +} from "../effect_runtime_errors.ts"; +import { atomicWriteJson, withFileMutationLock } from "../effect_runtime_io.ts"; +import { + optionalNonEmptyString, + requireBoolean, + requireJsonObject, + requireNonEmptyString, + requireStringLiteral, +} from "../runtime_decode.ts"; +import { + OUTCOME_ROUTING_PLAN_SCHEMA_VERSION, + projectOutcomeRoutingPlan, +} from "./outcome_routing_plan.ts"; + +export const OUTCOME_ROUTING_STATE_STORE_REQUEST_SCHEMA = + "outcome_routing_state_store_request_v0"; +export const OUTCOME_ROUTING_STATE_STORE_SCHEMA = + "outcome_routing_state_store_v0"; +export const OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA = + "outcome_routing_state_store_result_v0"; +export const OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA = + "outcome_routing_state_reconcile_request_v0"; +export const OUTCOME_ROUTING_STATE_RECONCILIATION_SCHEMA = + "outcome_routing_state_reconciliation_v0"; +export const OUTCOME_ROUTING_STATE_BIND_REQUEST_SCHEMA = + "outcome_routing_state_bind_request_v0"; +export const OUTCOME_ROUTING_NEXT_CYCLE_REQUEST_SCHEMA = + "outcome_routing_next_cycle_request_v0"; +export const OUTCOME_ROUTING_NEXT_CYCLE_SCHEMA = + "outcome_routing_next_cycle_v0"; + +function stableValue(value: unknown): unknown { + if (Array.isArray(value)) return value.map(stableValue); + if (typeof value !== "object" || value === null) return value; + return Object.fromEntries( + Object.entries(value as JsonObject) + .sort(([left], [right]) => left.localeCompare(right)) + .map(([key, child]) => [key, stableValue(child)]), + ); +} + +function revision(projection: JsonObject): string { + return createHash("sha256") + .update(JSON.stringify(stableValue(projection)), "utf8") + .digest("hex"); +} + +function todoFeedbackId( + workItemId: string, + nextStatus: string, + todoId: string, + cycle: number, +): string { + const digest = createHash("sha256") + .update(`${workItemId}\u001f${nextStatus}\u001f${todoId}\u001f${cycle}`, "utf8") + .digest("hex") + .slice(0, 24); + return `todo_feedback_${digest}`; +} + +function safeGoalSegment(goalId: string): string { + const label = goalId + .toLowerCase() + .replace(/[^a-z0-9._-]+/g, "-") + .replace(/^-+|-+$/g, "") + .slice(0, 47) || "goal"; + const digest = createHash("sha256").update(goalId, "utf8").digest("hex").slice(0, 16); + return `${label}-${digest}`; +} + +export function outcomeRoutingStatePath(runtimeRoot: string, goalId: string): string { + if (!isAbsolute(runtimeRoot)) { + throw new EffectRuntimeRequestError("runtime_root must be absolute"); + } + return join( + runtimeRoot, + "goals", + safeGoalSegment(goalId), + "outcome-routing", + "state.json", + ); +} + +function storeRequest(value: unknown): { + request: JsonObject; + goalId: string; + path: string; +} { + const request = requireJsonObject(value, "outcome_routing_state_store params"); + if (request.schema_version !== OUTCOME_ROUTING_STATE_STORE_REQUEST_SCHEMA) { + throw new EffectRuntimeRequestError("outcome routing state store request schema mismatch"); + } + const runtimeRoot = requireNonEmptyString(request.runtime_root, "runtime_root"); + const goalId = requireNonEmptyString(request.goal_id, "goal_id"); + return { request, goalId, path: outcomeRoutingStatePath(runtimeRoot, goalId) }; +} + +function decodeStoredState(value: unknown, goalId: string): JsonObject { + const stored = requireJsonObject(value, "stored outcome routing state"); + if ( + stored.schema_version !== OUTCOME_ROUTING_STATE_STORE_SCHEMA || + stored.goal_id !== goalId || + typeof stored.revision !== "string" || + !/^[a-f0-9]{64}$/.test(stored.revision) + ) { + throw new EffectRuntimeRequestError("stored outcome routing state is invalid"); + } + const projection = requireJsonObject(stored.projection, "stored projection"); + if (projection.schema_version !== OUTCOME_ROUTING_PLAN_SCHEMA_VERSION) { + throw new EffectRuntimeRequestError("stored outcome routing projection schema is invalid"); + } + const bindings = stored.todo_bindings === undefined + ? [] + : requireBindings(stored.todo_bindings, projection); + const reconciliation = stored.reconciliation === undefined + ? null + : requireJsonObject(stored.reconciliation, "stored reconciliation"); + if ( + reconciliation !== null && + reconciliation.schema_version !== OUTCOME_ROUTING_STATE_RECONCILIATION_SCHEMA + ) { + throw new EffectRuntimeRequestError("stored outcome routing reconciliation is invalid"); + } + const revisionContent: JsonObject = { projection }; + if (bindings.length > 0) revisionContent.todo_bindings = bindings; + if (reconciliation !== null) revisionContent.reconciliation = reconciliation; + if (revision(revisionContent) !== stored.revision) { + throw new EffectRuntimeRequestError("stored outcome routing state revision does not match content"); + } + return stored; +} + +function requireBindings(value: unknown, projection: JsonObject): JsonObject[] { + if (!Array.isArray(value)) { + throw new EffectRuntimeRequestError("outcome routing Todo bindings must be an array"); + } + const workItems = projection.work_items; + if (!Array.isArray(workItems)) { + throw new EffectRuntimeRequestError("stored outcome work_items must be an array"); + } + const targets = new Map(workItems.map((value) => { + const work = requireJsonObject(value, "stored outcome work item"); + return [ + requireNonEmptyString(work.work_item_id, "stored work_item_id"), + requireNonEmptyString(work.target_key, "stored target_key"), + ]; + })); + const workIds = new Set(); + const todoIds = new Set(); + return value.map((value, index) => { + const binding = requireJsonObject(value, `todo_bindings[${index}]`); + const workItemId = requireNonEmptyString(binding.work_item_id, `todo_bindings[${index}].work_item_id`); + const targetKey = requireNonEmptyString(binding.target_key, `todo_bindings[${index}].target_key`); + const todoId = requireNonEmptyString(binding.todo_id, `todo_bindings[${index}].todo_id`); + if (targets.get(workItemId) !== targetKey) { + throw new EffectRuntimeRequestError("Todo binding must match a projected work item and target"); + } + if (workIds.has(workItemId) || todoIds.has(todoId)) { + throw new EffectRuntimeRequestError("Todo bindings must have unique work_item_id and todo_id values"); + } + workIds.add(workItemId); + todoIds.add(todoId); + return { + work_item_id: workItemId, + target_key: targetKey, + todo_id: todoId, + role: requireStringLiteral(binding.role, ["agent", "user"] as const, `todo_bindings[${index}].role`), + }; + }); +} + +async function readStoredState(path: string, goalId: string): Promise { + try { + return decodeStoredState(JSON.parse(await readFile(path, "utf8")), goalId); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return null; + throw error; + } +} + +export async function loadOutcomeRoutingState(value: unknown): Promise { + const { goalId, path } = storeRequest(value); + return { + schema_version: OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA, + operation: "load", + goal_id: goalId, + path, + state: await readStoredState(path, goalId), + }; +} + +export async function writeOutcomeRoutingState(value: unknown): Promise { + const { request, goalId, path } = storeRequest(value); + const expectedRevision = optionalNonEmptyString( + request.expected_revision, + "expected_revision", + ); + const projection = projectOutcomeRoutingPlan(request.state); + const nextRevision = revision({ projection }); + return await withFileMutationLock(path, async () => { + const existing = await readStoredState(path, goalId); + if (existing?.revision === nextRevision) { + return { + schema_version: OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA, + operation: "write", + goal_id: goalId, + path, + state: existing, + written: false, + replayed: true, + }; + } + if (existing && expectedRevision === null) { + throw new EffectRuntimeConflictError( + "expected_revision is required when outcome routing state already exists", + ); + } + if (expectedRevision !== (existing?.revision ?? null)) { + throw new EffectRuntimeConflictError("outcome routing state revision changed"); + } + const stored: JsonObject = { + schema_version: OUTCOME_ROUTING_STATE_STORE_SCHEMA, + goal_id: goalId, + revision: nextRevision, + updated_at: requireNonEmptyString(request.updated_at, "updated_at"), + projection, + }; + await atomicWriteJson(path, stored); + const readback = await readStoredState(path, goalId); + if (!readback || readback.revision !== nextRevision) { + throw new Error("outcome routing state readback failed"); + } + return { + schema_version: OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA, + operation: "write", + goal_id: goalId, + path, + state: readback, + written: true, + replayed: false, + }; + }); +} + +export async function bindOutcomeRoutingTodos(value: unknown): Promise { + const request = requireJsonObject(value, "outcome_routing_state_bind params"); + if (request.schema_version !== OUTCOME_ROUTING_STATE_BIND_REQUEST_SCHEMA) { + throw new EffectRuntimeRequestError("outcome routing Todo bind request schema mismatch"); + } + const runtimeRoot = requireNonEmptyString(request.runtime_root, "runtime_root"); + const goalId = requireNonEmptyString(request.goal_id, "goal_id"); + const path = outcomeRoutingStatePath(runtimeRoot, goalId); + const expectedRevision = requireNonEmptyString(request.expected_revision, "expected_revision"); + const updatedAt = requireNonEmptyString(request.updated_at, "updated_at"); + return await withFileMutationLock(path, async () => { + const existing = await readStoredState(path, goalId); + if (!existing) throw new EffectRuntimeRequestError("persisted outcome routing state does not exist"); + if (existing.revision !== expectedRevision) { + throw new EffectRuntimeConflictError("outcome routing state revision changed"); + } + const projection = requireJsonObject(existing.projection, "stored projection"); + const todoBindings = requireBindings(request.todo_bindings, projection); + const revisionContent: JsonObject = { projection }; + if (todoBindings.length > 0) revisionContent.todo_bindings = todoBindings; + if (existing.reconciliation !== undefined) { + revisionContent.reconciliation = requireJsonObject(existing.reconciliation, "stored reconciliation"); + } + const nextRevision = revision(revisionContent); + if (nextRevision === existing.revision) { + return { schema_version: OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA, operation: "bind", goal_id: goalId, path, state: existing, written: false, replayed: true }; + } + const nextState: JsonObject = { ...existing, ...revisionContent, revision: nextRevision, updated_at: updatedAt }; + await atomicWriteJson(path, nextState); + const readback = await readStoredState(path, goalId); + if (!readback || readback.revision !== nextRevision) throw new Error("outcome routing Todo binding readback failed"); + return { schema_version: OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA, operation: "bind", goal_id: goalId, path, state: readback, written: true, replayed: false }; + }); +} + +function todoReconciliation(value: unknown, projection: JsonObject): JsonObject { + const request = requireJsonObject(value, "outcome routing reconciliation"); + if (!Array.isArray(request.observations)) { + throw new EffectRuntimeRequestError("outcome routing observations must be an array"); + } + const workItems = projection.work_items; + if (!Array.isArray(workItems)) { + throw new EffectRuntimeRequestError("stored outcome work_items must be an array"); + } + const byTarget = new Map(); + for (const item of workItems) { + const work = requireJsonObject(item, "stored outcome work item"); + const target = requireNonEmptyString(work.target_key, "stored work target_key"); + if (byTarget.has(target)) { + throw new EffectRuntimeRequestError("stored outcome work target_key must be unique"); + } + byTarget.set(target, work); + } + const seenTargets = new Set(); + const observations = request.observations.map((value, index) => { + const raw = requireJsonObject(value, `observations[${index}]`); + const targetKey = requireNonEmptyString(raw.target_key, `observations[${index}].target_key`); + if (seenTargets.has(targetKey)) { + throw new EffectRuntimeRequestError("routed Todo observations must have unique target_key values"); + } + seenTargets.add(targetKey); + const work = byTarget.get(targetKey); + if (!work) { + throw new EffectRuntimeRequestError( + `routed Todo observation target_key ${JSON.stringify(targetKey)} is unknown`, + ); + } + const todoStatus = requireStringLiteral( + raw.status, + ["open", "done", "blocked", "deferred"] as const, + `observations[${index}].status`, + ); + const evidenceRef = optionalNonEmptyString( + raw.evidence_ref, + `observations[${index}].evidence_ref`, + ); + const priorStatus = requireNonEmptyString(work.status, "stored work status"); + const nextStatus = todoStatus === "done" + ? (evidenceRef === null ? "awaiting_evidence" : "done") + : todoStatus === "blocked" + ? "replanning" + : priorStatus; + return { + work_item_id: work.work_item_id, + target_key: targetKey, + todo_id: requireNonEmptyString(raw.todo_id, `observations[${index}].todo_id`), + todo_status: todoStatus, + prior_status: priorStatus, + next_status: nextStatus, + ...(evidenceRef === null ? {} : { evidence_ref: evidenceRef }), + changed: priorStatus !== nextStatus, + }; + }); + return { + schema_version: OUTCOME_ROUTING_STATE_RECONCILIATION_SCHEMA, + observations, + replan_required: observations.some((item) => + item.next_status === "replanning" || item.next_status === "awaiting_evidence" + ), + }; +} + +export async function reconcileOutcomeRoutingState(value: unknown): Promise { + const request = requireJsonObject(value, "outcome_routing_state_reconcile params"); + if (request.schema_version !== OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA) { + throw new EffectRuntimeRequestError("outcome routing reconciliation request schema mismatch"); + } + const runtimeRoot = requireNonEmptyString(request.runtime_root, "runtime_root"); + const goalId = requireNonEmptyString(request.goal_id, "goal_id"); + const path = outcomeRoutingStatePath(runtimeRoot, goalId); + const expectedRevision = requireNonEmptyString( + request.expected_revision, + "expected_revision", + ); + const execute = requireBoolean(request.execute, "execute"); + const updatedAt = requireNonEmptyString(request.updated_at, "updated_at"); + return await withFileMutationLock(path, async () => { + const existing = await readStoredState(path, goalId); + if (!existing) { + throw new EffectRuntimeRequestError("persisted outcome routing state does not exist"); + } + if (existing.revision !== expectedRevision) { + throw new EffectRuntimeConflictError("outcome routing state revision changed"); + } + const projection = requireJsonObject(existing.projection, "stored projection"); + const reconciliation = todoReconciliation(request, projection); + const revisionContent: JsonObject = { projection, reconciliation }; + if (Array.isArray(existing.todo_bindings) && existing.todo_bindings.length > 0) { + revisionContent.todo_bindings = existing.todo_bindings; + } + const nextRevision = revision(revisionContent); + const nextState: JsonObject = { + ...existing, + revision: nextRevision, + updated_at: updatedAt, + reconciliation, + }; + if (!execute) { + return { + schema_version: OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA, + operation: "reconcile", + goal_id: goalId, + path, + dry_run: true, + state: nextState, + written: false, + replayed: false, + }; + } + if (existing.revision === nextRevision) { + return { + schema_version: OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA, + operation: "reconcile", + goal_id: goalId, + path, + dry_run: false, + state: existing, + written: false, + replayed: true, + }; + } + await atomicWriteJson(path, nextState); + const readback = await readStoredState(path, goalId); + if (!readback || readback.revision !== nextRevision) { + throw new Error("outcome routing reconciliation readback failed"); + } + return { + schema_version: OUTCOME_ROUTING_STATE_STORE_RESULT_SCHEMA, + operation: "reconcile", + goal_id: goalId, + path, + dry_run: false, + state: readback, + written: true, + replayed: false, + }; + }); +} + +export function planOutcomeRoutingNextCycle(value: unknown): JsonObject { + const request = requireJsonObject(value, "outcome_routing_next_cycle params"); + if (request.schema_version !== OUTCOME_ROUTING_NEXT_CYCLE_REQUEST_SCHEMA) { + throw new EffectRuntimeRequestError("outcome routing next-cycle request schema mismatch"); + } + const goalId = requireNonEmptyString(request.goal_id, "goal_id"); + const stored = decodeStoredState(request.state, goalId); + const projection = requireJsonObject(stored.projection, "stored projection"); + const reconciliation = stored.reconciliation === undefined + ? null + : requireJsonObject(stored.reconciliation, "stored reconciliation"); + const observations = reconciliation?.observations; + if (!Array.isArray(observations) || observations.length === 0) { + throw new EffectRuntimeRequestError("outcome routing next cycle requires Todo reconciliation"); + } + const byTarget = new Map(); + for (const value of observations) { + const observation = requireJsonObject(value, "reconciliation observation"); + byTarget.set( + requireNonEmptyString(observation.target_key, "observation target_key"), + observation, + ); + } + const workItems = projection.work_items; + if (!Array.isArray(workItems)) { + throw new EffectRuntimeRequestError("stored outcome work_items must be an array"); + } + const nextWorkItems: JsonObject[] = []; + const feedback: JsonObject[] = Array.isArray(projection.feedback) + ? projection.feedback.map((item) => { + const prior = requireJsonObject(item, "stored routing feedback"); + const { disposition: _disposition, ...requestFeedback } = prior; + return requestFeedback; + }) + : []; + let convergedCount = 0; + for (const value of workItems) { + const work = requireJsonObject(value, "stored outcome work item"); + const targetKey = requireNonEmptyString(work.target_key, "stored work target_key"); + const observation = byTarget.get(targetKey); + const nextStatus = observation?.next_status; + if (nextStatus === "done") { + const todoId = requireNonEmptyString(observation?.todo_id, "observation todo_id"); + convergedCount += 1; + feedback.push({ + feedback_id: todoFeedbackId( + requireNonEmptyString(work.work_item_id, "stored work work_item_id"), + nextStatus, + todoId, + Number(projection.cycle), + ), + source: "loopx_todo", + subject: `Todo completed: ${work.title}`, + kind: "execution_result", + observed_at: requireNonEmptyString(stored.updated_at, "stored updated_at"), + evidence_ref: requireNonEmptyString( + observation?.evidence_ref, + "completion evidence_ref", + ), + affected_outcome_ids: [work.outcome_id], + }); + continue; + } + const { + route: _route, + route_reason: _routeReason, + status: _status, + todo_projection: _todoProjection, + ...nextWork + } = work; + nextWorkItems.push(nextWork); + if (nextStatus === "replanning" || nextStatus === "awaiting_evidence") { + const todoId = requireNonEmptyString(observation?.todo_id, "observation todo_id"); + feedback.push({ + feedback_id: todoFeedbackId( + requireNonEmptyString(work.work_item_id, "stored work work_item_id"), + nextStatus, + todoId, + Number(projection.cycle), + ), + source: "loopx_todo", + subject: nextStatus === "replanning" + ? `Todo blocked: ${work.title}` + : `Todo completion needs evidence: ${work.title}`, + kind: "risk", + observed_at: requireNonEmptyString(stored.updated_at, "stored updated_at"), + evidence_ref: `loopx-todo:${todoId}`, + affected_outcome_ids: [work.outcome_id], + }); + } + } + const nextState: JsonObject = { + schema_version: "outcome_routing_plan_request_v0", + direction: projection.direction, + cycle: Number(projection.cycle) + 1, + outcomes: projection.outcomes, + work_items: nextWorkItems, + feedback, + }; + const nextProjection = projectOutcomeRoutingPlan(nextState); + return { + schema_version: OUTCOME_ROUTING_NEXT_CYCLE_SCHEMA, + goal_id: goalId, + source_revision: stored.revision, + converged_work_item_count: convergedCount, + remaining_work_item_count: nextWorkItems.length, + goal_converged: nextWorkItems.length === 0, + replan_required: nextProjection.replan_required, + state: nextState, + projection: nextProjection, + }; +} diff --git a/tests/control_plane/test_company_control_loop_cli.py b/tests/control_plane/test_company_control_loop_cli.py new file mode 100644 index 0000000000..c2476aef1c --- /dev/null +++ b/tests/control_plane/test_company_control_loop_cli.py @@ -0,0 +1,458 @@ +from __future__ import annotations + +import json + +from loopx.cli import main +from loopx.cli_commands import company_control_loop + + +def _request() -> dict[str, object]: + return { + "schema_version": "outcome_routing_plan_request_v0", + "direction": "Improve durable customer value.", + "cycle": 1, + "outcomes": [ + { + "outcome_id": "outcome_activation", + "title": "Improve activation", + "metric": "seven day activation rate", + "target": ">= 40%", + "evidence_source": "activation analytics", + } + ], + "work_items": [], + "feedback": [], + } + + +def _stored_projection( + *work_items: dict[str, object], + bindings: list[dict[str, object]] | None = None, +) -> dict[str, object]: + return { + "schema_version": "outcome_routing_state_store_result_v0", + "operation": "load", + "goal_id": "company-goal", + "state": { + "revision": "a" * 64, + "projection": { + "schema_version": "outcome_routing_plan_v0", + "work_items": list(work_items), + }, + **({"todo_bindings": bindings} if bindings is not None else {}), + }, + } + + +def _routed_work( + work_item_id: str, + target_key: str, + *, + role: str = "agent", + task_class: str = "advancement_task", + action_kind: str = "ai_execute", +) -> dict[str, object]: + return { + "work_item_id": work_item_id, + "target_key": target_key, + "todo_projection": { + "role": role, + "task_class": task_class, + "action_kind": action_kind, + "target_key": target_key, + "text": f"Advance {work_item_id}", + "acceptance": f"Evidence for {work_item_id}", + }, + } + + +def test_outcome_routing_plan_cli_calls_typed_projection( + tmp_path, monkeypatch, capsys +) -> None: + state_path = tmp_path / "company.json" + state_path.write_text(json.dumps(_request()), encoding="utf-8") + calls: list[tuple[str, dict[str, object]]] = [] + + def project(method: str, params: dict[str, object]) -> dict[str, object]: + calls.append((method, params)) + return { + "schema_version": "outcome_routing_plan_v0", + "direction": params["direction"], + "cycle": params["cycle"], + "outcomes": params["outcomes"], + "work_items": [], + "feedback": [], + "replan_required": False, + } + + monkeypatch.setattr(company_control_loop, "effect_runtime_result", project) + + assert main([ + "--format", + "json", + "company-control-loop", + "project", + "--state-json", + str(state_path), + ]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["schema_version"] == "outcome_routing_plan_v0" + assert calls == [("work_item.outcome_routing_plan.project", _request())] + + +def test_outcome_routing_plan_cli_rejects_non_object_json(tmp_path, capsys) -> None: + state_path = tmp_path / "company.json" + state_path.write_text("[]", encoding="utf-8") + + assert main([ + "--format", + "json", + "company-control-loop", + "project", + "--state-json", + str(state_path), + ]) == 1 + payload = json.loads(capsys.readouterr().out) + assert payload["ok"] is False + assert "must contain an object" in payload["error"] + + +def test_outcome_routing_plan_save_previews_then_writes_with_revision( + tmp_path, monkeypatch, capsys +) -> None: + state_path = tmp_path / "company.json" + state_path.write_text(json.dumps(_request()), encoding="utf-8") + calls: list[str] = [] + + def runtime(method: str, params: dict[str, object]) -> dict[str, object]: + calls.append(method) + if method.endswith("project"): + return {"schema_version": "outcome_routing_plan_v0"} + assert params["expected_revision"] == "a" * 64 + return { + "schema_version": "outcome_routing_state_store_result_v0", + "operation": "write", + "written": True, + "replayed": False, + "state": {"revision": "b" * 64}, + } + + monkeypatch.setattr(company_control_loop, "effect_runtime_result", runtime) + common = [ + "--format", "json", "--runtime-root", str(tmp_path / "runtime"), + "company-control-loop", "save", "--goal-id", "company-goal", + "--state-json", str(state_path), + ] + assert main(common) == 0 + assert json.loads(capsys.readouterr().out)["dry_run"] is True + assert calls == ["work_item.outcome_routing_plan.project"] + + calls.clear() + assert main([*common, "--expected-revision", "a" * 64, "--execute"]) == 0 + assert json.loads(capsys.readouterr().out)["written"] is True + assert calls == [ + "work_item.outcome_routing_plan.project", + "work_item.outcome_routing_state.write", + ] + + +def test_outcome_routing_plan_show_reads_goal_state(tmp_path, monkeypatch, capsys) -> None: + calls: list[tuple[str, dict[str, object]]] = [] + + def runtime(method: str, params: dict[str, object]) -> dict[str, object]: + calls.append((method, params)) + return { + "schema_version": "outcome_routing_state_store_result_v0", + "operation": "load", + "goal_id": params["goal_id"], + "state": None, + } + + monkeypatch.setattr(company_control_loop, "effect_runtime_result", runtime) + assert main([ + "--format", "json", "--runtime-root", str(tmp_path / "runtime"), + "company-control-loop", "show", "--goal-id", "company-goal", + ]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["operation"] == "load" + assert calls[0][0] == "work_item.outcome_routing_state.load" + assert calls[0][1]["goal_id"] == "company-goal" + + +def test_outcome_routing_plan_sync_todos_previews_existing_and_missing( + tmp_path, monkeypatch, capsys +) -> None: + monkeypatch.setattr( + company_control_loop, + "effect_runtime_result", + lambda method, params: _stored_projection( + _routed_work("work_existing", "target_existing"), + _routed_work("work_missing", "target_missing"), + ), + ) + monkeypatch.setattr( + company_control_loop, + "list_goal_todos", + lambda **kwargs: { + "todos": [{ + "todo_id": "todo_existing", + "target_key": "target_existing", + "status": "open", + }] + }, + ) + writes: list[dict[str, object]] = [] + monkeypatch.setattr( + company_control_loop, + "add_goal_todo", + lambda **kwargs: writes.append(kwargs), + ) + + assert main([ + "--format", "json", "--registry", str(tmp_path / "registry.json"), + "--runtime-root", str(tmp_path / "runtime"), + "company-control-loop", "sync-todos", "--goal-id", "company-goal", + "--agent-id", "agent-ceo", "--project", str(tmp_path), + ]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["dry_run"] is True + assert payload["readback_verified"] is False + assert [item["action"] for item in payload["actions"]] == [ + "linked_existing", "would_create", + ] + assert writes == [] + + +def test_outcome_routing_plan_sync_todos_creates_and_verifies_readback( + tmp_path, monkeypatch, capsys +) -> None: + monkeypatch.setattr( + company_control_loop, + "effect_runtime_result", + lambda method, params: ( + {"state": {"revision": "b" * 64}} + if method.endswith(".bind") + else _stored_projection(_routed_work( + "work_decision", + "target_decision", + role="user", + task_class="user_gate", + action_kind="human_decide", + ), + _routed_work( + "work_watch", + "target_watch", + task_class="continuous_monitor", + action_kind="observe", + )) + ), + ) + listings = iter([ + {"todos": []}, + {"todos": [ + {"todo_id": "created_1"}, + {"todo_id": "created_2", "target_key": "target_watch"}, + ]}, + ]) + monkeypatch.setattr( + company_control_loop, "list_goal_todos", lambda **kwargs: next(listings) + ) + writes: list[dict[str, object]] = [] + + def add(**kwargs: object) -> dict[str, object]: + writes.append(kwargs) + return {"todo_id": f"created_{len(writes)}"} + + monkeypatch.setattr(company_control_loop, "add_goal_todo", add) + + assert main([ + "--format", "json", "--registry", str(tmp_path / "registry.json"), + "--runtime-root", str(tmp_path / "runtime"), + "company-control-loop", "sync-todos", "--goal-id", "company-goal", + "--agent-id", "agent-ceo", "--project", str(tmp_path), "--execute", + ]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["readback_verified"] is True + assert [item["todo_id"] for item in payload["actions"]] == [ + "created_1", "created_2", + ] + assert writes[0]["blocks_agent"] == "agent-ceo" + assert writes[0]["decision_scope"] == "direction:action:target_decision" + assert writes[0]["monitor_metadata"] == {} + assert writes[1]["monitor_metadata"] == { + "target_key": "target_watch", + "watch_only": "true", + } + + +def test_outcome_routing_plan_sync_todos_fails_on_missing_readback( + tmp_path, monkeypatch, capsys +) -> None: + monkeypatch.setattr( + company_control_loop, + "effect_runtime_result", + lambda method, params: _stored_projection( + _routed_work("work_missing", "target_missing") + ), + ) + listings = iter([{"todos": []}, {"todos": []}]) + monkeypatch.setattr( + company_control_loop, "list_goal_todos", lambda **kwargs: next(listings) + ) + monkeypatch.setattr( + company_control_loop, + "add_goal_todo", + lambda **kwargs: {"todo_id": "todo_unreadable"}, + ) + + assert main([ + "--format", "json", "--registry", str(tmp_path / "registry.json"), + "--runtime-root", str(tmp_path / "runtime"), + "company-control-loop", "sync-todos", "--goal-id", "company-goal", + "--agent-id", "agent-ceo", "--project", str(tmp_path), "--execute", + ]) == 1 + payload = json.loads(capsys.readouterr().out) + assert payload["ok"] is False + assert "readback missing ids" in payload["error"] + + +def test_outcome_routing_plan_reconcile_todos_sends_evidence_to_typed_owner( + tmp_path, monkeypatch, capsys +) -> None: + calls: list[tuple[str, dict[str, object]]] = [] + + def runtime(method: str, params: dict[str, object]) -> dict[str, object]: + calls.append((method, params)) + if method.endswith(".load"): + return _stored_projection( + _routed_work("work_activation", "target_activation"), + bindings=[{ + "work_item_id": "work_activation", + "target_key": "target_activation", + "todo_id": "todo_activation", + "role": "agent", + }], + ) + return { + "schema_version": "outcome_routing_state_store_result_v0", + "operation": "reconcile", + "dry_run": True, + "written": False, + } + + monkeypatch.setattr(company_control_loop, "effect_runtime_result", runtime) + monkeypatch.setattr( + company_control_loop, + "list_goal_todos", + lambda **kwargs: {"todos": [ + { + "todo_id": "todo_activation", + "target_key": "target_activation", + "status": "done", + "evidence": "artifact:activation-report", + }, + { + "todo_id": "todo_unrelated", + "target_key": "other_target", + "status": "done", + }, + ]}, + ) + + assert main([ + "--format", "json", "--registry", str(tmp_path / "registry.json"), + "--runtime-root", str(tmp_path / "runtime"), + "company-control-loop", "reconcile-todos", "--goal-id", "company-goal", + "--agent-id", "agent-ceo", "--project", str(tmp_path), + ]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["dry_run"] is True + assert calls[1][0] == "work_item.outcome_routing_state.reconcile" + assert calls[1][1]["expected_revision"] == "a" * 64 + assert calls[1][1]["execute"] is False + assert calls[1][1]["observations"] == [{ + "target_key": "target_activation", + "todo_id": "todo_activation", + "status": "done", + "evidence_ref": "artifact:activation-report", + }] + + +def test_outcome_routing_plan_reconcile_todos_uses_persisted_todo_identity( + tmp_path, monkeypatch, capsys +) -> None: + monkeypatch.setattr( + company_control_loop, + "effect_runtime_result", + lambda method, params: _stored_projection( + _routed_work("work_activation", "target_activation"), + bindings=[{ + "work_item_id": "work_activation", + "target_key": "target_activation", + "todo_id": "todo_first", + "role": "agent", + }], + ), + ) + monkeypatch.setattr( + company_control_loop, + "list_goal_todos", + lambda **kwargs: {"todos": [ + {"todo_id": "todo_first", "target_key": "target_activation", "status": "open"}, + {"todo_id": "todo_second", "target_key": "target_activation", "status": "done"}, + ]}, + ) + assert main([ + "--format", "json", "--registry", str(tmp_path / "registry.json"), + "--runtime-root", str(tmp_path / "runtime"), + "company-control-loop", "reconcile-todos", "--goal-id", "company-goal", + "--agent-id", "agent-ceo", "--execute", + ]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["ok"] is True + + +def test_outcome_routing_plan_next_cycle_uses_persisted_reconciliation( + tmp_path, monkeypatch, capsys +) -> None: + stored = _stored_projection( + _routed_work("work_activation", "target_activation") + ) + state = stored["state"] + assert isinstance(state, dict) + state["reconciliation"] = { + "schema_version": "outcome_routing_state_reconciliation_v0", + "observations": [{ + "work_item_id": "work_activation", + "target_key": "target_activation", + "todo_id": "todo_activation", + "todo_status": "done", + "prior_status": "ready", + "next_status": "done", + "evidence_ref": "artifact:activation-report", + "changed": True, + }], + "replan_required": False, + } + calls: list[tuple[str, dict[str, object]]] = [] + + def runtime(method: str, params: dict[str, object]) -> dict[str, object]: + calls.append((method, params)) + if method.endswith(".load"): + return stored + return { + "schema_version": "outcome_routing_next_cycle_v0", + "goal_id": "company-goal", + "goal_converged": True, + "remaining_work_item_count": 0, + } + + monkeypatch.setattr(company_control_loop, "effect_runtime_result", runtime) + assert main([ + "--format", "json", "--runtime-root", str(tmp_path / "runtime"), + "company-control-loop", "next-cycle", "--goal-id", "company-goal", + ]) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["goal_converged"] is True + assert calls[1][0] == "work_item.outcome_routing_state.next_cycle" + assert calls[1][1]["state"] == state diff --git a/tests/control_plane_ts/outcome_routing_e2e.test.ts b/tests/control_plane_ts/outcome_routing_e2e.test.ts new file mode 100644 index 0000000000..fb47339b49 --- /dev/null +++ b/tests/control_plane_ts/outcome_routing_e2e.test.ts @@ -0,0 +1,154 @@ +import assert from "node:assert/strict"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; + +import { + OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + OUTCOME_ROUTING_STATE_STORE_REQUEST_SCHEMA, + loadOutcomeRoutingState, + planOutcomeRoutingNextCycle, + reconcileOutcomeRoutingState, + writeOutcomeRoutingState, +} from "../../loopx/control_plane/work_items/outcome_routing_state.ts"; + +function state() { + return { + schema_version: "outcome_routing_plan_request_v0", + direction: "Improve durable customer value.", + cycle: 1, + outcomes: [{ + outcome_id: "outcome_activation", + title: "Improve activation", + metric: "seven day activation rate", + target: ">= 40%", + evidence_source: "activation analytics", + }], + work_items: [ + { + work_item_id: "work_ai_delivery", + outcome_id: "outcome_activation", + title: "Implement activation instrumentation", + acceptance: "instrumentation report", + authority_tier: "A", + ai_capable: true, + target_key: "ai_delivery", + }, + { + work_item_id: "work_human_decision", + outcome_id: "outcome_activation", + title: "Choose the activation threshold", + acceptance: "recorded threshold decision", + authority_tier: "B", + ai_capable: false, + material_decision: true, + target_key: "human_decision", + }, + { + work_item_id: "work_human_execution", + outcome_id: "outcome_activation", + title: "Interview the launch customer", + acceptance: "customer interview record", + authority_tier: "C", + ai_capable: false, + human_identity_required: true, + target_key: "human_execution", + }, + ], + feedback: [], + }; +} + +function storeRequest(runtimeRoot: string, extra: Record) { + return { + schema_version: OUTCOME_ROUTING_STATE_STORE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + ...extra, + }; +} + +test("outcome routing v0 closes AI, human, restart, escalation, and convergence scenarios", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-e2e-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + + const first = await writeOutcomeRoutingState(storeRequest(runtimeRoot, { + state: state(), + updated_at: "2026-09-17T00:00:00Z", + })); + const firstState = first.state as Record; + assert.deepEqual( + firstState.projection.work_items.map((item: Record) => item.route), + ["ai_execute", "human_decide", "human_execute"], + ); + + const restarted = await loadOutcomeRoutingState(storeRequest(runtimeRoot, {})); + assert.deepEqual(restarted.state, first.state); + + const failedCycle = await reconcileOutcomeRoutingState({ + schema_version: OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: firstState.revision, + updated_at: "2026-09-17T00:01:00Z", + execute: true, + observations: [ + { + target_key: "ai_delivery", + todo_id: "todo_ai_delivery", + status: "done", + evidence_ref: "artifact:instrumentation-report", + }, + { + target_key: "human_decision", + todo_id: "todo_human_decision", + status: "done", + evidence_ref: "decision:activation-threshold", + }, + { + target_key: "human_execution", + todo_id: "todo_human_execution", + status: "blocked", + }, + ], + }); + const replan = planOutcomeRoutingNextCycle({ + schema_version: "outcome_routing_next_cycle_request_v0", + goal_id: "company-goal", + state: failedCycle.state, + }); + assert.equal(replan.replan_required, true); + assert.equal(replan.converged_work_item_count, 2); + assert.equal(replan.remaining_work_item_count, 1); + const replannedState = replan.state as Record; + assert.equal(replannedState.feedback.at(-1).kind, "risk"); + assert.equal(replannedState.work_items[0].work_item_id, "work_human_execution"); + + const second = await writeOutcomeRoutingState(storeRequest(runtimeRoot, { + state: replannedState, + expected_revision: (failedCycle.state as Record).revision, + updated_at: "2026-09-17T00:02:00Z", + })); + const completedCycle = await reconcileOutcomeRoutingState({ + schema_version: OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: (second.state as Record).revision, + updated_at: "2026-09-17T00:03:00Z", + execute: true, + observations: [{ + target_key: "human_execution", + todo_id: "todo_human_execution_retry", + status: "done", + evidence_ref: "record:customer-interview", + }], + }); + const converged = planOutcomeRoutingNextCycle({ + schema_version: "outcome_routing_next_cycle_request_v0", + goal_id: "company-goal", + state: completedCycle.state, + }); + assert.equal(converged.goal_converged, true); + assert.equal(converged.remaining_work_item_count, 0); +}); diff --git a/tests/control_plane_ts/outcome_routing_plan.test.ts b/tests/control_plane_ts/outcome_routing_plan.test.ts new file mode 100644 index 0000000000..1f2e8886bd --- /dev/null +++ b/tests/control_plane_ts/outcome_routing_plan.test.ts @@ -0,0 +1,168 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { + OUTCOME_ROUTING_PLAN_REQUEST_SCHEMA_VERSION, + projectOutcomeRoutingPlan, +} from "../../loopx/control_plane/work_items/outcome_routing_plan.ts"; + +function request(overrides: Record = {}) { + return { + schema_version: OUTCOME_ROUTING_PLAN_REQUEST_SCHEMA_VERSION, + direction: "Improve durable customer value.", + cycle: 3, + outcomes: [{ + outcome_id: "outcome_activation", + title: "Improve activation", + metric: "seven day activation rate", + target: ">= 40%", + evidence_source: "activation analytics", + }], + work_items: [], + feedback: [], + ...overrides, + }; +} + +function work(overrides: Record = {}) { + return { + work_item_id: "work_activation_analysis", + outcome_id: "outcome_activation", + title: "Analyze the activation funnel.", + acceptance: "Baseline every stage and propose three measurable experiments.", + authority_tier: "A", + ai_capable: true, + target_key: "activation_funnel_analysis", + ...overrides, + }; +} + +test("outcome routing loop routes AI work into an advancement Todo", () => { + const result = projectOutcomeRoutingPlan(request({ work_items: [work()] })); + const item = (result.work_items as Record[])[0]; + + assert.equal(result.schema_version, "outcome_routing_plan_v0"); + assert.equal(item.route, "ai_execute"); + assert.equal(item.status, "ready"); + assert.deepEqual(item.todo_projection, { + role: "agent", + task_class: "advancement_task", + action_kind: "ai_execute", + target_key: "activation_funnel_analysis", + text: "Analyze the activation funnel.", + acceptance: "Baseline every stage and propose three measurable experiments.", + }); +}); + +test("routing precedence preserves authority, waiting, and human boundaries", () => { + const result = projectOutcomeRoutingPlan(request({ + work_items: [ + work({ work_item_id: "work_rejected", target_key: "target_rejected", prohibited: true }), + work({ work_item_id: "work_observe", target_key: "target_observe", wait_for: "provider result" }), + work({ work_item_id: "work_decide", target_key: "target_decide", material_decision: true }), + work({ work_item_id: "work_execute", target_key: "target_execute", human_identity_required: true }), + work({ work_item_id: "work_incomplete", target_key: "target_incomplete", ai_capable: false }), + ], + })); + + assert.deepEqual( + (result.work_items as Record[]).map((item) => item.route), + ["reject", "observe", "human_decide", "human_execute", "human_decide"], + ); + assert.deepEqual( + (result.work_items as Record[]).map((item) => + (item.todo_projection as Record).task_class + ), + ["blocker", "continuous_monitor", "user_gate", "user_action", "user_gate"], + ); +}); + +test("material feedback creates an explicit replan signal", () => { + const result = projectOutcomeRoutingPlan(request({ + feedback: [ + { + feedback_id: "feedback_metric_change", + source: "analytics", + subject: "activation", + kind: "metric_change", + observed_at: "2026-09-17T00:00:00Z", + evidence_ref: "report:activation-2026-09-17", + affected_outcome_ids: ["outcome_activation"], + }, + { + feedback_id: "feedback_comment", + source: "support", + subject: "onboarding copy", + kind: "comment", + observed_at: "2026-09-17T00:01:00Z", + evidence_ref: "ticket:123", + affected_outcome_ids: ["outcome_activation"], + }, + ], + })); + + assert.equal(result.replan_required, true); + assert.deepEqual( + (result.feedback as Record[]).map((item) => item.disposition), + ["replan", "recorded"], + ); +}); + +test("outcome routing loop rejects dangling outcome references and unsafe ids", () => { + assert.throws( + () => projectOutcomeRoutingPlan(request({ + work_items: [work({ outcome_id: "outcome_missing" })], + })), + /outcome_id must reference an outcome/, + ); + assert.throws( + () => projectOutcomeRoutingPlan(request({ + work_items: [work({ target_key: "../../private" })], + })), + /target_key must be a public-safe id/, + ); +}); + +test("outcome routing loop rejects ambiguous identifiers and unsafe cycle integers", () => { + assert.throws( + () => projectOutcomeRoutingPlan(request({ + outcomes: [request().outcomes[0], request().outcomes[0]], + })), + /outcome_id values must be unique/, + ); + assert.throws( + () => projectOutcomeRoutingPlan(request({ + work_items: [ + work({ work_item_id: "work_first" }), + work({ work_item_id: "work_second" }), + ], + })), + /target_key values must be unique/, + ); + assert.throws( + () => projectOutcomeRoutingPlan(request({ + work_items: [ + work({ target_key: "target_first" }), + work({ target_key: "target_second" }), + ], + })), + /work_item_id values must be unique/, + ); + const feedback = { + feedback_id: "feedback_duplicate", + source: "analytics", + subject: "activation", + kind: "fact", + observed_at: "2026-09-17T00:00:00Z", + evidence_ref: "report:activation", + affected_outcome_ids: ["outcome_activation"], + }; + assert.throws( + () => projectOutcomeRoutingPlan(request({ feedback: [feedback, feedback] })), + /feedback_id values must be unique/, + ); + assert.throws( + () => projectOutcomeRoutingPlan(request({ cycle: Number.MAX_SAFE_INTEGER + 1 })), + /non-negative safe integer/, + ); +}); diff --git a/tests/control_plane_ts/outcome_routing_state.test.ts b/tests/control_plane_ts/outcome_routing_state.test.ts new file mode 100644 index 0000000000..03b815afc0 --- /dev/null +++ b/tests/control_plane_ts/outcome_routing_state.test.ts @@ -0,0 +1,430 @@ +import assert from "node:assert/strict"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; + +import { + bindOutcomeRoutingTodos, + OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + OUTCOME_ROUTING_STATE_STORE_REQUEST_SCHEMA, + outcomeRoutingStatePath, + loadOutcomeRoutingState, + planOutcomeRoutingNextCycle, + reconcileOutcomeRoutingState, + writeOutcomeRoutingState, +} from "../../loopx/control_plane/work_items/outcome_routing_state.ts"; + +function state(direction = "Improve durable customer value.") { + return { + schema_version: "outcome_routing_plan_request_v0", + direction, + cycle: 1, + outcomes: [{ + outcome_id: "outcome_activation", + title: "Improve activation", + metric: "seven day activation rate", + target: ">= 40%", + evidence_source: "activation analytics", + }], + work_items: [], + feedback: [], + }; +} + +function stateWithWork() { + const value = state(); + return { + ...value, + work_items: [{ + work_item_id: "work_activation", + outcome_id: "outcome_activation", + title: "Ship activation improvement", + acceptance: "validated activation evidence", + authority_tier: "A", + ai_capable: true, + target_key: "activation_delivery", + }], + }; +} + +function request(runtimeRoot: string, extra: Record = {}) { + return { + schema_version: OUTCOME_ROUTING_STATE_STORE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + ...extra, + }; +} + +test("outcome routing state writes atomically and reads back exact revision", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-state-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: state(), + updated_at: "2026-09-17T00:00:00Z", + })); + assert.equal(first.written, true); + const stored = first.state as Record; + assert.match(String(stored.revision), /^[a-f0-9]{64}$/); + + const loaded = await loadOutcomeRoutingState(request(runtimeRoot)); + assert.deepEqual(loaded.state, first.state); + assert.equal(loaded.path, outcomeRoutingStatePath(runtimeRoot, "company-goal")); + + const replay = await writeOutcomeRoutingState(request(runtimeRoot, { + state: state(), + updated_at: "2026-09-17T00:01:00Z", + })); + assert.equal(replay.written, false); + assert.equal(replay.replayed, true); +}); + +test("outcome routing state requires revision matching for updates", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-state-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: state(), + updated_at: "2026-09-17T00:00:00Z", + })); + const stored = first.state as Record; + + await assert.rejects( + writeOutcomeRoutingState(request(runtimeRoot, { + state: state("Changed direction."), + updated_at: "2026-09-17T00:01:00Z", + })), + /expected_revision is required/, + ); + await assert.rejects( + writeOutcomeRoutingState(request(runtimeRoot, { + state: state("Changed direction."), + expected_revision: "0".repeat(64), + updated_at: "2026-09-17T00:01:00Z", + })), + /revision changed/, + ); + const updated = await writeOutcomeRoutingState(request(runtimeRoot, { + state: state("Changed direction."), + expected_revision: stored.revision, + updated_at: "2026-09-17T00:01:00Z", + })); + assert.equal(updated.written, true); + assert.notEqual( + (updated.state as Record).revision, + stored.revision, + ); +}); + +test("Todo bindings are revisioned profile state with exact work identity", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-outcome-bind-")); + t.after(async () => await rm(runtimeRoot, { recursive: true, force: true })); + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: { + ...stateWithWork(), + direction: "Route human work without widening shared Todo identity.", + }, + updated_at: "2026-09-17T00:00:00Z", + })); + const stored = first.state as Record; + const bound = await bindOutcomeRoutingTodos({ + schema_version: "outcome_routing_state_bind_request_v0", + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: stored.revision, + updated_at: "2026-09-17T00:01:00Z", + todo_bindings: [{ + work_item_id: "work_activation", + target_key: "activation_delivery", + todo_id: "todo_human_decision", + role: "user", + }], + }); + assert.notEqual((bound.state as Record).revision, stored.revision); + assert.deepEqual((bound.state as Record).todo_bindings, [{ + work_item_id: "work_activation", + target_key: "activation_delivery", + todo_id: "todo_human_decision", + role: "user", + }]); + await assert.rejects( + bindOutcomeRoutingTodos({ + schema_version: "outcome_routing_state_bind_request_v0", + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: (bound.state as Record).revision, + updated_at: "2026-09-17T00:02:00Z", + todo_bindings: [{ + work_item_id: "work_activation", + target_key: "wrong_target", + todo_id: "todo_human_decision", + role: "user", + }], + }), + /must match a projected work item and target/, + ); +}); + +test("outcome routing state path is bounded and rejects relative runtime roots", () => { + const left = outcomeRoutingStatePath("/runtime", "company goal"); + const right = outcomeRoutingStatePath("/runtime", "company-goal"); + assert.notEqual(left, right); + assert.match(left, /outcome-routing\/state\.json$/); + assert.throws( + () => outcomeRoutingStatePath("relative", "company-goal"), + /runtime_root must be absolute/, + ); +}); + +test("outcome routing reconciliation previews and persists evidence-gated Todo status", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-state-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: stateWithWork(), + updated_at: "2026-09-17T00:00:00Z", + })); + const original = first.state as Record; + const reconcileRequest = { + schema_version: OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: original.revision, + updated_at: "2026-09-17T00:01:00Z", + execute: false, + observations: [{ + target_key: "activation_delivery", + todo_id: "todo_activation", + status: "done", + evidence_ref: "artifact:activation-report", + }], + }; + const preview = await reconcileOutcomeRoutingState(reconcileRequest); + assert.equal(preview.dry_run, true); + assert.equal(preview.written, false); + assert.equal( + ((preview.state as Record).reconciliation.observations[0]).next_status, + "done", + ); + assert.deepEqual((await loadOutcomeRoutingState(request(runtimeRoot))).state, first.state); + + const written = await reconcileOutcomeRoutingState({ + ...reconcileRequest, + execute: true, + }); + assert.equal(written.written, true); + const reconciled = written.state as Record; + assert.notEqual(reconciled.revision, original.revision); + assert.equal(reconciled.reconciliation.replan_required, false); + assert.equal(reconciled.reconciliation.observations[0].evidence_ref, "artifact:activation-report"); + + const replay = await reconcileOutcomeRoutingState({ + ...reconcileRequest, + expected_revision: reconciled.revision, + updated_at: "2026-09-17T00:02:00Z", + execute: true, + }); + assert.equal(replay.written, false); + assert.equal(replay.replayed, true); +}); + +test("outcome routing reconciliation requests replanning for blocked or unproven completion", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-state-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: stateWithWork(), + updated_at: "2026-09-17T00:00:00Z", + })); + const revision = (first.state as Record).revision; + + for (const [status, expected] of [ + ["blocked", "replanning"], + ["done", "awaiting_evidence"], + ] as const) { + const result = await reconcileOutcomeRoutingState({ + schema_version: OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: revision, + updated_at: "2026-09-17T00:01:00Z", + execute: false, + observations: [{ + target_key: "activation_delivery", + todo_id: "todo_activation", + status, + }], + }); + const reconciliation = (result.state as Record).reconciliation; + assert.equal(reconciliation.replan_required, true); + assert.equal(reconciliation.observations[0].next_status, expected); + } +}); + +test("outcome routing reconciliation rejects stale revisions and unknown targets", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-state-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: stateWithWork(), + updated_at: "2026-09-17T00:00:00Z", + })); + const revision = (first.state as Record).revision; + const base = { + schema_version: OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: revision, + updated_at: "2026-09-17T00:01:00Z", + execute: false, + }; + await assert.rejects( + reconcileOutcomeRoutingState({ + ...base, + expected_revision: "0".repeat(64), + observations: [], + }), + /revision changed/, + ); + await assert.rejects( + reconcileOutcomeRoutingState({ + ...base, + observations: [{ + target_key: "unknown_target", + todo_id: "todo_unknown", + status: "open", + }], + }), + /is unknown/, + ); +}); + +test("next company cycle converts evidence and blockers into feedback and replanning", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-state-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + const input = stateWithWork(); + input.work_items.push({ + work_item_id: "work_retention", + outcome_id: "outcome_activation", + title: "Resolve retention risk", + acceptance: "risk is cleared", + authority_tier: "A", + ai_capable: true, + target_key: "retention_risk", + }); + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: input, + updated_at: "2026-09-17T00:00:00Z", + })); + const reconciled = await reconcileOutcomeRoutingState({ + schema_version: OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: (first.state as Record).revision, + updated_at: "2026-09-17T00:01:00Z", + execute: true, + observations: [ + { + target_key: "activation_delivery", + todo_id: "todo_activation", + status: "done", + evidence_ref: "artifact:activation-report", + }, + { + target_key: "retention_risk", + todo_id: "todo_retention", + status: "blocked", + }, + ], + }); + const next = planOutcomeRoutingNextCycle({ + schema_version: "outcome_routing_next_cycle_request_v0", + goal_id: "company-goal", + state: reconciled.state, + }); + assert.equal(next.goal_converged, false); + assert.equal(next.converged_work_item_count, 1); + assert.equal(next.remaining_work_item_count, 1); + assert.equal(next.replan_required, true); + const nextState = next.state as Record; + assert.equal(nextState.cycle, 2); + assert.deepEqual( + nextState.work_items.map((item: Record) => item.work_item_id), + ["work_retention"], + ); + assert.deepEqual( + nextState.feedback.map((item: Record) => item.kind), + ["execution_result", "risk"], + ); +}); + +test("next company cycle reports goal convergence after all work has evidence", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-state-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: stateWithWork(), + updated_at: "2026-09-17T00:00:00Z", + })); + const reconciled = await reconcileOutcomeRoutingState({ + schema_version: OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: (first.state as Record).revision, + updated_at: "2026-09-17T00:01:00Z", + execute: true, + observations: [{ + target_key: "activation_delivery", + todo_id: "todo_activation", + status: "done", + evidence_ref: "artifact:activation-report", + }], + }); + const next = planOutcomeRoutingNextCycle({ + schema_version: "outcome_routing_next_cycle_request_v0", + goal_id: "company-goal", + state: reconciled.state, + }); + assert.equal(next.goal_converged, true); + assert.equal(next.remaining_work_item_count, 0); + assert.equal(next.replan_required, true); +}); + +test("next company cycle keeps derived feedback ids valid for maximum-length work ids", async (t) => { + const runtimeRoot = await mkdtemp(join(tmpdir(), "loopx-company-state-")); + t.after(() => rm(runtimeRoot, { recursive: true, force: true })); + const input = stateWithWork(); + input.work_items[0].work_item_id = `w${"a".repeat(127)}`; + const first = await writeOutcomeRoutingState(request(runtimeRoot, { + state: input, + updated_at: "2026-09-17T00:00:00Z", + })); + const reconciled = await reconcileOutcomeRoutingState({ + schema_version: OUTCOME_ROUTING_STATE_RECONCILE_REQUEST_SCHEMA, + runtime_root: runtimeRoot, + goal_id: "company-goal", + expected_revision: (first.state as Record).revision, + updated_at: "2026-09-17T00:01:00Z", + execute: true, + observations: [{ + target_key: "activation_delivery", + todo_id: "todo_activation", + status: "done", + evidence_ref: "artifact:activation-report", + }], + }); + + const next = planOutcomeRoutingNextCycle({ + schema_version: "outcome_routing_next_cycle_request_v0", + goal_id: "company-goal", + state: reconciled.state, + }); + const feedback = (next.state as Record).feedback[0]; + assert.match(feedback.feedback_id, /^todo_feedback_[a-f0-9]{24}$/); + assert.ok(feedback.feedback_id.length <= 128); + assert.deepEqual( + planOutcomeRoutingNextCycle({ + schema_version: "outcome_routing_next_cycle_request_v0", + goal_id: "company-goal", + state: reconciled.state, + }).state, + next.state, + ); +}); diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index bcf19da66a..7e56caabab 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -70,6 +70,8 @@ "loopx/control_plane/turn_driver/turn_journal_effects.ts", "loopx/control_plane/turn_driver/delivery_continuity.ts", "loopx/control_plane/work_items/delivery_outcome.ts", + "loopx/control_plane/work_items/company_control_loop.ts", + "loopx/control_plane/work_items/company_control_state.ts", "loopx/control_plane/work_items/interaction_contract.ts", "loopx/control_plane/work_items/task_lease_acquire.ts", "loopx/control_plane/work_items/task_lease_lifecycle.ts", @@ -113,6 +115,8 @@ "tests/control_plane_ts/postgresql_authority_service_fixture.ts", "tests/control_plane_ts/delivery_continuity.test.ts", "tests/control_plane_ts/delivery_history.test.ts", + "tests/control_plane_ts/company_control_loop.test.ts", + "tests/control_plane_ts/company_control_state.test.ts", "tests/control_plane_ts/delivery_workspace.test.ts", "tests/control_plane_ts/settlement_workspace_causality.test.ts", "tests/control_plane_ts/quota_settlement_readback.test.ts",