From 8378f166d54f36961174e20e95697b8ea8c1f689 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 15 Sep 2026 17:44:09 +0800 Subject: [PATCH 1/6] fix(heartbeat): keep an unapplied automation prompt adoption on disk Update-time reconciliation can only report a prompt migration it cannot write while the Codex App is running; the pending request lived in that one command's stdout. Record it per lane under the runtime root, with the same reviewed prompt-only request update-time already builds, so the obligation survives the report that discovered it. - share automation_update_request between the plan/reconcile paths instead of rebuilding the same App request inline - resolve only this host home's records, so one runtime root fronting several Codex homes keeps its other pending adoptions - drop a record once its reviewed body is installed, without re-classifying any lane Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../heartbeat/automation_upgrade.py | 21 ++- .../heartbeat/installed_prompt_update.py | 146 ++++++++++++++++-- .../test_automation_prompt_upgrade.py | 119 ++++++++++++-- 3 files changed, 267 insertions(+), 19 deletions(-) diff --git a/loopx/control_plane/heartbeat/automation_upgrade.py b/loopx/control_plane/heartbeat/automation_upgrade.py index c1f8993fd6..86fca270d3 100644 --- a/loopx/control_plane/heartbeat/automation_upgrade.py +++ b/loopx/control_plane/heartbeat/automation_upgrade.py @@ -15,7 +15,7 @@ import sqlite3 import tempfile import tomllib -from typing import Any +from typing import Any, Mapping from types import SimpleNamespace from .bootstrap_prompt import ( @@ -229,6 +229,25 @@ def build_plan(*, registry: Path, home: Path | None = None, "policy": "Discovery is not adoption authority. Review each replacement; use the App API first."} +def automation_update_request(*, automation_id: str, manifest: Mapping[str, Any], + expected_prompt_sha256: str, desired_prompt: str) -> dict[str, Any]: + """Complete prompt-only App request; scheduling and thread binding are preserved. + + The CLI cannot call an in-App tool itself, so every entrypoint that hands a + stale installed body to its host sends this exact reviewed request. + """ + return {"tool": "automation_update", + "expected_prompt_sha256": expected_prompt_sha256, + "precondition": "View the same automation; verify this prompt hash and all " + "preserved fields before update; read back afterward.", + "arguments": {"mode": "update", "id": automation_id, "kind": "heartbeat", + "name": manifest["name"], "status": manifest["status"], + "rrule": manifest["rrule"], + "targetThreadId": manifest["target_thread_id"], + "notificationPolicy": manifest.get("notification_policy"), + "prompt": desired_prompt}} + + def apply_offline(*, home: Path, automation_id: str, expected_prompt_sha256: str, desired_prompt: str, expected_source_sha256: str | None = None) -> dict[str, Any]: """Journaled prompt-only write, also used by qualified update-time migration. diff --git a/loopx/control_plane/heartbeat/installed_prompt_update.py b/loopx/control_plane/heartbeat/installed_prompt_update.py index 94569f89f3..f465b4cc80 100644 --- a/loopx/control_plane/heartbeat/installed_prompt_update.py +++ b/loopx/control_plane/heartbeat/installed_prompt_update.py @@ -6,6 +6,7 @@ """ from __future__ import annotations +from datetime import datetime, timezone import json import os from pathlib import Path @@ -13,10 +14,117 @@ import subprocess import sys import tempfile +import tomllib +from typing import Any -from .automation_upgrade import SCHEMA, _atomic, apply_offline, bootstrap_binding, build_plan +from .automation_upgrade import ( + SCHEMA, + _atomic, + apply_offline, + automation_update_request, + bootstrap_binding, + build_plan, + digest, +) from .bootstrap_prompt import host_bootstrap_binding +PENDING_SCHEMA_VERSION = "loopx_automation_prompt_adoption_pending_v0" +PENDING_HOST_ACTION = "adopt_managed_bootstrap" +PENDING_HOST_ACTION_CONTRACT = "codex_app_automation_prompt_adoption" +PENDING_SPEND_POLICY = "no_spend_for_automation_prompt_adoption" +PENDING_STORE_FILENAME = "app-automation-prompt-adoptions.json" + + +def pending_store_path(runtime_root: str | Path) -> Path: + return Path(runtime_root).expanduser().resolve() / PENDING_STORE_FILENAME + + +def _read_pending_store(runtime_root: str | Path) -> dict[str, dict[str, Any]]: + try: + payload = json.loads(pending_store_path(runtime_root).read_text(encoding="utf-8")) + except (OSError, ValueError): + return {} + entries = payload.get("entries") if isinstance(payload, dict) else None + if not isinstance(entries, dict): + return {} + return {str(key): dict(value) for key, value in entries.items() if isinstance(value, dict)} + + +def update_pending_adoptions(*, runtime_root: str | Path, codex_home: str, + pending: list[dict[str, Any]], resolved: list[str]) -> Path | None: + """Replace the adoptions this reconciliation left pending, per automation. + + The store records an update-time observation, never adoption authority: it + only names entries this run already reviewed, so a later turn can apply the + reviewed prompt-only request instead of re-classifying the host store. + One runtime root can front several Codex homes, so only this home's records + are resolved; another home's pending obligations are left alone. + """ + path = pending_store_path(runtime_root) + entries = _read_pending_store(runtime_root) + for automation_id in resolved: + if str(entries.get(automation_id, {}).get("codex_home") or "") == codex_home: + entries.pop(automation_id, None) + for record in pending: + entries[str(record["automation_id"])] = record + if not entries: + if path.is_file(): + path.unlink() + return None + _atomic(path, json.dumps({"schema_version": PENDING_SCHEMA_VERSION, "entries": entries}, + ensure_ascii=False, indent=2) + "\n") + return path + + +def load_pending_adoption(*, runtime_root: str | Path | None, goal_id: str, + agent_id: str | None) -> dict[str, Any] | None: + """Project this lane's recorded adoption, or None once the host already applied it. + + The record was reviewed at update time; this read only checks that its + expected body is still not installed, so a satisfied obligation stops + projecting itself without re-classifying any lane. + """ + if not runtime_root or not agent_id: + return None + records = [record for record in _read_pending_store(runtime_root).values() + if record.get("goal_id") == goal_id and record.get("agent_id") == agent_id] + if len(records) != 1: + # No record, or several lanes sharing one identifier, is not a lane + # obligation this turn may act on. + return None + record = records[0] + request = record.get("api_update_request") + if not isinstance(request, dict) or not record.get("desired_sha256"): + return None + if _installed_prompt_matches_digest(record): + return None + return {"schema_version": PENDING_SCHEMA_VERSION, "status": "adoption_required", + "goal_id": goal_id, "agent_id": agent_id, + "automation_id": str(record.get("automation_id") or ""), + "automation_status": record.get("automation_status"), + "host_action": PENDING_HOST_ACTION, "host_action_contract": PENDING_HOST_ACTION_CONTRACT, + "spend_policy": PENDING_SPEND_POLICY, + "expected_prompt_sha256": record.get("expected_prompt_sha256"), + "desired_sha256": record.get("desired_sha256"), + "source": record.get("source"), "recorded_at": record.get("recorded_at"), + "reason": "the installed automation body is not the current managed loader; " + "review this prompt-only request through the App", + "api_update_request": request} + + +def _installed_prompt_matches_digest(record: dict[str, Any]) -> bool: + """True only when this exact automation now carries the desired body.""" + codex_home = str(record.get("codex_home") or "").strip() + automation_id = str(record.get("automation_id") or "") + if not codex_home or not automation_id: + return False + try: + manifest = tomllib.loads((Path(codex_home) / "automations" / automation_id + / "automation.toml").read_text(encoding="utf-8")) + except (OSError, UnicodeError, tomllib.TOMLDecodeError): + return False + return digest(str(manifest.get("prompt") or "")) == record.get("desired_sha256") + def require_closed_app() -> None: if sys.platform != "darwin": @@ -80,10 +188,14 @@ def reconcile(*, before: dict, registry: Path, home: Path, build_plan(registry=registry, home=home, runtime_root=runtime_root, cli_bin=cli_bin)["entries"]} results = [] api_updates = [] + pending_records = [] + resolved = [] + recorded_at = datetime.now(timezone.utc).isoformat() for old in before["entries"]: identifier = old["automation_id"] now = current.get(identifier) result = {"automation_id": identifier, "status": "review_required"} + pending_here = False if now is None: result["status"] = "missing" elif now["status"] in {"current", "unmanaged", "blocked"}: @@ -107,22 +219,38 @@ def reconcile(*, before: dict, registry: Path, home: Path, result.update(status="deferred", reason=str(error)) # The CLI cannot call an in-App tool itself. Give its host a # complete prompt-only request, plus a precondition to re-view. - import tomllib try: manifest = tomllib.loads((home / "automations" / identifier / "automation.toml").read_text()) except (OSError, ValueError): manifest = {} required = {"name", "status", "rrule", "target_thread_id"} if required <= manifest.keys() and manifest.get("prompt") == now["current_prompt"]: - api_updates.append({"tool": "automation_update", + request = automation_update_request(automation_id=identifier, + manifest=manifest, expected_prompt_sha256=now["prompt_sha256"], + desired_prompt=now["desired_prompt"]) + api_updates.append(request) + # A completed report cannot carry the obligation into the + # next turn, and nothing else re-observes the installed + # body, so record what the host still has to adopt. + pending_here = True + pending_records.append({"automation_id": identifier, + "goal_id": now["goal_id"], "agent_id": now["agent_id"], + "automation_status": manifest["status"], "codex_home": str(home), "expected_prompt_sha256": now["prompt_sha256"], - "precondition": "View the same automation; verify this prompt hash and all preserved fields before update; read back afterward.", - "arguments": {"mode": "update", "id": identifier, "kind": "heartbeat", - "name": manifest["name"], "status": manifest["status"], - "rrule": manifest["rrule"], "targetThreadId": manifest["target_thread_id"], - "notificationPolicy": manifest.get("notification_policy"), - "prompt": now["desired_prompt"]}}) + "desired_sha256": now["desired_sha256"], "recorded_at": recorded_at, + "source": "update_time_reconciliation", + "api_update_request": request}) + if not pending_here: + resolved.append(identifier) results.append(result) + if pending_records or resolved: + from loopx.control_plane.coordination.local_authority_shadow_adapter import ( + effective_runtime_root, + ) + update_pending_adoptions( + runtime_root=effective_runtime_root(registry.resolve(), runtime_root), + codex_home=str(home), + pending=pending_records, resolved=resolved) pending = any(result["status"] not in {"current", "updated", "unmanaged", "missing"} for result in results) return {"ok": not pending, "status": "attention_required" if pending else "current", "results": results, "api_updates": api_updates, diff --git a/tests/control_plane/test_automation_prompt_upgrade.py b/tests/control_plane/test_automation_prompt_upgrade.py index a771c8df6c..9d7ae1a9e7 100644 --- a/tests/control_plane/test_automation_prompt_upgrade.py +++ b/tests/control_plane/test_automation_prompt_upgrade.py @@ -13,6 +13,12 @@ from loopx.control_plane.heartbeat import automation_upgrade as upgrade +def runtime_root(tmp_path: Path) -> Path: + """One explicit runtime root per test: reconciliation records deferred adoptions.""" + + return tmp_path / "runtime" + + def fixture(tmp_path: Path, backing_kind="heartbeat"): home = tmp_path / "host" path = home / "automations/watch/automation.toml" @@ -218,7 +224,7 @@ def test_failed_install_never_attempts_prompt_writes(tmp_path, monkeypatch): def forbidden(*args, **kwargs): raise AssertionError("failed installer must not invoke prompt writer") monkeypatch.setattr(lifecycle.subprocess, "run", forbidden) - result = lifecycle.update_with_prompts({}, registry=registry, runtime_root=None, + result = lifecycle.update_with_prompts({}, registry=registry, runtime_root=str(runtime_root(tmp_path)), timeout_seconds=1, runtime_update=lambda payload, **_: {"ok": False}) assert result["automation_prompt_upgrade"]["status"] == "skipped_runtime_update_failed" assert path.read_bytes() == original @@ -249,7 +255,7 @@ def run(command, **kwargs): "next_action": {"kind": "review_or_rollback"}, "recommended_action": "Inspect optional extensions"} result = lifecycle.update_with_prompts( {"install_lifecycle": {"execution_driver": "python_pip"}}, registry=registry, - runtime_root=None, timeout_seconds=60, runtime_update=lambda *_, **__: runtime_result) + runtime_root=str(runtime_root(tmp_path)), timeout_seconds=60, runtime_update=lambda *_, **__: runtime_result) assert result["automation_prompt_upgrade"]["status"] == expected_status assert len(calls) == (1 if expected_status == "attention_required" else 0) assert not result["ok"] and not result["upgrade_complete"] @@ -270,12 +276,12 @@ def test_update_identifies_owned_legacy_body_and_migrates_without_changing_sched metadata = tomllib.loads(path.read_text()) monkeypatch.setattr(lifecycle, "require_closed_app", lambda: None) monkeypatch.setattr(lifecycle.sys, "platform", "darwin") - result = lifecycle.reconcile(before=before, registry=registry, home=home) + result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path)) assert result["ok"] and result["results"][0]["status"] == "updated" after = tomllib.loads(path.read_text()) assert {k: v for k, v in after.items() if k != "prompt"} == {k: v for k, v in metadata.items() if k != "prompt"} assert upgrade.bootstrap_binding(after["prompt"]) is not None - assert lifecycle.reconcile(before=before, registry=registry, home=home)["results"][0]["status"] == "current" + assert lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path))["results"][0]["status"] == "current" @pytest.mark.parametrize("scenario", ["custom", "unsupported_host", "race", "wrong_home", "canary"]) @@ -298,7 +304,7 @@ def test_update_does_not_overwrite_custom_changed_or_foreign_hosts(tmp_path, mon with pytest.raises(ValueError, match="another host"): lifecycle.reconcile(before=before, registry=registry, home=tmp_path / "other") else: - result = lifecycle.reconcile(before=before, registry=registry, home=home) + result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path)) expected = {"custom": "review_required", "canary": "review_required", "unsupported_host": "deferred", "race": "changed_since_snapshot"}[scenario] assert result["results"][0]["status"] == expected @@ -337,11 +343,12 @@ def fail_mirror(target, text): def test_running_app_defers_to_native_api_without_touching_cached_scheduler(tmp_path, monkeypatch): from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle home, path, database, registry, _ = fixture(tmp_path) - prompt = upgrade.bootstrap_prompt(registry=registry, goal_id="fixture-goal", agent_id="agent-a") + prompt = upgrade.bootstrap_prompt(registry=registry, goal_id="fixture-goal", agent_id="agent-a", + runtime_root=str(runtime_root(tmp_path))) legacy = prompt.replace(upgrade.BOOTSTRAP, upgrade._LEGACY_BOOTSTRAP, 1).removesuffix( upgrade._BOOTSTRAP_INSTRUCTION) + upgrade._LEGACY_INSTRUCTION _set_fixture_prompt(path, database, legacy) - before = lifecycle.snapshot(registry=registry, home=home) + before = lifecycle.snapshot(registry=registry, home=home, runtime_root=str(runtime_root(tmp_path))) monkeypatch.setattr(lifecycle.sys, "platform", "darwin") probes = [] def running(command, **kwargs): @@ -351,7 +358,7 @@ def running(command, **kwargs): original = path.read_bytes() with sqlite3.connect(database) as observer: row = observer.execute("SELECT * FROM automations").fetchone() - result = lifecycle.reconcile(before=before, registry=registry, home=home) + result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path)) assert not result["ok"] and result["results"][0]["status"] == "deferred" assert observer.execute("SELECT * FROM automations").fetchone() == row assert path.read_bytes() == original @@ -531,9 +538,103 @@ def test_exact_legacy_host_loader_upgrades_to_v2_without_dropping_explicit_polic assert entry["desired_prompt"].startswith("LoopX managed heartbeat bootstrap v2\n") assert host_bootstrap_binding(entry["desired_prompt"])["permission_rule"] == "Read only" monkeypatch.setattr(lifecycle, "require_closed_app", lambda: None) - assert lifecycle.reconcile(before=before, registry=registry, home=home)["ok"] + assert lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path))["ok"] assert lifecycle.snapshot(registry=registry, home=home)["entries"][0]["status"] == "current" malformed = legacy.replace("heartbeat-prompt ", "") _set_fixture_prompt(path, database, malformed) rejected = lifecycle.snapshot(registry=registry, home=home)["entries"][0] assert rejected["status"] != "current" and not rejected["automatic_eligible"] + + +def managed_legacy_fixture(path, database, registry, root) -> str: + """A managed v1 wrapper this host may still adopt, not a hand-written body.""" + + desired = upgrade.bootstrap_prompt(registry=registry, goal_id="fixture-goal", + agent_id="agent-a", runtime_root=str(root)) + legacy = desired.replace(upgrade.BOOTSTRAP, upgrade._LEGACY_BOOTSTRAP, 1).removesuffix( + upgrade._BOOTSTRAP_INSTRUCTION) + upgrade._LEGACY_INSTRUCTION + _set_fixture_prompt(path, database, legacy) + return legacy + + +def test_deferred_adoption_is_recorded_and_stops_projecting_once_installed(tmp_path, monkeypatch): + from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle + home, path, database, registry, _ = fixture(tmp_path) + root = runtime_root(tmp_path) + legacy = managed_legacy_fixture(path, database, registry, root) + before = lifecycle.snapshot(registry=registry, home=home, runtime_root=str(root)) + assert before["entries"][0]["automatic_eligible"] is True + # An unqualified host cannot take the write, so the reviewed request has to + # outlive this one report. + monkeypatch.setattr(lifecycle.sys, "platform", "linux") + result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=root) + assert not result["ok"] and result["results"][0]["status"] == "deferred" + assert lifecycle.pending_store_path(root).is_file() + pending = lifecycle.load_pending_adoption( + runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") + assert pending["status"] == "adoption_required" + assert pending["host_action"] == lifecycle.PENDING_HOST_ACTION + assert pending["host_action_contract"] == lifecycle.PENDING_HOST_ACTION_CONTRACT + assert pending["spend_policy"] == lifecycle.PENDING_SPEND_POLICY + # Same reviewed request as the update report, scheduling and binding intact. + assert pending["api_update_request"] == result["api_updates"][0] + assert pending["api_update_request"]["arguments"]["targetThreadId"] == "thread-a" + assert pending["api_update_request"]["arguments"]["status"] == "PAUSED" + # Another lane, an absent store, and no lane identity project nothing. + assert lifecycle.load_pending_adoption( + runtime_root=root, goal_id="fixture-goal", agent_id="agent-b") is None + assert lifecycle.load_pending_adoption( + runtime_root=tmp_path / "absent", goal_id="fixture-goal", agent_id="agent-a") is None + assert lifecycle.load_pending_adoption( + runtime_root=root, goal_id="fixture-goal", agent_id=None) is None + # The exact reviewed body satisfies the obligation without re-reading the lane. + upgrade.apply_offline(home=home, automation_id="watch", expected_prompt_sha256=upgrade.digest(legacy), + desired_prompt=pending["api_update_request"]["arguments"]["prompt"]) + assert lifecycle.load_pending_adoption( + runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") is None + + +def test_applied_adoption_clears_its_record(tmp_path, monkeypatch): + from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle + home, path, database, registry, _ = fixture(tmp_path) + root = runtime_root(tmp_path) + managed_legacy_fixture(path, database, registry, root) + before = lifecycle.snapshot(registry=registry, home=home, runtime_root=str(root)) + monkeypatch.setattr(lifecycle.sys, "platform", "linux") + assert lifecycle.reconcile( + before=before, registry=registry, home=home, runtime_root=root)["api_updates"] + assert lifecycle.pending_store_path(root).is_file() + monkeypatch.setattr(lifecycle, "require_closed_app", lambda: None) + monkeypatch.setattr(lifecycle.sys, "platform", "darwin") + result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=root) + assert result["ok"] and result["results"][0]["status"] == "updated" + assert not lifecycle.pending_store_path(root).is_file() + assert lifecycle.load_pending_adoption( + runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") is None + + +def test_pending_store_serves_only_an_unambiguous_reviewed_record(tmp_path): + from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle + root = runtime_root(tmp_path) + request = upgrade.automation_update_request(automation_id="watch", + manifest={"name": "Fixture watch", "status": "ACTIVE", "rrule": "FREQ=HOURLY", + "target_thread_id": "thread-a"}, + expected_prompt_sha256="a" * 64, desired_prompt="LoopX managed heartbeat bootstrap v2") + record = {"automation_id": "watch", "goal_id": "fixture-goal", "agent_id": "agent-a", + "automation_status": "ACTIVE", "codex_home": str(tmp_path / "host"), + "expected_prompt_sha256": "a" * 64, "desired_sha256": "b" * 64, + "source": "update_time_reconciliation", "api_update_request": request} + lifecycle.update_pending_adoptions(runtime_root=root, codex_home=str(tmp_path / "host"), + pending=[record], resolved=[]) + assert lifecycle.load_pending_adoption( + runtime_root=root, goal_id="fixture-goal", agent_id="agent-a")["automation_id"] == "watch" + # Two installed automations claiming one lane stay a review, not a host action. + lifecycle.update_pending_adoptions(runtime_root=root, codex_home=str(tmp_path / "host"), + pending=[{**record, "automation_id": "watch-2"}], resolved=[]) + assert lifecycle.load_pending_adoption( + runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") is None + # A record without its reviewed request is never handed to a host. + lifecycle.update_pending_adoptions(runtime_root=root, codex_home=str(tmp_path / "host"), + pending=[{**record, "api_update_request": None}], resolved=["watch-2"]) + assert lifecycle.load_pending_adoption( + runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") is None From 80ca91509edc0e348a3eb9023618b7e3dadcd8e5 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 15 Sep 2026 17:44:16 +0800 Subject: [PATCH 2/6] feat(quota): project a recorded automation prompt adoption every turn The heartbeat contract has to describe the automation the host actually runs. A recorded pending adoption is projected as scheduler_hint.app_automation.prompt_adoption with the reviewed prompt-only automation_update request, its no-spend policy, and both prompt digests, so the turn applies it once and reads the automation back without spending quota. The projection reuses the existing payload obligation channel rather than adding a scheduler-hint parameter, and a satisfied record stops projecting itself once the exact body is installed. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/heartbeat-automation-prompt.md | 7 +- docs/reference/automation-prompt-upgrades.md | 26 ++++ .../control_plane/heartbeat-prompt-smoke.py | 2 + loopx/control_plane/heartbeat/rules.py | 10 +- .../control_plane/quota/should_run_packet.py | 4 +- .../control_plane/quota/should_run_prepare.py | 34 +++++ .../control_plane/scheduler/scheduler_hint.py | 5 + skills/loopx-project/SKILL.md | 6 +- .../test_prompt_adoption_projection.py | 135 ++++++++++++++++++ 9 files changed, 222 insertions(+), 7 deletions(-) create mode 100644 tests/control_plane/test_prompt_adoption_projection.py diff --git a/docs/heartbeat-automation-prompt.md b/docs/heartbeat-automation-prompt.md index 463b7e64ec..e7a913cca8 100644 --- a/docs/heartbeat-automation-prompt.md +++ b/docs/heartbeat-automation-prompt.md @@ -637,7 +637,12 @@ heartbeats should search/use `automation_update` when available. If `scheduler_hint.app_automation.host_action=pause_or_delete_current_heartbeat`: in that terminal case, call `automation_update` once to pause the current heartbeat (delete only if pause is unavailable), verify the host result, spend -no quota, and end the turn without a scheduler ACK. Otherwise call it only when +no quota, and end the turn without a scheduler ACK. When the same lane reports +`scheduler_hint.app_automation.prompt_adoption.host_action=adopt_managed_bootstrap`, +its installed body is not the current managed loader: apply that reported +prompt-only `api_update_request` once through `automation_update` after checking +the reported prompt hash, read the automation back, and spend no quota. Otherwise +call it only when `scheduler_hint.app_automation.stateful_backoff.apply_needed=true` and `scheduler_hint.app_automation.recommended_rrule` is present. After a successful RRULE update, run `loopx` with diff --git a/docs/reference/automation-prompt-upgrades.md b/docs/reference/automation-prompt-upgrades.md index 22d4edbab9..93fb25c189 100644 --- a/docs/reference/automation-prompt-upgrades.md +++ b/docs/reference/automation-prompt-upgrades.md @@ -68,6 +68,22 @@ read-only preview, not the upgrade executor. Do not infer a manual-only policy from its `adoption_required` status. Custom or inconsistent entries still need review; automatic prompt migration never grants scheduler or thread authority. +## Deferred adoption obligations + +Reconciliation cannot write while the App runs, and a completed report cannot +carry an unapplied migration into the next turn. An entry that is reviewable but +not written is therefore recorded per lane in the runtime root's +`app-automation-prompt-adoptions.json`, holding the reviewed prompt-only +`automation_update` request, both prompt digests, and the host that was read. +`quota should-run` projects that record as +`scheduler_hint.app_automation.prompt_adoption` with +`host_action=adopt_managed_bootstrap` and the no-spend policy, so the obligation +survives the update that discovered it. The turn applies the request once +through `automation_update`, reads the automation back, and spends no quota. +The projection is not adoption authority and not delivery permission: it names +only an entry this host already reviewed, it disappears as soon as the exact +body is installed, and nothing re-classifies installed automations per turn. + On the qualified macOS heartbeat schema, direct migration requires the App closed. The adapter holds a SQLite writer transaction through TOML delivery, compares the entire previewed manifest, preserves every non-prompt field, and @@ -202,3 +218,13 @@ gh 登录,仍失败则明确要求已核验 SHA,不切换分支或静默覆 日程、暂停状态、模型、线程、通知偏好和历史均不迁移。 不支持的存储仍需原生 API;运行中的本轮不热切换。普通测试不消耗模型 token, 真实模型发布资格仍需独立评测,不能由迁移成功推断。 + +App 运行中对账无法写入,而一次性报告也带不走未应用的迁移。因此可审阅但未写入的 +条目会按 lane 记录到 runtime root 的 `app-automation-prompt-adoptions.json`,保存 +已审阅的、仅改 prompt 的 `automation_update` 请求、两个 prompt 摘要以及读取来源。 +`quota should-run` 将该记录投影为 +`scheduler_hint.app_automation.prompt_adoption`,携带 +`host_action=adopt_managed_bootstrap` 与 no-spend 策略,使义务不被一次性报告带走; +turn 通过 `automation_update` 应用一次并读回,不消耗额度。该投影不是采纳权威、 +也不授予交付权限:它只描述本 host 已审阅过的条目,目标 body 一旦安装即自动消失, +且不按轮重新分类已安装的 automation。 diff --git a/examples/control_plane/heartbeat-prompt-smoke.py b/examples/control_plane/heartbeat-prompt-smoke.py index 71167f43a2..ae25c82a39 100644 --- a/examples/control_plane/heartbeat-prompt-smoke.py +++ b/examples/control_plane/heartbeat-prompt-smoke.py @@ -594,6 +594,7 @@ def main() -> int: "具体user todo未投影", "Observed capabilities -> `--available-capability`; never user gates", "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend)", + "prompt_adoption=adopt(no-spend);", "else RRULE/projected-fallback_hint/ack/fail", "no-change=`surface_only`/no spend", "unchanged->`--vision-unchanged-reason`", @@ -696,6 +697,7 @@ def main() -> int: "NOTIFY缺动作→", "具体user todo未投影", "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend)", + "prompt_adoption=adopt(no-spend);", "else RRULE/projected-fallback_hint/ack/fail", "no-change=`surface_only`/no spend", "unchanged->`--vision-unchanged-reason`", diff --git a/loopx/control_plane/heartbeat/rules.py b/loopx/control_plane/heartbeat/rules.py index 264bd5f2b9..7df419569b 100644 --- a/loopx/control_plane/heartbeat/rules.py +++ b/loopx/control_plane/heartbeat/rules.py @@ -41,19 +41,23 @@ ) SCHEDULER_HINT_APPLICATION_RULE = ( "`scheduler_hint` no-spend. host_action=pause_or_delete_current_heartbeat -> " - "automation_update stop once, verify, end; else apply_needed -> RRULE via " - "automation_update; unavailable -> use fallback_hint.cli_args only when projected " - "(SQLite/app API " + "automation_update stop once, verify, end; prompt_adoption.host_action=" + "adopt_managed_bootstrap -> its prompt-only api_update_request once, no spend, " + "then readback; else apply_needed -> RRULE via automation_update; unavailable -> " + "use fallback_hint.cli_args only when projected (SQLite/app API " "bypass - fallback only), then ack; further failure -> failure_hint; " "ack_needed -> ack." ) SCHEDULER_HINT_COMPACT_RULE = ( "host_action=pause_or_delete_current_heartbeat: automation_update stop; " + "prompt_adoption.host_action=adopt_managed_bootstrap: prompt-only " + "api_update_request once, no spend; " "else RRULE apply via automation_update, projected fallback_hint when unavailable, " "then ack/fail. No spend." ) SCHEDULER_HINT_THIN_RULE = ( "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend); " + "prompt_adoption=adopt(no-spend); " "else RRULE/projected-fallback_hint/ack/fail." ) RUNTIME_CAPABILITY_PROJECTION_THIN_RULE = ( diff --git a/loopx/control_plane/quota/should_run_packet.py b/loopx/control_plane/quota/should_run_packet.py index ab7eb108d0..bbc2e6eadf 100644 --- a/loopx/control_plane/quota/should_run_packet.py +++ b/loopx/control_plane/quota/should_run_packet.py @@ -1317,6 +1317,7 @@ def _build_quota_should_run_payload( agent_scope_frontier=route.agent_scope_frontier, workspace_guard=prepared.workspace_guard, automation_prompt_upgrade=prepared.automation_prompt_upgrade, + automation_prompt_adoption=prepared.automation_prompt_adoption, ) if prepared.agent_scoped_user_todo_override: payload[str(prepared.agent_scoped_user_todo_override["kind"])] = ( @@ -1471,8 +1472,7 @@ def _build_quota_should_run_payload( _load_app_automation_scheduler_state( prepared.status_payload, goal_id=prepared.safe_goal_id, - agent_id=quota_decision_agent_id(payload) - or prepared.requested_agent_id, + agent_id=quota_decision_agent_id(payload) or prepared.requested_agent_id, surface=prepared.resolved_scheduler_context.context.host_surface.value, ) if prepared.resolved_scheduler_context.ok diff --git a/loopx/control_plane/quota/should_run_prepare.py b/loopx/control_plane/quota/should_run_prepare.py index 159e938af6..6ef2e70c23 100644 --- a/loopx/control_plane/quota/should_run_prepare.py +++ b/loopx/control_plane/quota/should_run_prepare.py @@ -31,6 +31,7 @@ build_goal_frontier_projection_context_from_status, ) from ..quota.error_codes import HeartbeatReceiptIdentityConflictError +from ..quota.decision_summary import quota_decision_agent_id from ..agents.capability_memory import resolve_agent_capabilities from ..quota.goal_boundary import ( goal_boundary as _goal_boundary, @@ -140,6 +141,7 @@ class _QuotaDecisionPreparation: goal_boundary: dict[str, Any] | None automation_prompt_upgrade: dict[str, Any] | None automation_prompt_upgrade_required: bool + automation_prompt_adoption: dict[str, Any] | None blocked_priority_fallback: dict[str, Any] | None stall_self_repair: dict[str, Any] | None self_repair_allowed: bool @@ -296,6 +298,35 @@ def _automation_prompt_upgrade( ) +def _automation_prompt_adoption( + status_payload: dict[str, Any], + *, + goal_id: str, + requested_agent_id: str | None, + resolved_scheduler_context: SchedulerExecutionContextResolution, +) -> dict[str, Any] | None: + """The lane's recorded prompt-only adoption, when an App automation drives it. + + Update-time reconciliation cannot apply a prompt change while the App runs, + so it records the request it reviewed. Project that record instead of + re-classifying the installed body every turn. + """ + context = resolved_scheduler_context.context + if ( + not resolved_scheduler_context.ok + or context is None + or not context.app_automation_applicable + ): + return None + from ..heartbeat.installed_prompt_update import load_pending_adoption + + return load_pending_adoption( + runtime_root=status_payload.get("runtime_root"), + goal_id=goal_id, + agent_id=quota_decision_agent_id(status_payload) or requested_agent_id, + ) + + def _todo_task_class(item: dict[str, Any]) -> str: return projection_todo_item_task_class(item) @@ -826,6 +857,9 @@ def _prepare_quota_should_run_item( goal_boundary=goal_boundary, automation_prompt_upgrade=automation_prompt_upgrade, automation_prompt_upgrade_required=automation_prompt_upgrade_required, + automation_prompt_adoption=_automation_prompt_adoption(status_payload, + goal_id=safe_goal_id, requested_agent_id=requested_agent_id, + resolved_scheduler_context=resolved_scheduler_context), blocked_priority_fallback=blocked_priority_fallback, stall_self_repair=stall_self_repair, self_repair_allowed=self_repair_allowed, diff --git a/loopx/control_plane/scheduler/scheduler_hint.py b/loopx/control_plane/scheduler/scheduler_hint.py index f385ce5f05..6d65f58c37 100644 --- a/loopx/control_plane/scheduler/scheduler_hint.py +++ b/loopx/control_plane/scheduler/scheduler_hint.py @@ -881,6 +881,11 @@ def build( }, "no_spend_for_cadence_change": True, } + prompt_adoption = self.payload.get("automation_prompt_adoption") + if isinstance(prompt_adoption, dict): + # The installed body decides which rules later wakes follow, so a + # recorded pending adoption is part of the observed App automation. + app_automation["prompt_adoption"] = dict(prompt_adoption) stateful_backoff = app_automation["stateful_backoff"] if host_update_failures: stateful_backoff["host_update_failures"] = [ diff --git a/skills/loopx-project/SKILL.md b/skills/loopx-project/SKILL.md index 5e13884959..d0bdbfe696 100644 --- a/skills/loopx-project/SKILL.md +++ b/skills/loopx-project/SKILL.md @@ -582,7 +582,11 @@ search/use `automation_update` when available. If `automation_update` once to pause the current heartbeat (delete only when the host cannot pause), verify the host result, spend no quota, and end the turn. This terminal host action takes precedence over RRULE handling and requires no -scheduler ACK. Otherwise use `automation_update` only when +scheduler ACK. When the same lane reports +`scheduler_hint.app_automation.prompt_adoption.host_action=adopt_managed_bootstrap`, +apply that prompt-only `api_update_request` once through `automation_update` +after checking the reported prompt hash, read the automation back, and spend no +quota. Otherwise use `automation_update` only when `scheduler_hint.app_automation.stateful_backoff.apply_needed=true` and `scheduler_hint.app_automation.recommended_rrule` is present. After a successful RRULE update, run `loopx` with diff --git a/tests/control_plane/test_prompt_adoption_projection.py b/tests/control_plane/test_prompt_adoption_projection.py new file mode 100644 index 0000000000..86788fafbb --- /dev/null +++ b/tests/control_plane/test_prompt_adoption_projection.py @@ -0,0 +1,135 @@ +from __future__ import annotations + +import json +from pathlib import Path +import sqlite3 + +from loopx.control_plane.heartbeat import automation_upgrade as upgrade +from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle +from loopx.control_plane.quota.should_run import build_quota_should_run +from loopx.control_plane.scheduler.execution_context import ( + scheduler_execution_context_for_runtime_profile, +) +from loopx.control_plane.testing.quota_fixtures import ( + quota_status_payload, + quota_todo_item, +) + + +APP_CONTEXT = scheduler_execution_context_for_runtime_profile("codex_app_heartbeat") +CLI_CONTEXT = scheduler_execution_context_for_runtime_profile("codex_cli") +GOAL_ID = "fixture-goal" +AGENT_ID = "agent-a" +THREAD_ID = "thread-a" +FROZEN_BODY = f"Advance `{GOAL_ID}` from registry. --agent-id {AGENT_ID}" +DESIRED_BODY = "LoopX managed heartbeat bootstrap v2" + + +def _host_with_frozen_body(tmp_path: Path) -> tuple[Path, Path]: + home = tmp_path / "host" + path = home / "automations/watch/automation.toml" + path.parent.mkdir(parents=True) + path.write_text('version = 1\nid = "watch"\nname = "Fixture watch"\nkind = "heartbeat"\n' + 'status = "PAUSED"\ntarget_thread_id = "' + THREAD_ID + '"\n' + 'rrule = "FREQ=MINUTELY;INTERVAL=3"\n' + 'prompt = ' + json.dumps(FROZEN_BODY) + "\n", encoding="utf-8") + database = home / "sqlite/codex-dev.db" + database.parent.mkdir() + with sqlite3.connect(database) as connection: + connection.execute("CREATE TABLE automations (id TEXT PRIMARY KEY, kind TEXT, prompt TEXT," + " status TEXT, target_thread_id TEXT, rrule TEXT, model TEXT," + " updated_at INTEGER, next_run_at INTEGER)") + connection.execute("INSERT INTO automations VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", + ("watch", "heartbeat", FROZEN_BODY, "PAUSED", THREAD_ID, + "FREQ=MINUTELY;INTERVAL=3", "fixture-model", 123, 456)) + registry = tmp_path / "registry.json" + state = tmp_path / "STATE.md" + state.write_text("# Fixture\n", encoding="utf-8") + registry.write_text(json.dumps({"goals": [{"id": GOAL_ID, "repo": str(tmp_path), + "state_file": str(state), "registered_agents": [AGENT_ID]}]}), encoding="utf-8") + return home, registry + + +def _record_deferred_adoption(tmp_path: Path, home: Path, desired: str) -> None: + """Store the reviewed request an update-time reconciliation could not apply.""" + + lifecycle.update_pending_adoptions(runtime_root=tmp_path / "runtime", codex_home=str(home), + resolved=[], + pending=[{"automation_id": "watch", "goal_id": GOAL_ID, "agent_id": AGENT_ID, + "automation_status": "PAUSED", "codex_home": str(home), + "expected_prompt_sha256": upgrade.digest(FROZEN_BODY), + "desired_sha256": upgrade.digest(desired), + "source": "update_time_reconciliation", + "api_update_request": upgrade.automation_update_request(automation_id="watch", + manifest={"name": "Fixture watch", "status": "PAUSED", + "rrule": "FREQ=MINUTELY;INTERVAL=3", "target_thread_id": THREAD_ID}, + expected_prompt_sha256=upgrade.digest(FROZEN_BODY), desired_prompt=desired)}]) + + +def _lane_payload(tmp_path: Path, registry: Path) -> dict: + """One runnable lane whose packet has to carry the recorded obligation.""" + + return { + **quota_status_payload( + goal_id=GOAL_ID, + status="active", + agent_todo_items=[ + quota_todo_item( + todo_id="todo_current001", + index=1, + priority="P1", + title="Advance the reviewed slice.", + claimed_by=AGENT_ID, + ) + ], + recommended_action="Advance the reviewed slice.", + coordination={"agent_model": "peer_v1", "registered_agents": [AGENT_ID]}, + claim_scope_agent_id=AGENT_ID, + ), + "registry": str(registry), + "runtime_root": str(tmp_path / "runtime"), + } + + +def _packet(tmp_path: Path, registry: Path, context): + return build_quota_should_run( + _lane_payload(tmp_path, registry), + goal_id=GOAL_ID, + agent_id=AGENT_ID, + scheduler_execution_context=context, + ) + + +def test_recorded_adoption_is_projected_into_the_lane_packet(tmp_path): + home, registry = _host_with_frozen_body(tmp_path) + _record_deferred_adoption(tmp_path, home, DESIRED_BODY) + hint = _packet(tmp_path, registry, APP_CONTEXT)["scheduler_hint"] + adoption = hint["app_automation"]["prompt_adoption"] + assert adoption["status"] == "adoption_required" + assert adoption["automation_id"] == "watch" + assert adoption["host_action"] == lifecycle.PENDING_HOST_ACTION + assert adoption["host_action_contract"] == lifecycle.PENDING_HOST_ACTION_CONTRACT + assert adoption["spend_policy"] == lifecycle.PENDING_SPEND_POLICY + assert adoption["api_update_request"]["arguments"] == { + "mode": "update", "id": "watch", "kind": "heartbeat", "name": "Fixture watch", + "status": "PAUSED", "rrule": "FREQ=MINUTELY;INTERVAL=3", "targetThreadId": THREAD_ID, + "notificationPolicy": None, "prompt": DESIRED_BODY} + # The legacy Codex App projection carries the same recorded obligation. + assert hint["codex_app"]["prompt_adoption"] == adoption + + +def test_applied_body_and_non_app_hosts_project_no_adoption(tmp_path): + home, registry = _host_with_frozen_body(tmp_path) + _record_deferred_adoption(tmp_path, home, DESIRED_BODY) + assert "prompt_adoption" not in json.dumps(_packet(tmp_path, registry, CLI_CONTEXT)["scheduler_hint"]) + # A different installed body than the reviewed one stays pending. + assert _packet(tmp_path, registry, APP_CONTEXT)["scheduler_hint"]["app_automation"]["prompt_adoption"] + path = home / "automations/watch/automation.toml" + path.write_text(upgrade._replace_prompt(path.read_text(), DESIRED_BODY), encoding="utf-8") + assert "prompt_adoption" not in json.dumps(_packet(tmp_path, registry, APP_CONTEXT)["scheduler_hint"]) + + +def test_lane_without_a_record_projects_no_adoption(tmp_path): + home, registry = _host_with_frozen_body(tmp_path) + hint = _packet(tmp_path, registry, APP_CONTEXT)["scheduler_hint"] + assert "prompt_adoption" not in json.dumps(hint) From 4d11f774c227b599d5a14c36f23b5b6b76e655d0 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 15 Sep 2026 18:06:29 +0800 Subject: [PATCH 3/6] fix(heartbeat): keep the thin prompt adoption clause inside its budget The agent-facing CLI output differential measures the thin heartbeat prompt markdown against a 32-character hot-path allowance. Spelling the new obligation as `prompt_adoption=adopt(no-spend); ` grew those rows by 33 characters in all three thin variants (small, multi_agent, crowded), which fails `examples/control_plane/cli-output-budget-regression-smoke.py` and the kernel static checks job. Keep the clause and the no-spend marker, but drop the `prompt_` prefix inside the thin variant so the obligation costs 27 characters and leaves headroom. The full `prompt_adoption` field name stays visible in the compact and application rule variants and in the `scheduler_hint` payload. Validation: heartbeat-prompt-smoke.py, cli-output-budget-regression-smoke.py and cli-output-base-head-differential-smoke.py (base=102 candidate=102) pass; test_prompt_adoption_projection.py, test_automation_prompt_upgrade.py and test_cli_output_budget.py report 70 passed. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- examples/control_plane/heartbeat-prompt-smoke.py | 4 ++-- loopx/control_plane/heartbeat/rules.py | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/examples/control_plane/heartbeat-prompt-smoke.py b/examples/control_plane/heartbeat-prompt-smoke.py index ae25c82a39..e116811497 100644 --- a/examples/control_plane/heartbeat-prompt-smoke.py +++ b/examples/control_plane/heartbeat-prompt-smoke.py @@ -594,7 +594,7 @@ def main() -> int: "具体user todo未投影", "Observed capabilities -> `--available-capability`; never user gates", "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend)", - "prompt_adoption=adopt(no-spend);", + "adoption=adopt(no-spend);", "else RRULE/projected-fallback_hint/ack/fail", "no-change=`surface_only`/no spend", "unchanged->`--vision-unchanged-reason`", @@ -697,7 +697,7 @@ def main() -> int: "NOTIFY缺动作→", "具体user todo未投影", "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend)", - "prompt_adoption=adopt(no-spend);", + "adoption=adopt(no-spend);", "else RRULE/projected-fallback_hint/ack/fail", "no-change=`surface_only`/no spend", "unchanged->`--vision-unchanged-reason`", diff --git a/loopx/control_plane/heartbeat/rules.py b/loopx/control_plane/heartbeat/rules.py index 7df419569b..0f53a13979 100644 --- a/loopx/control_plane/heartbeat/rules.py +++ b/loopx/control_plane/heartbeat/rules.py @@ -57,7 +57,7 @@ ) SCHEDULER_HINT_THIN_RULE = ( "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend); " - "prompt_adoption=adopt(no-spend); " + "adoption=adopt(no-spend); " "else RRULE/projected-fallback_hint/ack/fail." ) RUNTIME_CAPABILITY_PROJECTION_THIN_RULE = ( From a4bfdb00ca8c87bc75bd3a585ee4727cc3e30ce6 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 15 Sep 2026 18:55:45 +0800 Subject: [PATCH 4/6] refactor(heartbeat): inject deferred upgrade hints through the existing turn-start hook Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../control_plane/heartbeat-prompt-smoke.py | 2 - .../heartbeat/automation_upgrade.py | 21 +- .../heartbeat/installed_prompt_update.py | 149 ++----------- .../heartbeat/prompt_upgrade_hook.py | 128 +++++++++++ loopx/control_plane/heartbeat/rules.py | 10 +- loopx/control_plane/quota/live_decision.py | 7 + .../control_plane/quota/should_run_packet.py | 4 +- .../control_plane/quota/should_run_prepare.py | 34 --- loopx/control_plane/quota/turn_envelope.ts | 3 +- .../control_plane/scheduler/scheduler_hint.py | 5 - .../work_items/interaction_contract.py | 3 +- .../test_automation_prompt_upgrade.py | 121 +---------- .../test_prompt_adoption_projection.py | 135 ------------ .../control_plane/test_prompt_upgrade_hook.py | 202 ++++++++++++++++++ 14 files changed, 369 insertions(+), 455 deletions(-) create mode 100644 loopx/control_plane/heartbeat/prompt_upgrade_hook.py delete mode 100644 tests/control_plane/test_prompt_adoption_projection.py create mode 100644 tests/control_plane/test_prompt_upgrade_hook.py diff --git a/examples/control_plane/heartbeat-prompt-smoke.py b/examples/control_plane/heartbeat-prompt-smoke.py index e116811497..71167f43a2 100644 --- a/examples/control_plane/heartbeat-prompt-smoke.py +++ b/examples/control_plane/heartbeat-prompt-smoke.py @@ -594,7 +594,6 @@ def main() -> int: "具体user todo未投影", "Observed capabilities -> `--available-capability`; never user gates", "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend)", - "adoption=adopt(no-spend);", "else RRULE/projected-fallback_hint/ack/fail", "no-change=`surface_only`/no spend", "unchanged->`--vision-unchanged-reason`", @@ -697,7 +696,6 @@ def main() -> int: "NOTIFY缺动作→", "具体user todo未投影", "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend)", - "adoption=adopt(no-spend);", "else RRULE/projected-fallback_hint/ack/fail", "no-change=`surface_only`/no spend", "unchanged->`--vision-unchanged-reason`", diff --git a/loopx/control_plane/heartbeat/automation_upgrade.py b/loopx/control_plane/heartbeat/automation_upgrade.py index 86fca270d3..c1f8993fd6 100644 --- a/loopx/control_plane/heartbeat/automation_upgrade.py +++ b/loopx/control_plane/heartbeat/automation_upgrade.py @@ -15,7 +15,7 @@ import sqlite3 import tempfile import tomllib -from typing import Any, Mapping +from typing import Any from types import SimpleNamespace from .bootstrap_prompt import ( @@ -229,25 +229,6 @@ def build_plan(*, registry: Path, home: Path | None = None, "policy": "Discovery is not adoption authority. Review each replacement; use the App API first."} -def automation_update_request(*, automation_id: str, manifest: Mapping[str, Any], - expected_prompt_sha256: str, desired_prompt: str) -> dict[str, Any]: - """Complete prompt-only App request; scheduling and thread binding are preserved. - - The CLI cannot call an in-App tool itself, so every entrypoint that hands a - stale installed body to its host sends this exact reviewed request. - """ - return {"tool": "automation_update", - "expected_prompt_sha256": expected_prompt_sha256, - "precondition": "View the same automation; verify this prompt hash and all " - "preserved fields before update; read back afterward.", - "arguments": {"mode": "update", "id": automation_id, "kind": "heartbeat", - "name": manifest["name"], "status": manifest["status"], - "rrule": manifest["rrule"], - "targetThreadId": manifest["target_thread_id"], - "notificationPolicy": manifest.get("notification_policy"), - "prompt": desired_prompt}} - - def apply_offline(*, home: Path, automation_id: str, expected_prompt_sha256: str, desired_prompt: str, expected_source_sha256: str | None = None) -> dict[str, Any]: """Journaled prompt-only write, also used by qualified update-time migration. diff --git a/loopx/control_plane/heartbeat/installed_prompt_update.py b/loopx/control_plane/heartbeat/installed_prompt_update.py index f465b4cc80..a0304d186c 100644 --- a/loopx/control_plane/heartbeat/installed_prompt_update.py +++ b/loopx/control_plane/heartbeat/installed_prompt_update.py @@ -6,7 +6,6 @@ """ from __future__ import annotations -from datetime import datetime, timezone import json import os from pathlib import Path @@ -14,117 +13,10 @@ import subprocess import sys import tempfile -import tomllib -from typing import Any -from .automation_upgrade import ( - SCHEMA, - _atomic, - apply_offline, - automation_update_request, - bootstrap_binding, - build_plan, - digest, -) +from .automation_upgrade import SCHEMA, _atomic, apply_offline, bootstrap_binding, build_plan from .bootstrap_prompt import host_bootstrap_binding -PENDING_SCHEMA_VERSION = "loopx_automation_prompt_adoption_pending_v0" -PENDING_HOST_ACTION = "adopt_managed_bootstrap" -PENDING_HOST_ACTION_CONTRACT = "codex_app_automation_prompt_adoption" -PENDING_SPEND_POLICY = "no_spend_for_automation_prompt_adoption" -PENDING_STORE_FILENAME = "app-automation-prompt-adoptions.json" - - -def pending_store_path(runtime_root: str | Path) -> Path: - return Path(runtime_root).expanduser().resolve() / PENDING_STORE_FILENAME - - -def _read_pending_store(runtime_root: str | Path) -> dict[str, dict[str, Any]]: - try: - payload = json.loads(pending_store_path(runtime_root).read_text(encoding="utf-8")) - except (OSError, ValueError): - return {} - entries = payload.get("entries") if isinstance(payload, dict) else None - if not isinstance(entries, dict): - return {} - return {str(key): dict(value) for key, value in entries.items() if isinstance(value, dict)} - - -def update_pending_adoptions(*, runtime_root: str | Path, codex_home: str, - pending: list[dict[str, Any]], resolved: list[str]) -> Path | None: - """Replace the adoptions this reconciliation left pending, per automation. - - The store records an update-time observation, never adoption authority: it - only names entries this run already reviewed, so a later turn can apply the - reviewed prompt-only request instead of re-classifying the host store. - One runtime root can front several Codex homes, so only this home's records - are resolved; another home's pending obligations are left alone. - """ - path = pending_store_path(runtime_root) - entries = _read_pending_store(runtime_root) - for automation_id in resolved: - if str(entries.get(automation_id, {}).get("codex_home") or "") == codex_home: - entries.pop(automation_id, None) - for record in pending: - entries[str(record["automation_id"])] = record - if not entries: - if path.is_file(): - path.unlink() - return None - _atomic(path, json.dumps({"schema_version": PENDING_SCHEMA_VERSION, "entries": entries}, - ensure_ascii=False, indent=2) + "\n") - return path - - -def load_pending_adoption(*, runtime_root: str | Path | None, goal_id: str, - agent_id: str | None) -> dict[str, Any] | None: - """Project this lane's recorded adoption, or None once the host already applied it. - - The record was reviewed at update time; this read only checks that its - expected body is still not installed, so a satisfied obligation stops - projecting itself without re-classifying any lane. - """ - if not runtime_root or not agent_id: - return None - records = [record for record in _read_pending_store(runtime_root).values() - if record.get("goal_id") == goal_id and record.get("agent_id") == agent_id] - if len(records) != 1: - # No record, or several lanes sharing one identifier, is not a lane - # obligation this turn may act on. - return None - record = records[0] - request = record.get("api_update_request") - if not isinstance(request, dict) or not record.get("desired_sha256"): - return None - if _installed_prompt_matches_digest(record): - return None - return {"schema_version": PENDING_SCHEMA_VERSION, "status": "adoption_required", - "goal_id": goal_id, "agent_id": agent_id, - "automation_id": str(record.get("automation_id") or ""), - "automation_status": record.get("automation_status"), - "host_action": PENDING_HOST_ACTION, "host_action_contract": PENDING_HOST_ACTION_CONTRACT, - "spend_policy": PENDING_SPEND_POLICY, - "expected_prompt_sha256": record.get("expected_prompt_sha256"), - "desired_sha256": record.get("desired_sha256"), - "source": record.get("source"), "recorded_at": record.get("recorded_at"), - "reason": "the installed automation body is not the current managed loader; " - "review this prompt-only request through the App", - "api_update_request": request} - - -def _installed_prompt_matches_digest(record: dict[str, Any]) -> bool: - """True only when this exact automation now carries the desired body.""" - codex_home = str(record.get("codex_home") or "").strip() - automation_id = str(record.get("automation_id") or "") - if not codex_home or not automation_id: - return False - try: - manifest = tomllib.loads((Path(codex_home) / "automations" / automation_id - / "automation.toml").read_text(encoding="utf-8")) - except (OSError, UnicodeError, tomllib.TOMLDecodeError): - return False - return digest(str(manifest.get("prompt") or "")) == record.get("desired_sha256") - def require_closed_app() -> None: if sys.platform != "darwin": @@ -188,14 +80,10 @@ def reconcile(*, before: dict, registry: Path, home: Path, build_plan(registry=registry, home=home, runtime_root=runtime_root, cli_bin=cli_bin)["entries"]} results = [] api_updates = [] - pending_records = [] - resolved = [] - recorded_at = datetime.now(timezone.utc).isoformat() for old in before["entries"]: identifier = old["automation_id"] now = current.get(identifier) result = {"automation_id": identifier, "status": "review_required"} - pending_here = False if now is None: result["status"] = "missing" elif now["status"] in {"current", "unmanaged", "blocked"}: @@ -219,38 +107,25 @@ def reconcile(*, before: dict, registry: Path, home: Path, result.update(status="deferred", reason=str(error)) # The CLI cannot call an in-App tool itself. Give its host a # complete prompt-only request, plus a precondition to re-view. + import tomllib try: manifest = tomllib.loads((home / "automations" / identifier / "automation.toml").read_text()) except (OSError, ValueError): manifest = {} required = {"name", "status", "rrule", "target_thread_id"} if required <= manifest.keys() and manifest.get("prompt") == now["current_prompt"]: - request = automation_update_request(automation_id=identifier, - manifest=manifest, expected_prompt_sha256=now["prompt_sha256"], - desired_prompt=now["desired_prompt"]) - api_updates.append(request) - # A completed report cannot carry the obligation into the - # next turn, and nothing else re-observes the installed - # body, so record what the host still has to adopt. - pending_here = True - pending_records.append({"automation_id": identifier, - "goal_id": now["goal_id"], "agent_id": now["agent_id"], - "automation_status": manifest["status"], "codex_home": str(home), + api_updates.append({"tool": "automation_update", "expected_prompt_sha256": now["prompt_sha256"], - "desired_sha256": now["desired_sha256"], "recorded_at": recorded_at, - "source": "update_time_reconciliation", - "api_update_request": request}) - if not pending_here: - resolved.append(identifier) + "precondition": "View the same automation; verify this prompt hash and all preserved fields before update; read back afterward.", + "arguments": {"mode": "update", "id": identifier, "kind": "heartbeat", + "name": manifest["name"], "status": manifest["status"], + "rrule": manifest["rrule"], "targetThreadId": manifest["target_thread_id"], + "notificationPolicy": manifest.get("notification_policy"), + "prompt": now["desired_prompt"]}}) results.append(result) - if pending_records or resolved: - from loopx.control_plane.coordination.local_authority_shadow_adapter import ( - effective_runtime_root, - ) - update_pending_adoptions( - runtime_root=effective_runtime_root(registry.resolve(), runtime_root), - codex_home=str(home), - pending=pending_records, resolved=resolved) + from .prompt_upgrade_hook import record_deferred_upgrades + record_deferred_upgrades(registry=registry, home=home, runtime_root=runtime_root, + cli_bin=cli_bin, entries=current, results=results) pending = any(result["status"] not in {"current", "updated", "unmanaged", "missing"} for result in results) return {"ok": not pending, "status": "attention_required" if pending else "current", "results": results, "api_updates": api_updates, diff --git a/loopx/control_plane/heartbeat/prompt_upgrade_hook.py b/loopx/control_plane/heartbeat/prompt_upgrade_hook.py new file mode 100644 index 0000000000..21c1e34a6a --- /dev/null +++ b/loopx/control_plane/heartbeat/prompt_upgrade_hook.py @@ -0,0 +1,128 @@ +"""Content-free update receipts and a read-only, conditional turn-start hint. + +Installed-prompt reconciliation owns discovery. The existing typed hook owns +observation admission; neither this receipt nor its hint authorizes adoption. +""" +from __future__ import annotations + +from contextlib import closing +import json +from pathlib import Path +import shlex +from typing import Any, Mapping + +from ...file_lock import exclusive_file_lock +from ...history import load_registry +from ...paths import resolve_runtime_root +from ...upgrade import codex_home +from ..capability_hooks import ( + TURN_START_HOOK_RESULT_SCHEMA_VERSION, + TurnStartHookRegistration, + dispatch_turn_start_hooks, +) +from .automation_upgrade import _atomic, _connect, _read, digest + +_SCHEMA = "loopx_deferred_prompt_upgrade_v0" +_HOOK = "heartbeat.prompt_upgrade" + + +def receipt_path(runtime_root: Path, registry: Path, home: Path) -> Path: + scope = json.dumps([str(registry.resolve()), str(home.resolve())]) + return runtime_root / "automation-prompt-upgrades" / (digest(scope) + ".json") + + +def _read_receipts(path: Path) -> dict[str, Any]: + if not path.exists(): + return {} + payload = json.loads(path.read_text(encoding="utf-8")) + if (not isinstance(payload, dict) or payload.get("schema_version") != _SCHEMA + or not isinstance(payload.get("entries"), dict)): + raise ValueError("invalid deferred prompt upgrade receipt") + return dict(payload["entries"]) + + +def record_deferred_upgrades(*, registry: Path, home: Path, runtime_root: str | None, + cli_bin: str, entries: dict[str, Any], results: list[dict[str, Any]]) -> None: + root = resolve_runtime_root(load_registry(registry), runtime_root, registry_path=registry) + path = receipt_path(root, registry, home) + if not path.exists() and not any(result["status"] == "deferred" for result in results): + return + # Independent sync-installed subsets cannot erase each other's reminders. + with exclusive_file_lock(path): + pending = _read_receipts(path) + for result in results: + identifier = result["automation_id"] + if result["status"] == "deferred": + entry = entries[identifier] + pending[identifier] = {key: entry[key] for key in ( + "goal_id", "agent_id", "prompt_sha256", "target_thread_id", + )} + pending[identifier].update(cli_bin=cli_bin, runtime_root=runtime_root) + else: + pending.pop(identifier, None) + if pending: + _atomic(path, json.dumps({"schema_version": _SCHEMA, "entries": pending})) + elif path.exists(): + path.unlink() + + +def prompt_upgrade_hook(*, registry: Path, runtime_root: Path, goal_id: str, + agent_id: str) -> TurnStartHookRegistration | None: + home = codex_home().expanduser().resolve() + pending = _read_receipts(receipt_path(runtime_root, registry, home)) + candidates = [(key, value) for key, value in pending.items() + if isinstance(value, dict) and value.get("goal_id") == goal_id + and value.get("agent_id") == agent_id] + if len(candidates) != 1: + return None + identifier, record = candidates[0] + command = [record["cli_bin"], "--format", "json", "--registry", str(registry)] + if record.get("runtime_root"): + command += ["--runtime-root", record["runtime_root"]] + command += ["automation-prompts", "plan", "--codex-home", str(home), + "--automation-id", identifier, "--cli-bin", record["cli_bin"]] + required_read = { + "kind": "automation_prompt_upgrade", + "command": shlex.join(command), + "reason": "A managed prompt upgrade is pending. Read the fresh plan; review and apply only the prompt through automation_update, then read back. Preserve other fields; no quota spend for repair. Continue normal work under its existing decision.", + "ordering": "before_work", + } + + def produce() -> dict[str, Any]: + with closing(_connect(home)) as connection: + _, item, _ = _read(home, identifier, connection) + applicable = (digest(item["prompt"]) == record.get("prompt_sha256") + and item["target_thread_id"] == record.get("target_thread_id")) + return { + "schema_version": TURN_START_HOOK_RESULT_SCHEMA_VERSION, + "hook_id": _HOOK, "capability_id": "automation-prompt-upgrade", + "phase": "turn_start", "status": "observed" if applicable else "empty", + "observation_count": int(applicable), "agent_read_required": applicable, + "external_reads_performed": False, "external_writes_performed": False, + "local_private_state_mutated": False, "private_content_returned": False, + "provider_payload_returned": False, "error_code": None, + } + + return TurnStartHookRegistration( + hook_id=_HOOK, capability_id="automation-prompt-upgrade", + requested_read_scope=("deferred_prompt_upgrade", "installed_automation"), + requested_write_scope=(), producer=produce, required_read=required_read, + ) + + +def extend_prompt_upgrade_reads( + dispatch: Mapping[str, Any] | None, **kwargs: Any, +) -> Mapping[str, Any] | None: + """Fixed hook loading; healthy lanes preserve their exact existing projection.""" + try: + hook = prompt_upgrade_hook(**kwargs) + if hook is None: + return dispatch + extra = dispatch_turn_start_hooks((hook,)) + except (OSError, ValueError, KeyError, TypeError): + return dispatch + if not extra["required_reads"]: + return dispatch + return {**(dispatch or {}), "required_reads": [ + *(dispatch or {}).get("required_reads", []), *extra["required_reads"], + ]} diff --git a/loopx/control_plane/heartbeat/rules.py b/loopx/control_plane/heartbeat/rules.py index 0f53a13979..264bd5f2b9 100644 --- a/loopx/control_plane/heartbeat/rules.py +++ b/loopx/control_plane/heartbeat/rules.py @@ -41,23 +41,19 @@ ) SCHEDULER_HINT_APPLICATION_RULE = ( "`scheduler_hint` no-spend. host_action=pause_or_delete_current_heartbeat -> " - "automation_update stop once, verify, end; prompt_adoption.host_action=" - "adopt_managed_bootstrap -> its prompt-only api_update_request once, no spend, " - "then readback; else apply_needed -> RRULE via automation_update; unavailable -> " - "use fallback_hint.cli_args only when projected (SQLite/app API " + "automation_update stop once, verify, end; else apply_needed -> RRULE via " + "automation_update; unavailable -> use fallback_hint.cli_args only when projected " + "(SQLite/app API " "bypass - fallback only), then ack; further failure -> failure_hint; " "ack_needed -> ack." ) SCHEDULER_HINT_COMPACT_RULE = ( "host_action=pause_or_delete_current_heartbeat: automation_update stop; " - "prompt_adoption.host_action=adopt_managed_bootstrap: prompt-only " - "api_update_request once, no spend; " "else RRULE apply via automation_update, projected fallback_hint when unavailable, " "then ack/fail. No spend." ) SCHEDULER_HINT_THIN_RULE = ( "host_action=pause_or_delete_current_heartbeat->automation_update stop(no-spend); " - "adoption=adopt(no-spend); " "else RRULE/projected-fallback_hint/ack/fail." ) RUNTIME_CAPABILITY_PROJECTION_THIN_RULE = ( diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index b08c8a9ed1..180be5b02e 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -484,6 +484,13 @@ def build_live_quota_should_run_decision( available_capabilities = remembered_runtime if route_source.startswith("loopx_turn_"): payload["runtime_root"] = str(runtime_root) + if codex_app_host and agent_id: + from ..heartbeat.prompt_upgrade_hook import extend_prompt_upgrade_reads + + turn_start_hook_dispatch = extend_prompt_upgrade_reads( + turn_start_hook_dispatch, registry=registry_path, runtime_root=runtime_root, + goal_id=goal_id, agent_id=agent_id, + ) _project_turn_start_required_reads( payload, turn_start_hook_dispatch, diff --git a/loopx/control_plane/quota/should_run_packet.py b/loopx/control_plane/quota/should_run_packet.py index bbc2e6eadf..ab7eb108d0 100644 --- a/loopx/control_plane/quota/should_run_packet.py +++ b/loopx/control_plane/quota/should_run_packet.py @@ -1317,7 +1317,6 @@ def _build_quota_should_run_payload( agent_scope_frontier=route.agent_scope_frontier, workspace_guard=prepared.workspace_guard, automation_prompt_upgrade=prepared.automation_prompt_upgrade, - automation_prompt_adoption=prepared.automation_prompt_adoption, ) if prepared.agent_scoped_user_todo_override: payload[str(prepared.agent_scoped_user_todo_override["kind"])] = ( @@ -1472,7 +1471,8 @@ def _build_quota_should_run_payload( _load_app_automation_scheduler_state( prepared.status_payload, goal_id=prepared.safe_goal_id, - agent_id=quota_decision_agent_id(payload) or prepared.requested_agent_id, + agent_id=quota_decision_agent_id(payload) + or prepared.requested_agent_id, surface=prepared.resolved_scheduler_context.context.host_surface.value, ) if prepared.resolved_scheduler_context.ok diff --git a/loopx/control_plane/quota/should_run_prepare.py b/loopx/control_plane/quota/should_run_prepare.py index 6ef2e70c23..159e938af6 100644 --- a/loopx/control_plane/quota/should_run_prepare.py +++ b/loopx/control_plane/quota/should_run_prepare.py @@ -31,7 +31,6 @@ build_goal_frontier_projection_context_from_status, ) from ..quota.error_codes import HeartbeatReceiptIdentityConflictError -from ..quota.decision_summary import quota_decision_agent_id from ..agents.capability_memory import resolve_agent_capabilities from ..quota.goal_boundary import ( goal_boundary as _goal_boundary, @@ -141,7 +140,6 @@ class _QuotaDecisionPreparation: goal_boundary: dict[str, Any] | None automation_prompt_upgrade: dict[str, Any] | None automation_prompt_upgrade_required: bool - automation_prompt_adoption: dict[str, Any] | None blocked_priority_fallback: dict[str, Any] | None stall_self_repair: dict[str, Any] | None self_repair_allowed: bool @@ -298,35 +296,6 @@ def _automation_prompt_upgrade( ) -def _automation_prompt_adoption( - status_payload: dict[str, Any], - *, - goal_id: str, - requested_agent_id: str | None, - resolved_scheduler_context: SchedulerExecutionContextResolution, -) -> dict[str, Any] | None: - """The lane's recorded prompt-only adoption, when an App automation drives it. - - Update-time reconciliation cannot apply a prompt change while the App runs, - so it records the request it reviewed. Project that record instead of - re-classifying the installed body every turn. - """ - context = resolved_scheduler_context.context - if ( - not resolved_scheduler_context.ok - or context is None - or not context.app_automation_applicable - ): - return None - from ..heartbeat.installed_prompt_update import load_pending_adoption - - return load_pending_adoption( - runtime_root=status_payload.get("runtime_root"), - goal_id=goal_id, - agent_id=quota_decision_agent_id(status_payload) or requested_agent_id, - ) - - def _todo_task_class(item: dict[str, Any]) -> str: return projection_todo_item_task_class(item) @@ -857,9 +826,6 @@ def _prepare_quota_should_run_item( goal_boundary=goal_boundary, automation_prompt_upgrade=automation_prompt_upgrade, automation_prompt_upgrade_required=automation_prompt_upgrade_required, - automation_prompt_adoption=_automation_prompt_adoption(status_payload, - goal_id=safe_goal_id, requested_agent_id=requested_agent_id, - resolved_scheduler_context=resolved_scheduler_context), blocked_priority_fallback=blocked_priority_fallback, stall_self_repair=stall_self_repair, self_repair_allowed=self_repair_allowed, diff --git a/loopx/control_plane/quota/turn_envelope.ts b/loopx/control_plane/quota/turn_envelope.ts index f2b71642e1..af43a05e52 100644 --- a/loopx/control_plane/quota/turn_envelope.ts +++ b/loopx/control_plane/quota/turn_envelope.ts @@ -294,7 +294,8 @@ function requiredReads(interaction: JsonObject, payload: JsonObject): JsonObject const result: JsonObject[] = []; for (const value of raw.slice(0, 5)) { const item = object(value); - const command = text(item.command, 360); + // Preserve routes admitted by the turn-start hook command budget. + const command = text(item.command, 1024); if (!command) continue; const compact: JsonObject = { command }; for (const field of ["kind", "reason", "source"]) { diff --git a/loopx/control_plane/scheduler/scheduler_hint.py b/loopx/control_plane/scheduler/scheduler_hint.py index 6d65f58c37..f385ce5f05 100644 --- a/loopx/control_plane/scheduler/scheduler_hint.py +++ b/loopx/control_plane/scheduler/scheduler_hint.py @@ -881,11 +881,6 @@ def build( }, "no_spend_for_cadence_change": True, } - prompt_adoption = self.payload.get("automation_prompt_adoption") - if isinstance(prompt_adoption, dict): - # The installed body decides which rules later wakes follow, so a - # recorded pending adoption is part of the observed App automation. - app_automation["prompt_adoption"] = dict(prompt_adoption) stateful_backoff = app_automation["stateful_backoff"] if host_update_failures: stateful_backoff["host_update_failures"] = [ diff --git a/loopx/control_plane/work_items/interaction_contract.py b/loopx/control_plane/work_items/interaction_contract.py index 68c8ebfd06..c633cd4da4 100644 --- a/loopx/control_plane/work_items/interaction_contract.py +++ b/loopx/control_plane/work_items/interaction_contract.py @@ -959,7 +959,8 @@ def _interaction_required_reads(payload: dict[str, Any]) -> list[dict[str, Any]] for item in reads: if not isinstance(item, dict): continue - command = protocol_action_text(item.get("command"), limit=360) + # Match the admitted turn-start command budget; never clip a valid route. + command = protocol_action_text(item.get("command"), limit=1024) if not command: continue result.append({**item, "command": command}) diff --git a/tests/control_plane/test_automation_prompt_upgrade.py b/tests/control_plane/test_automation_prompt_upgrade.py index 9d7ae1a9e7..2b8ae714df 100644 --- a/tests/control_plane/test_automation_prompt_upgrade.py +++ b/tests/control_plane/test_automation_prompt_upgrade.py @@ -13,12 +13,6 @@ from loopx.control_plane.heartbeat import automation_upgrade as upgrade -def runtime_root(tmp_path: Path) -> Path: - """One explicit runtime root per test: reconciliation records deferred adoptions.""" - - return tmp_path / "runtime" - - def fixture(tmp_path: Path, backing_kind="heartbeat"): home = tmp_path / "host" path = home / "automations/watch/automation.toml" @@ -40,7 +34,7 @@ def fixture(tmp_path: Path, backing_kind="heartbeat"): registry = tmp_path / "registry.json" state = tmp_path / "STATE.md" state.write_text("# Fixture\n", encoding="utf-8") - registry.write_text(json.dumps({"goals": [{"id": "fixture-goal", "repo": str(tmp_path), + registry.write_text(json.dumps({"common_runtime_root": str(tmp_path / "runtime"), "goals": [{"id": "fixture-goal", "repo": str(tmp_path), "state_file": str(state), "registered_agents": ["agent-a"]}]}), encoding="utf-8") return home, path, database, registry, prompt @@ -224,7 +218,7 @@ def test_failed_install_never_attempts_prompt_writes(tmp_path, monkeypatch): def forbidden(*args, **kwargs): raise AssertionError("failed installer must not invoke prompt writer") monkeypatch.setattr(lifecycle.subprocess, "run", forbidden) - result = lifecycle.update_with_prompts({}, registry=registry, runtime_root=str(runtime_root(tmp_path)), + result = lifecycle.update_with_prompts({}, registry=registry, runtime_root=None, timeout_seconds=1, runtime_update=lambda payload, **_: {"ok": False}) assert result["automation_prompt_upgrade"]["status"] == "skipped_runtime_update_failed" assert path.read_bytes() == original @@ -255,7 +249,7 @@ def run(command, **kwargs): "next_action": {"kind": "review_or_rollback"}, "recommended_action": "Inspect optional extensions"} result = lifecycle.update_with_prompts( {"install_lifecycle": {"execution_driver": "python_pip"}}, registry=registry, - runtime_root=str(runtime_root(tmp_path)), timeout_seconds=60, runtime_update=lambda *_, **__: runtime_result) + runtime_root=None, timeout_seconds=60, runtime_update=lambda *_, **__: runtime_result) assert result["automation_prompt_upgrade"]["status"] == expected_status assert len(calls) == (1 if expected_status == "attention_required" else 0) assert not result["ok"] and not result["upgrade_complete"] @@ -276,12 +270,12 @@ def test_update_identifies_owned_legacy_body_and_migrates_without_changing_sched metadata = tomllib.loads(path.read_text()) monkeypatch.setattr(lifecycle, "require_closed_app", lambda: None) monkeypatch.setattr(lifecycle.sys, "platform", "darwin") - result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path)) + result = lifecycle.reconcile(before=before, registry=registry, home=home) assert result["ok"] and result["results"][0]["status"] == "updated" after = tomllib.loads(path.read_text()) assert {k: v for k, v in after.items() if k != "prompt"} == {k: v for k, v in metadata.items() if k != "prompt"} assert upgrade.bootstrap_binding(after["prompt"]) is not None - assert lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path))["results"][0]["status"] == "current" + assert lifecycle.reconcile(before=before, registry=registry, home=home)["results"][0]["status"] == "current" @pytest.mark.parametrize("scenario", ["custom", "unsupported_host", "race", "wrong_home", "canary"]) @@ -304,7 +298,7 @@ def test_update_does_not_overwrite_custom_changed_or_foreign_hosts(tmp_path, mon with pytest.raises(ValueError, match="another host"): lifecycle.reconcile(before=before, registry=registry, home=tmp_path / "other") else: - result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path)) + result = lifecycle.reconcile(before=before, registry=registry, home=home) expected = {"custom": "review_required", "canary": "review_required", "unsupported_host": "deferred", "race": "changed_since_snapshot"}[scenario] assert result["results"][0]["status"] == expected @@ -343,12 +337,11 @@ def fail_mirror(target, text): def test_running_app_defers_to_native_api_without_touching_cached_scheduler(tmp_path, monkeypatch): from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle home, path, database, registry, _ = fixture(tmp_path) - prompt = upgrade.bootstrap_prompt(registry=registry, goal_id="fixture-goal", agent_id="agent-a", - runtime_root=str(runtime_root(tmp_path))) + prompt = upgrade.bootstrap_prompt(registry=registry, goal_id="fixture-goal", agent_id="agent-a") legacy = prompt.replace(upgrade.BOOTSTRAP, upgrade._LEGACY_BOOTSTRAP, 1).removesuffix( upgrade._BOOTSTRAP_INSTRUCTION) + upgrade._LEGACY_INSTRUCTION _set_fixture_prompt(path, database, legacy) - before = lifecycle.snapshot(registry=registry, home=home, runtime_root=str(runtime_root(tmp_path))) + before = lifecycle.snapshot(registry=registry, home=home) monkeypatch.setattr(lifecycle.sys, "platform", "darwin") probes = [] def running(command, **kwargs): @@ -358,7 +351,7 @@ def running(command, **kwargs): original = path.read_bytes() with sqlite3.connect(database) as observer: row = observer.execute("SELECT * FROM automations").fetchone() - result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path)) + result = lifecycle.reconcile(before=before, registry=registry, home=home) assert not result["ok"] and result["results"][0]["status"] == "deferred" assert observer.execute("SELECT * FROM automations").fetchone() == row assert path.read_bytes() == original @@ -538,103 +531,9 @@ def test_exact_legacy_host_loader_upgrades_to_v2_without_dropping_explicit_polic assert entry["desired_prompt"].startswith("LoopX managed heartbeat bootstrap v2\n") assert host_bootstrap_binding(entry["desired_prompt"])["permission_rule"] == "Read only" monkeypatch.setattr(lifecycle, "require_closed_app", lambda: None) - assert lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=runtime_root(tmp_path))["ok"] + assert lifecycle.reconcile(before=before, registry=registry, home=home)["ok"] assert lifecycle.snapshot(registry=registry, home=home)["entries"][0]["status"] == "current" malformed = legacy.replace("heartbeat-prompt ", "") _set_fixture_prompt(path, database, malformed) rejected = lifecycle.snapshot(registry=registry, home=home)["entries"][0] assert rejected["status"] != "current" and not rejected["automatic_eligible"] - - -def managed_legacy_fixture(path, database, registry, root) -> str: - """A managed v1 wrapper this host may still adopt, not a hand-written body.""" - - desired = upgrade.bootstrap_prompt(registry=registry, goal_id="fixture-goal", - agent_id="agent-a", runtime_root=str(root)) - legacy = desired.replace(upgrade.BOOTSTRAP, upgrade._LEGACY_BOOTSTRAP, 1).removesuffix( - upgrade._BOOTSTRAP_INSTRUCTION) + upgrade._LEGACY_INSTRUCTION - _set_fixture_prompt(path, database, legacy) - return legacy - - -def test_deferred_adoption_is_recorded_and_stops_projecting_once_installed(tmp_path, monkeypatch): - from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle - home, path, database, registry, _ = fixture(tmp_path) - root = runtime_root(tmp_path) - legacy = managed_legacy_fixture(path, database, registry, root) - before = lifecycle.snapshot(registry=registry, home=home, runtime_root=str(root)) - assert before["entries"][0]["automatic_eligible"] is True - # An unqualified host cannot take the write, so the reviewed request has to - # outlive this one report. - monkeypatch.setattr(lifecycle.sys, "platform", "linux") - result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=root) - assert not result["ok"] and result["results"][0]["status"] == "deferred" - assert lifecycle.pending_store_path(root).is_file() - pending = lifecycle.load_pending_adoption( - runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") - assert pending["status"] == "adoption_required" - assert pending["host_action"] == lifecycle.PENDING_HOST_ACTION - assert pending["host_action_contract"] == lifecycle.PENDING_HOST_ACTION_CONTRACT - assert pending["spend_policy"] == lifecycle.PENDING_SPEND_POLICY - # Same reviewed request as the update report, scheduling and binding intact. - assert pending["api_update_request"] == result["api_updates"][0] - assert pending["api_update_request"]["arguments"]["targetThreadId"] == "thread-a" - assert pending["api_update_request"]["arguments"]["status"] == "PAUSED" - # Another lane, an absent store, and no lane identity project nothing. - assert lifecycle.load_pending_adoption( - runtime_root=root, goal_id="fixture-goal", agent_id="agent-b") is None - assert lifecycle.load_pending_adoption( - runtime_root=tmp_path / "absent", goal_id="fixture-goal", agent_id="agent-a") is None - assert lifecycle.load_pending_adoption( - runtime_root=root, goal_id="fixture-goal", agent_id=None) is None - # The exact reviewed body satisfies the obligation without re-reading the lane. - upgrade.apply_offline(home=home, automation_id="watch", expected_prompt_sha256=upgrade.digest(legacy), - desired_prompt=pending["api_update_request"]["arguments"]["prompt"]) - assert lifecycle.load_pending_adoption( - runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") is None - - -def test_applied_adoption_clears_its_record(tmp_path, monkeypatch): - from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle - home, path, database, registry, _ = fixture(tmp_path) - root = runtime_root(tmp_path) - managed_legacy_fixture(path, database, registry, root) - before = lifecycle.snapshot(registry=registry, home=home, runtime_root=str(root)) - monkeypatch.setattr(lifecycle.sys, "platform", "linux") - assert lifecycle.reconcile( - before=before, registry=registry, home=home, runtime_root=root)["api_updates"] - assert lifecycle.pending_store_path(root).is_file() - monkeypatch.setattr(lifecycle, "require_closed_app", lambda: None) - monkeypatch.setattr(lifecycle.sys, "platform", "darwin") - result = lifecycle.reconcile(before=before, registry=registry, home=home, runtime_root=root) - assert result["ok"] and result["results"][0]["status"] == "updated" - assert not lifecycle.pending_store_path(root).is_file() - assert lifecycle.load_pending_adoption( - runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") is None - - -def test_pending_store_serves_only_an_unambiguous_reviewed_record(tmp_path): - from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle - root = runtime_root(tmp_path) - request = upgrade.automation_update_request(automation_id="watch", - manifest={"name": "Fixture watch", "status": "ACTIVE", "rrule": "FREQ=HOURLY", - "target_thread_id": "thread-a"}, - expected_prompt_sha256="a" * 64, desired_prompt="LoopX managed heartbeat bootstrap v2") - record = {"automation_id": "watch", "goal_id": "fixture-goal", "agent_id": "agent-a", - "automation_status": "ACTIVE", "codex_home": str(tmp_path / "host"), - "expected_prompt_sha256": "a" * 64, "desired_sha256": "b" * 64, - "source": "update_time_reconciliation", "api_update_request": request} - lifecycle.update_pending_adoptions(runtime_root=root, codex_home=str(tmp_path / "host"), - pending=[record], resolved=[]) - assert lifecycle.load_pending_adoption( - runtime_root=root, goal_id="fixture-goal", agent_id="agent-a")["automation_id"] == "watch" - # Two installed automations claiming one lane stay a review, not a host action. - lifecycle.update_pending_adoptions(runtime_root=root, codex_home=str(tmp_path / "host"), - pending=[{**record, "automation_id": "watch-2"}], resolved=[]) - assert lifecycle.load_pending_adoption( - runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") is None - # A record without its reviewed request is never handed to a host. - lifecycle.update_pending_adoptions(runtime_root=root, codex_home=str(tmp_path / "host"), - pending=[{**record, "api_update_request": None}], resolved=["watch-2"]) - assert lifecycle.load_pending_adoption( - runtime_root=root, goal_id="fixture-goal", agent_id="agent-a") is None diff --git a/tests/control_plane/test_prompt_adoption_projection.py b/tests/control_plane/test_prompt_adoption_projection.py deleted file mode 100644 index 86788fafbb..0000000000 --- a/tests/control_plane/test_prompt_adoption_projection.py +++ /dev/null @@ -1,135 +0,0 @@ -from __future__ import annotations - -import json -from pathlib import Path -import sqlite3 - -from loopx.control_plane.heartbeat import automation_upgrade as upgrade -from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle -from loopx.control_plane.quota.should_run import build_quota_should_run -from loopx.control_plane.scheduler.execution_context import ( - scheduler_execution_context_for_runtime_profile, -) -from loopx.control_plane.testing.quota_fixtures import ( - quota_status_payload, - quota_todo_item, -) - - -APP_CONTEXT = scheduler_execution_context_for_runtime_profile("codex_app_heartbeat") -CLI_CONTEXT = scheduler_execution_context_for_runtime_profile("codex_cli") -GOAL_ID = "fixture-goal" -AGENT_ID = "agent-a" -THREAD_ID = "thread-a" -FROZEN_BODY = f"Advance `{GOAL_ID}` from registry. --agent-id {AGENT_ID}" -DESIRED_BODY = "LoopX managed heartbeat bootstrap v2" - - -def _host_with_frozen_body(tmp_path: Path) -> tuple[Path, Path]: - home = tmp_path / "host" - path = home / "automations/watch/automation.toml" - path.parent.mkdir(parents=True) - path.write_text('version = 1\nid = "watch"\nname = "Fixture watch"\nkind = "heartbeat"\n' - 'status = "PAUSED"\ntarget_thread_id = "' + THREAD_ID + '"\n' - 'rrule = "FREQ=MINUTELY;INTERVAL=3"\n' - 'prompt = ' + json.dumps(FROZEN_BODY) + "\n", encoding="utf-8") - database = home / "sqlite/codex-dev.db" - database.parent.mkdir() - with sqlite3.connect(database) as connection: - connection.execute("CREATE TABLE automations (id TEXT PRIMARY KEY, kind TEXT, prompt TEXT," - " status TEXT, target_thread_id TEXT, rrule TEXT, model TEXT," - " updated_at INTEGER, next_run_at INTEGER)") - connection.execute("INSERT INTO automations VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", - ("watch", "heartbeat", FROZEN_BODY, "PAUSED", THREAD_ID, - "FREQ=MINUTELY;INTERVAL=3", "fixture-model", 123, 456)) - registry = tmp_path / "registry.json" - state = tmp_path / "STATE.md" - state.write_text("# Fixture\n", encoding="utf-8") - registry.write_text(json.dumps({"goals": [{"id": GOAL_ID, "repo": str(tmp_path), - "state_file": str(state), "registered_agents": [AGENT_ID]}]}), encoding="utf-8") - return home, registry - - -def _record_deferred_adoption(tmp_path: Path, home: Path, desired: str) -> None: - """Store the reviewed request an update-time reconciliation could not apply.""" - - lifecycle.update_pending_adoptions(runtime_root=tmp_path / "runtime", codex_home=str(home), - resolved=[], - pending=[{"automation_id": "watch", "goal_id": GOAL_ID, "agent_id": AGENT_ID, - "automation_status": "PAUSED", "codex_home": str(home), - "expected_prompt_sha256": upgrade.digest(FROZEN_BODY), - "desired_sha256": upgrade.digest(desired), - "source": "update_time_reconciliation", - "api_update_request": upgrade.automation_update_request(automation_id="watch", - manifest={"name": "Fixture watch", "status": "PAUSED", - "rrule": "FREQ=MINUTELY;INTERVAL=3", "target_thread_id": THREAD_ID}, - expected_prompt_sha256=upgrade.digest(FROZEN_BODY), desired_prompt=desired)}]) - - -def _lane_payload(tmp_path: Path, registry: Path) -> dict: - """One runnable lane whose packet has to carry the recorded obligation.""" - - return { - **quota_status_payload( - goal_id=GOAL_ID, - status="active", - agent_todo_items=[ - quota_todo_item( - todo_id="todo_current001", - index=1, - priority="P1", - title="Advance the reviewed slice.", - claimed_by=AGENT_ID, - ) - ], - recommended_action="Advance the reviewed slice.", - coordination={"agent_model": "peer_v1", "registered_agents": [AGENT_ID]}, - claim_scope_agent_id=AGENT_ID, - ), - "registry": str(registry), - "runtime_root": str(tmp_path / "runtime"), - } - - -def _packet(tmp_path: Path, registry: Path, context): - return build_quota_should_run( - _lane_payload(tmp_path, registry), - goal_id=GOAL_ID, - agent_id=AGENT_ID, - scheduler_execution_context=context, - ) - - -def test_recorded_adoption_is_projected_into_the_lane_packet(tmp_path): - home, registry = _host_with_frozen_body(tmp_path) - _record_deferred_adoption(tmp_path, home, DESIRED_BODY) - hint = _packet(tmp_path, registry, APP_CONTEXT)["scheduler_hint"] - adoption = hint["app_automation"]["prompt_adoption"] - assert adoption["status"] == "adoption_required" - assert adoption["automation_id"] == "watch" - assert adoption["host_action"] == lifecycle.PENDING_HOST_ACTION - assert adoption["host_action_contract"] == lifecycle.PENDING_HOST_ACTION_CONTRACT - assert adoption["spend_policy"] == lifecycle.PENDING_SPEND_POLICY - assert adoption["api_update_request"]["arguments"] == { - "mode": "update", "id": "watch", "kind": "heartbeat", "name": "Fixture watch", - "status": "PAUSED", "rrule": "FREQ=MINUTELY;INTERVAL=3", "targetThreadId": THREAD_ID, - "notificationPolicy": None, "prompt": DESIRED_BODY} - # The legacy Codex App projection carries the same recorded obligation. - assert hint["codex_app"]["prompt_adoption"] == adoption - - -def test_applied_body_and_non_app_hosts_project_no_adoption(tmp_path): - home, registry = _host_with_frozen_body(tmp_path) - _record_deferred_adoption(tmp_path, home, DESIRED_BODY) - assert "prompt_adoption" not in json.dumps(_packet(tmp_path, registry, CLI_CONTEXT)["scheduler_hint"]) - # A different installed body than the reviewed one stays pending. - assert _packet(tmp_path, registry, APP_CONTEXT)["scheduler_hint"]["app_automation"]["prompt_adoption"] - path = home / "automations/watch/automation.toml" - path.write_text(upgrade._replace_prompt(path.read_text(), DESIRED_BODY), encoding="utf-8") - assert "prompt_adoption" not in json.dumps(_packet(tmp_path, registry, APP_CONTEXT)["scheduler_hint"]) - - -def test_lane_without_a_record_projects_no_adoption(tmp_path): - home, registry = _host_with_frozen_body(tmp_path) - hint = _packet(tmp_path, registry, APP_CONTEXT)["scheduler_hint"] - assert "prompt_adoption" not in json.dumps(hint) diff --git a/tests/control_plane/test_prompt_upgrade_hook.py b/tests/control_plane/test_prompt_upgrade_hook.py new file mode 100644 index 0000000000..dab3d3e495 --- /dev/null +++ b/tests/control_plane/test_prompt_upgrade_hook.py @@ -0,0 +1,202 @@ +"""Deferred upgrade -> conditional read -> quiet, using isolated real host stores.""" +from __future__ import annotations + +from datetime import datetime, timezone +import json +import shlex +import sqlite3 +import subprocess +import sys + +import pytest + +from loopx.control_plane.heartbeat import automation_upgrade as upgrade +from loopx.control_plane.heartbeat import installed_prompt_update as lifecycle +from loopx.control_plane.heartbeat import prompt_upgrade_hook as hooks +from loopx.control_plane.quota.cli_projection import compact_quota_should_run_cli_payload +from loopx.control_plane.quota.turn_envelope import build_turn_envelope +from loopx.control_plane.quota.live_decision import build_live_quota_should_run_decision +from loopx.control_plane.scheduler.execution_context import scheduler_execution_context_for_runtime_profile +from loopx.control_plane.testing.quota_fixtures import quota_status_payload +from test_automation_prompt_upgrade import fixture, _set_fixture_prompt + + +def _deferred(tmp_path, monkeypatch, *, cli_bin="loopx", explicit_root=True): + home, path, database, registry, _ = fixture(tmp_path) + registration = json.loads(registry.read_text()) + registration["goals"][0]["coordination"] = {"registered_agents": ["agent-a"]} + registry.write_text(json.dumps(registration)) + root = tmp_path / "runtime" + root_arg = str(root) if explicit_root else None + desired = upgrade.bootstrap_prompt(registry=registry, goal_id="fixture-goal", + agent_id="agent-a", runtime_root=root_arg, cli_bin=cli_bin) + legacy = desired.replace(upgrade.BOOTSTRAP, upgrade._LEGACY_BOOTSTRAP, 1).removesuffix( + upgrade.BOOTSTRAP_INSTRUCTION) + upgrade._LEGACY_INSTRUCTION + _set_fixture_prompt(path, database, legacy) + before = lifecycle.snapshot(registry=registry, home=home, runtime_root=root_arg, cli_bin=cli_bin) + def unavailable(): + raise ValueError("host is running") + monkeypatch.setattr(lifecycle, "require_closed_app", unavailable) + monkeypatch.setenv("CODEX_HOME", str(home)) + result = lifecycle.reconcile(before=before, registry=registry, home=home, + runtime_root=root_arg, cli_bin=cli_bin) + assert result["results"] == [{"automation_id": "watch", "status": "deferred", "reason": "host is running"}] + return home, path, database, registry, root, legacy, desired, before + + +def _reads(registry, root, *, agent_id="agent-a", dispatch=None): + return hooks.extend_prompt_upgrade_reads(dispatch, registry=registry, runtime_root=root, + goal_id="fixture-goal", agent_id=agent_id) + + +@pytest.mark.parametrize("explicit_root", [False, True]) +def test_pending_read_uses_fresh_real_cli_plan_and_stops_after_adoption(tmp_path, monkeypatch, explicit_root): + home, path, database, registry, root, legacy, desired, before = _deferred( + tmp_path, monkeypatch, explicit_root=explicit_root) + receipt = hooks.receipt_path(root, registry, home) + contents = receipt.read_bytes() + assert legacy not in contents.decode() and desired not in contents.decode() + assert "api_update_request" not in contents.decode() + assert receipt.stat().st_mode & 0o077 == 0 + host_before = path.read_bytes(), database.read_bytes() + hint = _reads(registry, root)["required_reads"][0] + assert hint["source"] == "turn_start_capability_hook" + command = shlex.split(hint["command"]) + assert command[command.index("--automation-id") + 1] == "watch" + result = subprocess.run([sys.executable, "-m", "loopx.cli", *command[1:]], + capture_output=True, text=True) + assert result.returncode == 0, result.stdout + result.stderr + plan = json.loads(result.stdout) + assert plan["entries"][0]["desired_prompt"] == desired + assert (path.read_bytes(), database.read_bytes()) == host_before + assert receipt.read_bytes() == contents + # Both host mirrors acknowledge adoption; no receipt mutation/ACK is needed. + _set_fixture_prompt(path, database, desired) + existing = {"required_reads": [{"kind": "existing"}]} + assert _reads(registry, root, dispatch=existing) is existing + assert receipt.read_bytes() == contents + lifecycle.reconcile(before=before, registry=registry, home=home, + runtime_root=str(root) if explicit_root else None) + assert not receipt.exists() + + +@pytest.mark.parametrize("change", ["prompt", "thread", "divergent", "deleted", "missing", "ambiguous", "schema", "invalid_id"]) +def test_stale_ambiguous_or_unreadable_receipts_cannot_inject_a_repair(tmp_path, monkeypatch, change): + home, path, database, registry, root, legacy, _, _ = _deferred(tmp_path, monkeypatch) + if change == "prompt": + _set_fixture_prompt(path, database, legacy + "\nOwner customization") + elif change in {"thread", "deleted"}: + field, value = ("target_thread_id", "thread-b") if change == "thread" else ("status", "DELETED") + old = "thread-a" if change == "thread" else "PAUSED" + path.write_text(path.read_text().replace(f'{field} = "{old}"', f'{field} = "{value}"')) + with sqlite3.connect(database) as connection: + connection.execute(f"UPDATE automations SET {field}=?", (value,)) + elif change == "divergent": + with sqlite3.connect(database) as connection: + connection.execute("UPDATE automations SET prompt='different'") + elif change == "missing": + path.unlink() + else: + receipt = hooks.receipt_path(root, registry, home) + payload = json.loads(receipt.read_text()) + if change == "schema": + payload["schema_version"] = "unknown" + else: + payload["entries"]["other" if change == "ambiguous" else "../outside"] = payload["entries"]["watch"] + if change == "invalid_id": + del payload["entries"]["watch"] + receipt.write_text(json.dumps(payload)) + existing = {"required_reads": [{"kind": "existing"}]} + assert _reads(registry, root, dispatch=existing) is existing + + +def test_scope_is_exact_and_schedule_changes_do_not_hide_pending_prompt(tmp_path, monkeypatch): + home, path, database, registry, root, _, _, _ = _deferred(tmp_path, monkeypatch, cli_bin="loopx-canary") + assert _reads(registry, root, agent_id="other") is None + assert _reads(registry.with_name("other.json"), root) is None + monkeypatch.setenv("CODEX_HOME", str(tmp_path / "other-host")) + assert _reads(registry, root) is None + monkeypatch.setenv("CODEX_HOME", str(home)) + path.write_text(path.read_text().replace("FREQ=HOURLY", "FREQ=DAILY")) + with sqlite3.connect(database) as connection: + connection.execute("UPDATE automations SET rrule='FREQ=DAILY'") + hint = _reads(registry, root)["required_reads"][0] + command = shlex.split(hint["command"]) + assert command[0] == "loopx-canary" + assert command[command.index("--cli-bin") + 1] == "loopx-canary" + assert "FREQ=" not in hint["command"] + + +def test_partial_reconciliation_preserves_other_pending_entries(tmp_path, monkeypatch): + home, _, _, registry, root, _, _, before = _deferred(tmp_path, monkeypatch) + entries = {"other": {**before["entries"][0], "agent_id": "agent-b"}} + hooks.record_deferred_upgrades(registry=registry, home=home, runtime_root=str(root), + cli_bin="loopx", entries=entries, results=[{"automation_id": "other", "status": "deferred"}]) + hooks.record_deferred_upgrades(registry=registry, home=home, runtime_root=str(root), + cli_bin="loopx", entries={}, results=[{"automation_id": "watch", "status": "current"}]) + assert set(json.loads(hooks.receipt_path(root, registry, home).read_text())["entries"]) == {"other"} + + +@pytest.mark.parametrize("route_source", ["quota_cli_invocation", "loopx_turn_run_once"]) +def test_live_decision_adds_only_existing_required_read_channel(tmp_path, monkeypatch, route_source): + home, path, database, registry, root, _, desired, _ = _deferred(tmp_path, monkeypatch) + monkeypatch.setattr("loopx.control_plane.scheduler.scheduler_hint.now_utc", + lambda: datetime(2026, 1, 1, tzinfo=timezone.utc)) + status = quota_status_payload(goal_id="fixture-goal", status="active", recommended_action="Advance the selected work", + coordination={"registered_agents": ["agent-a"]}, + agent_todo_items=[{"todo_id": "todo_ordinary_work", "index": 1, + "text": "[P1] Advance the selected work", "role": "agent", "status": "open", + "priority": "P1", "task_class": "advancement_task", "claimed_by": "agent-a"}]) + kwargs = dict(goal_id="fixture-goal", agent_id="agent-a", available_capabilities=["shell"], + include_scheduler_detail=False, codex_app_current_rrule="FREQ=HOURLY", + registry_path=registry, runtime_root=root, route_source=route_source, + scheduler_execution_context=scheduler_execution_context_for_runtime_profile("codex_app_heartbeat")) + receipt = hooks.receipt_path(root, registry, home) + contents = receipt.read_bytes() + receipt.unlink() + baseline = build_live_quota_should_run_decision(status, **kwargs) + receipt.write_bytes(contents) + pending = build_live_quota_should_run_decision(status, **kwargs) + assert pending["required_reads"][-1]["kind"] == "automation_prompt_upgrade" + assert pending["interaction_contract"]["agent_channel"]["required_reads"] == pending["required_reads"] + hint = pending["required_reads"][-1] + assert len(hint["command"]) > 360 + assert compact_quota_should_run_cli_payload(pending)["required_reads"][-1] == hint + assert build_turn_envelope(pending)["required_reads"][-1]["command"] == hint["command"] + for key in baseline.keys() | pending.keys(): + if key not in {"required_reads", "interaction_contract", "protocol_action_packet"}: + assert pending.get(key) == baseline.get(key), key + _set_fixture_prompt(path, database, desired) + assert build_live_quota_should_run_decision(status, **kwargs) == baseline + # Other hosts never load the Codex receipt hook at all. + monkeypatch.setattr(hooks, "extend_prompt_upgrade_reads", lambda *a, **k: pytest.fail("unexpected Codex hook")) + for profile in ("generic_cli", "codex_cli", "trae_app"): + kwargs["scheduler_execution_context"] = scheduler_execution_context_for_runtime_profile(profile) + build_live_quota_should_run_decision(status, **kwargs) + + +def test_healthy_lane_does_not_dispatch_or_read_host_store(tmp_path, monkeypatch): + monkeypatch.setenv("CODEX_HOME", str(tmp_path / "host")) + monkeypatch.setattr(hooks, "dispatch_turn_start_hooks", lambda *a: pytest.fail("healthy lane dispatched")) + assert _reads(tmp_path / "registry.json", tmp_path / "runtime") is None + assert not list(tmp_path.iterdir()) + + +def test_real_quota_cli_exposes_hint_without_bypassing_health_gate(tmp_path, monkeypatch): + _, path, database, registry, root, _, desired, _ = _deferred(tmp_path, monkeypatch) + # This minimal lifecycle fixture deliberately lacks a healthy active state. + command = [sys.executable, "-m", "loopx.cli", "--format", "json", + "--registry", str(registry), "--runtime-root", str(root), + "quota", "should-run", "--goal-id", "fixture-goal", "--agent-id", "agent-a", + "--codex-app", "--scan-path", str(tmp_path / "STATE.md")] + pending_result = subprocess.run(command, capture_output=True, text=True) + pending = json.loads(pending_result.stdout) + assert pending["required_reads"][-1]["kind"] == "automation_prompt_upgrade" + assert pending["should_run"] is False + _set_fixture_prompt(path, database, desired) + current_result = subprocess.run(command, capture_output=True, text=True) + current = json.loads(current_result.stdout) + assert not current.get("required_reads") + assert current_result.returncode == pending_result.returncode == 1 + for key in ("should_run", "decision", "reason", "state"): + assert current[key] == pending[key] From 4479d82674db4fef32f28e812fc8625b5907fb88 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 15 Sep 2026 18:55:45 +0800 Subject: [PATCH 5/6] docs(heartbeat): describe conditional prompt upgrade hints Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/heartbeat-automation-prompt.md | 7 +-- docs/reference/automation-prompt-upgrades.md | 63 ++++++++++++-------- skills/loopx-project/SKILL.md | 6 +- 3 files changed, 41 insertions(+), 35 deletions(-) diff --git a/docs/heartbeat-automation-prompt.md b/docs/heartbeat-automation-prompt.md index e7a913cca8..463b7e64ec 100644 --- a/docs/heartbeat-automation-prompt.md +++ b/docs/heartbeat-automation-prompt.md @@ -637,12 +637,7 @@ heartbeats should search/use `automation_update` when available. If `scheduler_hint.app_automation.host_action=pause_or_delete_current_heartbeat`: in that terminal case, call `automation_update` once to pause the current heartbeat (delete only if pause is unavailable), verify the host result, spend -no quota, and end the turn without a scheduler ACK. When the same lane reports -`scheduler_hint.app_automation.prompt_adoption.host_action=adopt_managed_bootstrap`, -its installed body is not the current managed loader: apply that reported -prompt-only `api_update_request` once through `automation_update` after checking -the reported prompt hash, read the automation back, and spend no quota. Otherwise -call it only when +no quota, and end the turn without a scheduler ACK. Otherwise call it only when `scheduler_hint.app_automation.stateful_backoff.apply_needed=true` and `scheduler_hint.app_automation.recommended_rrule` is present. After a successful RRULE update, run `loopx` with diff --git a/docs/reference/automation-prompt-upgrades.md b/docs/reference/automation-prompt-upgrades.md index 93fb25c189..7757c697b6 100644 --- a/docs/reference/automation-prompt-upgrades.md +++ b/docs/reference/automation-prompt-upgrades.md @@ -68,21 +68,30 @@ read-only preview, not the upgrade executor. Do not infer a manual-only policy from its `adoption_required` status. Custom or inconsistent entries still need review; automatic prompt migration never grants scheduler or thread authority. -## Deferred adoption obligations - -Reconciliation cannot write while the App runs, and a completed report cannot -carry an unapplied migration into the next turn. An entry that is reviewable but -not written is therefore recorded per lane in the runtime root's -`app-automation-prompt-adoptions.json`, holding the reviewed prompt-only -`automation_update` request, both prompt digests, and the host that was read. -`quota should-run` projects that record as -`scheduler_hint.app_automation.prompt_adoption` with -`host_action=adopt_managed_bootstrap` and the no-spend policy, so the obligation -survives the update that discovered it. The turn applies the request once -through `automation_update`, reads the automation back, and spends no quota. -The projection is not adoption authority and not delivery permission: it names -only an entry this host already reviewed, it disappears as soon as the exact -body is installed, and nothing re-classifies installed automations per turn. +## Deferred upgrade hint + +When an automatically eligible migration is deferred, reconciliation records +only its identity, old prompt digest and CLI route under the private runtime root's +`automation-prompt-upgrades/` directory. Records are scoped to the registry and +Codex home; they contain no prompt body or saved host update request. + +Codex App heartbeat decisions load a fixed, read-only turn-start hook from the +existing heartbeat lifecycle. It uses the existing typed capability-hook +observation and required-read contracts, without adding a standalone capability, +provider package, prompt template or scheduler action. For one unambiguous pending +entry whose old prompt and thread still match both host stores, the hook inserts +an `automation-prompts plan --automation-id ...` read into the existing Agent/CLI +channel. Read that fresh plan, review the prompt-only adoption through the App, +and read back the result. Repair alone spends no quota; normal work keeps its +existing decision and permission boundaries. + +No pending entry, an adopted prompt, customization, a changed thread, ambiguous +identity or unavailable host evidence produces no adoption hint. With no pending receipt +the hook does not open the host database or dispatch a capability call. +An unrelated RRULE change does not hide a pending prompt. The next reconciliation +removes resolved records; the turn never needs a new ACK or state write. +Detection remains update-time: edits made outside LoopX between updates are not +new migration candidates. Re-run `automation-prompts plan` for explicit review. On the qualified macOS heartbeat schema, direct migration requires the App closed. The adapter holds a SQLite writer transaction through TOML delivery, @@ -219,12 +228,18 @@ gh 登录,仍失败则明确要求已核验 SHA,不切换分支或静默覆 不支持的存储仍需原生 API;运行中的本轮不热切换。普通测试不消耗模型 token, 真实模型发布资格仍需独立评测,不能由迁移成功推断。 -App 运行中对账无法写入,而一次性报告也带不走未应用的迁移。因此可审阅但未写入的 -条目会按 lane 记录到 runtime root 的 `app-automation-prompt-adoptions.json`,保存 -已审阅的、仅改 prompt 的 `automation_update` 请求、两个 prompt 摘要以及读取来源。 -`quota should-run` 将该记录投影为 -`scheduler_hint.app_automation.prompt_adoption`,携带 -`host_action=adopt_managed_bootstrap` 与 no-spend 策略,使义务不被一次性报告带走; -turn 通过 `automation_update` 应用一次并读回,不消耗额度。该投影不是采纳权威、 -也不授予交付权限:它只描述本 host 已审阅过的条目,目标 body 一旦安装即自动消失, -且不按轮重新分类已安装的 automation。 +自动升级候选未能应用时,对账只在 runtime root 的 +`automation-prompt-upgrades/` 中记录按 registry 和 Codex home 隔离的身份、 +prompt 摘要与 CLI 路由,不保存 prompt 正文或宿主更新请求。Codex App heartbeat +固定加载现有 heartbeat 生命周期内的只读 turn-start hook,复用已有的类型化 +capability-hook 观察与 required-read 契约,不新增独立 capability、provider 包、 +prompt 模板或 scheduler action。只有唯一未完成项的旧 prompt 和线程仍与两个 +宿主存储一致时,才在现有 Agent/CLI 通道插入指定 automation 的最新 plan 读取提示。 +按最新计划审阅、通过 App 仅更新 prompt 并读回;修复本身不消耗额度,正常工作仍按 +原有决策和权限执行。 + +无未完成项、已升级、自定义修改、线程变化、身份歧义或无法核验时不注入采纳提示; +无未完成记录时不打开宿主数据库、不调用 capability 分发器。单独的 RRULE 变化不影响提示。 +完成升级即停止提示,下次对账清除记录,无需新的 ACK 或按轮状态写入。发现仍发生在 +升级时;两次升级之间的外部修改不会自动成为迁移候选,可显式运行 +`automation-prompts plan` 审阅。 diff --git a/skills/loopx-project/SKILL.md b/skills/loopx-project/SKILL.md index d0bdbfe696..5e13884959 100644 --- a/skills/loopx-project/SKILL.md +++ b/skills/loopx-project/SKILL.md @@ -582,11 +582,7 @@ search/use `automation_update` when available. If `automation_update` once to pause the current heartbeat (delete only when the host cannot pause), verify the host result, spend no quota, and end the turn. This terminal host action takes precedence over RRULE handling and requires no -scheduler ACK. When the same lane reports -`scheduler_hint.app_automation.prompt_adoption.host_action=adopt_managed_bootstrap`, -apply that prompt-only `api_update_request` once through `automation_update` -after checking the reported prompt hash, read the automation back, and spend no -quota. Otherwise use `automation_update` only when +scheduler ACK. Otherwise use `automation_update` only when `scheduler_hint.app_automation.stateful_backoff.apply_needed=true` and `scheduler_hint.app_automation.recommended_rrule` is present. After a successful RRULE update, run `loopx` with From 7dd4332b86356faab8ce463dd812cf890408e2b8 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Tue, 15 Sep 2026 19:07:24 +0800 Subject: [PATCH 6/6] fix(hooks): carry bounded prompt budget only with active hint reads Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- docs/reference/automation-prompt-upgrades.md | 12 +++++++++ loopx/control_plane/capability_hooks.ts | 24 +++++++++++++---- .../heartbeat/prompt_upgrade_hook.py | 1 + loopx/control_plane/quota/live_decision.py | 2 +- loopx/control_plane/quota/turn_envelope.ts | 11 +++++--- .../quota/turn_envelope_budget.ts | 26 +++++++++++++----- .../work_items/interaction_contract.py | 4 +-- .../control_plane/test_prompt_upgrade_hook.py | 6 ++++- .../control_plane_ts/capability_hooks.test.ts | 16 +++++++++++ tests/control_plane_ts/turn_envelope.test.ts | 27 +++++++++++++++++++ 10 files changed, 110 insertions(+), 19 deletions(-) diff --git a/docs/reference/automation-prompt-upgrades.md b/docs/reference/automation-prompt-upgrades.md index 7757c697b6..d2f9f8535b 100644 --- a/docs/reference/automation-prompt-upgrades.md +++ b/docs/reference/automation-prompt-upgrades.md @@ -85,6 +85,13 @@ channel. Read that fresh plan, review the prompt-only adoption through the App, and read back the result. Repair alone spends no quota; normal work keeps its existing decision and permission boundaries. +The active hint declares `prompt_budget_bytes=1536` in its required read. +The typed hook validates this optional allowance (at most 2048 bytes per read). +Only emitted hook reads extend the envelope's 8192-byte budget and their command +projection allowance; inactive hooks contribute zero. Existing unbudgeted reads +retain their 360-character projection and the normal envelope budget. This is +prompt capacity, not execution, quota or adoption authority. + No pending entry, an adopted prompt, customization, a changed thread, ambiguous identity or unavailable host evidence produces no adoption hint. With no pending receipt the hook does not open the host database or dispatch a capability call. @@ -243,3 +250,8 @@ prompt 模板或 scheduler action。只有唯一未完成项的旧 prompt 和线 完成升级即停止提示,下次对账清除记录,无需新的 ACK 或按轮状态写入。发现仍发生在 升级时;两次升级之间的外部修改不会自动成为迁移候选,可显式运行 `automation-prompts plan` 审阅。 + +激活的 hint 同时声明 `prompt_budget_bytes=1536`,由类型化 hook 校验(单条最多 +2048 字节)。只有实际输出的 hook read 才增加 envelope 原有 8192 字节预算及该条 +命令的投影空间;未激活时增加量为零。未声明预算的 read 保留原有 360 字符投影和 +默认 envelope 预算。这仅增加 prompt 容量,不增加执行、额度或采纳权限。 diff --git a/loopx/control_plane/capability_hooks.ts b/loopx/control_plane/capability_hooks.ts index af11dbd789..a9e8644361 100644 --- a/loopx/control_plane/capability_hooks.ts +++ b/loopx/control_plane/capability_hooks.ts @@ -351,6 +351,15 @@ export function validateInteractionProjectionHookInvocation(input: { }; } +/** Optional per-read prompt allowance; never execution or effect authority. */ +export function turnStartPromptBudgetBytes(value: unknown): number { + if (value === undefined) return 0; + if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 1 || value > 2_048) { + throw new Error("turn-start prompt budget must be an integer from 1 to 2048 bytes"); + } + return value; +} + export function validateTurnStartHookRegistration( value: unknown, ): JsonObject & { @@ -413,11 +422,10 @@ export function validateTurnStartHookRegistration( registration.required_read, "turn-start hook required_read", ); - requireExactFields( - candidate, - TURN_START_REQUIRED_READ_FIELDS, - "turn-start hook required_read", - ); + const readFields = new Set(TURN_START_REQUIRED_READ_FIELDS); + if ("prompt_budget_bytes" in candidate) readFields.add("prompt_budget_bytes"); + requireExactFields(candidate, readFields, "turn-start hook required_read"); + const promptBudget = turnStartPromptBudgetBytes(candidate.prompt_budget_bytes); const kind = requiredString(candidate.kind, "turn-start hook required_read kind"); const command = requiredString( candidate.command, @@ -443,6 +451,12 @@ export function validateTurnStartHookRegistration( throw new Error("turn-start hook required_read ordering is invalid"); } requiredRead = { kind, command, reason, ordering: "before_work" }; + if (promptBudget) { + requiredRead.prompt_budget_bytes = promptBudget; + if (Buffer.byteLength(JSON.stringify(requiredRead), "utf8") > promptBudget) { + throw new Error("turn-start required read exceeds its declared prompt budget"); + } + } } return { ...registration, diff --git a/loopx/control_plane/heartbeat/prompt_upgrade_hook.py b/loopx/control_plane/heartbeat/prompt_upgrade_hook.py index 21c1e34a6a..c8e5b9f131 100644 --- a/loopx/control_plane/heartbeat/prompt_upgrade_hook.py +++ b/loopx/control_plane/heartbeat/prompt_upgrade_hook.py @@ -86,6 +86,7 @@ def prompt_upgrade_hook(*, registry: Path, runtime_root: Path, goal_id: str, "command": shlex.join(command), "reason": "A managed prompt upgrade is pending. Read the fresh plan; review and apply only the prompt through automation_update, then read back. Preserve other fields; no quota spend for repair. Continue normal work under its existing decision.", "ordering": "before_work", + "prompt_budget_bytes": 1536, } def produce() -> dict[str, Any]: diff --git a/loopx/control_plane/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index 180be5b02e..c9ac6427b6 100644 --- a/loopx/control_plane/quota/live_decision.py +++ b/loopx/control_plane/quota/live_decision.py @@ -73,7 +73,7 @@ def _turn_start_required_reads( projected.append( { key: read[key] - for key in ("kind", "command", "reason", "source", "ordering") + for key in ("kind", "command", "reason", "source", "ordering", "prompt_budget_bytes") if key in read } ) diff --git a/loopx/control_plane/quota/turn_envelope.ts b/loopx/control_plane/quota/turn_envelope.ts index af43a05e52..e3d3c0f453 100644 --- a/loopx/control_plane/quota/turn_envelope.ts +++ b/loopx/control_plane/quota/turn_envelope.ts @@ -5,9 +5,10 @@ import { type JsonObject, } from "../effect_program.ts"; import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { turnStartPromptBudgetBytes } from "../capability_hooks.ts"; import { requireJsonObject } from "../runtime_decode.ts"; import { projectPendingCapabilityIntent } from "../work_items/pending_capability_intent.ts"; -import { measureTurnEnvelope, TURN_ENVELOPE_BUDGET_BYTES } from "./turn_envelope_budget.ts"; +import { measureTurnEnvelope, turnEnvelopeBudgetBytes } from "./turn_envelope_budget.ts"; export { TURN_ENVELOPE_BUDGET_BYTES } from "./turn_envelope_budget.ts"; export const TURN_ENVELOPE_SCHEMA_VERSION = "loopx_turn_envelope_v0"; @@ -294,10 +295,12 @@ function requiredReads(interaction: JsonObject, payload: JsonObject): JsonObject const result: JsonObject[] = []; for (const value of raw.slice(0, 5)) { const item = object(value); - // Preserve routes admitted by the turn-start hook command budget. - const command = text(item.command, 1024); + const promptBudget = item.source === "turn_start_capability_hook" + ? turnStartPromptBudgetBytes(item.prompt_budget_bytes) : 0; + const command = text(item.command, promptBudget || 360); if (!command) continue; const compact: JsonObject = { command }; + if (promptBudget) compact.prompt_budget_bytes = promptBudget; for (const field of ["kind", "reason", "source"]) { const rendered = text(item[field], 240); if (rendered) compact[field] = rendered; @@ -682,7 +685,7 @@ function turnActionProjection(payload: JsonObject, protocolActionFields: JsonObj // Guidance must not crowd out the actionable contract. Preserve a signed // content reference to the existing full-decision route under budget pressure. if (Object.keys(context).length > 0 - && Buffer.byteLength(JSON.stringify(projection), "utf8") > TURN_ENVELOPE_BUDGET_BYTES - 1_400) { + && Buffer.byteLength(JSON.stringify(projection), "utf8") > turnEnvelopeBudgetBytes(projection) - 1_400) { projection.agent_context = { schema_version: context.schema_version, phase: context.phase, scope: context.scope, target: "coordinator", authority: "guidance_only", delivery: "projected", diff --git a/loopx/control_plane/quota/turn_envelope_budget.ts b/loopx/control_plane/quota/turn_envelope_budget.ts index 94fd60638a..2fe48f6995 100644 --- a/loopx/control_plane/quota/turn_envelope_budget.ts +++ b/loopx/control_plane/quota/turn_envelope_budget.ts @@ -1,4 +1,5 @@ /** Performance diagnostics, never Turn admission or execution authority. */ +import { turnStartPromptBudgetBytes } from "../capability_hooks.ts"; import type { JsonObject } from "../effect_program.ts"; export const TURN_ENVELOPE_BUDGET_BYTES = 8_192; @@ -30,13 +31,25 @@ function sectionBytes(envelope: JsonObject): Record { return sizes; } +export function turnEnvelopeBudgetBytes(envelope: JsonObject): number { + const reads = Array.isArray(envelope.required_reads) ? envelope.required_reads : []; + return TURN_ENVELOPE_BUDGET_BYTES + reads.slice(0, 5).reduce((total: number, value: unknown) => { + if (!value || typeof value !== "object" || Array.isArray(value)) return total; + const read = value as JsonObject; + return total + (read.source === "turn_start_capability_hook" + ? turnStartPromptBudgetBytes(read.prompt_budget_bytes) : 0); + }, 0); +} + export function measureTurnEnvelope(envelope: JsonObject, source: JsonObject): void { // Keep v0 *_json_bytes code-point metrics for compatibility. New diagnostics // and the performance target use actual compact JSON UTF-8 bytes. const sourceChars = [...JSON.stringify(source)].length; + const budgetBytes = turnEnvelopeBudgetBytes(envelope); + const hookBudget = budgetBytes - TURN_ENVELOPE_BUDGET_BYTES; envelope.compaction = { source_json_bytes: sourceChars, envelope_json_bytes: 0, - byte_reduction_ratio: 0, budget_bytes: TURN_ENVELOPE_BUDGET_BYTES, + byte_reduction_ratio: 0, budget_bytes: budgetBytes, within_budget: true, envelope_utf8_bytes: 0, }; // Measurements include their own serialized metadata. Recompute to a fixed @@ -56,18 +69,19 @@ export function measureTurnEnvelope(envelope: JsonObject, source: JsonObject): v byte_reduction_ratio: ratioLocked ? (envelope.compaction as JsonObject).byte_reduction_ratio : sourceChars ? Math.round((1 - chars / sourceChars) * 10_000) / 10_000 : 0, - budget_bytes: TURN_ENVELOPE_BUDGET_BYTES, - within_budget: bytes <= TURN_ENVELOPE_BUDGET_BYTES, + budget_bytes: budgetBytes, + within_budget: bytes <= budgetBytes, envelope_utf8_bytes: bytes, }; - if (bytes > TURN_ENVELOPE_BUDGET_BYTES) { + if (hookBudget) metric.hook_prompt_budget_bytes = hookBudget; + if (bytes > budgetBytes) { const sections = sectionBytes(envelope); metric.warning = { code: "turn_envelope_budget_exceeded", severity: "warning", - excess_bytes: bytes - TURN_ENVELOPE_BUDGET_BYTES, + excess_bytes: bytes - budgetBytes, section_bytes: sections, over_target_sections: (Object.keys(sections) as Section[]) - .filter((key) => sections[key] > TURN_ENVELOPE_SECTION_TARGETS[key]), + .filter((key) => sections[key] > TURN_ENVELOPE_SECTION_TARGETS[key] + (key === "action" ? hookBudget : 0)), }; } envelope.compaction = metric; diff --git a/loopx/control_plane/work_items/interaction_contract.py b/loopx/control_plane/work_items/interaction_contract.py index c633cd4da4..fccf313982 100644 --- a/loopx/control_plane/work_items/interaction_contract.py +++ b/loopx/control_plane/work_items/interaction_contract.py @@ -959,8 +959,8 @@ def _interaction_required_reads(payload: dict[str, Any]) -> list[dict[str, Any]] for item in reads: if not isinstance(item, dict): continue - # Match the admitted turn-start command budget; never clip a valid route. - command = protocol_action_text(item.get("command"), limit=1024) + command = protocol_action_text(item.get("command"), limit=( + item.get("prompt_budget_bytes", 360) if item.get("source") == "turn_start_capability_hook" else 360)) if not command: continue result.append({**item, "command": command}) diff --git a/tests/control_plane/test_prompt_upgrade_hook.py b/tests/control_plane/test_prompt_upgrade_hook.py index dab3d3e495..a8ade637cb 100644 --- a/tests/control_plane/test_prompt_upgrade_hook.py +++ b/tests/control_plane/test_prompt_upgrade_hook.py @@ -162,7 +162,11 @@ def test_live_decision_adds_only_existing_required_read_channel(tmp_path, monkey hint = pending["required_reads"][-1] assert len(hint["command"]) > 360 assert compact_quota_should_run_cli_payload(pending)["required_reads"][-1] == hint - assert build_turn_envelope(pending)["required_reads"][-1]["command"] == hint["command"] + envelope = build_turn_envelope(pending) + assert envelope["required_reads"][-1]["command"] == hint["command"] + assert envelope["compaction"]["budget_bytes"] == 8192 + 1536 + assert envelope["compaction"]["hook_prompt_budget_bytes"] == 1536 + assert build_turn_envelope(baseline)["compaction"]["budget_bytes"] == 8192 for key in baseline.keys() | pending.keys(): if key not in {"required_reads", "interaction_contract", "protocol_action_packet"}: assert pending.get(key) == baseline.get(key), key diff --git a/tests/control_plane_ts/capability_hooks.test.ts b/tests/control_plane_ts/capability_hooks.test.ts index bcdcfd5ec2..d4737ba34c 100644 --- a/tests/control_plane_ts/capability_hooks.test.ts +++ b/tests/control_plane_ts/capability_hooks.test.ts @@ -476,6 +476,22 @@ test("turn-start read commands retain explicit routes with a bounded byte budget } }); +test("turn-start reads can opt into a bounded prompt budget without changing defaults", () => { + const registration = turnStartRegistration(); + const read = registration.required_read; + assert.deepEqual(validateTurnStartHookRegistration(registration).required_read, read); + const budgeted = { ...read, prompt_budget_bytes: 1_536 }; + assert.deepEqual(validateTurnStartHookRegistration({ ...registration, required_read: budgeted }).required_read, budgeted); + for (const value of [0, -1, 2_049, 1.5, true, "1536", null]) { + assert.throws(() => validateTurnStartHookRegistration({ ...registration, + required_read: { ...read, prompt_budget_bytes: value }, + }), /prompt budget/); + } + assert.throws(() => validateTurnStartHookRegistration({ ...registration, + required_read: { ...read, prompt_budget_bytes: 64 }, + }), /exceeds its declared prompt budget/); +}); + test("turn-start observations require Agent reading without returning private content", () => { const observed = validateTurnStartHookInvocation({ registration: turnStartRegistration(), diff --git a/tests/control_plane_ts/turn_envelope.test.ts b/tests/control_plane_ts/turn_envelope.test.ts index 9e3e076b4e..a55fff3edc 100644 --- a/tests/control_plane_ts/turn_envelope.test.ts +++ b/tests/control_plane_ts/turn_envelope.test.ts @@ -72,6 +72,33 @@ const protocolActionFields = { agent_action: "advance one bounded segment", }; +test("only active hook reads carry additive prompt budget through the envelope", () => { + const source = payload(); + const baseline = buildTurnEnvelope({ payload: source, protocol_action_fields: protocolActionFields, scheduler_execution_args: "" }); + const command = "loopx inspect --registry /" + "route/".repeat(80) + "registry.json"; + const read = { kind: "fixture_read", command, reason: "Read the pending observation", + source: "turn_start_capability_hook" }; + source.required_reads = [read]; + const ordinary = buildTurnEnvelope({ payload: source, protocol_action_fields: protocolActionFields, scheduler_execution_args: "" }); + assert.equal((ordinary.compaction as JsonObject).budget_bytes, 8_192); + assert.equal(((ordinary.required_reads as JsonObject[])[0].command as string).length, 360); + assert.equal((ordinary.compaction as JsonObject).hook_prompt_budget_bytes, undefined); + source.required_reads = [{ ...read, prompt_budget_bytes: 1_536 }]; + const active = buildTurnEnvelope({ payload: source, protocol_action_fields: protocolActionFields, scheduler_execution_args: "" }); + assert.equal((active.required_reads as JsonObject[])[0].command, command); + assert.equal((active.compaction as JsonObject).budget_bytes, 8_192 + 1_536); + assert.equal((active.compaction as JsonObject).hook_prompt_budget_bytes, 1_536); + assert.equal((active.compaction as JsonObject).envelope_utf8_bytes, Buffer.byteLength(JSON.stringify(active))); + for (const field of ["action", "user", "scheduler", "execution_policy", "writeback"]) { + assert.deepEqual(active[field], baseline[field]); + } + source.required_reads = [{ ...read, source: "other", prompt_budget_bytes: 1_536 }]; + const unrelated = buildTurnEnvelope({ payload: source, protocol_action_fields: protocolActionFields, scheduler_execution_args: "" }); + assert.equal((unrelated.compaction as JsonObject).budget_bytes, 8_192); + delete source.required_reads; + assert.deepEqual(buildTurnEnvelope({ payload: source, protocol_action_fields: protocolActionFields, scheduler_execution_args: "" }), baseline); +}); + test("pending capability action outranks stale replan commands and remains signed", () => { const source = payload(); const command = "loopx periodic-report consume-pending --goal-id goal-turn-envelope --agent-id agent-ts --execute";