diff --git a/benchmark/LHTB/README.md b/benchmark/LHTB/README.md index fe1b81086f..da6c3075ca 100644 --- a/benchmark/LHTB/README.md +++ b/benchmark/LHTB/README.md @@ -113,8 +113,8 @@ not replayed; bootstrap and Todo lifecycle follow the installed product version. ## Shared execution configuration `LOOPX_EXECUTION_MODE` selects `plain`, `native-goal`, `heartbeat` (default), -`turn` or `loopx-goal`. `LOOPX_ITERATION_CONTEXT` defaults to `fresh`; only Turn -accepts `resume-if-available`. Turn also requires `LOOPX_VALIDATION_COMMAND_JSON`, +`turn` or `loopx-goal`. `LOOPX_ITERATION_CONTEXT` defaults to `fresh`; heartbeat and Turn +accept `resume`, sharing the same Goal/Agent session across planning, wakes and Todos. Turn also requires `LOOPX_VALIDATION_COMMAND_JSON`, an argv array for the independently protected task validator. No generic benchmark scoring or hidden-verifier feedback is introduced. diff --git a/benchmark/runtime/RUNTIME.md b/benchmark/runtime/RUNTIME.md index 50b7935f95..595bc27925 100644 --- a/benchmark/runtime/RUNTIME.md +++ b/benchmark/runtime/RUNTIME.md @@ -30,15 +30,28 @@ agents: | --- | --- | --- | | `plain` | One Codex exec, native Goals disabled | Absent | | `native-goal` | Installed native Goal transport; objective `Finish the task.` | Absent | -| `heartbeat` | Product thin heartbeat + external scheduler, fresh each wake | Present | +| `heartbeat` | Product thin heartbeat + external scheduler, fresh or same-session resume | Present | | `turn` | Public Turn CLI, typed result, independent validation, settlement | Present | | `loopx-goal` | Product Goal body + installed native Goal transport | Present | -Only `turn` accepts `iteration_context: resume-if-available`. Core session -compatibility determines whether it actually resumes, including after changing -Todo. Native Goal continuation stays with Codex; blocked Goals are not -automatically unblocked. Plain exec versus Goal app-server also changes -transport; it does not isolate the continuation effect alone. +`heartbeat` and `turn` accept `iteration_context: resume`. The first invocation +creates one Codex conversation; later planning checkpoints and execution wakes +resume that exact native session ID for the trial's Goal and Agent, including +when the selected Todo changes. Both drivers share the product's agent-scoped +Codex session store. `fresh` remains the runner default and starts a new session +on each invocation. The former context name is rejected, with no alias. + +Resume never falls back to a new conversation when a binding is corrupt, the +trial home/workspace/model/settings change, or Codex returns a different ID. +Repair the configuration or explicitly select `fresh`; the next observed fresh +session replaces the binding. A timeout preserves an observed ID without +claiming progress. Wakes remain serialized by the outer controller. Private +wake receipts record the requested action and confirmed native session ID; +aggregate trajectories copy each native session once. + +Native Goal continuation stays with Codex; blocked Goals are not automatically +unblocked. Plain exec versus Goal app-server also changes transport; it does +not isolate the continuation effect alone. `turn` requires `validation_command`, an argv list for an independently protected validator available inside the task environment. It receives the @@ -67,10 +80,13 @@ planning process or missing result fails the entry; it never falls back to a generic Todo. A blocked entry retains the referenced blockers and starts no execution driver. Readback proves state and ownership, not semantic plan quality. -Planning uses a separate fresh `codex exec` session with native Goals disabled -for that call. Its session is not inserted into core Turn session bindings or -resumed by the subsequent execution. This is a planning-contract ablation, not -an exact reproduction of same-conversation interactive `$loopx` startup. +Planning follows the chosen context policy for heartbeat and Turn. With +`resume`, planning and execution share the same conversation, including later +phase planning; each checkpoint still renders fresh public task inputs and +validates actual Todo readback. With `fresh`, each planning/execution invocation +starts a new conversation. LoopX Goal planning uses a separate exec conversation +because native Goal execution owns its app-server thread lifecycle. The runner +does not claim exact equivalence to interactive `$loopx` startup. The default `planning_timeout_sec` is 300; planning and preparation consume the same `scheduler_timeout_sec` phase budget as execution. Planning sessions are included in native session/token aggregation. No planning checkpoint is counted diff --git a/benchmark/runtime/codex.py b/benchmark/runtime/codex.py index 4da2ff4fe7..e323b1d61d 100644 --- a/benchmark/runtime/codex.py +++ b/benchmark/runtime/codex.py @@ -10,7 +10,7 @@ MODES = ("plain", "native-goal", "heartbeat", "turn", "loopx-goal") -CONTEXTS = ("fresh", "resume-if-available") +CONTEXTS = ("fresh", "resume") TASK_ENTRIES = ("seeded-todo", "loopx-planned") SANDBOXES = ("read-only", "workspace-write", "danger-full-access") @@ -31,8 +31,8 @@ def __post_init__(self) -> None: raise ValueError("unsupported task entry") if self.task_entry == "loopx-planned" and not self.uses_loopx: raise ValueError("loopx-planned requires a LoopX execution mode") - if self.context != "fresh" and self.mode != "turn": - raise ValueError("resume-if-available currently requires mode=turn") + if self.context != "fresh" and self.mode not in {"turn", "heartbeat"}: + raise ValueError("resume requires mode=turn or mode=heartbeat") if self.sandbox not in SANDBOXES: raise ValueError("unsupported Codex sandbox") if not math.isfinite(self.timeout_seconds) or self.timeout_seconds <= 0: diff --git a/benchmark/runtime/sessions.py b/benchmark/runtime/sessions.py new file mode 100644 index 0000000000..87eb73a774 --- /dev/null +++ b/benchmark/runtime/sessions.py @@ -0,0 +1,117 @@ +"""Trial-local conversation continuity through the Codex session owner.""" + +from __future__ import annotations + +import json +import shutil +from pathlib import Path +from typing import Any + +from .codex import Execution + +from loopx.control_plane.goals.first_party_host_admission import ( + FirstPartyHostGoalAdmission, + capture_first_party_host_goal_ref, +) +from loopx.control_plane.turn_driver.codex_cli import codex_cli_event_session_id +from loopx.control_plane.turn_driver.codex_sessions import ( + _store_codex_cli_session, + codex_session_profile_digest, + require_codex_session_profile, + select_codex_cli_session, +) + + +class BenchmarkSessionWake: + """Planning and heartbeat use the same Goal/Agent binding as governed Turns. + + The outer controller serializes wakes. A native ID must be observed even on + timeout; neither missing history nor a failed resume permits a fresh fork. + """ + + def __init__( + self, env: dict[str, str], execution: Execution, receipt: dict[str, Any], + ) -> None: + self.root = Path(env["LOOPX_RUNTIME_ROOT"]) + self.lineage = { + "goal_id": env["LOOPX_GOAL_ID"], + "agent_id": env["LOOPX_AGENT_ID"], + } + self.goal_ref = capture_first_party_host_goal_ref( + registry_path=Path(env["LOOPX_REGISTRY"]), + goal_id=self.lineage["goal_id"], + ) + self.admission = FirstPartyHostGoalAdmission.for_plan( + registry_path=Path(env["LOOPX_REGISTRY"]), + goal_id=self.lineage["goal_id"], + planned_goal_ref=self.goal_ref, + ) + if execution.context == "fresh": + self.admission.require_current() + binary = shutil.which(env["CODEX_BIN"]) + if binary is None: + raise ValueError("Codex CLI executable is unavailable") + self.digest = codex_session_profile_digest( + project=Path(env["LOOPX_PROJECT"]), + codex_bin=binary, + home=Path(env["CODEX_HOME"]), + model=env["MODEL_NAME"], + reasoning_effort=env["REASONING_EFFORT"], + sandbox=execution.sandbox, + ) + binding = ( + select_codex_cli_session( + self.root, + lineage=self.lineage, + session_scope="agent", + goal_admission=self.admission, + ) + if execution.context == "resume" + else None + ) + if binding: + require_codex_session_profile(binding, self.digest) + if binding.get("operation_transport"): + raise ValueError("session requires its original managed transport") + self.session_id = binding["session_id"] if binding else None + self.receipt = receipt + receipt["session"] = { + "binding_scope": "agent", + "action": "resume" if binding else "start_new", + "session_id": self.session_id, + } + + def observe(self, path: Path) -> None: + if not path.exists(): + return # No process was launched. + ids = set() + with path.open() as stream: + for line in stream: + try: + event = json.loads(line) + except ValueError: + continue + candidate = ( + codex_cli_event_session_id(event) + if isinstance(event, dict) + else None + ) + if candidate: + ids.add(candidate) + if len(ids) != 1 or ( + self.session_id is not None and self.session_id not in ids + ): + self.receipt.update(ok=False, error_kind="session_identity_unconfirmed") + raise RuntimeError("Codex did not confirm the expected session identity") + observed = ids.pop() + self.admission.accept_result( + lambda: _store_codex_cli_session( + self.root, + lineage=self.lineage, + session_scope="agent", + session_id=observed, + goal_ref=self.goal_ref, + session_profile_digest=self.digest, + ) + ) + self.receipt["session"]["session_id"] = observed diff --git a/benchmark/runtime/worker.py b/benchmark/runtime/worker.py index afd9fafd64..bd336e987c 100644 --- a/benchmark/runtime/worker.py +++ b/benchmark/runtime/worker.py @@ -134,6 +134,10 @@ def turn_command( ), "--iteration-context", execution.context, + "--session-scope", + "agent", + "--codex-reasoning-effort", + env["REASONING_EFFORT"], "--codex-bin", env["CODEX_BIN"], "--codex-model", @@ -200,6 +204,32 @@ def run_native_goal( receipt["native_goal"] = compact_native_goal_receipt(observed[0]) +def native_command(env, execution, stage, wake, session_wake): + if session_wake is None: + # Preserve the baseline/Goal-planning entrypoint and its transport inputs. + return [ + env["CODEX_BIN"], "exec", "--json", "--skip-git-repo-check", + "--sandbox", execution.sandbox, "--cd", env["LOOPX_PROJECT"], + *(["-c", "features.goals=false", "--output-schema", str(wake / "planning-schema.json"), + "--output-last-message", str(wake / "planning-result.json")] if stage == "plan" else []), + "-", + ] + from loopx.control_plane.turn_driver.codex_cli import _codex_command + + command = _codex_command( + codex_bin=env["CODEX_BIN"], project=Path(env["LOOPX_PROJECT"]), + schema_path=wake / "planning-schema.json" if stage == "plan" else None, + output_path=wake / "planning-result.json" if stage == "plan" else None, + sandbox=execution.sandbox, model=env["MODEL_NAME"], + reasoning_effort=env["REASONING_EFFORT"], + session_id=session_wake.session_id if session_wake else None, + mcp_server=None, + ) + if stage == "plan": + command[-1:-1] = ["-c", "features.goals=false"] + return command + + def run_once(env: dict[str, str]) -> dict: execution = Execution( mode=env.get("LOOPX_EXECUTION_MODE", "heartbeat"), @@ -229,6 +259,7 @@ def run_once(env: dict[str, str]) -> dict: "ok": False, "timed_out": False, } + session_wake = None pending_path = ( Path(env.get("LOOPX_RUNTIME_ROOT", str(home))) / "benchmark-pending-turn.json" ) @@ -266,6 +297,11 @@ def run_once(env: dict[str, str]) -> dict: body = heartbeat_body(env, turn_id, native_goal=execution.native_goal) elif execution.mode == "plain": body = Path(env["LOOPX_TASK_DOC"]).read_text(encoding="utf-8") + session_wake = None + if execution.mode in {"heartbeat", "turn"} and not (execution.mode == "turn" and stage == "execute"): + from benchmark.runtime.sessions import BenchmarkSessionWake + + session_wake = BenchmarkSessionWake(env, execution, receipt) with (wake / "stderr.log").open("w") as stderr: if execution.native_goal and stage == "execute": run_native_goal(env, execution, body, receipt, stderr) @@ -286,22 +322,7 @@ def run_once(env: dict[str, str]) -> dict: pending.get("resume_turn_key"), ) if execution.mode == "turn" and stage == "execute" - else [ - env["CODEX_BIN"], - "exec", - "--json", - "--skip-git-repo-check", - "--sandbox", - execution.sandbox, - "--cd", - env["LOOPX_PROJECT"], - *([ - "-c", "features.goals=false", - "--output-schema", str(wake / "planning-schema.json"), - "--output-last-message", str(wake / "planning-result.json"), - ] if stage == "plan" else []), - "-", - ] + else native_command(env, execution, stage, wake, session_wake) ) with (wake / "stdout.jsonl").open("w") as stdout: with child_process( @@ -347,13 +368,19 @@ def run_once(env: dict[str, str]) -> dict: receipt["error_kind"] = type(exc).__name__ raise finally: - if (home / "sessions").is_dir(): - # One authoritative copy per native session; resume must not count - # the same prefix again in every wake's aggregate trajectory. - shutil.copytree( - home / "sessions", log_root.parent / "sessions", dirs_exist_ok=True - ) - (wake / "receipt.json").write_text(json.dumps(receipt, indent=2) + "\n") + try: + if session_wake is not None: + session_wake.observe(wake / "stdout.jsonl") + except BaseException as exc: + receipt.update(ok=False, error_kind=type(exc).__name__) + raise + finally: + if (home / "sessions").is_dir(): + # One authoritative copy; resumed prefixes are never double counted. + shutil.copytree( + home / "sessions", log_root.parent / "sessions", dirs_exist_ok=True + ) + (wake / "receipt.json").write_text(json.dumps(receipt, indent=2) + "\n") return receipt diff --git a/benchmark/tests/test_resume_sessions.py b/benchmark/tests/test_resume_sessions.py new file mode 100644 index 0000000000..5d98f734a0 --- /dev/null +++ b/benchmark/tests/test_resume_sessions.py @@ -0,0 +1,247 @@ +"""Conversation continuity contracts, including the native exec transport.""" + +from __future__ import annotations + +import json +import os +import threading +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path + +import pytest + +from benchmark.runtime.worker import run_once +from benchmark.tests.test_shared_codex_runtime import worker_env +from loopx.control_plane.turn_driver.codex_sessions import load_codex_cli_session + + +def resume_env(tmp_path, monkeypatch): + env = worker_env(tmp_path) + skills = tmp_path / "skills" + skills.mkdir() + env.update( + LOOPX_EXECUTION_MODE="heartbeat", + LOOPX_ITERATION_CONTEXT="resume", + LOOPX_RUNTIME_ROOT=str(tmp_path / "runtime"), + LOOPX_REGISTRY=str(tmp_path / "registry.json"), + LOOPX_GOAL_ID="fixture-goal", + LOOPX_AGENT_ID="fixture-agent", + LOOPX_SHARED_SKILLS=str(skills), + LOOPX_CLI="loopx", + ) + monkeypatch.setattr( + "benchmark.runtime.worker.heartbeat_body", + lambda *a, **kw: "Current Todo input.", + ) + Path(env["CODEX_BIN"]).write_text( + f"#!{__import__('sys').executable}\n" + + """ +import json, os, pathlib, sys, time, uuid +args = sys.argv[1:] +body = sys.stdin.read() +session = args[-2] if "resume" in args else str(uuid.uuid4()) +if os.environ.get("FORK_SESSION"): + session = str(uuid.uuid4()) +home = pathlib.Path(os.environ["CODEX_HOME"]) / "sessions" +home.mkdir(exist_ok=True) +with (home / (session + ".jsonl")).open("a") as stream: + stream.write(json.dumps({"body": body}) + "\\n") +print(json.dumps({"type": "thread.started", "thread_id": session}), flush=True) +print(json.dumps({"argv": args}), flush=True) +if os.environ.get("HOLD_SESSION"): + time.sleep(60) +if "--output-last-message" in args: + pathlib.Path(args[args.index("--output-last-message") + 1]).write_text('{"marker":"planning-ack"}') +""" + ) + return env + + +def plan_stage(env, monkeypatch): + from benchmark.runtime import planning + + schema = { + "type": "object", + "properties": {"marker": {"type": "string"}}, + "required": ["marker"], + "additionalProperties": False, + } + monkeypatch.setattr( + planning, + "task_plan_packet", + lambda *a: {"result_schema": schema, "marker": "planning-input"}, + ) + monkeypatch.setattr(planning, "validate_plan_readback", lambda *a: {"ok": True}) + return env | { + "LOOPX_TASK_ENTRY": "loopx-planned", + "LOOPX_TASK_STAGE": "plan", + "LOOPX_PLANNING_TIMEOUT_SEC": "30", + "LOOPX_PLANNING_RESULT": str(Path(env["LOOPX_PROJECT"]) / "planning.json"), + } + + +def binding(env): + return load_codex_cli_session( + Path(env["LOOPX_RUNTIME_ROOT"]), + lineage={"goal_id": env["LOOPX_GOAL_ID"], "agent_id": env["LOOPX_AGENT_ID"]}, + session_scope="agent", + ) + + +@pytest.mark.parametrize("planning_mode", ["heartbeat", "turn"]) +def test_planning_heartbeat_and_turn_share_one_binding( + tmp_path, monkeypatch, planning_mode +): + env = resume_env(tmp_path, monkeypatch) + plan_env = plan_stage(env, monkeypatch) | {"LOOPX_EXECUTION_MODE": planning_mode} + if planning_mode == "turn": + plan_env["LOOPX_VALIDATION_COMMAND_JSON"] = '["true"]' + first = run_once(plan_env) + second = run_once(env) + third = run_once(env) + assert first["ok"] and second["ok"] and third["ok"] + assert first["session"]["action"] == "start_new" + assert second["session"]["action"] == third["session"]["action"] == "resume" + assert len({r["session"]["session_id"] for r in (first, second, third)}) == 1 + assert binding(env)["session_id"] == first["session"]["session_id"] + assert len(list((tmp_path / "logs/sessions").glob("*.jsonl"))) == 1 + + +def test_timeout_keeps_the_observed_session_for_next_heartbeat(tmp_path, monkeypatch): + env = resume_env(tmp_path, monkeypatch) + first = run_once(env | {"HOLD_SESSION": "1", "LOOPX_CODEX_TURN_TIMEOUT_SEC": "1"}) + assert first["timed_out"] and not first["ok"] + second = run_once(env) + assert second["ok"] and second["session"]["action"] == "resume" + assert second["session"]["session_id"] == first["session"]["session_id"] + + +@pytest.mark.parametrize("mutation", ["corrupt", "home", "fork"]) +def test_resume_refuses_silent_replacement(tmp_path, monkeypatch, mutation): + env = resume_env(tmp_path, monkeypatch) + first = run_once(env) + if mutation == "corrupt": + next((tmp_path / "runtime").glob("goals/*/turn-sessions/*.json")).write_text( + "invalid" + ) + elif mutation == "home": + env["LOOPX_CODEX_HOME"] = str(tmp_path / "other-home") + else: + env["FORK_SESSION"] = "1" + with pytest.raises((ValueError, RuntimeError)): + run_once(env) + receipt = json.loads( + sorted((tmp_path / "logs/wakes").glob("*/receipt.json"))[-1].read_text() + ) + assert not receipt["ok"] + if mutation == "fork": + assert binding(env)["session_id"] == first["session"]["session_id"] + + +@pytest.mark.parametrize("field", ["LOOPX_GOAL_ID", "LOOPX_AGENT_ID"]) +def test_another_identity_starts_an_independent_conversation( + tmp_path, monkeypatch, field +): + env = resume_env(tmp_path, monkeypatch) + first = run_once(env) + second = run_once(env | {field: "another-identity"}) + assert second["session"]["action"] == "start_new" + assert second["session"]["session_id"] != first["session"]["session_id"] + + +def test_explicit_fresh_replaces_corrupt_shared_binding(tmp_path, monkeypatch): + env = resume_env(tmp_path, monkeypatch) + first = run_once(env) + next((tmp_path / "runtime").glob("goals/*/turn-sessions/*.json")).write_text( + "invalid" + ) + second = run_once(env | {"LOOPX_ITERATION_CONTEXT": "fresh"}) + third = run_once(env) + assert second["session"]["session_id"] != first["session"]["session_id"] + assert third["session"]["session_id"] == second["session"]["session_id"] + + +@pytest.mark.skipif( + not os.environ.get("LOOPX_TEST_CODEX_BIN"), reason="native Codex binary required" +) +def test_native_codex_resume_keeps_planning_history(tmp_path, monkeypatch): + """Exercise installed Codex against a disposable local Responses endpoint.""" + requests = [] + + class Handler(BaseHTTPRequestHandler): + def do_POST(self): + raw = self.rfile.read(int(self.headers["Content-Length"])) + requests.append(json.loads(raw)) + text = ( + '{"marker":"planning-ack"}' if len(requests) == 1 else "execution-ack" + ) + message = { + "id": "msg_fixture", + "type": "message", + "role": "assistant", + "status": "completed", + "content": [{"type": "output_text", "text": text, "annotations": []}], + } + response = { + "id": f"resp_{len(requests)}", + "object": "response", + "model": "gpt-5.4", + "status": "completed", + "output": [message], + "usage": {"input_tokens": 10, "output_tokens": 5, "total_tokens": 15}, + } + events = [ + { + "type": "response.created", + "response": response | {"status": "in_progress", "output": []}, + }, + { + "type": "response.output_item.added", + "output_index": 0, + "item": message, + }, + { + "type": "response.output_text.delta", + "output_index": 0, + "content_index": 0, + "item_id": message["id"], + "delta": text, + }, + { + "type": "response.output_item.done", + "output_index": 0, + "item": message, + }, + {"type": "response.completed", "response": response}, + ] + self.send_response(200) + self.send_header("Content-Type", "text/event-stream") + self.end_headers() + for event in events: + self.wfile.write(("data: " + json.dumps(event) + "\n\n").encode()) + self.wfile.flush() + + def log_message(self, *args): + pass + + server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) + threading.Thread(target=server.serve_forever, daemon=True).start() + try: + env = resume_env(tmp_path, monkeypatch) | { + "CODEX_BIN": os.environ["LOOPX_TEST_CODEX_BIN"], + "MODEL_NAME": "gpt-5.4", + "OPENAI_BASE_URL": f"http://127.0.0.1:{server.server_port}/v1", + "LOOPX_CODEX_TURN_TIMEOUT_SEC": "30", + } + first = run_once(plan_stage(env, monkeypatch)) + second = run_once(env) + assert first["ok"] and second["ok"] + assert first["session"]["session_id"] == second["session"]["session_id"] + assert len(requests) == 2 + history = json.dumps(requests[1]["input"]) + assert "planning-input" in history and "planning-ack" in history + assert "Current Todo input." in history + assert len(list((tmp_path / "logs/sessions").rglob("*.jsonl"))) == 1 + finally: + server.shutdown() + server.server_close() diff --git a/benchmark/tests/test_shared_codex_runtime.py b/benchmark/tests/test_shared_codex_runtime.py index 144797ff55..83e6d2bcf9 100644 --- a/benchmark/tests/test_shared_codex_runtime.py +++ b/benchmark/tests/test_shared_codex_runtime.py @@ -239,18 +239,30 @@ def test_turn_uses_public_cli_and_core_session_policy(tmp_path): } execution = Execution( mode="turn", - context="resume-if-available", + context="resume", validation_command=("python", "trusted-validator.py"), ) command = turn_command(env, execution, "wake-fixture") from loopx.cli import build_parser parsed = build_parser().parse_args(command[1:]) - assert parsed.iteration_context == "resume-if-available" + assert parsed.iteration_context == "resume" + assert parsed.session_scope == "agent" assert parsed.validation_command_json == '["python", "trusted-validator.py"]' assert parsed.codex_sandbox == "danger-full-access" +def test_turn_defaults_and_retired_context_rejection(): + from loopx.cli import build_parser + + parser = build_parser() + argv = ["turn", "plan", "--goal-id", "fixture-goal", "--agent-id", "fixture-agent"] + parsed = parser.parse_args(argv) + assert parsed.iteration_context == "resume" and parsed.session_scope == "agent" + with pytest.raises(SystemExit): + parser.parse_args([*argv, "--iteration-context", "resume-if-available"]) + + def test_failed_turn_restarts_same_transaction_until_core_recovers(tmp_path): env = worker_env(tmp_path) skills = tmp_path / "skills" diff --git a/benchmark/tests/test_task_entry.py b/benchmark/tests/test_task_entry.py index 72b073b061..3fde565f02 100644 --- a/benchmark/tests/test_task_entry.py +++ b/benchmark/tests/test_task_entry.py @@ -64,6 +64,7 @@ def planning_env(tmp_path): f"#!{sys.executable}\n" + """ import json, os, pathlib, subprocess, sys +print(json.dumps({"type": "thread.started", "thread_id": "planning-fixture-session"}), flush=True) packet = json.loads(sys.stdin.read().split("Host-supplied planning checkpoint:\\n", 1)[1]) command = [os.environ["LOOPX_CLI"], "--format", "json", "--registry", os.environ["LOOPX_REGISTRY"], "--runtime-root", os.environ["LOOPX_RUNTIME_ROOT"], "todo", "add", "--goal-id", "planning-goal", diff --git a/docs/architecture/rfcs/long-horizon-harness-benchmark-research-program-v0.md b/docs/architecture/rfcs/long-horizon-harness-benchmark-research-program-v0.md index 6c233f5d41..29f54f3b9b 100644 --- a/docs/architecture/rfcs/long-horizon-harness-benchmark-research-program-v0.md +++ b/docs/architecture/rfcs/long-horizon-harness-benchmark-research-program-v0.md @@ -958,8 +958,11 @@ model planning through the product's `todo plan` checkpoint. Planning runs before the selected driver, uses the shared Goal planner/Todo-delta contract, and consumes the phase budget without counting as advancement. Qualification must read back task identity and actual Todos, preserve blocked state and -disclose the separate planning session; synthetic task success alone does not -prove equivalence to interactive `$loopx` startup or planning effectiveness. +disclose the context policy: heartbeat/Turn resume shares the Goal/Agent +conversation across planning and execution, while fresh uses separate sessions. +Exact native session ID/history conformance qualifies continuity, not task +quality or matched-budget gains; synthetic task success alone does not prove +equivalence to interactive `$loopx` startup or planning effectiveness. ### 11.3 Required delivery slice diff --git a/docs/integrations/deepseek-harness-connector.md b/docs/integrations/deepseek-harness-connector.md index 076dcedc7e..1048239b7e 100644 --- a/docs/integrations/deepseek-harness-connector.md +++ b/docs/integrations/deepseek-harness-connector.md @@ -122,7 +122,7 @@ resume behavior for the newly selected id remains owned by the dsh composition. `--iteration-context fresh` selects a session scoped to the current Turn key; retrying that same transaction keeps its identity. The default -`resume-if-available` retains the Goal/Agent/Todo lineage behavior described +`resume` retains the Goal/Agent/Todo lineage behavior described above. A durable Agent identity is not a reason to reuse another Todo's task packet or chat context. Independent questions should select fresh context; continuations must keep the exact task lineage and refresh explicit inputs. @@ -139,7 +139,7 @@ additional permissions. Domain task packages and artifact validators must still verify their own task identity and revision. `fresh` 按本次 Turn 选择新上下文,同一事务重试保持身份;默认 -`resume-if-available` 按 Goal/Agent/Todo 延续。Agent 可以长期存在,但独立 +`resume` 按 Goal/Agent/Todo 延续。Agent 可以长期存在,但独立 问题应使用新上下文,不能把另一个 Todo 的旧任务包当成交接。宿主将上述四个 本次调用身份变量传给运行时工具,覆盖陈旧值;无 Todo 时为空。这是覆盖映射, 不是环境隔离边界:当前 pin 的 SDK 先继承父进程环境,再应用映射,因此父进程 diff --git a/docs/reference/protocols/loopx-turn-v0.md b/docs/reference/protocols/loopx-turn-v0.md index a1d45d461c..7b28e209ce 100644 --- a/docs/reference/protocols/loopx-turn-v0.md +++ b/docs/reference/protocols/loopx-turn-v0.md @@ -208,12 +208,13 @@ Before wiring Trae CLI, Codex CLI, or another host, answer these five questions: dedicated result file. Do not scrape arbitrary conversation text as the completion contract. 3. **What is its resume handle?** Keep the opaque handle in local adapter - state, keyed by `(goal_id, agent_id, todo_id)`. Never put it in LoopX state - or public evidence. + state, keyed by the declared scope: Goal/Agent for Codex exec by default, + or Goal/Agent/Todo for isolated Todo context. Keep it out of public Goal + narratives and evidence. 4. **Which failures may resume?** A bounded timeout or lost transport may - preserve an observed session. A rejected startup contract, incompatible - host version, or missing session invalidates it so the next Turn starts - cleanly. + preserve an observed session. For agent-scoped Codex exec, a rejected + startup contract, incompatible host version or missing session requires + repair or explicit fresh context; it cannot silently fork. 5. **What proves the work independently?** Name a command that checks the real repository, artifact, service readback, document revision, or other postcondition without trusting the agent CLI's own claim. @@ -232,7 +233,8 @@ A thin adapter can be implemented with this host-neutral algorithm: ```text request = read_one_json(stdin) todo = request.turn_envelope.action.selected_todo -session = load_local_session(goal_id, agent_id, todo.todo_id) +scope = request.session.context_policy.get("binding_scope", "todo") +session = load_local_session(goal_id, agent_id, scope, todo.todo_id) prompt = render_bounded_prompt(todo, request.result_contract, temporary_result_path) invoke_agent_cli(prompt, workspace, session, explicit_timeout) candidate = read_and_shape_temporary_result(temporary_result_path) @@ -762,10 +764,10 @@ Session recovery is fail-closed: | Host observation | Session disposition | Next Turn | | --- | --- | --- | -| Typed result returned | Keep the opaque session eligible. | Resume when the same todo remains selected. | +| Typed result returned | Keep the opaque session eligible. | Resume within the selected conversation scope; refresh current Todo inputs. | | Timeout or transport loss after a session was observed | Keep it eligible, but do not infer progress. | Retry the side-effect-safe host phase. | -| Incompatible host version or rejected startup/output contract | Invalidate it. | Start a fresh session after repair. | -| Host reports the session is missing | Invalidate it. | Start a fresh session if policy still allows execution. | +| Incompatible host version or rejected startup/output contract | Agent-scoped Codex exec retains the binding and fails closed; Todo scope invalidates it. | Repair, then explicitly select fresh for agent scope. | +| Host reports the session is missing | Agent-scoped Codex exec retains the binding and fails closed; Todo scope invalidates it. | Repair history or explicitly select fresh for agent scope. | | Failure before any session was observed | Store nothing. | Re-decide, then start fresh only if allowed. | Session eligibility is recovery metadata, not evidence that work happened. It @@ -777,13 +779,26 @@ writeback ordering. Each `turn plan` or `turn run-once` invocation declares an iteration context policy independently from the Todo and Goal lifecycle: -- `resume-if-available` preserves the existing behavior and resumes a compatible - opaque Host Session for the same Goal, Agent, and Todo; -- `fresh` ignores a compatible saved session for this invocation and starts a - clean Host Session. Selecting `fresh` does not itself delete the prior binding; - after a successful host start, the newly observed session becomes the eligible - binding for later iterations. It does not imply a new Todo, successor, retry, - or Goal. +- `resume` is the CLI default. The first invocation starts a session; later + invocations reuse its exact native ID. The retired context name is rejected. +- Codex exec defaults to `--session-scope agent`: the same Goal and Agent keep + their conversation when Todos change. `--session-scope todo` explicitly + isolates conversations per Todo. These scopes use separate persistence keys; + existing Todo-scoped bindings are not automatically adopted into agent scope. +- `fresh` starts a clean session. After a native ID is observed, that session + replaces the binding for the selected scope, including on timeout. It implies + no new Todo, successor, retry, progress or Goal. + +The Codex exec binding includes the exact Goal lifetime where available and a +profile digest for agent scope: workspace, Codex home/settings, executable, +model, effort, sandbox and MCP configuration. A corrupt/incompatible binding, +missing native history or unexpected resumed ID fails closed; it cannot silently +fork. Explicit `fresh` is the recovery operation. Current Turn selection, task +lease, validation and settlement remain Todo-bound. Session scope grants no +additional tool or effect authority. Managed operation-equipped app-server +sessions retain their existing Todo-bound approval/handoff contract; this CLI +option applies to exec conversations. Other hosts retain their adapter's +Todo-based session contract. Use a new `turn_instance_id` for each new iteration. Reuse the same id only for an explicit replay or failed-Turn recovery. The context policy controls Host diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index 7ce1473462..80b4d9b79b 100644 --- a/loopx/cli_commands/turn.py +++ b/loopx/cli_commands/turn.py @@ -157,11 +157,16 @@ def handle_turn_command( decision, scheduler_execution_context=scheduler_context, ) + # Operation-equipped app-server sessions retain their Todo-bound + # approval/handoff contract; this option controls exec conversation reuse. + session_scope = (args.session_scope if args.host == "codex-cli" + and not getattr(args, "codex_operation_tools", False) else "todo") if ( args.turn_command == "run-once" and args.host == "codex-cli" and not resume_requested and not args.resume_turn_key + and args.iteration_context != "fresh" and turn_envelope.get("effective_action") != EffectiveAction.GOVERNED_CAPABILITY_INTENT.value ): session_binding = ( @@ -169,9 +174,10 @@ def handle_turn_command( runtime_root, turn_envelope, goal_admission=strict_goal_admission, + session_scope=session_scope, ) if strict_goal_admission is not None - else codex_cli_session_binding(runtime_root, turn_envelope) + else codex_cli_session_binding(runtime_root, turn_envelope, session_scope=session_scope) ) payload = build_loopx_turn_plan( turn_envelope, @@ -180,7 +186,8 @@ def handle_turn_command( scheduler_owner=args.scheduler_owner, session_binding=session_binding, turn_instance_id=args.turn_instance_id, - iteration_context_policy=args.iteration_context.replace("-", "_"), + iteration_context_policy=args.iteration_context, + session_scope=session_scope, goal_ref=goal_ref, ) # Resolve machine authentication once for the readback and host launch. diff --git a/loopx/cli_commands/turn_inspection.py b/loopx/cli_commands/turn_inspection.py index 42d16b5186..4ce4f94639 100644 --- a/loopx/cli_commands/turn_inspection.py +++ b/loopx/cli_commands/turn_inspection.py @@ -11,6 +11,7 @@ LOOPX_TURN_JOURNAL_INSPECTION_SCHEMA_VERSION, codex_cli_session_binding, inspect_loopx_turn_journal, + load_loopx_turn_plan_from_journal, ) from .turn_rendering import render_loopx_turn_journal_inspection_markdown @@ -36,6 +37,12 @@ def handle_turn_journal_inspection( registry_path=registry_path, runtime_root_override=runtime_root_arg, ) + scope = "todo" + if args.retry_failed_turn: + plan = load_loopx_turn_plan_from_journal( + runtime_root, goal_id=args.goal_id, turn_key=args.turn_key, + ) + scope = str(((plan.get("session") or {}).get("context_policy") or {}).get("binding_scope") or "todo") payload = inspect_loopx_turn_journal( runtime_root, goal_id=args.goal_id, @@ -46,6 +53,7 @@ def handle_turn_journal_inspection( lambda turn_envelope: codex_cli_session_binding( runtime_root, turn_envelope, + session_scope=scope, ) ), ) diff --git a/loopx/cli_commands/turn_registration.py b/loopx/cli_commands/turn_registration.py index ab857990d6..da9369406f 100644 --- a/loopx/cli_commands/turn_registration.py +++ b/loopx/cli_commands/turn_registration.py @@ -7,6 +7,7 @@ from collections.abc import Callable from ..control_plane.operator_provider import operator_provider_environ +from ..control_plane.turn_driver.driver import SessionBindingScope, SUPPORTED_ITERATION_CONTEXT_POLICIES from ..control_plane.turn_driver.host_binding import ( MANAGED_TURN_HOST, resolve_default_turn_host, @@ -383,14 +384,20 @@ def _add_turn_decision_arguments( ) parser.add_argument( "--iteration-context", - choices=["fresh", "resume-if-available"], - default="resume-if-available", + choices=sorted(SUPPORTED_ITERATION_CONTEXT_POLICIES), + default="resume", help=( "Host context policy for this iteration. fresh starts a clean " "session even when a compatible prior session exists; " - "resume-if-available preserves the existing continuation behavior." + "resume continues the recorded session." ), ) + parser.add_argument( + "--session-scope", + choices=[scope.value for scope in SessionBindingScope], + default=SessionBindingScope.AGENT.value, + help="Codex exec conversation binding: todo isolates each Todo; agent continues across Todos in the same Goal/Agent. Turn authority and settlement stay Todo-scoped.", + ) parser.add_argument( "--resume-goal-id", help="Goal identity bound to an available opaque host session.", diff --git a/loopx/cli_commands/turn_run_once.py b/loopx/cli_commands/turn_run_once.py index 355c3c6ecd..653bb708db 100644 --- a/loopx/cli_commands/turn_run_once.py +++ b/loopx/cli_commands/turn_run_once.py @@ -801,16 +801,19 @@ def run_built_in_host( def resolve_built_in_session_binding( turn_envelope: Mapping[str, Any], ) -> dict[str, str] | None: + session_scope = str((payload.get("session", {}).get("context_policy") or {}).get("binding_scope") or "todo") return ( codex_cli_session_binding( runtime_root, turn_envelope, goal_admission=goal_admission, + session_scope=session_scope, ) if goal_admission is not None else codex_cli_session_binding( runtime_root, turn_envelope, + session_scope=session_scope, ) ) diff --git a/loopx/control_plane/collaboration/operation_handoff.py b/loopx/control_plane/collaboration/operation_handoff.py index 2bada69a3f..0b20532e06 100644 --- a/loopx/control_plane/collaboration/operation_handoff.py +++ b/loopx/control_plane/collaboration/operation_handoff.py @@ -219,7 +219,7 @@ def agent_operation_action( # Session writes and operation commits share Goal -> registry -> session # -> action-store order. No lock is held across a domain effect. from ...file_lock import exclusive_file_lock - from ..turn_driver.codex_cli import _session_path + from ..turn_driver.codex_sessions import _session_path executor = parameters["executor"] session_paths = set() diff --git a/loopx/control_plane/turn_driver/codex_cli.py b/loopx/control_plane/turn_driver/codex_cli.py index cdd0610ff9..d2e5ab370c 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -2,7 +2,6 @@ from __future__ import annotations -import hashlib import json import os import re @@ -12,13 +11,23 @@ from pathlib import Path from typing import Any -from ...runtime import validate_goal_id_path_segment -from ...file_lock import exclusive_file_lock from ..goals.first_party_host_admission import FirstPartyHostGoalAdmission from .subagent_execution_topology import ( child_execution_receipts_json_schema, ) -from .driver import selected_turn_todo +from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES +from .codex_sessions import ( + CODEX_CLI_SESSION_SCHEMA_VERSION as CODEX_CLI_SESSION_SCHEMA_VERSION, + _discard_codex_cli_session, + _lineage, + _store_codex_cli_session, + _valid_session_id, + codex_cli_session_binding as codex_cli_session_binding, + load_codex_cli_session as load_codex_cli_session, + select_codex_cli_session, + codex_session_profile_digest, + require_codex_session_profile, +) from .executor import ( HOST_AGENT_VISION_JSON_MAX_CHARS, HOST_REWARD_MEMORY_REFLECTION_JSON_MAX_CHARS, @@ -31,7 +40,6 @@ from .transaction import LOOPX_TURN_RESULT_SCHEMA_VERSION, TRANSACTION_PHASES -CODEX_CLI_SESSION_SCHEMA_VERSION = "loopx_codex_cli_session_v1" CODEX_STDIO_MCP_SERVER_SCHEMA_VERSION = "codex_stdio_mcp_server_v0" CODEX_CLI_RESULT_KINDS = ( "validated_progress", @@ -42,7 +50,6 @@ "iteration_failed", ) CODEX_CLI_SANDBOXES = ("read-only", "workspace-write", "danger-full-access") -SESSION_ID_MAX_CHARS = 256 OUTPUT_DRAIN_TIMEOUT_SECONDS = 2.0 SESSION_INVALIDATING_FAILURE_CATEGORIES = frozenset( { @@ -173,219 +180,6 @@ def _codex_mcp_config_arguments( return [item for pair in pairs for item in ("-c", pair)] -def _lineage(request: Mapping[str, Any]) -> dict[str, str]: - envelope = _mapping(request.get("turn_envelope")) - todo = selected_turn_todo(envelope) - lineage = { - "goal_id": str(envelope.get("goal_id") or "").strip(), - "agent_id": str(envelope.get("agent_id") or "").strip(), - "todo_id": str(todo.get("todo_id") or "").strip(), - } - if not all(lineage.values()): - raise ValueError("Codex CLI host request has incomplete turn lineage") - lineage["goal_id"] = validate_goal_id_path_segment(lineage["goal_id"]) - return lineage - - -def _session_path(runtime_root: Path, lineage: Mapping[str, str]) -> Path: - digest = hashlib.sha256( - json.dumps( - dict(lineage), - ensure_ascii=False, - sort_keys=True, - separators=(",", ":"), - ).encode("utf-8") - ).hexdigest() - return ( - runtime_root - / "goals" - / validate_goal_id_path_segment(lineage["goal_id"]) - / "turn-sessions" - / f"{digest}.json" - ) - - -def _valid_session_id(value: Any) -> str | None: - session_id = str(value or "").strip() - if not session_id or len(session_id) > SESSION_ID_MAX_CHARS: - return None - if any(character in session_id for character in ("\x00", "\r", "\n")): - return None - return session_id - - -def load_codex_cli_session( - 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 None - if not isinstance(value, dict): - return None - if value.get("schema_version") != CODEX_CLI_SESSION_SCHEMA_VERSION: - return None - if any(value.get(field) != lineage[field] for field in lineage): - return None - session_id = _valid_session_id(value.get("session_id")) - if not session_id: - return None - 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 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", - **lineage, - } - - -def _store_codex_cli_session( - runtime_root: Path, - *, - lineage: Mapping[str, str], - session_id: str, - goal_ref: Mapping[str, Any] | None = None, - operation_profile_digest: str | None = None, - operation_model: str | None = None, - operation_reasoning_effort: str | None = None, -) -> None: - with exclusive_file_lock(_session_path(runtime_root, lineage)): - _write_codex_cli_session( - runtime_root, - lineage=lineage, - session_id=session_id, - goal_ref=goal_ref, - operation_profile_digest=operation_profile_digest, - operation_model=operation_model, - operation_reasoning_effort=operation_reasoning_effort, - ) - - -def _write_codex_cli_session( - runtime_root: Path, - *, - lineage: Mapping[str, str], - session_id: str, - goal_ref: Mapping[str, Any] | None = None, - operation_profile_digest: str | None = None, - operation_model: str | None = None, - operation_reasoning_effort: str | None = None, -) -> None: - normalized_session_id = _valid_session_id(session_id) - if not normalized_session_id: - raise ValueError("Codex CLI returned an invalid session id") - path = _session_path(runtime_root, lineage) - path.parent.mkdir(parents=True, exist_ok=True) - descriptor, temporary_name = tempfile.mkstemp( - prefix=f".{path.name}.", dir=path.parent - ) - temporary = Path(temporary_name) - try: - if hasattr(os, "fchmod"): - os.fchmod(descriptor, 0o600) - 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) - if operation_profile_digest is not None: - payload["operation_transport"] = "app-server-operation-tools-v0" - payload["operation_profile_digest"] = operation_profile_digest - payload["operation_model"] = operation_model - payload["operation_reasoning_effort"] = operation_reasoning_effort - json.dump( - payload, - handle, - ensure_ascii=False, - indent=2, - sort_keys=True, - ) - handle.write("\n") - os.replace(temporary, path) - path.chmod(0o600) - finally: - if descriptor >= 0: - os.close(descriptor) - temporary.unlink(missing_ok=True) - - -def _discard_codex_cli_session( - runtime_root: Path, - *, - lineage: Mapping[str, str], -) -> None: - path = _session_path(runtime_root, lineage) - with exclusive_file_lock(path): - path.unlink(missing_ok=True) - - def _has_subagent_topology(request: Mapping[str, Any] | None) -> bool: return bool( isinstance(request, Mapping) @@ -789,8 +583,8 @@ def _codex_command( *, codex_bin: str, project: Path, - schema_path: Path, - output_path: Path, + schema_path: Path | None, + output_path: Path | None, sandbox: str, model: str | None, reasoning_effort: str | None, @@ -807,10 +601,6 @@ def _codex_command( "resume", "-c", f'sandbox_mode="{sandbox}"', - "--output-schema", - str(schema_path), - "--output-last-message", - str(output_path), "--json", ] else: @@ -822,12 +612,12 @@ def _codex_command( sandbox, "-C", str(project), - "--output-schema", - str(schema_path), - "--output-last-message", - str(output_path), "--json", ] + if schema_path is not None: + command.extend(["--output-schema", str(schema_path)]) + if output_path is not None: + command.extend(["--output-last-message", str(output_path)]) if model: command.extend(["--model", model]) if reasoning_effort: @@ -871,25 +661,23 @@ def run_codex_cli_host( planned_session = _mapping(request.get("session")) planned_action = str(planned_session.get("action") or "") context_policy = _mapping(planned_session.get("context_policy")) + if context_policy.get("mode") is not None and context_policy["mode"] not in SUPPORTED_ITERATION_CONTEXT_POLICIES: + raise ValueError("iteration context policy must be fresh or resume") fresh_iteration = context_policy.get("mode") == "fresh" - 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 + session_scope = str(context_policy.get("binding_scope") or "todo") + if fresh_iteration and goal_admission is not None: + goal_admission.require_current() + binding = None if fresh_iteration else select_codex_cli_session( + runtime_root, lineage=lineage, session_scope=session_scope, + goal_admission=goal_admission, + ) + profile_digest = (codex_session_profile_digest( + project=project, codex_bin=str(resolved), + home=Path(os.environ.get("CODEX_HOME", "~/.codex")).expanduser(), + model=model, reasoning_effort=reasoning_effort, sandbox=sandbox, mcp_server=mcp_server, + ) if session_scope == "agent" else None) + if binding and profile_digest: + require_codex_session_profile(binding, profile_digest) 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: @@ -910,6 +698,8 @@ def commit() -> None: runtime_root, lineage=lineage, session_id=observed_session_id, + session_scope=session_scope, + session_profile_digest=profile_digest, goal_ref=exact_goal_ref, ) @@ -923,6 +713,7 @@ def commit() -> None: _discard_codex_cli_session( runtime_root, lineage=lineage, + session_scope=session_scope, ) if goal_admission is None: @@ -964,7 +755,7 @@ def observe_event(line: str) -> None: return if isinstance(event, dict): candidate = codex_cli_event_session_id(event) - if candidate and not observed_session: + if candidate and candidate not in observed_session: observed_session.append(candidate) structured, diagnostic = _event_failure_categories(event) if structured: @@ -989,6 +780,12 @@ def observe_stderr(line: str) -> None: output_observation_incomplete = not (observed["output_complete"] and events.complete and diagnostics.complete) if observed["outcome"] not in {"exited", "timeout"}: raise BuiltInHostError("codex_cli_process_" + observed["outcome"]) + if session_scope == "agent" and ( + len(observed_session) > 1 or (returncode == 0 and not observed_session) + ): + raise BuiltInHostError("codex_cli_session_identity_unconfirmed") + if session_id and any(candidate != session_id for candidate in observed_session): + raise BuiltInHostError("codex_cli_resume_session_identity_changed") if timed_out: if observed_session: store_session(observed_session[0]) @@ -1006,7 +803,7 @@ def observe_stderr(line: str) -> None: or "exit_nonzero" ) ) - if returncode != 0 and category in SESSION_INVALIDATING_FAILURE_CATEGORIES: + if session_scope == "todo" and returncode != 0 and category in SESSION_INVALIDATING_FAILURE_CATEGORIES: discard_session() if observed_session and ( returncode == 0 or category not in SESSION_INVALIDATING_FAILURE_CATEGORIES diff --git a/loopx/control_plane/turn_driver/codex_operation_host.py b/loopx/control_plane/turn_driver/codex_operation_host.py index 1643a2536c..5135d93d5d 100644 --- a/loopx/control_plane/turn_driver/codex_operation_host.py +++ b/loopx/control_plane/turn_driver/codex_operation_host.py @@ -26,14 +26,16 @@ from ..goals.first_party_host_admission import FirstPartyHostGoalAdmission from ..effect_runtime import EffectRuntimeRejected from .codex_cli import ( - _lineage, _prompt, + normalize_codex_stdio_mcp_server, + codex_cli_result_schema, +) +from .codex_sessions import ( + _lineage, _read_codex_cli_session_document, _codex_session_goal_ref, _store_codex_cli_session, load_codex_cli_session, - normalize_codex_stdio_mcp_server, - codex_cli_result_schema, ) from .executor import LOOPX_TURN_HOST_REQUEST_SCHEMA_VERSION from .host_failure import BuiltInHostError diff --git a/loopx/control_plane/turn_driver/codex_sessions.py b/loopx/control_plane/turn_driver/codex_sessions.py new file mode 100644 index 0000000000..fdf6664244 --- /dev/null +++ b/loopx/control_plane/turn_driver/codex_sessions.py @@ -0,0 +1,335 @@ +"""Codex transport session persistence; context scope never grants Turn authority.""" + +from __future__ import annotations + +import hashlib +import json +import os +import tempfile +from collections.abc import Mapping +from pathlib import Path +from typing import Any + +from ...file_lock import exclusive_file_lock +from ...runtime import validate_goal_id_path_segment +from ..goals.first_party_host_admission import FirstPartyHostGoalAdmission +from .driver import selected_turn_todo, session_identity_fields + +CODEX_CLI_SESSION_SCHEMA_VERSION = "loopx_codex_cli_session_v1" +SESSION_ID_MAX_CHARS = 256 + + +def _lineage(request: Mapping[str, Any]) -> dict[str, str]: + envelope = request.get("turn_envelope") or {} + todo = selected_turn_todo(envelope) + lineage = { + "goal_id": str(envelope.get("goal_id") or "").strip(), + "agent_id": str(envelope.get("agent_id") or "").strip(), + "todo_id": str(todo.get("todo_id") or "").strip(), + } + if not all(lineage.values()): + raise ValueError("Codex CLI host request has incomplete turn lineage") + lineage["goal_id"] = validate_goal_id_path_segment(lineage["goal_id"]) + return lineage + + +def _session_path( + runtime_root: Path, + lineage: Mapping[str, str], + session_scope: str = "todo", +) -> Path: + fields = session_identity_fields(session_scope) + identity = ( + dict(lineage) + if session_scope == "todo" + else { + **{field: lineage[field] for field in fields}, + "session_scope": session_scope, + } + ) + digest = hashlib.sha256( + json.dumps( + identity, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + ).encode("utf-8") + ).hexdigest() + goal_id: str = validate_goal_id_path_segment(lineage["goal_id"]) + return runtime_root / "goals" / goal_id / "turn-sessions" / f"{digest}.json" + + +def _valid_session_id(value: Any) -> str | None: + session_id = str(value or "").strip() + if not session_id or len(session_id) > SESSION_ID_MAX_CHARS: + return None + if any(character in session_id for character in ("\x00", "\r", "\n")): + return None + return session_id + + +def load_codex_cli_session( + runtime_root: Path, + *, + lineage: Mapping[str, str], + session_scope: str = "todo", +) -> dict[str, Any] | None: + path = _session_path(runtime_root, lineage, session_scope) + if not path.exists(): + return None + try: + value = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + return None + if not isinstance(value, dict): + return None + if value.get("schema_version") != CODEX_CLI_SESSION_SCHEMA_VERSION: + return None + if value.get("session_scope", "todo") != session_scope: + return None + if any( + value.get(field) != lineage[field] + for field in session_identity_fields(session_scope) + ): + return None + session_id = _valid_session_id(value.get("session_id")) + if not session_id: + return None + return {**value, "session_id": session_id} + + +def _codex_session_goal_ref( + value: Mapping[str, Any], + *, + lineage: Mapping[str, str], + session_scope: str = "todo", +) -> object: + if ( + value.get("schema_version") != CODEX_CLI_SESSION_SCHEMA_VERSION + or value.get("session_scope", "todo") != session_scope + or any( + value.get(field) != lineage[field] + for field in session_identity_fields(session_scope) + ) + 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], + session_scope: str = "todo", +) -> dict[str, Any] | None: + path = _session_path(runtime_root, lineage, session_scope) + 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, + session_scope: str = "todo", +) -> dict[str, str] | None: + lineage = _lineage({"turn_envelope": dict(turn_envelope)}) + session = select_codex_cli_session( + runtime_root, + lineage=lineage, + session_scope=session_scope, + goal_admission=goal_admission, + ) + if session is None: + return None + return { + "schema_version": "loopx_turn_session_binding_v0", + **lineage, + } + + +def _store_codex_cli_session( + runtime_root: Path, + *, + lineage: Mapping[str, str], + session_scope: str = "todo", + session_id: str, + goal_ref: Mapping[str, Any] | None = None, + session_profile_digest: str | None = None, + operation_profile_digest: str | None = None, + operation_model: str | None = None, + operation_reasoning_effort: str | None = None, +) -> None: + with exclusive_file_lock(_session_path(runtime_root, lineage, session_scope)): + _write_codex_cli_session( + runtime_root, + lineage=lineage, + session_scope=session_scope, + session_profile_digest=session_profile_digest, + session_id=session_id, + goal_ref=goal_ref, + operation_profile_digest=operation_profile_digest, + operation_model=operation_model, + operation_reasoning_effort=operation_reasoning_effort, + ) + + +def _write_codex_cli_session( + runtime_root: Path, + *, + lineage: Mapping[str, str], + session_scope: str = "todo", + session_id: str, + goal_ref: Mapping[str, Any] | None = None, + session_profile_digest: str | None = None, + operation_profile_digest: str | None = None, + operation_model: str | None = None, + operation_reasoning_effort: str | None = None, +) -> None: + normalized_session_id = _valid_session_id(session_id) + if not normalized_session_id: + raise ValueError("Codex CLI returned an invalid session id") + path = _session_path(runtime_root, lineage, session_scope) + path.parent.mkdir(parents=True, exist_ok=True) + descriptor, temporary_name = tempfile.mkstemp( + prefix=f".{path.name}.", dir=path.parent + ) + temporary = Path(temporary_name) + try: + if hasattr(os, "fchmod"): + os.fchmod(descriptor, 0o600) + handle = os.fdopen(descriptor, "w", encoding="utf-8") + descriptor = -1 + with handle: + payload: dict[str, Any] = { + "schema_version": CODEX_CLI_SESSION_SCHEMA_VERSION, + **lineage, + "host": "codex-cli", + "session_id": normalized_session_id, + } + if session_scope != "todo": + payload["session_scope"] = session_scope + if session_profile_digest is not None: + payload["session_profile_digest"] = session_profile_digest + if goal_ref is not None: + payload["goal_ref"] = dict(goal_ref) + if operation_profile_digest is not None: + payload["operation_transport"] = "app-server-operation-tools-v0" + payload["operation_profile_digest"] = operation_profile_digest + payload["operation_model"] = operation_model + payload["operation_reasoning_effort"] = operation_reasoning_effort + json.dump( + payload, + handle, + ensure_ascii=False, + indent=2, + sort_keys=True, + ) + handle.write("\n") + os.replace(temporary, path) + path.chmod(0o600) + finally: + if descriptor >= 0: + os.close(descriptor) + temporary.unlink(missing_ok=True) + + +def _discard_codex_cli_session( + runtime_root: Path, + *, + lineage: Mapping[str, str], + session_scope: str = "todo", +) -> None: + path = _session_path(runtime_root, lineage, session_scope) + with exclusive_file_lock(path): + path.unlink(missing_ok=True) + + +def select_codex_cli_session( + runtime_root: Path, + *, + lineage: Mapping[str, str], + session_scope: str = "todo", + goal_admission: FirstPartyHostGoalAdmission | None = None, +) -> dict[str, Any] | None: + """Select through the existing Goal lifetime fence, then validate identity. + + A present but invalid agent-scoped record is an error, never permission to + silently create a replacement conversation. + """ + session_identity_fields(session_scope) + + def read() -> dict[str, Any] | None: + return _read_codex_cli_session_document( + runtime_root, + lineage=lineage, + session_scope=session_scope, + ) + + value = ( + goal_admission.select_state( + read_state=read, + goal_ref_of=lambda item: _codex_session_goal_ref( + item, + lineage=lineage, + session_scope=session_scope, + ), + ) + if goal_admission is not None + else read() + ) + if value is None: + return None + if _codex_session_goal_ref(value, lineage=lineage, session_scope=session_scope) == { + "malformed": True + }: + if session_scope == "agent": + raise ValueError( + "invalid agent session binding; explicitly select fresh to replace it" + ) + return None + return dict(value) + + +def codex_session_profile_digest( + *, + project: Path, + codex_bin: str, + home: Path, + model: str | None = None, + reasoning_effort: str | None = None, + sandbox: str | None = None, + mcp_server: Mapping[str, Any] | None = None, +) -> str: + """Bind a shared conversation to its workspace, provider home and settings.""" + config = home / "config.toml" + profile = { + "workspace": str(project.resolve()), + "codex_home": str(home.expanduser().resolve()), + "codex_bin": str(Path(codex_bin).resolve()), + "model": model, + "reasoning_effort": reasoning_effort, + "sandbox": sandbox, + "mcp_server": mcp_server, + "config_digest": hashlib.sha256( + config.read_bytes() if config.exists() else b"" + ).hexdigest(), + } + return hashlib.sha256(json.dumps(profile, sort_keys=True).encode()).hexdigest() + + +def require_codex_session_profile(binding: Mapping[str, Any], digest: str) -> None: + if binding.get("session_profile_digest") != digest: + raise ValueError( + "Codex session profile changed; explicitly select fresh or use a new runtime" + ) diff --git a/loopx/control_plane/turn_driver/driver.py b/loopx/control_plane/turn_driver/driver.py index 9ffa6f1a4f..5d6c643adc 100644 --- a/loopx/control_plane/turn_driver/driver.py +++ b/loopx/control_plane/turn_driver/driver.py @@ -4,6 +4,7 @@ import json from collections.abc import Mapping +from enum import StrEnum from hashlib import sha256 from typing import Any @@ -34,7 +35,22 @@ LOOPX_CHILD_HOST_OPERATION_SCHEMA_VERSION = "loopx_child_host_operation_v0" SUPPORTED_HOSTS = {"codex-cli", "claude-code", "dsh", "generic-cli"} SUPPORTED_EXECUTION_MODES = {"interactive-visible", "isolated-headless"} -SUPPORTED_ITERATION_CONTEXT_POLICIES = {"fresh", "resume_if_available"} +SUPPORTED_ITERATION_CONTEXT_POLICIES = {"fresh", "resume"} + + +class SessionBindingScope(StrEnum): + TODO = "todo" + AGENT = "agent" + + +def session_identity_fields(scope: str) -> tuple[str, ...]: + """Conversation identity is independent of a Turn's settlement identity.""" + selected = SessionBindingScope(scope) + if selected is SessionBindingScope.TODO: + return ("goal_id", "agent_id", "todo_id") + return ("goal_id", "agent_id") + + REPLAN_ACTIONS = { "autonomous_replan", "autonomous_replan_required", @@ -221,7 +237,15 @@ def _session_plan( lineage: Mapping[str, str], session_binding: Mapping[str, Any] | None, iteration_context_policy: str, + session_scope: str = "agent", ) -> tuple[dict[str, Any], str | None]: + identity_fields = session_identity_fields(session_scope) + context_policy = { + "schema_version": LOOPX_ITERATION_CONTEXT_POLICY_SCHEMA_VERSION, + "mode": iteration_context_policy, + "scope": "iteration", + **({"binding_scope": session_scope} if session_scope != "todo" else {}), + } host_route = route in { LoopXTurnRoute.READY_FOR_HOST, LoopXTurnRoute.REPAIR_REQUIRED, @@ -232,11 +256,7 @@ def _session_plan( "schema_version": LOOPX_TURN_SESSION_BINDING_SCHEMA_VERSION, "action": "none", "binding_status": "not_applicable", - "context_policy": { - "schema_version": LOOPX_ITERATION_CONTEXT_POLICY_SCHEMA_VERSION, - "mode": iteration_context_policy, - "scope": "iteration", - }, + "context_policy": context_policy, }, None if not all(lineage.values()): return { @@ -252,11 +272,7 @@ def _session_plan( "binding_status": ( "existing_binding_ignored" if session_binding else "not_found" ), - "context_policy": { - "schema_version": LOOPX_ITERATION_CONTEXT_POLICY_SCHEMA_VERSION, - "mode": "fresh", - "scope": "iteration", - }, + "context_policy": context_policy, }, None binding = dict(session_binding or {}) @@ -264,11 +280,7 @@ def _session_plan( return { "schema_version": LOOPX_TURN_SESSION_BINDING_SCHEMA_VERSION, "action": "start_new", - "context_policy": { - "schema_version": LOOPX_ITERATION_CONTEXT_POLICY_SCHEMA_VERSION, - "mode": "resume_if_available", - "scope": "iteration", - }, + "context_policy": context_policy, }, None if binding.get("schema_version") != LOOPX_TURN_SESSION_BINDING_SCHEMA_VERSION: return { @@ -277,7 +289,6 @@ def _session_plan( "binding_status": "unsupported_schema", }, "unsupported LoopX Turn session binding schema" - identity_fields = ("goal_id", "agent_id", "todo_id") actual = {field: str(binding.get(field) or "") for field in identity_fields} expected = {field: lineage[field] for field in identity_fields} if actual != expected: @@ -285,16 +296,12 @@ def _session_plan( "schema_version": LOOPX_TURN_SESSION_BINDING_SCHEMA_VERSION, "action": "reject", "binding_status": "identity_mismatch", - }, "session binding does not match the current goal, agent, and todo" + }, "session binding does not match the current " + ", ".join(identity_fields) return { "schema_version": LOOPX_TURN_SESSION_BINDING_SCHEMA_VERSION, "action": "resume", "binding_status": "compatible", - "context_policy": { - "schema_version": LOOPX_ITERATION_CONTEXT_POLICY_SCHEMA_VERSION, - "mode": "resume_if_available", - "scope": "iteration", - }, + "context_policy": context_policy, }, None @@ -308,11 +315,13 @@ def reconcile_failed_turn_session_request( envelope = _mapping(request.get("turn_envelope")) selected_todo = selected_turn_todo(envelope) lineage = _turn_lineage(envelope, selected_todo=selected_todo) + planned_session = _mapping(request.get("session")) session, session_error = _session_plan( route=LoopXTurnRoute.READY_FOR_HOST, lineage=lineage, session_binding=session_binding, - iteration_context_policy="resume_if_available", + iteration_context_policy="resume", + session_scope=str(_mapping(planned_session.get("context_policy")).get("binding_scope") or "todo"), ) if session_error: status = str(session.get("binding_status") or "") @@ -331,7 +340,6 @@ def reconcile_failed_turn_session_request( "failed-Turn recovery requires a compatible persisted session binding", ) - planned_session = _mapping(request.get("session")) planned_action = str(planned_session.get("action") or "") if planned_action not in {"start_new", "resume"}: raise FailedTurnSessionRecoveryError( @@ -434,7 +442,8 @@ def build_loopx_turn_plan( scheduler_owner: str | None = None, session_binding: Mapping[str, Any] | None = None, turn_instance_id: str | None = None, - iteration_context_policy: str = "resume_if_available", + iteration_context_policy: str = "resume", + session_scope: str | None = None, goal_ref: Mapping[str, str] | None = None, ) -> dict[str, Any]: """Project a TurnEnvelope into a typed, side-effect-free host decision.""" @@ -445,8 +454,13 @@ def build_loopx_turn_plan( raise ValueError(f"unsupported LoopX Turn execution mode: {execution_mode}") if iteration_context_policy not in SUPPORTED_ITERATION_CONTEXT_POLICIES: raise ValueError( - "iteration context policy must be fresh or resume_if_available" + "iteration context policy must be fresh or resume" ) + if session_scope is None: + session_scope = "agent" if host == "codex-cli" else "todo" + session_identity_fields(session_scope) + if session_scope != "todo" and host != "codex-cli": + raise ValueError("agent session scope requires the codex-cli host") execution_context = scheduler_execution_context_for_turn( host=host, @@ -471,6 +485,7 @@ def build_loopx_turn_plan( lineage=lineage, session_binding=session_binding, iteration_context_policy=iteration_context_policy, + session_scope=session_scope, ) if session_error: route = LoopXTurnRoute.CONTRACT_ERROR diff --git a/packages/loopx-ark-turn/tests/test_host.py b/packages/loopx-ark-turn/tests/test_host.py index fef9cc108d..5db02b46e1 100644 --- a/packages/loopx-ark-turn/tests/test_host.py +++ b/packages/loopx-ark-turn/tests/test_host.py @@ -194,7 +194,7 @@ def test_invalid_request_never_starts_provider(tmp_path, mutation): if mutation == "signature": req["turn_envelope"]["action"]["primary_action"] = "Forged change" elif mutation == "context": - req["session"]["context_policy"]["mode"] = "resume-if-available" + req["session"]["context_policy"]["mode"] = "resume" else: req["turn_key"] = "not-a-turn-key" provider = Provider([]) diff --git a/tests/extensions/test_lark_goal_channel_operation.py b/tests/extensions/test_lark_goal_channel_operation.py index dc4d6f2948..6d241e911e 100644 --- a/tests/extensions/test_lark_goal_channel_operation.py +++ b/tests/extensions/test_lark_goal_channel_operation.py @@ -76,7 +76,7 @@ def _prepare_agent_handoff( } parameters["operation_kind"] = "fixture.submit" if managed: - from loopx.control_plane.turn_driver.codex_cli import _store_codex_cli_session + from loopx.control_plane.turn_driver.codex_sessions import _store_codex_cli_session _store_codex_cli_session( store.root.parent.parent, diff --git a/tests/test_chat_operation_actions.py b/tests/test_chat_operation_actions.py index 9e00f62a29..8089e816d2 100644 --- a/tests/test_chat_operation_actions.py +++ b/tests/test_chat_operation_actions.py @@ -322,8 +322,8 @@ def test_managed_replacement_has_evidence_only_access_and_never_inherits_executi from contextlib import contextmanager from threading import Event, current_thread from loopx.control_plane.collaboration import operation_handoff - from loopx.control_plane.turn_driver import codex_cli - from loopx.control_plane.turn_driver.codex_cli import _discard_codex_cli_session + from loopx.control_plane.turn_driver import codex_sessions + from loopx.control_plane.turn_driver.codex_sessions import _discard_codex_cli_session service, store = _service(tmp_path) original = _managed_handler(service, store) @@ -394,7 +394,7 @@ def test_managed_replacement_has_evidence_only_access_and_never_inherits_executi "reconciles_outcome_digest": reported["outcome_digest"], } attempted, acquired = Event(), Event() - original_lock = codex_cli.exclusive_file_lock + original_lock = codex_sessions.exclusive_file_lock original_binding = operation_handoff._binding original_write = ChatActionStore._write commits = [] @@ -424,7 +424,7 @@ def record_report(self, payload): original_write(self, payload) commits.append("report") - monkeypatch.setattr(codex_cli, "exclusive_file_lock", observed_lock) + monkeypatch.setattr(codex_sessions, "exclusive_file_lock", observed_lock) monkeypatch.setattr(ChatActionStore, "_write", record_report) with ThreadPoolExecutor( max_workers=1, thread_name_prefix="managed-revoker" @@ -506,7 +506,7 @@ def test_managed_tool_rejects_actor_injection_native_mismatch_and_revoked_profil decision="confirm", confirmation=_confirmation(delivered), ) - from loopx.control_plane.turn_driver.codex_cli import _store_codex_cli_session + from loopx.control_plane.turn_driver.codex_sessions import _store_codex_cli_session _store_codex_cli_session( store.root.parent.parent, diff --git a/tests/test_codex_operation_host.py b/tests/test_codex_operation_host.py index e73ce30545..381352cf6a 100644 --- a/tests/test_codex_operation_host.py +++ b/tests/test_codex_operation_host.py @@ -13,10 +13,10 @@ import pytest from loopx.control_plane.turn_driver.codex_cli import ( - _lineage, load_codex_cli_session, run_codex_cli_host, ) +from loopx.control_plane.turn_driver.codex_sessions import _lineage from loopx.control_plane.turn_driver.codex_operation_host import ( run_codex_operation_host, ) @@ -167,7 +167,7 @@ def test_callback_launch_fence_rejects_drift_before_native_process_start( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, drift: str, ) -> None: from loopx.control_plane.turn_driver import codex_operation_host as owner - from loopx.control_plane.turn_driver.codex_cli import _store_codex_cli_session + from loopx.control_plane.turn_driver.codex_sessions import _store_codex_cli_session store, claimed, request, options = _claimed_native_fixture(tmp_path) if drift == "session": diff --git a/tests/test_loopx_turn_codex_cli.py b/tests/test_loopx_turn_codex_cli.py index b466e35dd5..e80ed6783d 100644 --- a/tests/test_loopx_turn_codex_cli.py +++ b/tests/test_loopx_turn_codex_cli.py @@ -512,6 +512,66 @@ def test_codex_cli_host_starts_then_resumes_opaque_session( assert "private_material" not in persisted +@pytest.mark.parametrize("scope", ["todo", "agent"]) +def test_codex_session_scope_across_todos(tmp_path, monkeypatch, scope): + executable, log_path = _fake_codex(tmp_path) + monkeypatch.setenv("FAKE_CODEX_LOG", str(log_path)) + project = tmp_path / "project" + project.mkdir() + runtime_root = tmp_path / "runtime" + first = _request() + first["session"]["context_policy"] = {"mode": "resume", "binding_scope": scope} + run_codex_cli_host(first, runtime_root=runtime_root, project=project, + codex_bin=str(executable), timeout_seconds=5) + second = _request(turn_key="sha256:" + "b" * 64, + session_action="resume" if scope == "agent" else "start_new") + second["turn_envelope"]["action"]["selected_todo"]["todo_id"] = "todo_successor" + second["session"]["context_policy"] = {"mode": "resume", "binding_scope": scope} + run_codex_cli_host(second, runtime_root=runtime_root, project=project, + codex_bin=str(executable), timeout_seconds=5) + calls = [json.loads(line) for line in log_path.read_text().splitlines()] + assert ("resume" in calls[1]) is (scope == "agent") + if scope == "agent": + assert "session-fixture-0001" in calls[1] + assert len(list(runtime_root.glob("goals/*/turn-sessions/*.json"))) == (1 if scope == "agent" else 2) + + +def test_agent_resume_rejects_changed_profile_before_launch(tmp_path, monkeypatch): + executable, log = _fake_codex(tmp_path) + monkeypatch.setenv("FAKE_CODEX_LOG", str(log)) + request = _request() + request["session"]["context_policy"] = {"mode": "resume", "binding_scope": "agent"} + options = dict(runtime_root=tmp_path / "runtime", project=tmp_path, codex_bin=str(executable)) + run_codex_cli_host(request, model="fixture-model", **options) + request["session"]["action"] = "resume" + with pytest.raises(ValueError, match="profile changed"): + run_codex_cli_host(request, model="another-model", **options) + assert len(log.read_text().splitlines()) == 1 + + +def test_agent_binding_still_requires_exact_goal_lifetime(tmp_path, monkeypatch): + executable, log = _fake_codex(tmp_path) + monkeypatch.setenv("FAKE_CODEX_LOG", str(log)) + admission = _source_admission(tmp_path) + request = _request() + request["goal_ref"] = SOURCE_GOAL_REF + request["session"]["context_policy"] = {"mode": "resume", "binding_scope": "agent"} + run_codex_cli_host(request, runtime_root=tmp_path / "runtime", project=tmp_path, + codex_bin=str(executable), goal_admission=admission) + with source_session_registry_transaction(admission.registry_path, operation="fixture_new_lifetime") as transaction: + registry = transaction.payload_copy() + registry["goals"][0]["goal_instance_id"] = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" + transaction.commit(registry) + with pytest.raises(RuntimeError, match="first-party Host runtime rejected"): + codex_cli_session_binding(tmp_path / "runtime", request["turn_envelope"], + session_scope="agent", goal_admission=admission) + request["session"]["context_policy"]["mode"] = "fresh" + with pytest.raises(RuntimeError, match="first-party Host runtime rejected"): + run_codex_cli_host(request, runtime_root=tmp_path / "runtime", project=tmp_path, + codex_bin=str(executable), goal_admission=admission) + assert len(log.read_text().splitlines()) == 1 + + def test_codex_source_session_descriptor_persists_exact_goal_ref( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, diff --git a/tests/test_loopx_turn_driver.py b/tests/test_loopx_turn_driver.py index 531060444e..862f9e3bc6 100644 --- a/tests/test_loopx_turn_driver.py +++ b/tests/test_loopx_turn_driver.py @@ -32,7 +32,7 @@ reward_memory_reflection_digest, run_loopx_turn_once, ) -from loopx.control_plane.turn_driver.codex_cli import _store_codex_cli_session +from loopx.control_plane.turn_driver.codex_sessions import _store_codex_cli_session from loopx.control_plane.turn_driver.subagent_execution_topology import ( CHILD_FALLBACK_ACTIONS, ) @@ -100,8 +100,9 @@ def test_turn_plan_projects_ready_route_without_side_effects() -> None: "action": "start_new", "context_policy": { "schema_version": "loopx_iteration_context_policy_v0", - "mode": "resume_if_available", + "mode": "resume", "scope": "iteration", + "binding_scope": "agent", }, } assert payload["transaction"]["status"] == "planned" @@ -896,6 +897,25 @@ def test_turn_plan_resumes_only_a_matching_session_binding() -> None: assert payload["boundary"]["opaque_session_handle_omitted"] is True +@pytest.mark.parametrize("field", ["goal_id", "agent_id", "todo_id"]) +def test_agent_session_scope_preserves_current_turn_identity(field) -> None: + binding = { + "schema_version": LOOPX_TURN_SESSION_BINDING_SCHEMA_VERSION, + "goal_id": "fixture-goal", "agent_id": "codex-fixture", + "todo_id": "todo_fixture0001", + } + binding[field] = "different-identity" + payload = build_loopx_turn_plan( + _envelope(), host="codex-cli", execution_mode="interactive-visible", + session_binding=binding, session_scope="agent", + ) + assert payload["ok"] is (field == "todo_id") + assert payload["session"]["action"] == ("resume" if field == "todo_id" else "reject") + if field == "todo_id": + assert payload["session"]["context_policy"]["binding_scope"] == "agent" + assert payload["transaction"]["settlement_plan"]["identity"]["todo_id"] == "todo_fixture0001" + + def test_turn_plan_rejects_session_binding_identity_drift() -> None: payload = build_loopx_turn_plan( _envelope(), @@ -918,6 +938,59 @@ def test_turn_plan_rejects_session_binding_identity_drift() -> None: assert payload["effects"]["host_invoked"] is False +def test_default_agent_session_reuses_planning_across_completed_todos(tmp_path, monkeypatch): + from benchmark.runtime.codex import Execution + from benchmark.runtime.sessions import BenchmarkSessionWake + from tests.test_loopx_turn_codex_cli import _fake_codex + + project, runtime, registry = _write_live_fixture(tmp_path, extra_agent_todo_lines=( + "- [ ] [P1] Complete the successor fixture.", + " ", + )) + binary, log = _fake_codex(tmp_path) + monkeypatch.setenv("FAKE_CODEX_LOG", str(log)) + home = tmp_path / "codex-home" + home.mkdir() + monkeypatch.setenv("CODEX_HOME", str(home)) + # A planning exec observes the native ID before either Todo is executed. + env = {"LOOPX_RUNTIME_ROOT": str(runtime), "LOOPX_REGISTRY": str(registry), + "LOOPX_GOAL_ID": "loopx-turn-fixture", "LOOPX_AGENT_ID": "codex-fixture", + "LOOPX_PROJECT": str(project), "CODEX_HOME": str(home), "CODEX_BIN": str(binary), + "MODEL_NAME": "fixture-model", "REASONING_EFFORT": "high"} + wake = BenchmarkSessionWake(env, Execution(context="resume", sandbox="read-only"), {}) + observed = tmp_path / "planning-stream.jsonl" + observed.write_text('{"type":"thread.started","thread_id":"session-fixture-0001"}\n') + wake.observe(observed) + binary.write_text(binary.read_text().replace('"validated_progress"', '"validated_completion"').replace( + 'output_path = ', 'pathlib.Path("completed-todo.txt").write_text("complete")\noutput_path = ', + )) + argv = ["--registry", str(registry), "--runtime-root", str(runtime), "--format", "json", + "turn", "run-once", "--goal-id", "loopx-turn-fixture", "--agent-id", "codex-fixture", + "--host", "codex-cli", "--project", str(project), "--scan-root", str(project), + "--codex-bin", str(binary), "--codex-model", "fixture-model", + "--codex-reasoning-effort", "high", "--codex-sandbox", "read-only", + "--no-global-sync", "--validation-command-json", + json.dumps([sys.executable, "-c", "from pathlib import Path; assert Path('completed-todo.txt').read_text() == 'complete'"]), + "--execute"] + results = [] + for number in range(2): + output = io.StringIO() + with contextlib.redirect_stdout(output): + code = cli_main([*argv, "--turn-instance-id", f"completion-{number}"]) + result = json.loads(output.getvalue()) + assert code == 0, result + assert result["status"] == "committed", result + results.append(result) + calls = [json.loads(line) for line in log.read_text().splitlines()] + assert len(calls) == 2 and all("resume" in call and "session-fixture-0001" in call for call in calls) + assert results[0]["resume_turn_key"] != results[1]["resume_turn_key"] + journals = [json.loads(path.read_text()) for path in (runtime / "goals/loopx-turn-fixture/turns").glob("*.json")] + assert {journal["writeback"]["completion"]["todo_id"] for journal in journals} == { + "todo_fixture0001", "todo_fixture0002", + } + + def test_turn_plan_transaction_key_is_stable_and_todo_scoped() -> None: first = build_loopx_turn_plan( _envelope(), @@ -1036,6 +1109,7 @@ def test_turn_plan_fresh_iteration_ignores_compatible_session_binding() -> None: "schema_version": "loopx_iteration_context_policy_v0", "mode": "fresh", "scope": "iteration", + "binding_scope": "agent", }, } assert payload["transaction"]["turn_instance_id"] == "cycle-2:iteration-1" @@ -1685,6 +1759,7 @@ def test_turn_cli_projects_explicit_fresh_iteration_context( "schema_version": "loopx_iteration_context_policy_v0", "mode": "fresh", "scope": "iteration", + "binding_scope": "agent", } @@ -2180,8 +2255,8 @@ def run(turn_instance_id: str, context: str) -> dict[str, object]: run("fresh-fixture-1", "fresh") run("fresh-fixture-2", "fresh") - run("resume-fixture-1", "resume-if-available") - run("resume-fixture-2", "resume-if-available") + run("resume-fixture-1", "resume") + run("resume-fixture-2", "resume") session_ids = (host_project / "dsh-session-ids.txt").read_text( encoding="utf-8" @@ -3523,6 +3598,7 @@ def test_turn_run_once_cli_resumes_session_from_recoverable_failed_turn( def fake_session_binding( _runtime_root: Path, _turn_envelope: dict[str, object], + **_kwargs: object, ) -> dict[str, str] | None: nonlocal session_binding_calls session_binding_calls += 1