diff --git a/docs/reference/automation-prompt-upgrades.md b/docs/reference/automation-prompt-upgrades.md index 22d4edbab9..d2f9f8535b 100644 --- a/docs/reference/automation-prompt-upgrades.md +++ b/docs/reference/automation-prompt-upgrades.md @@ -68,6 +68,38 @@ 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 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. + +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. +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, compares the entire previewed manifest, preserves every non-prompt field, and @@ -202,3 +234,24 @@ gh 登录,仍失败则明确要求已核验 SHA,不切换分支或静默覆 日程、暂停状态、模型、线程、通知偏好和历史均不迁移。 不支持的存储仍需原生 API;运行中的本轮不热切换。普通测试不消耗模型 token, 真实模型发布资格仍需独立评测,不能由迁移成功推断。 + +自动升级候选未能应用时,对账只在 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` 审阅。 + +激活的 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/installed_prompt_update.py b/loopx/control_plane/heartbeat/installed_prompt_update.py index 94569f89f3..a0304d186c 100644 --- a/loopx/control_plane/heartbeat/installed_prompt_update.py +++ b/loopx/control_plane/heartbeat/installed_prompt_update.py @@ -123,6 +123,9 @@ def reconcile(*, before: dict, registry: Path, home: Path, "notificationPolicy": manifest.get("notification_policy"), "prompt": now["desired_prompt"]}}) results.append(result) + 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..c8e5b9f131 --- /dev/null +++ b/loopx/control_plane/heartbeat/prompt_upgrade_hook.py @@ -0,0 +1,129 @@ +"""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", + "prompt_budget_bytes": 1536, + } + + 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/quota/live_decision.py b/loopx/control_plane/quota/live_decision.py index b08c8a9ed1..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 } ) @@ -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/turn_envelope.ts b/loopx/control_plane/quota/turn_envelope.ts index f2b71642e1..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,9 +295,12 @@ 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); + 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; @@ -681,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 68c8ebfd06..fccf313982 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) + 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_automation_prompt_upgrade.py b/tests/control_plane/test_automation_prompt_upgrade.py index a771c8df6c..2b8ae714df 100644 --- a/tests/control_plane/test_automation_prompt_upgrade.py +++ b/tests/control_plane/test_automation_prompt_upgrade.py @@ -34,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 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..a8ade637cb --- /dev/null +++ b/tests/control_plane/test_prompt_upgrade_hook.py @@ -0,0 +1,206 @@ +"""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 + 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 + _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] 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";