From deee231c8664f8d2cb96e4339eb06e52667cf11e Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sun, 27 Sep 2026 01:46:41 +0800 Subject: [PATCH] feat(runtime): fence host state by goal instance Signed-off-by: duanjialing.777 --- ...nstance-identity-and-orphan-recovery-v0.md | 15 + ...e-identity-and-orphan-recovery-v0.zh-CN.md | 13 + loopx/cli_commands/turn.py | 75 ++++- .../control_plane/effect_runtime_handlers.ts | 2 + .../goals/first_party_host_admission.py | 212 ++++++++++++++ .../goals/first_party_host_runtime.ts | 171 +++++++++++ loopx/control_plane/turn_driver/codex_cli.py | 133 +++++++-- loopx/control_plane/turn_driver/driver.py | 4 + loopx/control_plane/turn_driver/executor.py | 37 ++- .../control_plane/turn_driver/transaction.py | 5 + loopx/dsh_goal_mode/turn_host_adapter.py | 6 + loopx/kunluncode_goal_mode/runtime.py | 101 +++++-- .../goal_instance_binding_inventory_v1.json | 12 +- .../project_registry_io_manifest_v1.json | 24 ++ .../test_goal_instance_binding_inventory.py | 7 + .../test_source_session_registry_denial.py | 7 +- .../test_first_party_host_runtime.py | 268 ++++++++++++++++++ .../first_party_host_runtime.test.ts | 168 +++++++++++ tests/test_dsh_goal_mode.py | 25 ++ tests/test_kunluncode_goal_mode.py | 134 +++++++++ tests/test_loopx_turn_codex_cli.py | 122 ++++++++ tests/test_loopx_turn_driver.py | 38 +++ tests/test_loopx_turn_executor.py | 220 ++++++++++++++ 23 files changed, 1736 insertions(+), 63 deletions(-) create mode 100644 loopx/control_plane/goals/first_party_host_admission.py create mode 100644 loopx/control_plane/goals/first_party_host_runtime.ts create mode 100644 tests/control_plane/test_first_party_host_runtime.py create mode 100644 tests/control_plane_ts/first_party_host_runtime.test.ts diff --git a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md index 89e77cb82d..12b9af5f4e 100644 --- a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md +++ b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md @@ -769,6 +769,21 @@ promotion retain their own acceptance. No new paid cohort or soak is authorized. qualify the remaining effect owners before existing-project activation or global routing can open. +### 2026-09-27: first-party Host runtime partial enforcement + +- **Baseline:** `fd96e5e2574272262b9ea604a96581a0d20e94d1` +- **Delivered:** A TypeScript-owned exact GoalRef decision and alias-scoped + lifecycle guard for source-profile Turn journals, Codex descriptors, DSH + session identity, and the Kunlun native runtime journal. +- **Evidence:** Negative tests cover Goal A results returning after same-alias + Goal B publication, cached Turn-result recovery, legacy Host state, + cross-instance session selection, and serialized result/recreation commits. + Non-source plans, paths, schemas, and persisted bytes retain legacy behavior. +- **Remaining hold:** This is partial M3 enforcement. Accepted-before-retirement + downstream drain, unsupported/warm binaries, and the remaining inventory + owners are not qualified. `execution_authority: false` and the M3 activation + hold remain unchanged. + ## Appendix B: Decision log | Date | Decision | Owner / approval | Alternatives | Normative sections changed | diff --git a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md index a767398c5a..68c47ed23c 100644 --- a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md +++ b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md @@ -698,6 +698,19 @@ service adoption、D1–D3 provider promotion 保留各自验收。不授权付 - **剩余 hold:** 所有结果均为 `execution_authority: false`。M3 必须先完成其余 effect owner 资格化,才能开放既有项目 activation 或 global routing。 +### 2026-09-27:第一方 Host runtime 部分 enforcement + +- **基线:** `fd96e5e2574272262b9ea604a96581a0d20e94d1` +- **已交付:** 为 source profile 的 Turn journal、Codex descriptor、DSH session + identity 和 Kunlun native runtime journal 增加 TypeScript-owned exact GoalRef + 决策与 alias-scoped lifecycle guard。 +- **证据:** 负向测试覆盖同名 Goal B 发布后 Goal A 结果迟到、缓存 Turn result + 恢复、legacy Host state、跨实例 session selection,以及 result/recreation + commit 串行化。非 source plan、路径、schema 与持久化字节保持 legacy 行为。 +- **剩余 hold:** 这只是 M3 的部分 enforcement。accepted-before-retirement 的 + downstream drain、不支持的旧/常驻二进制和其余 inventory owner 尚未 + qualified;`execution_authority: false` 与 M3 activation hold 保持不变。 + ## 附录 B:决策日志 | 日期 | 决策 | Owner/批准 | 替代方案 | 变更的规范章节 | diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index 23eef2a5f8..ef9ba12819 100644 --- a/loopx/cli_commands/turn.py +++ b/loopx/cli_commands/turn.py @@ -21,6 +21,10 @@ from ..capabilities.periodic_report.cadence_runtime import extend_cadence_turn_start_dispatch from ..control_plane.quota.live_decision import build_live_quota_should_run_decision from ..control_plane.agents.workspace_guard import capture_delivery_workspace +from ..control_plane.goals.first_party_host_admission import ( + FirstPartyHostGoalAdmission, + capture_first_party_host_goal_ref, +) from ..control_plane.quota.heartbeat_receipt import ( ensure_turn_heartbeat_settlement_receipt, ) @@ -124,6 +128,18 @@ def handle_turn_command( registry_path=registry_path, runtime_root_override=runtime_root_arg, ) + goal_ref = capture_first_party_host_goal_ref( + registry_path=registry_path, + goal_id=args.goal_id, + ) + goal_admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry_path, + goal_id=args.goal_id, + planned_goal_ref=goal_ref, + ) + strict_goal_admission = ( + goal_admission if goal_admission.enabled else None + ) # Planning and dry-run execution inspect existing admitted intents. # Only an executing wake may sync inboxes or reserve a calendar window. turn_start_hook_dispatch = {} @@ -175,7 +191,15 @@ def handle_turn_command( and not args.resume_turn_key and turn_envelope.get("effective_action") != EffectiveAction.GOVERNED_CAPABILITY_INTENT.value ): - session_binding = codex_cli_session_binding(runtime_root, turn_envelope) + session_binding = ( + codex_cli_session_binding( + runtime_root, + turn_envelope, + goal_admission=strict_goal_admission, + ) + if strict_goal_admission is not None + else codex_cli_session_binding(runtime_root, turn_envelope) + ) payload = build_loopx_turn_plan( turn_envelope, host=args.host, @@ -184,6 +208,7 @@ def handle_turn_command( session_binding=session_binding, turn_instance_id=args.turn_instance_id, iteration_context_policy=args.iteration_context.replace("-", "_"), + goal_ref=goal_ref, ) # The executor readback names where this Turn's model work runs and # whether that host can launch here, so a caller never has to infer it @@ -304,6 +329,16 @@ def handle_turn_command( raise ValueError( "LoopX Turn resume journal belongs to another agent" ) + goal_admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry_path, + goal_id=args.goal_id, + planned_goal_ref=payload.get("goal_ref"), + ) + strict_goal_admission = ( + goal_admission if goal_admission.enabled else None + ) + if strict_goal_admission is not None: + strict_goal_admission.require_current() if payload.get("route", {}).get("kind") == "capability_action_required": # The normal host transaction forbids Core mutations. A # capability may prepare artifacts and require authored input; @@ -994,24 +1029,37 @@ def scheduler(_spend_payload: dict[str, object]) -> dict[str, object]: def run_built_in_host( request: Mapping[str, Any], ) -> dict[str, Any]: - return run_codex_cli_host( - request, - runtime_root=runtime_root, - project=project, - codex_bin=args.codex_bin, - sandbox=args.codex_sandbox, - model=args.codex_model, - reasoning_effort=args.codex_reasoning_effort, - mcp_server=args.codex_mcp_server_json, - timeout_seconds=max(1.0, args.timeout_seconds - 5.0), - ) + options = { + "runtime_root": runtime_root, + "project": project, + "codex_bin": args.codex_bin, + "sandbox": args.codex_sandbox, + "model": args.codex_model, + "reasoning_effort": args.codex_reasoning_effort, + "mcp_server": args.codex_mcp_server_json, + "timeout_seconds": max(1.0, args.timeout_seconds - 5.0), + } + if strict_goal_admission is not None: + options["goal_admission"] = strict_goal_admission + return run_codex_cli_host(request, **options) host_runner = run_built_in_host def resolve_built_in_session_binding( turn_envelope: Mapping[str, Any], ) -> dict[str, str] | None: - return codex_cli_session_binding(runtime_root, turn_envelope) + return ( + codex_cli_session_binding( + runtime_root, + turn_envelope, + goal_admission=strict_goal_admission, + ) + if strict_goal_admission is not None + else codex_cli_session_binding( + runtime_root, + turn_envelope, + ) + ) session_binding_resolver = resolve_built_in_session_binding elif args.host == "dsh": @@ -1086,6 +1134,7 @@ def on_managed_start_admitted() -> None: ), admit_start=managed_cadence.admit if args.execute else None, confirm_start=managed_cadence.confirm if args.execute else None, + goal_admission=strict_goal_admission, ) else: raise ValueError("turn requires the `plan` or `run-once` subcommand") diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index a774751aff..5299948af5 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -129,6 +129,7 @@ import { decideProjectSessionBind, decideProjectSessionUnbind, } from "./goals/source_session_lifetime.ts"; +import { decideFirstPartyHostRuntime } from "./goals/first_party_host_runtime.ts"; import { evaluateDeliveryRoute, } from "./turn_driver/delivery_continuity.ts"; @@ -532,6 +533,7 @@ export function createEffectRuntimeHandlers( ["goal.source_session.bind.decide", decideProjectSessionBind], ["goal.source_session.unbind.decide", decideProjectSessionUnbind], ["goal.source_session.recreate.decide", decideGoalRecreation], + ["goal.first_party_host_runtime.decide", decideFirstPartyHostRuntime], ["goal.acceptance.inspect", inspectLocalGoalAcceptance], ["goal.acceptance.configure", commitLocalGoalAcceptance], ["goal.acceptance.verify.commit", commitLocalGoalAcceptanceVerification], diff --git a/loopx/control_plane/goals/first_party_host_admission.py b/loopx/control_plane/goals/first_party_host_admission.py new file mode 100644 index 0000000000..0d2b45e034 --- /dev/null +++ b/loopx/control_plane/goals/first_party_host_admission.py @@ -0,0 +1,212 @@ +from __future__ import annotations + +from collections.abc import Callable, Mapping +from dataclasses import dataclass +from pathlib import Path +from typing import Any, TypeVar + +from ...file_lock import exclusive_cross_runtime_file_lock +from ..effect_runtime import effect_runtime_result +from ..projects.registry_codec import ( + SOURCE_SESSION_PROFILE_ID, + load_project_registry, +) +from .source_session_registry_state import ( + exact_goal_ref, + guard_path, + require_goal_id, +) + + +T = TypeVar("T") + + +class FirstPartyHostRuntimeRejected(RuntimeError): + """Refusal to reuse Host state or accept a stale Host result.""" + + def __init__(self, code: str) -> None: + super().__init__(f"first-party Host runtime rejected: {code}") + self.code = code + + +def _source_authority(registry_path: Path, goal_id: str) -> dict[str, Any]: + if not registry_path.is_file(): + return {"kind": "unavailable", "reason": "registry_missing"} + try: + registry = load_project_registry(registry_path) + except (OSError, TypeError, ValueError): + return {"kind": "unavailable", "reason": "registry_unreadable"} + if registry.get("profile_id") != SOURCE_SESSION_PROFILE_ID: + return {"kind": "unavailable", "reason": "registry_unreadable"} + goals = registry.get("goals") + if not isinstance(goals, list) or any( + not isinstance(item, Mapping) for item in goals + ): + return {"kind": "unavailable", "reason": "registry_unreadable"} + matches = [ + item + for item in goals + if item.get("id") == goal_id and item.get("status") == "active" + ] + if len(matches) != 1: + return {"kind": "absent"} + goal = matches[0] + goal_ref = {"goal_id": goal.get("id")} + if "goal_instance_id" in goal: + goal_ref["goal_instance_id"] = goal.get("goal_instance_id") + return {"kind": "present", "goal_ref": goal_ref} + + +def capture_first_party_host_goal_ref( + *, + registry_path: Path, + goal_id: str, +) -> dict[str, str] | None: + """Capture the current exact GoalRef, or preserve a legacy registry.""" + + requested_registry = registry_path.expanduser().resolve() + require_goal_id(goal_id) + if not requested_registry.is_file(): + return None + registry = load_project_registry(requested_registry) + if registry.get("profile_id") != SOURCE_SESSION_PROFILE_ID: + return None + with exclusive_cross_runtime_file_lock( + guard_path(requested_registry, goal_id), + operation="first_party_host_goal_capture", + ): + authority = _source_authority(requested_registry, goal_id) + if authority.get("kind") != "present": + raise FirstPartyHostRuntimeRejected( + "goal_not_registered" + if authority.get("kind") == "absent" + else "goal_authority_unavailable" + ) + goal_ref = authority.get("goal_ref") + if ( + not isinstance(goal_ref, Mapping) + or not isinstance(goal_ref.get("goal_id"), str) + or not isinstance(goal_ref.get("goal_instance_id"), str) + ): + raise FirstPartyHostRuntimeRejected("goal_instance_id_missing") + return exact_goal_ref( + goal_ref["goal_id"], + goal_ref["goal_instance_id"], + ) + + +@dataclass(frozen=True, slots=True) +class FirstPartyHostGoalAdmission: + """Keep exact Host-state I/O inside one source Goal lifetime guard.""" + + registry_path: Path + goal_id: str + planned_goal_ref: object | None + source_profile: bool + + @classmethod + def for_plan( + cls, + *, + registry_path: Path, + goal_id: str, + planned_goal_ref: object | None, + ) -> FirstPartyHostGoalAdmission: + requested_registry = registry_path.expanduser().resolve() + require_goal_id(goal_id) + source_profile = planned_goal_ref is not None + if not source_profile and requested_registry.is_file(): + source_profile = ( + load_project_registry(requested_registry).get("profile_id") + == SOURCE_SESSION_PROFILE_ID + ) + return cls( + registry_path=requested_registry, + goal_id=goal_id, + planned_goal_ref=planned_goal_ref, + source_profile=source_profile, + ) + + @property + def enabled(self) -> bool: + return self.source_profile + + def _decision( + self, + operation: str, + *, + host_state: Mapping[str, Any] | None = None, + ) -> dict[str, Any]: + decision = effect_runtime_result( + "goal.first_party_host_runtime.decide", + { + "profile_id": ( + SOURCE_SESSION_PROFILE_ID if self.source_profile else None + ), + "operation": operation, + "planned_goal_ref": self.planned_goal_ref, + "authority": _source_authority( + self.registry_path, + self.goal_id, + ), + **({"host_state": host_state} if host_state is not None else {}), + }, + ) + if not isinstance(decision, dict): + raise RuntimeError("first-party Host runtime decision must be an object") + if decision.get("kind") == "reject": + raise FirstPartyHostRuntimeRejected(str(decision.get("code") or "invalid")) + return decision + + def select_state( + self, + *, + read_state: Callable[[], T | None], + goal_ref_of: Callable[[T], object], + initialize_state: Callable[[dict[str, str]], T] | None = None, + ) -> T | None: + if not self.source_profile: + state = read_state() + return state + with exclusive_cross_runtime_file_lock( + guard_path(self.registry_path, self.goal_id), + operation="first_party_host_state_select", + ): + state = read_state() + host_state = ( + {"kind": "absent"} + if state is None + else {"kind": "present", "goal_ref": goal_ref_of(state)} + ) + decision = self._decision("select_state", host_state=host_state) + if decision.get("kind") == "start_new" and initialize_state is not None: + goal_ref = decision.get("goal_ref") + if not isinstance(goal_ref, dict): + raise RuntimeError( + "first-party Host start decision omitted goal_ref" + ) + return initialize_state(goal_ref) + return state + + def require_current(self) -> None: + if not self.source_profile: + return + with exclusive_cross_runtime_file_lock( + guard_path(self.registry_path, self.goal_id), + operation="first_party_host_current_check", + ): + self._decision("require_current") + + def accept_result(self, commit_result: Callable[[], T]) -> T: + if not self.source_profile: + return commit_result() + with exclusive_cross_runtime_file_lock( + guard_path(self.registry_path, self.goal_id), + operation="first_party_host_result_accept", + ): + decision = self._decision("accept_result") + if decision.get("kind") != "accept_result": + raise RuntimeError( + "first-party Host result decision did not accept a result" + ) + return commit_result() diff --git a/loopx/control_plane/goals/first_party_host_runtime.ts b/loopx/control_plane/goals/first_party_host_runtime.ts new file mode 100644 index 0000000000..d1e6fc2af7 --- /dev/null +++ b/loopx/control_plane/goals/first_party_host_runtime.ts @@ -0,0 +1,171 @@ +import type { JsonObject } from "../effect_program.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { assertNever, jsonObject } from "../runtime_decode.ts"; +import { + parseExactGoalRef, + type ExactGoalRef, +} from "./goal_instance_identity.ts"; + +const SOURCE_SESSION_PROFILE_ID = "source_session_v1"; + +type WireGoalRef = Readonly<{ + goal_id: string; + goal_instance_id: string; +}>; + +type Operation = "select_state" | "require_current" | "accept_result"; + +type RejectionCode = + | "goal_not_registered" + | "goal_authority_unavailable" + | "goal_instance_id_missing" + | "stale_goal_instance" + | "legacy_host_state" + | "host_state_goal_instance_mismatch" + | "host_state_malformed"; + +export type FirstPartyHostRuntimeDecision = + | Readonly<{ kind: "legacy" }> + | Readonly<{ kind: "start_new"; goal_ref: WireGoalRef }> + | Readonly<{ kind: "resume"; goal_ref: WireGoalRef }> + | Readonly<{ kind: "accept_result"; goal_ref: WireGoalRef }> + | Readonly<{ kind: "reject"; code: RejectionCode }>; + +type Authority = + | Readonly<{ kind: "present"; goalRef: ExactGoalRef }> + | Readonly<{ kind: "absent" }> + | Readonly<{ kind: "unavailable" }> + | Readonly<{ kind: "missing_instance" }>; + +function requiredOperation(value: unknown): Operation { + if ( + value === "select_state" + || value === "require_current" + || value === "accept_result" + ) { + return value; + } + throw new EffectRuntimeRequestError( + "first-party host runtime operation is unsupported", + ); +} + +function wireGoalRef(value: ExactGoalRef): WireGoalRef { + return { + goal_id: value.goalId.value, + goal_instance_id: value.goalInstanceId.value, + }; +} + +function sameGoalRef(left: ExactGoalRef, right: ExactGoalRef): boolean { + return left.goalId.value === right.goalId.value + && left.goalInstanceId.value === right.goalInstanceId.value; +} + +function authority(value: unknown): Authority { + const raw = jsonObject(value); + if (!raw) return { kind: "unavailable" }; + switch (raw.kind) { + case "absent": + return { kind: "absent" }; + case "unavailable": + return { kind: "unavailable" }; + case "present": { + const parsed = parseExactGoalRef(raw.goal_ref); + if (parsed.kind === "parsed") { + return { kind: "present", goalRef: parsed.value }; + } + return parsed.issue === "missing_goal_instance_id" + ? { kind: "missing_instance" } + : { kind: "unavailable" }; + } + default: + return { kind: "unavailable" }; + } +} + +function selectHostState( + value: unknown, + plannedGoalRef: ExactGoalRef, +): FirstPartyHostRuntimeDecision { + const state = jsonObject(value); + if (!state) { + return { kind: "reject", code: "host_state_malformed" }; + } + if (state.kind === "absent") { + return { kind: "start_new", goal_ref: wireGoalRef(plannedGoalRef) }; + } + if (state.kind !== "present") { + return { kind: "reject", code: "host_state_malformed" }; + } + const parsed = parseExactGoalRef(state.goal_ref); + if (parsed.kind === "invalid") { + return { + kind: "reject", + code: parsed.issue === "missing_goal_instance_id" + ? "legacy_host_state" + : "host_state_malformed", + }; + } + if (!sameGoalRef(parsed.value, plannedGoalRef)) { + return { + kind: "reject", + code: "host_state_goal_instance_mismatch", + }; + } + return { kind: "resume", goal_ref: wireGoalRef(plannedGoalRef) }; +} + +/** + * Decide whether one first-party Host may select state or accept a result. + * + * Non-source profiles retain their historical behavior. The source-session + * profile is exact-only: neither an alias-only plan nor alias-only Host state + * can be upgraded by inference. + */ +export function decideFirstPartyHostRuntime( + value: unknown, +): FirstPartyHostRuntimeDecision { + const facts = jsonObject(value); + if (!facts) { + throw new EffectRuntimeRequestError( + "first-party host runtime facts must be an object", + ); + } + if (facts.profile_id !== SOURCE_SESSION_PROFILE_ID) { + return { kind: "legacy" }; + } + + const operation = requiredOperation(facts.operation); + const planned = parseExactGoalRef(facts.planned_goal_ref); + if (planned.kind === "invalid") { + return { kind: "reject", code: "goal_instance_id_missing" }; + } + const current = authority(facts.authority); + switch (current.kind) { + case "absent": + return { kind: "reject", code: "goal_not_registered" }; + case "unavailable": + return { kind: "reject", code: "goal_authority_unavailable" }; + case "missing_instance": + return { kind: "reject", code: "goal_instance_id_missing" }; + case "present": + if (!sameGoalRef(planned.value, current.goalRef)) { + return { kind: "reject", code: "stale_goal_instance" }; + } + break; + default: + return assertNever(current, "unsupported Goal authority"); + } + + switch (operation) { + case "select_state": + return selectHostState(facts.host_state, planned.value); + case "require_current": + return { kind: "resume", goal_ref: wireGoalRef(planned.value) }; + case "accept_result": + return { kind: "accept_result", goal_ref: wireGoalRef(planned.value) }; + default: + return assertNever(operation, "unsupported Host runtime operation"); + } +} diff --git a/loopx/control_plane/turn_driver/codex_cli.py b/loopx/control_plane/turn_driver/codex_cli.py index 4b95984412..dcc15f0af5 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -16,6 +16,7 @@ from typing import Any from ...runtime import validate_goal_id_path_segment +from ..goals.first_party_host_admission import FirstPartyHostGoalAdmission from .subagent_execution_topology import ( child_execution_receipts_json_schema, ) @@ -238,13 +239,61 @@ def load_codex_cli_session( return {**value, "session_id": session_id} +def _codex_session_goal_ref( + value: Mapping[str, Any], + *, + lineage: Mapping[str, str], +) -> object: + if ( + value.get("schema_version") != CODEX_CLI_SESSION_SCHEMA_VERSION + or any(value.get(field) != lineage[field] for field in lineage) + or _valid_session_id(value.get("session_id")) is None + ): + return {"malformed": True} + goal_ref = value.get("goal_ref") + if goal_ref is not None: + return goal_ref + return {"goal_id": value.get("goal_id")} + + +def _read_codex_cli_session_document( + runtime_root: Path, + *, + lineage: Mapping[str, str], +) -> dict[str, Any] | None: + path = _session_path(runtime_root, lineage) + if not path.exists(): + return None + try: + value = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + return {"malformed": True} + return value if isinstance(value, dict) else {"malformed": True} + + def codex_cli_session_binding( runtime_root: Path, turn_envelope: Mapping[str, Any], + *, + goal_admission: FirstPartyHostGoalAdmission | None = None, ) -> dict[str, str] | None: request = {"turn_envelope": dict(turn_envelope)} lineage = _lineage(request) - if load_codex_cli_session(runtime_root, lineage=lineage) is None: + if goal_admission is None: + session = load_codex_cli_session(runtime_root, lineage=lineage) + else: + selected = goal_admission.select_state( + read_state=lambda: _read_codex_cli_session_document( + runtime_root, + lineage=lineage, + ), + goal_ref_of=lambda value: _codex_session_goal_ref( + value, + lineage=lineage, + ), + ) + session = dict(selected) if selected is not None else None + if session is None: return None return { "schema_version": "loopx_turn_session_binding_v0", @@ -257,6 +306,7 @@ def _store_codex_cli_session( *, lineage: Mapping[str, str], session_id: str, + goal_ref: Mapping[str, Any] | None = None, ) -> None: normalized_session_id = _valid_session_id(session_id) if not normalized_session_id: @@ -273,13 +323,16 @@ def _store_codex_cli_session( handle = os.fdopen(descriptor, "w", encoding="utf-8") descriptor = -1 with handle: + payload = { + "schema_version": CODEX_CLI_SESSION_SCHEMA_VERSION, + **lineage, + "host": "codex-cli", + "session_id": normalized_session_id, + } + if goal_ref is not None: + payload["goal_ref"] = dict(goal_ref) json.dump( - { - "schema_version": CODEX_CLI_SESSION_SCHEMA_VERSION, - **lineage, - "host": "codex-cli", - "session_id": normalized_session_id, - }, + payload, handle, ensure_ascii=False, indent=2, @@ -787,6 +840,7 @@ def run_codex_cli_host( reasoning_effort: str | None = None, mcp_server: Mapping[str, Any] | None = None, timeout_seconds: float = 115.0, + goal_admission: FirstPartyHostGoalAdmission | None = None, ) -> dict[str, Any]: if request.get("schema_version") != LOOPX_TURN_HOST_REQUEST_SCHEMA_VERSION: raise ValueError("unsupported LoopX Turn host request schema") @@ -803,11 +857,24 @@ def run_codex_cli_host( planned_action = str(planned_session.get("action") or "") context_policy = _mapping(planned_session.get("context_policy")) fresh_iteration = context_policy.get("mode") == "fresh" - binding = ( - None - if fresh_iteration - else load_codex_cli_session(runtime_root, lineage=lineage) - ) + if goal_admission is None: + binding = ( + None + if fresh_iteration + else load_codex_cli_session(runtime_root, lineage=lineage) + ) + else: + selected = goal_admission.select_state( + read_state=lambda: _read_codex_cli_session_document( + runtime_root, + lineage=lineage, + ), + goal_ref_of=lambda value: _codex_session_goal_ref( + value, + lineage=lineage, + ), + ) + binding = None if fresh_iteration else selected if planned_action == "resume" and binding is None: raise RuntimeError("Codex CLI resume binding disappeared after planning") if planned_action == "start_new" and binding is not None: @@ -815,6 +882,34 @@ def run_codex_cli_host( if planned_action not in {"resume", "start_new"}: raise ValueError("Codex CLI host request has no executable session action") session_id = str(binding.get("session_id")) if binding else None + goal_ref = request.get("goal_ref") + exact_goal_ref = dict(goal_ref) if isinstance(goal_ref, Mapping) else None + + def store_session(observed_session_id: str) -> None: + def commit() -> None: + _store_codex_cli_session( + runtime_root, + lineage=lineage, + session_id=observed_session_id, + goal_ref=exact_goal_ref, + ) + + if goal_admission is None: + commit() + else: + goal_admission.accept_result(commit) + + def discard_session() -> None: + def commit() -> None: + _discard_codex_cli_session( + runtime_root, + lineage=lineage, + ) + + if goal_admission is None: + commit() + else: + goal_admission.accept_result(commit) with tempfile.TemporaryDirectory(prefix="loopx-turn-codex-") as directory: temporary = Path(directory) @@ -902,11 +997,7 @@ def discard_stderr() -> None: output_observation_incomplete = reader.is_alive() or stderr_reader.is_alive() if timed_out: if observed_session: - _store_codex_cli_session( - runtime_root, - lineage=lineage, - session_id=observed_session[0], - ) + store_session(observed_session[0]) raise BuiltInHostError( "codex_cli_timeout", failure_kind="executor_timeout", @@ -922,15 +1013,11 @@ def discard_stderr() -> None: ) ) if returncode != 0 and category in SESSION_INVALIDATING_FAILURE_CATEGORIES: - _discard_codex_cli_session(runtime_root, lineage=lineage) + discard_session() if observed_session and ( returncode == 0 or category not in SESSION_INVALIDATING_FAILURE_CATEGORIES ): - _store_codex_cli_session( - runtime_root, - lineage=lineage, - session_id=observed_session[0], - ) + store_session(observed_session[0]) if returncode != 0: raise BuiltInHostError( f"codex_cli_{category}", diff --git a/loopx/control_plane/turn_driver/driver.py b/loopx/control_plane/turn_driver/driver.py index 0672847c97..a68fb48726 100644 --- a/loopx/control_plane/turn_driver/driver.py +++ b/loopx/control_plane/turn_driver/driver.py @@ -435,6 +435,7 @@ def build_loopx_turn_plan( session_binding: Mapping[str, Any] | None = None, turn_instance_id: str | None = None, iteration_context_policy: str = "resume_if_available", + goal_ref: Mapping[str, str] | None = None, ) -> dict[str, Any]: """Project a TurnEnvelope into a typed, side-effect-free host decision.""" @@ -485,6 +486,7 @@ def build_loopx_turn_plan( scheduler_owner=str(context_projection.get("scheduler_owner") or ""), session_action=str(session.get("action") or "none"), turn_instance_id=turn_instance_id, + goal_ref=goal_ref, ) execution_topology = build_subagent_execution_topology( turn_envelope=envelope, @@ -550,6 +552,8 @@ def build_loopx_turn_plan( payload["delegation_context"] = context if execution_topology: payload["subagent_execution_topology"] = execution_topology + if goal_ref is not None: + payload["goal_ref"] = dict(goal_ref) if route is LoopXTurnRoute.CAPABILITY_ACTION_REQUIRED: payload["capability_action"] = { "status": "required", diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index bb0fbd1c64..a9245556f4 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -16,6 +16,10 @@ interpret_turn_result_packet, settlement_result_payload, ) +from ..goals.first_party_host_admission import ( + FirstPartyHostGoalAdmission, + FirstPartyHostRuntimeRejected, +) from ..goals.goal_vision import normalize_goal_vision_packet from ..work_items.delivery_batch_scale import require_delivery_batch_scale from ..work_items.delivery_outcome import require_delivery_outcome @@ -147,6 +151,9 @@ def build_loopx_turn_host_request(plan: Mapping[str, Any]) -> dict[str, Any]: "stdout": "one public-safe JSON object", }, } + goal_ref = plan.get("goal_ref") + if isinstance(goal_ref, Mapping): + request["goal_ref"] = dict(goal_ref) reward_memory_recall = plan.get("reward_memory_recall") if isinstance(reward_memory_recall, Mapping): request["reward_memory_recall"] = dict(reward_memory_recall) @@ -712,6 +719,8 @@ def _run_host_runner( else {} ), } + except FirstPartyHostRuntimeRejected: + raise except Exception as exc: # noqa: BLE001 - host adapters fail closed at boundary return {"ok": False, "reason": type(exc).__name__, "returncode": None} if not isinstance(value, dict): @@ -757,6 +766,7 @@ def _host_result_stage( confirm_start: Callable[[], None] | None = None, usage_runtime_root: Path | None = None, usage_goal_id: str = "", + goal_admission: FirstPartyHostGoalAdmission | None = None, ) -> tuple[dict[str, Any] | None, list[str], dict[str, Any] | None]: completed_phases = list(journal.get("completed_phases") or []) result = ( @@ -808,7 +818,12 @@ def _host_result_stage( journal["host_recovery"] = build_host_recovery_record(recovery_kind) else: journal.pop("host_recovery", None) - _write_journal(journal_path, journal) + if goal_admission is None: + _write_journal(journal_path, journal) + else: + goal_admission.accept_result( + lambda: _write_journal(journal_path, journal) + ) return ( None, [], @@ -846,7 +861,12 @@ def _host_result_stage( result_kind=LoopXTurnResultKind.VALIDATION_FAILED.value, validation_stage="host_result_contract", ) - _write_journal(journal_path, journal) + if goal_admission is None: + _write_journal(journal_path, journal) + else: + goal_admission.accept_result( + lambda: _write_journal(journal_path, journal) + ) return ( None, list(TRANSACTION_PHASES[:2]), @@ -867,7 +887,12 @@ def _host_result_stage( result_kind=normalized.get("result_kind"), completed_phases=completed_phases, ) - _write_journal(journal_path, journal) + if goal_admission is None: + _write_journal(journal_path, journal) + else: + goal_admission.accept_result( + lambda: _write_journal(journal_path, journal) + ) return normalized, completed_phases, None @@ -1234,6 +1259,7 @@ def run_loopx_turn_once( post_settlement: PostSettlement | None = None, admit_start: Callable[[Mapping[str, Any]], dict[str, Any]] | None = None, confirm_start: Callable[[], None] | None = None, + goal_admission: FirstPartyHostGoalAdmission | None = None, ) -> dict[str, Any]: if host_runner is not None and host_argv is not None: raise ValueError("run-once accepts either host_argv or host_runner, not both") @@ -1293,6 +1319,8 @@ def run_loopx_turn_once( journal = _load_journal(journal_path) recovery_decision: dict[str, Any] | None = None if journal is not None: + if goal_admission is not None: + goal_admission.require_current() envelope = ( plan.get("turn_envelope") if isinstance(plan.get("turn_envelope"), Mapping) @@ -1341,6 +1369,8 @@ def run_loopx_turn_once( needs_host = validation_reinvokes_host or journal is None or "typed_result" not in list( journal.get("completed_phases") or [] ) + if needs_host and goal_admission is not None: + goal_admission.require_current() admission = None if needs_host and admit_start is not None: admission = admit_start({ @@ -1443,6 +1473,7 @@ def finish_recovery(payload: dict[str, Any]) -> dict[str, Any]: if admission is not None and admission.get("reserved") is True else None ), + goal_admission=goal_admission, ) if terminal is not None: return finish_recovery(terminal) diff --git a/loopx/control_plane/turn_driver/transaction.py b/loopx/control_plane/turn_driver/transaction.py index 2079a46a72..2722a3940c 100644 --- a/loopx/control_plane/turn_driver/transaction.py +++ b/loopx/control_plane/turn_driver/transaction.py @@ -126,6 +126,7 @@ def build_loopx_turn_transaction_plan( session_action: str, scheduler_owner: str = "none", turn_instance_id: str | None = None, + goal_ref: Mapping[str, str] | None = None, ) -> dict[str, Any]: normalized_instance_id = normalize_turn_instance_id(turn_instance_id) identity = { @@ -137,6 +138,8 @@ def build_loopx_turn_transaction_plan( } if normalized_instance_id is not None: identity["turn_instance_id"] = normalized_instance_id + if goal_ref is not None: + identity["goal_ref"] = dict(goal_ref) turn_key = _canonical_hash(identity) settlement_identity = SettlementIdentity( goal_id=str(lineage.get("goal_id") or ""), @@ -198,6 +201,8 @@ def build_loopx_turn_transaction_plan( plan["settlement_plan"] = settlement_plan.as_dict() if normalized_instance_id is not None: plan["turn_instance_id"] = normalized_instance_id + if goal_ref is not None: + plan["goal_ref"] = dict(goal_ref) return plan diff --git a/loopx/dsh_goal_mode/turn_host_adapter.py b/loopx/dsh_goal_mode/turn_host_adapter.py index 53cbfabb49..d9672718cd 100644 --- a/loopx/dsh_goal_mode/turn_host_adapter.py +++ b/loopx/dsh_goal_mode/turn_host_adapter.py @@ -422,6 +422,12 @@ def _derive_session_id(request: Mapping[str, Any], turn_key: str) -> str: # values are not lineage identities, while their positions stay encoded. lineage = [value if value else None for value in lineage] if any(lineage): + goal_ref = _mapping(request.get("goal_ref")) + goal_instance_id = goal_ref.get("goal_instance_id") + if goal_instance_id: + return "dsh-lineage-v2-" + _canonical_hash( + [*lineage, goal_instance_id] + ).removeprefix("sha256:") return "dsh-lineage-v1-" + _canonical_hash(lineage).removeprefix("sha256:") return f"dsh-{turn_key.removeprefix('sha256:')[:24]}" diff --git a/loopx/kunluncode_goal_mode/runtime.py b/loopx/kunluncode_goal_mode/runtime.py index c293a21dee..6033124c99 100644 --- a/loopx/kunluncode_goal_mode/runtime.py +++ b/loopx/kunluncode_goal_mode/runtime.py @@ -12,6 +12,10 @@ from typing import Any from loopx.file_lock import exclusive_file_lock +from loopx.control_plane.goals.first_party_host_admission import ( + FirstPartyHostGoalAdmission, + capture_first_party_host_goal_ref, +) from loopx.kunluncode_goal_mode.app_server import ( NATIVE_GOAL_MODES, KunlunAppServerClient, @@ -234,16 +238,20 @@ def _new_state( todo_id: str, mode: str, objective_sha256: str, + goal_ref: Mapping[str, str] | None = None, ) -> dict[str, Any]: timestamp = _now() + binding = { + "goal_id": str(context.get("goal_id") or ""), + "agent_id": str(context.get("agent_id") or ""), + } + if goal_ref is not None: + binding["goal_instance_id"] = goal_ref["goal_instance_id"] return { "schema_version": RUNTIME_STATE_SCHEMA_VERSION, "created_at": timestamp, "updated_at": timestamp, - "binding": { - "goal_id": str(context.get("goal_id") or ""), - "agent_id": str(context.get("agent_id") or ""), - }, + "binding": binding, "todo_id": todo_id, "native": { "mode": mode, @@ -280,6 +288,24 @@ def _validate_binding(state: Mapping[str, Any], context: Mapping[str, Any]) -> N ) +def _runtime_goal_ref(state: Mapping[str, Any]) -> object: + binding = state.get("binding") + if not isinstance(binding, Mapping): + return {"malformed": True} + goal_ref = {"goal_id": binding.get("goal_id")} + if "goal_instance_id" in binding: + goal_ref["goal_instance_id"] = binding.get("goal_instance_id") + return goal_ref + + +def _commit_runtime_state( + project: Path, + state: dict[str, Any], + goal_admission: FirstPartyHostGoalAdmission, +) -> Path: + return goal_admission.accept_result(lambda: write_runtime_state(project, state)) + + def _compact_receipt( payload: Mapping[str, Any], fields: tuple[str, ...] ) -> dict[str, Any]: @@ -290,12 +316,14 @@ def _reconcile_writeback_evidence( project: Path, state: dict[str, Any], control: LoopXControlPlane, + goal_admission: FirstPartyHostGoalAdmission, ) -> None: native = state["native"] since = str(native.get("verified_at") or "") if not since: return todo_id = str(state.get("todo_id") or "") + goal_admission.require_current() payload = control.evidence_since(since, todo_id=todo_id) ledger = payload.get("ledger") if isinstance(payload.get("ledger"), list) else [] writeback = state["writeback"] @@ -315,7 +343,7 @@ def _reconcile_writeback_evidence( writeback["todo_completed"] = True if event_kind in {"quota_spend", "quota_slot_spent"}: writeback["quota_spent"] = True - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) def _verified_evidence(state: Mapping[str, Any]) -> str: @@ -335,11 +363,14 @@ def _closeout_writeback( project: Path, state: dict[str, Any], control: LoopXControlPlane, + goal_admission: FirstPartyHostGoalAdmission, ) -> None: - _reconcile_writeback_evidence(project, state, control) + goal_admission.require_current() + _reconcile_writeback_evidence(project, state, control, goal_admission) writeback = state["writeback"] mode = str(state["native"]["mode"]) if writeback.get("delivery_recorded") is not True: + goal_admission.require_current() payload = control.record_verified_delivery( mode=mode, todo_id=str(state["todo_id"]), @@ -352,8 +383,9 @@ def _closeout_writeback( writeback["delivery_receipt"] = _compact_receipt( payload, ("classification", "generated_at", "appended") ) - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) if writeback.get("todo_completed") is not True: + goal_admission.require_current() payload = control.complete( str(state["todo_id"]), evidence=_verified_evidence(state) ) @@ -365,8 +397,9 @@ def _closeout_writeback( writeback["todo_receipt"] = _compact_receipt( payload, ("todo_id", "status", "status_changed", "changed") ) - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) if writeback.get("quota_spent") is not True: + goal_admission.require_current() payload = control.spend(todo_id=str(state["todo_id"])) if payload.get("appended") is not True: raise KunlunNativeGoalRuntimeError( @@ -377,7 +410,7 @@ def _closeout_writeback( writeback["quota_receipt"] = _compact_receipt( payload, ("classification", "generated_at", "appended", "slots") ) - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) def _update_native_from_goal(state: dict[str, Any], goal: Mapping[str, Any]) -> None: @@ -427,6 +460,7 @@ def _run_native_host( token_budget: int | None, scenario: str | None, client_factory: Callable[..., Any], + goal_admission: FirstPartyHostGoalAdmission, ) -> None: command = build_app_server_command( kunluncode_bin, @@ -449,6 +483,7 @@ def record_event(message: Mapping[str, Any]) -> None: event_counts[method] = event_counts.get(method, 0) + 1 native = state["native"] + goal_admission.require_current() with client_factory( command, cwd=project, @@ -466,7 +501,7 @@ def record_event(message: Mapping[str, Any]) -> None: thread_id = client.start_thread(project) native["thread_id"] = thread_id native["status"] = "thread_started" - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) elif native.get("goal_id"): goal = client.get_goal(thread_id) _validate_native_identity(state, goal) @@ -486,10 +521,10 @@ def record_event(message: Mapping[str, Any]) -> None: ) _validate_native_identity(state, goal) _update_native_from_goal(state, goal) - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) else: _update_native_from_goal(state, goal) - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) if not native_goal_success(goal, mode=str(native["mode"])): status = str(goal.get("status") or "") @@ -510,14 +545,14 @@ def record_event(message: Mapping[str, Any]) -> None: _update_native_from_goal(state, goal) native["event_counts"] = dict(sorted(event_counts.items())) if not native_goal_success(goal, mode=str(native["mode"])): - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) raise KunlunNativeGoalRuntimeError( "KunlunCode native goal did not reach an accepted terminal state: " + json.dumps(compact_goal(goal), ensure_ascii=False, sort_keys=True) ) native["verified"] = True native["verified_at"] = native.get("verified_at") or _now() - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) def run_native_goal( @@ -556,20 +591,41 @@ def run_native_goal( if not binary: raise KunlunNativeGoalRuntimeError("kunluncode is not on PATH") control = control_plane or LoopXControlPlane(project, context) + goal_id = str(context.get("goal_id") or "") + registry_path = Path( + str(context.get("registry") or project / ".loopx" / "registry.json") + ) + goal_ref = capture_first_party_host_goal_ref( + registry_path=registry_path, + goal_id=goal_id, + ) + goal_admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry_path, + goal_id=goal_id, + planned_goal_ref=goal_ref, + ) journal_path = runtime_state_path(project) with exclusive_file_lock( journal_path, agent_id=str(context.get("agent_id") or ""), operation="kunluncode_native_goal", ): - state = read_runtime_state(project) + state = goal_admission.select_state( + read_state=lambda: read_runtime_state(project), + goal_ref_of=_runtime_goal_ref, + ) if state is not None: _validate_binding(state, context) if ( state["native"].get("verified") is True and runtime_phase(state) != "committed" ): - _closeout_writeback(project, state, control) + _closeout_writeback( + project, + state, + control, + goal_admission, + ) return { "ok": True, "status": "committed", @@ -619,9 +675,11 @@ def run_native_goal( todo_id=todo_id, mode=mode, objective_sha256=objective_sha256, + goal_ref=goal_ref, ) - write_runtime_state(project, state) + _commit_runtime_state(project, state, goal_admission) if not claimed_by: + goal_admission.require_current() control.claim(todo_id) _run_native_host( project, @@ -634,8 +692,15 @@ def run_native_goal( token_budget=token_budget, scenario=scenario, client_factory=client_factory, + goal_admission=goal_admission, + ) + goal_admission.require_current() + _closeout_writeback( + project, + state, + control, + goal_admission, ) - _closeout_writeback(project, state, control) return { "ok": True, "status": "committed", diff --git a/loopx/semantics/goal_instance_binding_inventory_v1.json b/loopx/semantics/goal_instance_binding_inventory_v1.json index 19ba6e412e..feb4c455bf 100644 --- a/loopx/semantics/goal_instance_binding_inventory_v1.json +++ b/loopx/semantics/goal_instance_binding_inventory_v1.json @@ -36,21 +36,23 @@ "locator": "/goals//sessions and host runtime state", "revision_signal": "turn_instance_id", "content_digest_signal": "turn_key", - "observed_reference": "goal_id", + "observed_reference": "goal_id + goal_instance_id for source_session_v1; goal_id otherwise", "producer_sites": [ + "loopx/control_plane/goals/first_party_host_admission.py::capture_first_party_host_goal_ref", "loopx/control_plane/turn_driver/codex_cli.py::codex_cli_session_binding", - "loopx/control_plane/turn_driver/host_binding.py::managed_executor_binding" + "loopx/control_plane/turn_driver/driver.py::build_loopx_turn_plan" ], "consumer_sites": [ + "loopx/control_plane/goals/first_party_host_runtime.ts::decideFirstPartyHostRuntime", "loopx/control_plane/turn_driver/executor.py::run_loopx_turn_once", "loopx/dsh_goal_mode/turn_host_adapter.py::run_dsh_host", "loopx/kunluncode_goal_mode/runtime.py::run_native_goal" ], - "effect_boundary": "loopx/control_plane/turn_driver/executor.py::run_loopx_turn_once", + "effect_boundary": "loopx/control_plane/goals/first_party_host_admission.py::FirstPartyHostGoalAdmission.accept_result", "authority_role": "host_execution_binding", - "current_identity_strength": "goal_alias_only", + "current_identity_strength": "source_exact_partial", "cleanup_support": "host_session_discard", - "m1_disposition": "alias_only_inventory", + "m1_disposition": "source_exact_partial_enforcement", "target_milestone": "M3" }, { diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 99b16fb2c1..1038062257 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1181,6 +1181,30 @@ "api": "atomic_write_json", "classification": "global_registry_io" }, + { + "site": "loopx/control_plane/goals/first_party_host_admission.py::.FirstPartyHostGoalAdmission.for_plan::codec_read:load_project_registry#1", + "line": 120, + "column": 17, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, + { + "site": "loopx/control_plane/goals/first_party_host_admission.py::._source_authority::codec_read:load_project_registry#1", + "line": 36, + "column": 20, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, + { + "site": "loopx/control_plane/goals/first_party_host_admission.py::.capture_first_party_host_goal_ref::codec_read:load_project_registry#1", + "line": 71, + "column": 16, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, { "site": "loopx/control_plane/goals/global_registry_health.py::.collect_global_registry_health::codec_read:load_registry#1", "line": 61, diff --git a/tests/architecture/test_goal_instance_binding_inventory.py b/tests/architecture/test_goal_instance_binding_inventory.py index 968ca37d78..497afb2203 100644 --- a/tests/architecture/test_goal_instance_binding_inventory.py +++ b/tests/architecture/test_goal_instance_binding_inventory.py @@ -50,6 +50,9 @@ "global_goal_projection", "project_registry_goal", } +PARTIALLY_ENFORCED_OWNER_IDS = { + "first_party_host_runtime", +} TYPESCRIPT_DECLARATION = re.compile( r"^(?:export\s+)?(?:async\s+)?(?:function|class)\s+([A-Za-z_$][A-Za-z0-9_$]*)", re.MULTILINE, @@ -156,6 +159,10 @@ def test_m1_observation_claims_are_bounded_to_the_selected_lifecycle() -> None: "build_goal_action_catalog" ) assert owner["target_milestone"] == "M2" + elif owner["owner_id"] in PARTIALLY_ENFORCED_OWNER_IDS: + assert owner["m1_disposition"] == "source_exact_partial_enforcement" + assert owner["current_identity_strength"] == "source_exact_partial" + assert owner["target_milestone"] == "M3" else: assert owner["m1_disposition"] == "alias_only_inventory" assert owner["current_identity_strength"] == "goal_alias_only" diff --git a/tests/architecture/test_source_session_registry_denial.py b/tests/architecture/test_source_session_registry_denial.py index d9d30e079e..d4887729f4 100644 --- a/tests/architecture/test_source_session_registry_denial.py +++ b/tests/architecture/test_source_session_registry_denial.py @@ -10,6 +10,7 @@ "loopx/bootstrap.py", "loopx/claude_goal_mode/scripts/connect.py", "loopx/configure_goal.py", + "loopx/control_plane/goals/first_party_host_admission.py", "loopx/control_plane/projects/registry.py", "loopx/kunluncode_goal_mode/cli.py", "loopx/state_migration.py", @@ -31,7 +32,11 @@ def test_direct_project_registry_loaders_have_source_session_denial() -> None: callers.add(path.relative_to(REPO_ROOT).as_posix()) assert callers == DIRECT_LOADER_ALLOWLIST - for relative in callers - {"loopx/control_plane/projects/registry.py"}: + source_session_owners = { + "loopx/control_plane/goals/first_party_host_admission.py", + "loopx/control_plane/projects/registry.py", + } + for relative in callers - source_session_owners: source = (REPO_ROOT / relative).read_text(encoding="utf-8") assert "require_runtime_compatible_project_registry(" in source, relative diff --git a/tests/control_plane/test_first_party_host_runtime.py b/tests/control_plane/test_first_party_host_runtime.py new file mode 100644 index 0000000000..2d1cccd726 --- /dev/null +++ b/tests/control_plane/test_first_party_host_runtime.py @@ -0,0 +1,268 @@ +from __future__ import annotations + +import threading +import time +from pathlib import Path + +import pytest + +from loopx.control_plane.goals.first_party_host_admission import ( + FirstPartyHostGoalAdmission, + FirstPartyHostRuntimeRejected, + capture_first_party_host_goal_ref, +) +from loopx.control_plane.goals.source_session_registry_state import guard_path +from loopx.control_plane.projects.registry_codec import ( + load_project_registry, + source_session_registry_transaction, +) +from loopx.file_lock import exclusive_cross_runtime_file_lock + + +INSTANCE_A = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" +INSTANCE_B = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" + + +def _source_registry(instance_id: str) -> dict[str, object]: + return { + "schema_version": "0.2", + "registry_role": "project-local", + "profile_id": "source_session_v1", + "common_runtime_root": "/tmp/loopx-runtime", + "projects": [], + "goals": [ + { + "id": "release", + "goal_instance_id": instance_id, + "status": "active", + "execution_authority": False, + } + ], + "session_bindings": [], + "session_receipts": [], + "lifetime_receipts": [], + "retired_goal_instances": [], + } + + +def _write_source_registry(path: Path, instance_id: str) -> None: + create = None if path.exists() else lambda: _source_registry(instance_id) + with source_session_registry_transaction( + path, + operation="first_party_host_runtime_test", + create=create, + ) as transaction: + payload = transaction.payload_copy() + payload["goals"] = _source_registry(instance_id)["goals"] + transaction.commit(payload) + + +def _replace_goal_instance(path: Path, instance_id: str) -> None: + with exclusive_cross_runtime_file_lock( + guard_path(path, "release"), + operation="first_party_host_runtime_test_recreate", + ): + _write_source_registry(path, instance_id) + + +def test_legacy_profile_preserves_callbacks_and_bytes( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + registry = tmp_path / ".loopx" / "registry.json" + registry.parent.mkdir(parents=True) + registry.write_text('{"schema_version":"0.1","goals":[]}\n', encoding="utf-8") + before = registry.read_bytes() + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="release", + planned_goal_ref=None, + ) + monkeypatch.setattr( + "loopx.control_plane.goals.first_party_host_admission.effect_runtime_result", + lambda *_args, **_kwargs: pytest.fail("legacy flow must not call TypeScript"), + ) + state = {"schema_version": "legacy", "goal_id": "release"} + callbacks: list[str] = [] + + assert ( + admission.select_state( + read_state=lambda: state, + goal_ref_of=lambda value: value, + ) + is state + ) + assert admission.accept_result(lambda: callbacks.append("commit")) is None + assert callbacks == ["commit"] + assert registry.read_bytes() == before + + +def test_source_profile_initializes_absent_state_with_exact_goal_ref( + tmp_path: Path, +) -> None: + registry = tmp_path / ".loopx" / "registry.json" + _write_source_registry(registry, INSTANCE_A) + goal_ref = capture_first_party_host_goal_ref( + registry_path=registry, + goal_id="release", + ) + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="release", + planned_goal_ref=goal_ref, + ) + + selected = admission.select_state( + read_state=lambda: None, + goal_ref_of=lambda value: value, + initialize_state=lambda exact: {"goal_ref": exact}, + ) + + assert selected == { + "goal_ref": { + "goal_id": "release", + "goal_instance_id": INSTANCE_A, + } + } + + +def test_source_profile_rejects_legacy_host_state_without_rewriting_it( + tmp_path: Path, +) -> None: + registry = tmp_path / ".loopx" / "registry.json" + _write_source_registry(registry, INSTANCE_A) + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="release", + planned_goal_ref={ + "goal_id": "release", + "goal_instance_id": INSTANCE_A, + }, + ) + state = {"goal_id": "release"} + initialized: list[dict[str, str]] = [] + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + admission.select_state( + read_state=lambda: state, + goal_ref_of=lambda value: value, + initialize_state=lambda exact: initialized.append(exact) or state, + ) + + assert exc_info.value.code == "legacy_host_state" + assert initialized == [] + + +def test_late_result_is_rejected_after_goal_recreation( + tmp_path: Path, +) -> None: + registry = tmp_path / ".loopx" / "registry.json" + _write_source_registry(registry, INSTANCE_A) + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="release", + planned_goal_ref=capture_first_party_host_goal_ref( + registry_path=registry, + goal_id="release", + ), + ) + _replace_goal_instance(registry, INSTANCE_B) + commits: list[str] = [] + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + admission.accept_result(lambda: commits.append("accepted")) + + assert exc_info.value.code == "stale_goal_instance" + assert commits == [] + + +@pytest.mark.parametrize("registry_state", ["missing", "unreadable"]) +def test_source_result_fails_closed_when_registry_is_unavailable( + tmp_path: Path, + registry_state: str, +) -> None: + registry = tmp_path / ".loopx" / "registry.json" + if registry_state == "unreadable": + registry.parent.mkdir(parents=True) + registry.write_text("{not-json", encoding="utf-8") + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="release", + planned_goal_ref={ + "goal_id": "release", + "goal_instance_id": INSTANCE_A, + }, + ) + commits: list[str] = [] + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + admission.accept_result(lambda: commits.append("accepted")) + + assert exc_info.value.code == "goal_authority_unavailable" + assert commits == [] + + +def test_old_source_plan_without_instance_id_fails_closed( + tmp_path: Path, +) -> None: + registry = tmp_path / ".loopx" / "registry.json" + _write_source_registry(registry, INSTANCE_A) + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="release", + planned_goal_ref=None, + ) + callbacks: list[str] = [] + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + admission.accept_result(lambda: callbacks.append("accepted")) + + assert exc_info.value.code == "goal_instance_id_missing" + assert callbacks == [] + + +def test_result_commit_and_recreation_are_serialized_by_the_lifetime_guard( + tmp_path: Path, +) -> None: + registry = tmp_path / ".loopx" / "registry.json" + _write_source_registry(registry, INSTANCE_A) + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="release", + planned_goal_ref=capture_first_party_host_goal_ref( + registry_path=registry, + goal_id="release", + ), + ) + commit_started = threading.Event() + allow_commit = threading.Event() + recreation_finished = threading.Event() + + def commit_result() -> None: + commit_started.set() + assert allow_commit.wait(timeout=5) + + commit_thread = threading.Thread( + target=lambda: admission.accept_result(commit_result), + ) + recreate_thread = threading.Thread( + target=lambda: ( + _replace_goal_instance(registry, INSTANCE_B), + recreation_finished.set(), + ), + ) + commit_thread.start() + assert commit_started.wait(timeout=5) + recreate_thread.start() + time.sleep(0.05) + assert recreation_finished.is_set() is False + + allow_commit.set() + commit_thread.join(timeout=5) + recreate_thread.join(timeout=5) + + assert commit_thread.is_alive() is False + assert recreate_thread.is_alive() is False + assert recreation_finished.is_set() is True + goal = load_project_registry(registry)["goals"][0] + assert goal["goal_instance_id"] == INSTANCE_B diff --git a/tests/control_plane_ts/first_party_host_runtime.test.ts b/tests/control_plane_ts/first_party_host_runtime.test.ts new file mode 100644 index 0000000000..d50bb65490 --- /dev/null +++ b/tests/control_plane_ts/first_party_host_runtime.test.ts @@ -0,0 +1,168 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { decideFirstPartyHostRuntime } from "../../loopx/control_plane/goals/first_party_host_runtime.ts"; + +const INSTANCE_A = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; +const INSTANCE_B = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; +const GOAL_A = { goal_id: "release", goal_instance_id: INSTANCE_A }; +const GOAL_B = { goal_id: "release", goal_instance_id: INSTANCE_B }; + +function facts( + operation: "select_state" | "require_current" | "accept_result", + overrides: Record = {}, +) { + return { + profile_id: "source_session_v1", + operation, + planned_goal_ref: GOAL_A, + authority: { kind: "present", goal_ref: GOAL_A }, + ...(operation === "select_state" + ? { host_state: { kind: "absent" } } + : {}), + ...overrides, + }; +} + +test("non-source profiles preserve legacy behavior without decoding strict facts", () => { + assert.deepEqual( + decideFirstPartyHostRuntime({ + profile_id: "legacy", + operation: "unsupported", + planned_goal_ref: null, + authority: null, + host_state: "malformed", + }), + { kind: "legacy" }, + ); +}); + +test("current exact GoalRef admits each host runtime checkpoint", () => { + assert.deepEqual(decideFirstPartyHostRuntime(facts("select_state")), { + kind: "start_new", + goal_ref: GOAL_A, + }); + assert.deepEqual( + decideFirstPartyHostRuntime( + facts("select_state", { + host_state: { kind: "present", goal_ref: GOAL_A }, + }), + ), + { kind: "resume", goal_ref: GOAL_A }, + ); + assert.deepEqual(decideFirstPartyHostRuntime(facts("require_current")), { + kind: "resume", + goal_ref: GOAL_A, + }); + assert.deepEqual(decideFirstPartyHostRuntime(facts("accept_result")), { + kind: "accept_result", + goal_ref: GOAL_A, + }); +}); + +test("Goal recreation rejects a stale plan at every checkpoint", () => { + for (const operation of [ + "select_state", + "require_current", + "accept_result", + ] as const) { + assert.deepEqual( + decideFirstPartyHostRuntime( + facts(operation, { + authority: { kind: "present", goal_ref: GOAL_B }, + }), + ), + { kind: "reject", code: "stale_goal_instance" }, + operation, + ); + } +}); + +test("strict mode rejects missing and unavailable Goal authority", () => { + const cases = [ + { + authority: { kind: "absent" }, + code: "goal_not_registered", + }, + { + authority: { kind: "unavailable", reason: "registry_missing" }, + code: "goal_authority_unavailable", + }, + { + authority: { kind: "unavailable", reason: "registry_unreadable" }, + code: "goal_authority_unavailable", + }, + { + authority: { kind: "present", goal_ref: { goal_id: "release" } }, + code: "goal_instance_id_missing", + }, + ]; + + for (const fixture of cases) { + assert.deepEqual( + decideFirstPartyHostRuntime( + facts("require_current", { authority: fixture.authority }), + ), + { kind: "reject", code: fixture.code }, + ); + } +}); + +test("strict mode never infers an exact binding from legacy host state", () => { + assert.deepEqual( + decideFirstPartyHostRuntime( + facts("select_state", { + host_state: { + kind: "present", + goal_ref: { goal_id: "release" }, + }, + }), + ), + { kind: "reject", code: "legacy_host_state" }, + ); +}); + +test("strict mode distinguishes stale and malformed host state", () => { + assert.deepEqual( + decideFirstPartyHostRuntime( + facts("select_state", { + host_state: { kind: "present", goal_ref: GOAL_B }, + }), + ), + { + kind: "reject", + code: "host_state_goal_instance_mismatch", + }, + ); + for (const hostState of [ + null, + {}, + { kind: "unknown" }, + { + kind: "present", + goal_ref: { + goal_id: "release", + goal_instance_id: "ginst_invalid", + }, + }, + ]) { + assert.deepEqual( + decideFirstPartyHostRuntime( + facts("select_state", { host_state: hostState }), + ), + { kind: "reject", code: "host_state_malformed" }, + ); + } +}); + +test("source mode rejects an unstamped plan before considering authority", () => { + assert.deepEqual( + decideFirstPartyHostRuntime( + facts("require_current", { + planned_goal_ref: { goal_id: "release" }, + authority: { kind: "absent" }, + }), + ), + { kind: "reject", code: "goal_instance_id_missing" }, + ); +}); diff --git a/tests/test_dsh_goal_mode.py b/tests/test_dsh_goal_mode.py index 493b535536..1a2cec43d7 100644 --- a/tests/test_dsh_goal_mode.py +++ b/tests/test_dsh_goal_mode.py @@ -135,6 +135,31 @@ def test_dsh_session_id_uses_a_versioned_lineage_digest() -> None: assert first_id == turn_host_adapter._derive_session_id(first, TURN_KEY) +def test_dsh_source_session_id_is_scoped_to_the_exact_goal_instance() -> None: + instance_a = _lineage_request( + goal_id="goal", + agent_id="agent", + todo_id="todo", + ) + instance_a["goal_ref"] = { + "goal_id": "goal", + "goal_instance_id": "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + } + instance_b = { + **instance_a, + "goal_ref": { + "goal_id": "goal", + "goal_instance_id": "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + }, + } + + first_id = turn_host_adapter._derive_session_id(instance_a, TURN_KEY) + second_id = turn_host_adapter._derive_session_id(instance_b, TURN_KEY) + + assert first_id.startswith("dsh-lineage-v2-") + assert first_id != second_id + + def test_dsh_session_id_preserves_missing_lineage_component_positions() -> None: missing_agent = _lineage_request(goal_id="goal", todo_id="todo") missing_todo = _lineage_request(goal_id="goal", agent_id="todo") diff --git a/tests/test_kunluncode_goal_mode.py b/tests/test_kunluncode_goal_mode.py index 3b685e7ae3..093aab60b4 100644 --- a/tests/test_kunluncode_goal_mode.py +++ b/tests/test_kunluncode_goal_mode.py @@ -18,6 +18,14 @@ create_fastmcp_server, ) from loopx.control_plane.effect_program import SettlementIdentity +from loopx.control_plane.goals.first_party_host_admission import ( + FirstPartyHostRuntimeRejected, +) +from loopx.control_plane.goals.source_session_registry_state import guard_path +from loopx.control_plane.projects.registry_codec import ( + source_session_registry_transaction, +) +from loopx.file_lock import exclusive_cross_runtime_file_lock from loopx.kunluncode_goal_mode import cli from loopx.kunluncode_goal_mode.app_server import ( NATIVE_GOAL_MODES, @@ -43,6 +51,7 @@ build_native_goal_objective, read_runtime_state, run_native_goal, + runtime_state_path, write_runtime_state, ) @@ -1090,6 +1099,48 @@ def _native_context(tmp_path: Path) -> dict[str, str]: } +def _write_source_native_registry(tmp_path: Path, instance_id: str) -> Path: + registry = tmp_path / ".loopx" / "registry.json" + payload = { + "schema_version": "0.2", + "registry_role": "project-local", + "profile_id": "source_session_v1", + "common_runtime_root": str(tmp_path / "runtime"), + "projects": [], + "goals": [ + { + "id": "shared-goal", + "goal_instance_id": instance_id, + "status": "active", + "execution_authority": False, + } + ], + "session_bindings": [], + "session_receipts": [], + "lifetime_receipts": [], + "retired_goal_instances": [], + } + create = None if registry.exists() else lambda: payload + with source_session_registry_transaction( + registry, + operation="kunlun_host_goal_instance_test", + create=create, + ) as transaction: + current = transaction.payload_copy() + current["goals"] = payload["goals"] + transaction.commit(current) + return registry + + +def _replace_source_native_goal(tmp_path: Path, instance_id: str) -> None: + registry = tmp_path / ".loopx" / "registry.json" + with exclusive_cross_runtime_file_lock( + guard_path(registry, "shared-goal"), + operation="kunlun_host_goal_instance_test_recreate", + ): + _write_source_native_registry(tmp_path, instance_id) + + def _native_todo(*, claimed: bool = False) -> dict[str, object]: return { "todo_id": "todo-native-1", @@ -1363,6 +1414,89 @@ def factory(command: list[str], **kwargs) -> _FakeAppServer: assert "Implement and verify the native Goal Pro adapter" not in journal +def test_source_native_goal_rejects_alias_only_runtime_state( + tmp_path: Path, +) -> None: + _write_source_native_registry( + tmp_path, + "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + ) + state = _native_state(_native_todo(claimed=True)) + write_runtime_state(tmp_path, state) + before = runtime_state_path(tmp_path).read_bytes() + control = _FakeControlPlane(_native_todo(claimed=True)) + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + run_native_goal( + tmp_path, + _native_context(tmp_path), + mode="goal-pro", + permission_mode="auto", + max_duration_secs=60, + controller_timeout_secs=120, + kunluncode_bin="/usr/bin/kunluncode", + control_plane=control, + client_factory=lambda *_args, **_kwargs: pytest.fail( + "legacy runtime state must not launch KunlunCode" + ), + ) + + assert exc_info.value.code == "legacy_host_state" + assert control.calls == [] + assert runtime_state_path(tmp_path).read_bytes() == before + + +def test_source_native_goal_rejects_result_after_recreation( + tmp_path: Path, +) -> None: + instance_a = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + instance_b = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" + _write_source_native_registry(tmp_path, instance_a) + control = _FakeControlPlane(_native_todo()) + + class RecreatingAppServer(_FakeAppServer): + def wait_for_goal_terminal( + self, + thread_id: str, + *, + timeout_seconds: float, + ) -> dict[str, object]: + result = super().wait_for_goal_terminal( + thread_id, + timeout_seconds=timeout_seconds, + ) + _replace_source_native_goal(tmp_path, instance_b) + return result + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + run_native_goal( + tmp_path, + _native_context(tmp_path), + mode="goal-pro", + permission_mode="auto", + max_duration_secs=60, + controller_timeout_secs=120, + kunluncode_bin="/usr/bin/kunluncode", + control_plane=control, + client_factory=lambda command, **kwargs: RecreatingAppServer( + command, + **kwargs, + ), + ) + + assert exc_info.value.code == "stale_goal_instance" + assert control.calls == ["should_run", "claim:todo-native-1"] + state = read_runtime_state(tmp_path) + assert state is not None + assert state["binding"]["goal_instance_id"] == instance_a + assert state["native"]["verified"] is False + assert state["writeback"] == { + "delivery_recorded": False, + "todo_completed": False, + "quota_spent": False, + } + + def test_native_goal_reconciles_crash_after_writeback_without_repeating_it( tmp_path: Path, ) -> None: diff --git a/tests/test_loopx_turn_codex_cli.py b/tests/test_loopx_turn_codex_cli.py index 07fc2aa54a..ac8a80f35f 100644 --- a/tests/test_loopx_turn_codex_cli.py +++ b/tests/test_loopx_turn_codex_cli.py @@ -8,6 +8,12 @@ import pytest +from loopx.control_plane.goals.first_party_host_admission import ( + FirstPartyHostGoalAdmission, +) +from loopx.control_plane.projects.registry_codec import ( + source_session_registry_transaction, +) from loopx.control_plane.turn_driver.codex_cli import ( CODEX_CLI_SESSION_SCHEMA_VERSION, CODEX_STDIO_MCP_SERVER_SCHEMA_VERSION, @@ -30,6 +36,44 @@ FAILURE_ENVELOPE_FIXTURES = ( Path(__file__).parent / "fixtures" / "codex_failure_envelopes.json" ) +SOURCE_GOAL_REF = { + "goal_id": "fixture-goal", + "goal_instance_id": "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", +} + + +def _source_admission(tmp_path: Path) -> FirstPartyHostGoalAdmission: + registry = tmp_path / "source" / ".loopx" / "registry.json" + payload = { + "schema_version": "0.2", + "registry_role": "project-local", + "profile_id": "source_session_v1", + "common_runtime_root": str(tmp_path / "runtime"), + "projects": [], + "goals": [ + { + "id": "fixture-goal", + "goal_instance_id": SOURCE_GOAL_REF["goal_instance_id"], + "status": "active", + "execution_authority": False, + } + ], + "session_bindings": [], + "session_receipts": [], + "lifetime_receipts": [], + "retired_goal_instances": [], + } + with source_session_registry_transaction( + registry, + operation="codex_host_goal_instance_test", + create=lambda: payload, + ) as transaction: + transaction.commit(payload) + return FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="fixture-goal", + planned_goal_ref=SOURCE_GOAL_REF, + ) def _request( @@ -448,6 +492,84 @@ def test_codex_cli_host_starts_then_resumes_opaque_session( assert "private_material" not in persisted +def test_codex_source_session_descriptor_persists_exact_goal_ref( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + executable, log_path = _fake_codex(tmp_path) + monkeypatch.setenv("FAKE_CODEX_LOG", str(log_path)) + runtime_root = tmp_path / "runtime" + project = tmp_path / "project" + project.mkdir() + admission = _source_admission(tmp_path) + request = _request() + request["goal_ref"] = SOURCE_GOAL_REF + + run_codex_cli_host( + request, + runtime_root=runtime_root, + project=project, + codex_bin=str(executable), + timeout_seconds=5, + goal_admission=admission, + ) + + envelope = request["turn_envelope"] + assert isinstance(envelope, dict) + binding = codex_cli_session_binding( + runtime_root, + envelope, + goal_admission=admission, + ) + assert binding is not None + stored = load_codex_cli_session( + runtime_root, + lineage={ + "goal_id": "fixture-goal", + "agent_id": "codex-fixture", + "todo_id": "todo_fixture0001", + }, + ) + assert stored is not None + assert stored["goal_ref"] == SOURCE_GOAL_REF + + +def test_codex_source_session_rejects_alias_only_descriptor( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + executable, log_path = _fake_codex(tmp_path) + monkeypatch.setenv("FAKE_CODEX_LOG", str(log_path)) + runtime_root = tmp_path / "runtime" + project = tmp_path / "project" + project.mkdir() + request = _request() + run_codex_cli_host( + request, + runtime_root=runtime_root, + project=project, + codex_bin=str(executable), + timeout_seconds=5, + ) + descriptor_path = next( + (runtime_root / "goals" / "fixture-goal" / "turn-sessions").glob("*.json") + ) + descriptor_before = descriptor_path.read_bytes() + admission = _source_admission(tmp_path) + + with pytest.raises( + RuntimeError, + match="first-party Host runtime rejected: legacy_host_state", + ): + codex_cli_session_binding( + runtime_root, + request["turn_envelope"], + goal_admission=admission, + ) + + assert descriptor_path.read_bytes() == descriptor_before + + def test_codex_cli_host_materializes_bound_mcp_tools_for_fresh_and_resume( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, diff --git a/tests/test_loopx_turn_driver.py b/tests/test_loopx_turn_driver.py index 0bcb53e19e..a6dfa8632e 100644 --- a/tests/test_loopx_turn_driver.py +++ b/tests/test_loopx_turn_driver.py @@ -949,6 +949,44 @@ def test_turn_plan_transaction_key_is_stable_and_todo_scoped() -> None: ) +def test_source_goal_ref_scopes_turn_identity_without_changing_legacy_plan() -> None: + legacy = build_loopx_turn_plan( + _envelope(), + host="generic-cli", + execution_mode="isolated-headless", + ) + instance_a = build_loopx_turn_plan( + _envelope(), + host="generic-cli", + execution_mode="isolated-headless", + goal_ref={ + "goal_id": "fixture-goal", + "goal_instance_id": "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + }, + ) + instance_b = build_loopx_turn_plan( + _envelope(), + host="generic-cli", + execution_mode="isolated-headless", + goal_ref={ + "goal_id": "fixture-goal", + "goal_instance_id": "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + }, + ) + + assert "goal_ref" not in legacy + assert "goal_ref" not in legacy["transaction"] + assert instance_a["goal_ref"] == instance_a["transaction"]["goal_ref"] + assert ( + instance_a["transaction"]["turn_key"] + != instance_b["transaction"]["turn_key"] + ) + assert ( + instance_a["transaction"]["turn_key"] + != legacy["transaction"]["turn_key"] + ) + + def test_turn_plan_instance_id_distinguishes_new_turns_from_retries() -> None: first = build_loopx_turn_plan( _envelope(), diff --git a/tests/test_loopx_turn_executor.py b/tests/test_loopx_turn_executor.py index bdefa2b32e..f4671360eb 100644 --- a/tests/test_loopx_turn_executor.py +++ b/tests/test_loopx_turn_executor.py @@ -10,6 +10,14 @@ from loopx.cli_commands import turn_cadence from loopx.cli_commands.turn_cadence import ManagedCadenceStart, managed_cadence_start from loopx.control_plane.effect_runtime import effect_runtime_result +from loopx.control_plane.goals.first_party_host_admission import ( + FirstPartyHostGoalAdmission, + FirstPartyHostRuntimeRejected, +) +from loopx.control_plane.goals.source_session_registry_state import guard_path +from loopx.control_plane.projects.registry_codec import ( + source_session_registry_transaction, +) from loopx.control_plane.turn_driver import executor as turn_executor from loopx.control_plane.turn_driver import ( LOOPX_TURN_RESULT_SCHEMA_VERSION, @@ -35,6 +43,50 @@ from loopx.control_plane.turn_driver.host_binding import managed_executor_binding from loopx.control_plane.turn_driver.settlement import execute_turn_driver_settlement from loopx.control_plane.turn_driver.transaction import TRANSACTION_PHASES +from loopx.file_lock import exclusive_cross_runtime_file_lock + + +INSTANCE_A = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" +INSTANCE_B = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" + + +def _write_source_registry(path: Path, instance_id: str) -> None: + payload = { + "schema_version": "0.2", + "registry_role": "project-local", + "profile_id": "source_session_v1", + "common_runtime_root": str(path.parent), + "projects": [], + "goals": [ + { + "id": "fixture-goal", + "goal_instance_id": instance_id, + "status": "active", + "execution_authority": False, + } + ], + "session_bindings": [], + "session_receipts": [], + "lifetime_receipts": [], + "retired_goal_instances": [], + } + create = None if path.exists() else lambda: payload + with source_session_registry_transaction( + path, + operation="turn_host_goal_instance_test", + create=create, + ) as transaction: + current = transaction.payload_copy() + current["goals"] = payload["goals"] + transaction.commit(current) + + +def _replace_source_goal(path: Path, instance_id: str) -> None: + with exclusive_cross_runtime_file_lock( + guard_path(path, "fixture-goal"), + operation="turn_host_goal_instance_test_recreate", + ): + _write_source_registry(path, instance_id) def _plan() -> dict[str, object]: @@ -74,6 +126,21 @@ def _plan() -> dict[str, object]: ) +def _source_plan() -> dict[str, object]: + plan = _plan() + envelope = plan["turn_envelope"] + assert isinstance(envelope, dict) + return build_loopx_turn_plan( + envelope, + host="generic-cli", + execution_mode="isolated-headless", + goal_ref={ + "goal_id": "fixture-goal", + "goal_instance_id": INSTANCE_A, + }, + ) + + def _codex_plan() -> dict[str, object]: plan = _plan() envelope = plan["turn_envelope"] @@ -463,6 +530,159 @@ def _passing_validator( } +def test_late_host_result_cannot_enter_a_recreated_goal( + tmp_path: Path, +) -> None: + registry = tmp_path / "project" / ".loopx" / "registry.json" + _write_source_registry(registry, INSTANCE_A) + plan = _source_plan() + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="fixture-goal", + planned_goal_ref=plan["goal_ref"], + ) + calls = {"writeback": 0, "spend": 0, "scheduler": 0} + writeback, spend, scheduler = _callbacks(calls) + runtime_root = tmp_path / "runtime" + + def stale_host_result(_request: Mapping[str, object]) -> dict[str, object]: + _replace_source_goal(registry, INSTANCE_B) + return _host_result(plan) + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + run_loopx_turn_once( + plan, + host_runner=stale_host_result, + project=tmp_path, + runtime_root=runtime_root, + goal_id="fixture-goal", + timeout_seconds=5, + execute=True, + task_validator=_passing_validator, + writeback=writeback, + spend=spend, + scheduler=scheduler, + goal_admission=admission, + ) + + assert exc_info.value.code == "stale_goal_instance" + assert calls == {"writeback": 0, "spend": 0, "scheduler": 0} + journal = _journal(runtime_root) + assert "host_result" not in journal + assert journal["completed_phases"] == [] + + +def test_stale_source_plan_is_rejected_before_host_start( + tmp_path: Path, +) -> None: + registry = tmp_path / "project" / ".loopx" / "registry.json" + _write_source_registry(registry, INSTANCE_A) + plan = _source_plan() + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="fixture-goal", + planned_goal_ref=plan["goal_ref"], + ) + _replace_source_goal(registry, INSTANCE_B) + calls = {"host": 0, "writeback": 0, "spend": 0, "scheduler": 0} + writeback, spend, scheduler = _callbacks(calls) + runtime_root = tmp_path / "runtime" + + def forbidden_host(_request: Mapping[str, object]) -> dict[str, object]: + calls["host"] += 1 + return _host_result(plan) + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + run_loopx_turn_once( + plan, + host_runner=forbidden_host, + project=tmp_path, + runtime_root=runtime_root, + goal_id="fixture-goal", + timeout_seconds=5, + execute=True, + task_validator=_passing_validator, + writeback=writeback, + spend=spend, + scheduler=scheduler, + goal_admission=admission, + ) + + assert exc_info.value.code == "stale_goal_instance" + assert calls == {"host": 0, "writeback": 0, "spend": 0, "scheduler": 0} + transaction = plan["transaction"] + assert isinstance(transaction, dict) + assert ( + turn_journal_path( + runtime_root, + goal_id="fixture-goal", + turn_key=str(transaction["turn_key"]), + ).exists() + is False + ) + + +def test_cached_host_result_cannot_resume_after_goal_recreation( + tmp_path: Path, +) -> None: + registry = tmp_path / "project" / ".loopx" / "registry.json" + _write_source_registry(registry, INSTANCE_A) + plan = _source_plan() + transaction = plan["transaction"] + assert isinstance(transaction, dict) + runtime_root = tmp_path / "runtime" + path = turn_journal_path( + runtime_root, + goal_id="fixture-goal", + turn_key=str(transaction["turn_key"]), + ) + journal = { + "schema_version": LOOPX_TURN_JOURNAL_SCHEMA_VERSION, + "turn_key": transaction["turn_key"], + "goal_id": "fixture-goal", + "status": "in_progress", + "host": {"kind": "generic-cli"}, + "completed_phases": [], + "plan": plan, + } + turn_executor._write_journal(path, journal) + journal.update( + completed_phases=list(TRANSACTION_PHASES[:2]), + host_result=_host_result(plan), + result_kind="validated_progress", + ) + turn_executor._write_journal(path, journal) + admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=registry, + goal_id="fixture-goal", + planned_goal_ref=plan["goal_ref"], + ) + _replace_source_goal(registry, INSTANCE_B) + calls = {"writeback": 0, "spend": 0, "scheduler": 0} + writeback, spend, scheduler = _callbacks(calls) + + with pytest.raises(FirstPartyHostRuntimeRejected) as exc_info: + run_loopx_turn_once( + plan, + host_runner=lambda _request: pytest.fail( + "cached host result must not relaunch the Host" + ), + project=tmp_path, + runtime_root=runtime_root, + goal_id="fixture-goal", + timeout_seconds=5, + execute=True, + task_validator=_passing_validator, + writeback=writeback, + spend=spend, + scheduler=scheduler, + goal_admission=admission, + ) + + assert exc_info.value.code == "stale_goal_instance" + assert calls == {"writeback": 0, "spend": 0, "scheduler": 0} + + def test_host_result_requires_bounded_public_material_fields() -> None: plan = _plan() result = _host_result(plan)