From 1068a5b3e45e202f0869f09066e21be0896254c1 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Wed, 23 Sep 2026 14:53:43 +0800 Subject: [PATCH 1/5] feat: admit managed Turn host starts against owner cadence Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/turn.py | 58 ++++++++++---- loopx/cli_commands/turn_registration.py | 8 ++ loopx/cli_commands/turn_rendering.py | 16 ++++ .../control_plane/effect_runtime_handlers.ts | 3 +- .../control_plane/quota/automation_cadence.ts | 76 +++++++++++++++++-- .../turn_driver/execution_readback.py | 2 + loopx/control_plane/turn_driver/executor.py | 45 +++++++++-- 7 files changed, 179 insertions(+), 29 deletions(-) diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index a43eff6d57..5dbfc65143 100644 --- a/loopx/cli_commands/turn.py +++ b/loopx/cli_commands/turn.py @@ -4,6 +4,7 @@ import argparse import json import shlex +import time from collections.abc import Callable, Mapping from pathlib import Path from typing import Any @@ -20,6 +21,7 @@ ) from ..capabilities.periodic_report.cadence_runtime import extend_cadence_turn_start_dispatch from ..control_plane.quota.live_decision import build_live_quota_should_run_decision +from ..control_plane.effect_runtime import effect_runtime_result from ..control_plane.agents.workspace_guard import capture_delivery_workspace from ..control_plane.quota.heartbeat_receipt import ( ensure_turn_heartbeat_settlement_receipt, @@ -395,22 +397,14 @@ def handle_turn_command( and persisted_effect_id != settlement_identity.effect_id ): raise ValueError("Turn settlement identity effect_id is inconsistent") - if args.execute: - ensure_turn_heartbeat_settlement_receipt( - runtime_root, - settlement_identity, - semantic_replan_guard_scoped=( - "replan_action_packet" in envelope - ), - semantic_replan_obligation_id=( - replan_obligation_id_from_packet( - envelope.get("replan_action_packet") - ) - if "replan_action_packet" in envelope - else None - ), - ) - + stable_envelope: Mapping[str, Any] = ( + envelope if isinstance(envelope, Mapping) else {} + ) + replan_guard_scoped = "replan_action_packet" in stable_envelope + replan_obligation_id = ( + replan_obligation_id_from_packet(stable_envelope.get("replan_action_packet")) + if replan_guard_scoped else None + ) def require_effect_ref( effect_ref: str, step_kind: SettlementStepKind, @@ -1049,6 +1043,37 @@ def post_settlement_reward_memory( settlement_evidence=settlement_evidence, ) + def admit_managed_start(identity: Mapping[str, Any]) -> dict[str, Any]: + now_ms = time.time_ns() // 1_000_000 + admitted = effect_runtime_result( + "quota.automation_cadence.admit", + { + "runtime_root": str(runtime_root), + "goal_id": args.goal_id, + "agent_id": args.agent_id, + "automation_id": args.automation_id, + "request_id": f"{identity['turn_key']}:{identity['attempt']}", + "trigger_at_ms": now_ms, + "now_ms": now_ms, + "manual_reason": args.manual_interval_bypass_reason, + }, + retry_safe=False, + ) + if admitted.get("admitted") is True: + ensure_turn_heartbeat_settlement_receipt( + runtime_root, + settlement_identity, + semantic_replan_guard_scoped=replan_guard_scoped, + semantic_replan_obligation_id=replan_obligation_id, + ) + return { + key: admitted.get(key) + for key in ( + "admitted", "reserved", "reason", "next_eligible_at_ms", + "min_interval_minutes", "pre_model_admission", + ) + } + payload = run_loopx_turn_once( payload, host_argv=raw_argv, @@ -1077,6 +1102,7 @@ def post_settlement_reward_memory( if args.execute and args.host == "codex-cli" else None ), + admit_start=admit_managed_start if args.execute else None, ) else: raise ValueError("turn requires the `plan` or `run-once` subcommand") diff --git a/loopx/cli_commands/turn_registration.py b/loopx/cli_commands/turn_registration.py index 7c86178ffc..d916f2fc85 100644 --- a/loopx/cli_commands/turn_registration.py +++ b/loopx/cli_commands/turn_registration.py @@ -176,6 +176,14 @@ def register_turn_commands( default_execution_mode="isolated-headless", ) run_once.add_argument("--project", required=True) + run_once.add_argument( + "--automation-id", + help="Stable automation identity for an automation-scoped execution interval.", + ) + run_once.add_argument( + "--manual-interval-bypass-reason", + help="Explicit manual intent; bypass only the interval and record this start.", + ) run_once.add_argument( "--host-command-json", "--host-adapter-command-json", diff --git a/loopx/cli_commands/turn_rendering.py b/loopx/cli_commands/turn_rendering.py index e0083e66b1..0ba3a8f4aa 100644 --- a/loopx/cli_commands/turn_rendering.py +++ b/loopx/cli_commands/turn_rendering.py @@ -1,5 +1,8 @@ from __future__ import annotations +from datetime import datetime, timezone +from typing import Any + from ..presentation.renderers.turn_envelope_markdown import ( turn_envelope_budget_warning_lines, ) @@ -33,6 +36,14 @@ def render_loopx_turn_plan_markdown(payload: dict[str, object]) -> str: def render_loopx_turn_execution_markdown(payload: dict[str, object]) -> str: effects = payload.get("effects") if isinstance(payload.get("effects"), dict) else {} + raw_admission = payload.get("admission") + admission: dict[str, Any] = raw_admission if isinstance(raw_admission, dict) else {} + next_ms = admission.get("next_eligible_at_ms") + next_at = ( + datetime.fromtimestamp(next_ms / 1000, tz=timezone.utc).isoformat().replace("+00:00", "Z") + if isinstance(next_ms, (int, float)) and not isinstance(next_ms, bool) + else None + ) receipt = payload.get("receipt") if isinstance(payload.get("receipt"), dict) else {} validation = ( payload.get("validation") if isinstance(payload.get("validation"), dict) else {} @@ -68,6 +79,11 @@ def render_loopx_turn_execution_markdown(payload: dict[str, object]) -> str: "# LoopX Turn Run Once", f"- status: {payload.get('status')}", f"- result_kind: {payload.get('result_kind')}", + *( + [f"- interval_reason: {admission.get('reason')}", + *([f"- next_eligible_at: {next_at}"] if next_at else [])] + if payload.get("status") == "interval_wait" else [] + ), *( [f"- execution_profile: {managed_executor['execution_profile']}"] if managed_executor.get("execution_profile") else [] diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 21970bb874..1faea5a3ad 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -1,4 +1,4 @@ -import {manageAutomationCadence, projectCadenceSchedule} from "./quota/automation_cadence.ts"; +import {admitAutomationStart, manageAutomationCadence, projectCadenceSchedule} from "./quota/automation_cadence.ts"; import {manageLocalAuthorityArchive} from "./coordination/local_authority_archive.ts"; import {selectPeriodicReportProgress, selectPeriodicReportApprovalRetry} from "./capabilities/periodic_report_progress.ts"; import {planIssueFixMonitorReconciliation} from "./capabilities/issue_fix_monitor_reconciliation.ts"; @@ -479,6 +479,7 @@ export function createEffectRuntimeHandlers( ["todo.external_wait.plan", planTodoExternalWaitTransition], ["scheduler.state_transition.evaluate", evaluateSchedulerStateTransition], ["quota.automation_cadence.manage", manageAutomationCadence], + ["quota.automation_cadence.admit", admitAutomationStart], ["quota.automation_cadence.schedule", projectCadenceSchedule], ["scheduler.state.evaluate", evaluateSchedulerStateOperation], ["scheduler.state.load", loadSchedulerState], diff --git a/loopx/control_plane/quota/automation_cadence.ts b/loopx/control_plane/quota/automation_cadence.ts index 73b15042c3..79d6c848e2 100644 --- a/loopx/control_plane/quota/automation_cadence.ts +++ b/loopx/control_plane/quota/automation_cadence.ts @@ -7,11 +7,13 @@ import {EffectRuntimeRequestError} from "../effect_runtime_errors.ts"; import {requireJsonObject, requireNonEmptyString, requireInteger, requireStringLiteral} from "../runtime_decode.ts"; import {schedulerStatePath} from "../scheduler/state_store.ts"; -const SCHEMA = "automation_cadence_store_v1"; +const LEGACY_SCHEMA = "automation_cadence_store_v1"; +const SCHEMA = "automation_cadence_store_v2"; const RESULT = "automation_cadence_result_v1"; type Scope = {agent_id: string | null; automation_id: string | null}; type Rule = Scope & {min_interval_minutes: number; revision: number; owner_reference: string}; -type Store = {schema_version: typeof SCHEMA; goal_id: string; revision: number; rules: Rule[]}; +type Start = Scope & {started_at_ms: number; trigger_at_ms: number; request_id: string; manual_reason: string | null}; +type Store = {schema_version: typeof SCHEMA; goal_id: string; revision: number; rules: Rule[]; starts: Start[]}; const fail = (message: string): never => {throw new EffectRuntimeRequestError(message, "automation_cadence_invalid");}; function text(value: unknown, name: string): string { const s = requireNonEmptyString(value, name).trim(); @@ -37,7 +39,8 @@ function scope(p: JsonObject): Scope { const key = (s: Scope): string => createHash("sha256").update(JSON.stringify([s.agent_id, s.automation_id])).digest("hex"); function decode(value: unknown, goal: string): Store { const p = requireJsonObject(value, "cadence store"); - if (p.schema_version !== SCHEMA || p.goal_id !== goal) fail("cadence store identity/schema mismatch"); + if ((p.schema_version !== SCHEMA && p.schema_version !== LEGACY_SCHEMA) || p.goal_id !== goal) + fail("cadence store identity/schema mismatch"); const revision = integer(p.revision, "revision"); if (!Array.isArray(p.rules)) fail("cadence rules must be an array"); const seen = new Set(); @@ -49,7 +52,19 @@ function decode(value: unknown, goal: string): Store { if (v > revision) fail("rule revision exceeds configuration revision"); return {...s, min_interval_minutes: minutes(r.min_interval_minutes), revision: v, owner_reference: text(r.owner_reference, "owner_reference")}; }); - return {schema_version: SCHEMA, goal_id: goal, revision, rules}; + if (p.starts !== undefined && !Array.isArray(p.starts)) fail("cadence starts must be an array"); + const startKeys = new Set(); + const starts = ((p.starts ?? []) as unknown[]).map(raw => { + const r = requireJsonObject(raw, "start"), s = scope(r), k = key(s); + if (!s.agent_id || startKeys.has(k)) fail("cadence start identity is invalid or duplicated"); + startKeys.add(k); + const started_at_ms = integer(r.started_at_ms, "started_at_ms"); + const trigger_at_ms = integer(r.trigger_at_ms, "trigger_at_ms"); + if (trigger_at_ms > started_at_ms) fail("cadence start trigger follows its start"); + return {...s, started_at_ms, trigger_at_ms, request_id: text(r.request_id, "request_id"), + manual_reason: r.manual_reason == null ? null : text(r.manual_reason, "manual_reason")}; + }); + return {schema_version: SCHEMA, goal_id: goal, revision, rules, starts}; } export function cadenceStorePath(runtimeRoot: string, goalId: string): string { return schedulerStatePath(runtimeRoot, {goalId, agentId: "owner-policy", surface: "quota", stateKey: "automation-cadence-v1"}); @@ -58,23 +73,70 @@ async function load(path: string, goal: string): Promise { try {return decode(JSON.parse(await readFile(path, "utf8")), goal);} catch (e) { if ((e as NodeJS.ErrnoException).code !== "ENOENT") throw e; - return {schema_version: SCHEMA, goal_id: goal, revision: 0, rules: []}; + return {schema_version: SCHEMA, goal_id: goal, revision: 0, rules: [], starts: []}; } } function applicable(r: Rule, s: Scope): boolean { return r.agent_id === null || (r.agent_id === s.agent_id && (r.automation_id === null || r.automation_id === s.automation_id)); } -function projection(store: Store, s: Scope): JsonObject { +function projection(store: Store, s: Scope, nowMs = Date.now()): JsonObject { const rules = store.rules.filter(r => applicable(r, s)); const floor = Math.max(0, ...rules.map(r => r.min_interval_minutes)); + const due = s.agent_id ? rules.filter(rule => rule.min_interval_minutes > 0).flatMap(rule => { + const start = store.starts.find(item => item.agent_id === s.agent_id && + item.automation_id === rule.automation_id); + return start ? [start.started_at_ms + rule.min_interval_minutes * 60_000] : []; + }) : []; + const next = due.length ? Math.max(...due) : null; return { schema_version: RESULT, ok: true, enabled: floor > 0, goal_id: store.goal_id, ...s, configuration_revision: store.revision, min_interval_minutes: floor, sources: rules, reason: floor === 0 ? "unconfigured" : "owner_minimum_interval", - enforcement: "scheduler_recommendation", pre_model_admission: "not_qualified", + next_eligible_at_ms: next, eligible_now: next === null || nowMs >= next, + enforcement: floor === 0 ? "scheduler_recommendation" : "managed_turn_atomic_admission_and_schedule_recommendation", + pre_model_admission: floor === 0 ? "not_qualified" : "managed_turn_only", }; } + +/** Reserve a new managed-host start under the same lock as policy changes. */ +export async function admitAutomationStart(p: JsonObject): Promise { + const goal = text(p.goal_id, "goal_id"), s = scope(p); + if (!s.agent_id) fail("admission requires agent_id"); + const path = cadenceStorePath(requireNonEmptyString(p.runtime_root, "runtime_root"), goal); + const now = integer(p.now_ms, "now_ms"), trigger = integer(p.trigger_at_ms, "trigger_at_ms"); + if (trigger > now) fail("trigger_at_ms cannot be in the future"); + const request = text(p.request_id, "request_id"); + const manual = p.manual_reason == null ? null : text(p.manual_reason, "manual_reason"); + // An absent policy does not create a store or lock file on the default-off path. + const snapshot = await load(path, goal); + const initial = projection(snapshot, s, now); + if (initial.enabled !== true) return {...initial, admitted: true, reserved: false}; + return withFileMutationLock(path, async () => { + const store = await load(path, goal), current = projection(store, s, now); + if (current.enabled !== true) return {...current, admitted: true, reserved: false}; + const activeScopes = new Set(store.rules.filter(rule => rule.min_interval_minutes > 0 && applicable(rule, s)) + .map(rule => rule.automation_id)); + const prior = store.starts.filter(start => start.agent_id === s.agent_id && + activeScopes.has(start.automation_id)); + if (prior.some(start => start.request_id === request || trigger < start.trigger_at_ms)) { + return {...current, admitted: false, reserved: false, reason: "duplicate_or_stale_trigger"}; + } + if (manual === null && current.eligible_now !== true) { + return {...current, admitted: false, reserved: false, reason: "minimum_interval_wait"}; + } + const scopes: Scope[] = [{agent_id: s.agent_id, automation_id: null}]; + if (s.automation_id && activeScopes.has(s.automation_id)) scopes.push(s); + for (const item of scopes) { + store.starts = store.starts.filter(start => key(start) !== key(item)); + store.starts.push({...item, started_at_ms: now, trigger_at_ms: trigger, request_id: request, + manual_reason: manual}); + } + await atomicWriteJson(path, store); + return {...projection(store, s, now), admitted: true, reserved: true, + reason: manual === null ? "admitted" : "explicit_manual_interval_bypass"}; + }); +} /** Pure typed calculation shared by scheduler adapters; no policy mutation. */ export function projectCadenceProgression(p: JsonObject): JsonObject { const floor = minutes(p.min_interval_minutes); diff --git a/loopx/control_plane/turn_driver/execution_readback.py b/loopx/control_plane/turn_driver/execution_readback.py index 9e19ab8742..34b35338ca 100644 --- a/loopx/control_plane/turn_driver/execution_readback.py +++ b/loopx/control_plane/turn_driver/execution_readback.py @@ -64,6 +64,8 @@ def execution_payload( "scheduler": journal.get("scheduler"), **subagent.subagent_execution_payload_projection(journal), "effects": dict(effects), + **({"admission": dict(journal["admission"])} + if isinstance(journal.get("admission"), Mapping) else {}), "quota_slot_spend_count": 1 if quota_spent else 0, **( {"settlement_result": journal["settlement_result"]} diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index d8e72b049d..54e20a23d4 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -1223,6 +1223,7 @@ def run_loopx_turn_once( terminal_closeout_resolver: TurnEffectResolver | None = None, scheduler: Scheduler | None = None, post_settlement: PostSettlement | None = None, + admit_start: Callable[[Mapping[str, Any]], dict[str, Any]] | None = None, ) -> dict[str, Any]: if host_runner is not None and host_argv is not None: raise ValueError("run-once accepts either host_argv or host_runner, not both") @@ -1315,11 +1316,42 @@ def run_loopx_turn_once( ) return payload request = require_turn_recovery_continuation(assessment) + + # Settlement recovery uses its cached host result and must not reserve + # another interval. A new host attempt, including failed-result retry, + # goes through admission before any host or journal attempt mutation. + receipt = journal.get("receipt") if isinstance(journal, Mapping) else None + validation_reinvokes_host = ( + isinstance(journal, Mapping) + and journal.get("status") == "failed" + and isinstance(receipt, Mapping) + and receipt.get("failed_phase") == "validation" + and journal.get("validation_stage") != "task_postcondition" + ) + needs_host = validation_reinvokes_host or journal is None or "typed_result" not in list( + journal.get("completed_phases") or [] + ) + admission = None + if needs_host and admit_start is not None: + admission = admit_start({ + "turn_key": turn_key, + "attempt": int(journal.get("host_attempt_count") or 0) + 1 + if journal is not None else 1, + }) + if admission.get("admitted") is not True: + waiting = { + "status": "interval_wait", + "host": host_projection, + "completed_phases": [], + "reason": str(admission.get("reason") or "automatic start not admitted"), + "admission": admission, + } + return execution_payload( + plan, waiting, execute=True, replayed=False, effects=empty_effects + ) + if journal is not None and recovery_decision is not None: journal["recovery_audit"] = build_turn_recovery_audit( - recovery_decision, - journal, - status="started", - host_invoked=None, + recovery_decision, journal, status="started", host_invoked=None, ) _write_journal(journal_path, journal) @@ -1330,7 +1362,7 @@ def run_loopx_turn_once( else {} ) if receipt.get("failed_phase") == "validation": - if journal.get("validation_stage") != "task_postcondition": + if validation_reinvokes_host: journal.pop("host_result", None) journal.pop("result_kind", None) journal["completed_phases"] = ( @@ -1357,6 +1389,9 @@ def run_loopx_turn_once( "plan": dict(plan), } _write_journal(journal_path, journal) + if admission is not None and admission.get("reserved") is True: + journal["admission"] = admission + _write_journal(journal_path, journal) effects = dict(empty_effects) From 0f15f9db2af89d3bd81362f4469888aa0f26313b Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Wed, 23 Sep 2026 14:53:55 +0800 Subject: [PATCH 2/5] test: cover cadence admission and managed Turn retries Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../automation_cadence.test.ts | 95 ++++++++++++++++++- tests/test_loopx_turn_driver.py | 34 +++++++ tests/test_loopx_turn_executor.py | 82 +++++++++++++++- 3 files changed, 205 insertions(+), 6 deletions(-) diff --git a/tests/control_plane_ts/automation_cadence.test.ts b/tests/control_plane_ts/automation_cadence.test.ts index 1b3820fc4b..4236ce2a14 100644 --- a/tests/control_plane_ts/automation_cadence.test.ts +++ b/tests/control_plane_ts/automation_cadence.test.ts @@ -1,9 +1,10 @@ import assert from "node:assert/strict"; import test from "node:test"; -import {mkdtemp, rm, readFile, writeFile, readdir} from "node:fs/promises"; +import {mkdtemp, mkdir, rm, readFile, writeFile, readdir} from "node:fs/promises"; import {tmpdir} from "node:os"; -import {join} from "node:path"; -import {cadenceStorePath, manageAutomationCadence as manage, projectCadenceProgression as progression} from "../../loopx/control_plane/quota/automation_cadence.ts"; +import {dirname, join} from "node:path"; +import {execFileSync} from "node:child_process"; +import {admitAutomationStart as admit, cadenceStorePath, manageAutomationCadence as manage, projectCadenceProgression as progression} from "../../loopx/control_plane/quota/automation_cadence.ts"; test("owner floor inherits without changing another agent; reductions and concurrent writes require authority", async () => { const root = await mkdtemp(join(tmpdir(), "cadence-")); @@ -34,12 +35,98 @@ test("owner floor inherits without changing another agent; reductions and concur const saved = JSON.parse(await readFile(cadenceStorePath(root, "fixture"), "utf8")); assert.equal(saved.revision, 5); assert.equal(saved.rules.length, 3); - assert.equal((await read()).enforcement, "scheduler_recommendation"); + assert.equal((await read()).enforcement, "managed_turn_atomic_admission_and_schedule_recommendation"); await writeFile(cadenceStorePath(root, "fixture"), '{"broken":true}'); await assert.rejects(read(), /identity\/schema mismatch/); } finally {await rm(root, {recursive: true, force: true});} }); +test("managed starts are atomic, durable, per agent, and due at the exact interval", async () => { + const root = await mkdtemp(join(tmpdir(), "cadence-admit-")); + const base = {runtime_root: root, goal_id: "fixture", agent_id: "a", automation_id: "daily"}; + const start = (request_id: string, now_ms: number, extra = {}) => admit({...base, request_id, + now_ms, trigger_at_ms: now_ms, ...extra}); + try { + assert.equal((await start("unconfigured", 1000)).reserved, false); + assert.deepEqual(await readdir(root), []); + await manage({...base, operation: "configure", expected_revision: 0, min_interval_minutes: 1440, + owner_reference: "owner-request", execute: true, agent_id: null, automation_id: null}); + const concurrent = await Promise.all([start("one", 1000), start("two", 1000)]); + assert.equal(concurrent.filter(r => r.admitted === true).length, 1); + assert.equal(concurrent.filter(r => r.admitted === false).length, 1); + assert.equal((await start("early", 1000 + 86_400_000 - 1)).reason, "minimum_interval_wait"); + assert.equal((await start("due", 1000 + 86_400_000)).admitted, true); + assert.equal((await start("due", 1000 + 2 * 86_400_000)).reason, "duplicate_or_stale_trigger"); + assert.equal((await start("stale", 1000 + 86_400_000, {trigger_at_ms: 999})).admitted, false); + assert.equal((await admit({...base, agent_id: "b", request_id: "other-agent", now_ms: 1001, trigger_at_ms: 1001})).admitted, true); + const manual = await start("manual", 1000 + 86_400_001, {manual_reason: "owner-request"}); + assert.equal(manual.reason, "explicit_manual_interval_bypass"); + assert.equal((await start("after-manual", 1000 + 2 * 86_400_000)).admitted, false); + const saved = JSON.parse(await readFile(cadenceStorePath(root, "fixture"), "utf8")); + assert.equal(saved.starts.length, 2); // one durable Goal-floor start per agent + assert.equal(saved.starts.find((row: {agent_id: string}) => row.agent_id === "a").manual_reason, "owner-request"); + await writeFile(cadenceStorePath(root, "fixture"), '{"broken":true}'); + await assert.rejects(start("corrupt", 1000 + 3 * 86_400_000), /identity\/schema mismatch/); + } finally {await rm(root, {recursive: true, force: true});} +}); + +test("automation-specific floors do not serialize unrelated automation lanes", async () => { + const root = await mkdtemp(join(tmpdir(), "cadence-scopes-")); + try { + for (const [automation_id, expected_revision] of [["first", 0], ["second", 1]] as const) { + await manage({runtime_root: root, goal_id: "fixture", agent_id: "a", automation_id, + operation: "configure", expected_revision, min_interval_minutes: 60, + owner_reference: "owner-request", execute: true}); + } + const start = (automation_id: string, request_id: string, now_ms: number) => admit({ + runtime_root: root, goal_id: "fixture", agent_id: "a", automation_id, + request_id, now_ms, trigger_at_ms: now_ms, + }); + assert.equal((await start("first", "one", 1000)).admitted, true); + assert.equal((await start("second", "two", 1001)).admitted, true); + assert.equal((await start("first", "three", 1002)).reason, "minimum_interval_wait"); + assert.equal((await start("second", "four", 1003)).reason, "minimum_interval_wait"); + const saved = JSON.parse(await readFile(cadenceStorePath(root, "fixture"), "utf8")); + assert.equal(saved.starts.length, 3); // one agent-wide record and two scoped records + } finally {await rm(root, {recursive: true, force: true});} +}); + +test("a restarted process observes the durable start before admitting another", async () => { + const root = await mkdtemp(join(tmpdir(), "cadence-restart-")); + try { + await manage({runtime_root: root, goal_id: "fixture", agent_id: "a", automation_id: null, + operation: "configure", expected_revision: 0, min_interval_minutes: 60, + owner_reference: "owner-request", execute: true}); + await admit({runtime_root: root, goal_id: "fixture", agent_id: "a", automation_id: null, + request_id: "first", now_ms: 1000, trigger_at_ms: 1000}); + const moduleUrl = new URL("../../loopx/control_plane/quota/automation_cadence.ts", import.meta.url).href; + const inChild = (request_id: string, now_ms: number) => JSON.parse(execFileSync(process.execPath, + ["--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", + `import {admitAutomationStart} from ${JSON.stringify(moduleUrl)}; console.log(JSON.stringify(await admitAutomationStart(${JSON.stringify({runtime_root: root, goal_id: "fixture", agent_id: "a", automation_id: null, request_id, now_ms, trigger_at_ms: now_ms})})));`], + {encoding: "utf8"})); + assert.equal(inChild("early", 1000 + 3_600_000 - 1).admitted, false); + assert.equal(inChild("due", 1000 + 3_600_000).admitted, true); + } finally {await rm(root, {recursive: true, force: true});} +}); + +test("an M1 policy file upgrades in place without losing its configured floor", async () => { + const root = await mkdtemp(join(tmpdir(), "cadence-migration-")); + try { + const path = cadenceStorePath(root, "fixture"); + await mkdir(dirname(path), {recursive: true}); + await writeFile(path, JSON.stringify({schema_version: "automation_cadence_store_v1", + goal_id: "fixture", revision: 1, rules: [{agent_id: "a", automation_id: null, + min_interval_minutes: 60, revision: 1, owner_reference: "owner-request"}]})); + const input = {runtime_root: root, goal_id: "fixture", agent_id: "a", automation_id: null}; + assert.equal((await manage({...input, operation: "read"})).min_interval_minutes, 60); + assert.equal((await admit({...input, request_id: "first", now_ms: 1000, trigger_at_ms: 1000})).reserved, true); + const saved = JSON.parse(await readFile(path, "utf8")); + assert.equal(saved.schema_version, "automation_cadence_store_v2"); + assert.equal(saved.rules[0].min_interval_minutes, 60); + assert.equal(saved.starts.length, 1); + } finally {await rm(root, {recursive: true, force: true});} +}); + test("every backoff/reset interval respects the owner floor without a 60-minute ceiling", () => { for (const sequence of [[3, 6, 12], [15, 30, 60], [120, 240, 480], [1440, 2880]]) { const result = progression({progression: sequence, min_interval_minutes: 1440}); diff --git a/tests/test_loopx_turn_driver.py b/tests/test_loopx_turn_driver.py index 17ef6c1618..0bcb53e19e 100644 --- a/tests/test_loopx_turn_driver.py +++ b/tests/test_loopx_turn_driver.py @@ -13,6 +13,7 @@ import pytest import loopx.cli_commands.turn as turn_command +from loopx.cli_commands.turn_rendering import render_loopx_turn_execution_markdown from tests.control_plane.canonical_authority_fixture import ( initialize_canonical_authority, ) @@ -1782,6 +1783,15 @@ def test_turn_run_once_cli_commits_validated_result_and_one_quota_slot( tmp_path: Path, ) -> None: project, runtime, registry = _write_live_fixture(tmp_path) + policy_output = io.StringIO() + with contextlib.redirect_stdout(policy_output): + policy_code = cli_main([ + "--registry", str(registry), "--runtime-root", str(runtime), "--format", "json", + "automation-cadence", "--goal-id", "loopx-turn-fixture", "--agent-id", "codex-fixture", + "--min-interval-minutes", "1440", "--expected-revision", "0", + "--owner-reference", "fixture-owner", "--execute", + ]) + assert policy_code == 0, policy_output.getvalue() host_project = tmp_path / "isolated-host-workspace" host_project.mkdir() host_script = """ @@ -1847,6 +1857,8 @@ def test_turn_run_once_cli_commits_validated_result_and_one_quota_slot( payload = json.loads(output.getvalue()) assert exit_code == 0, payload assert payload["status"] == "committed" + assert payload["admission"]["reserved"] is True + assert payload["admission"]["pre_model_admission"] == "managed_turn_only" assert payload["receipt"]["status"] == "committed" assert payload["receipt"]["next_phase"] is None assert payload["validation"]["status"] == "passed" @@ -1935,6 +1947,28 @@ def test_turn_run_once_cli_commits_validated_result_and_one_quota_slot( "quota_slot_spent", ] + next_output = io.StringIO() + with contextlib.redirect_stdout(next_output): + next_code = cli_main([ + "--registry", str(registry), "--runtime-root", str(runtime), "--format", "json", + "turn", "run-once", "--host", "generic-cli", "--goal-id", "loopx-turn-fixture", + "--agent-id", "codex-fixture", "--turn-instance-id", "next-cadence-fixture", + "--project", str(host_project), "--host-command-json", + json.dumps([sys.executable, "-c", host_script]), "--validation-command-json", + json.dumps([sys.executable, "-c", validation_script]), "--scan-root", str(project), + "--no-global-sync", "--execute", + ]) + waiting = json.loads(next_output.getvalue()) + assert next_code == 1, waiting + assert waiting["status"] == "interval_wait" + assert waiting["admission"]["next_eligible_at_ms"] > 0 + assert waiting["effects"]["host_invoked"] is False + assert waiting["effects"]["quota_spent"] is False + assert "- next_eligible_at: " in render_loopx_turn_execution_markdown(waiting) + assert [json.loads(line)["classification"] for line in index_path.read_text(encoding="utf-8").splitlines()] == [ + "fixture_progress", "quota_slot_spent", + ] + def test_turn_run_once_cli_dsh_fresh_sessions_and_iteration_failure( tmp_path: Path, diff --git a/tests/test_loopx_turn_executor.py b/tests/test_loopx_turn_executor.py index 958f4f5536..094e2b44f4 100644 --- a/tests/test_loopx_turn_executor.py +++ b/tests/test_loopx_turn_executor.py @@ -904,6 +904,47 @@ def test_run_once_preview_has_no_host_or_journal_effects(tmp_path: Path) -> None assert not (tmp_path / "runtime").exists() +def test_managed_start_waits_before_host_and_replay_needs_no_new_admission( + tmp_path: Path, +) -> None: + plan = _plan() + calls = {"host": 0, "admit": 0, "writeback": 0, "spend": 0, "scheduler": 0} + writeback, spend, scheduler = _callbacks(calls) + + def host(_request: object) -> dict[str, object]: + calls["host"] += 1 + return _host_result(plan) + + def wait(_identity: object) -> dict[str, object]: + calls["admit"] += 1 + return {"admitted": False, "reason": "minimum_interval_wait", "next_eligible_at_ms": 9000} + + common = dict( + host_runner=host, project=tmp_path, runtime_root=tmp_path / "runtime", + goal_id="fixture-goal", timeout_seconds=5, execute=True, + task_validator=_passing_validator, + writeback=writeback, spend=spend, scheduler=scheduler, + ) + denied = run_loopx_turn_once(plan, admit_start=wait, **common) + assert denied["status"] == "interval_wait" + assert denied["admission"]["next_eligible_at_ms"] == 9000 + assert denied["effects"]["host_invoked"] is False + assert calls == {"host": 0, "admit": 1, "writeback": 0, "spend": 0, "scheduler": 0} + assert not list((tmp_path / "runtime" / "goals" / "fixture-goal" / "turns").glob("*.json")) + + def allow(_identity: object) -> dict[str, object]: + calls["admit"] += 1 + return {"admitted": True, "reserved": True, "reason": "admitted"} + + committed = run_loopx_turn_once(plan, admit_start=allow, **common) + assert committed["status"] == "committed" + assert calls["host"] == 1 + replay = run_loopx_turn_once(plan, admit_start=wait, **common) + assert replay["replayed"] is True + assert calls["admit"] == 2 + assert calls["host"] == 1 + + def test_run_once_rejects_oversized_built_in_host_result(tmp_path: Path) -> None: plan = _plan() calls = {"writeback": 0, "spend": 0, "scheduler": 0} @@ -933,7 +974,7 @@ def test_run_once_explicitly_retries_failed_host_without_duplicate_effects( tmp_path: Path, ) -> None: plan = _plan() - calls = {"host": 0, "writeback": 0, "spend": 0, "scheduler": 0} + calls = {"host": 0, "admit": 0, "writeback": 0, "spend": 0, "scheduler": 0} writeback, spend, scheduler = _callbacks(calls) def host(_request: dict[str, object]) -> dict[str, object]: @@ -942,6 +983,11 @@ def host(_request: dict[str, object]) -> dict[str, object]: raise BuiltInHostError("codex_cli_model_requires_newer_codex") return _host_result(plan) + def admit(identity: dict[str, object]) -> dict[str, object]: + calls["admit"] += 1 + assert identity["attempt"] == (1 if calls["admit"] == 1 else 2) + return {"admitted": calls["admit"] != 2, "reason": "minimum_interval_wait"} + kwargs = { "host_runner": host, "project": tmp_path, @@ -953,9 +999,11 @@ def host(_request: dict[str, object]) -> dict[str, object]: "writeback": writeback, "spend": spend, "scheduler": scheduler, + "admit_start": admit, } failed = run_loopx_turn_once(plan, **kwargs) replayed = run_loopx_turn_once(plan, **kwargs) + waiting = run_loopx_turn_once(plan, retry_failed=True, **kwargs) recovered = run_loopx_turn_once(plan, retry_failed=True, **kwargs) assert failed["reason"] == "codex_cli_model_requires_newer_codex" @@ -963,8 +1011,38 @@ def host(_request: dict[str, object]) -> dict[str, object]: assert failed["receipt"]["result_kind"] == "host_failure" assert failed["receipt"]["failed_phase"] == "host_execute" assert replayed["replayed"] is True + assert waiting["status"] == "interval_wait" + assert waiting["effects"]["host_invoked"] is False assert recovered["status"] == "committed" - assert calls == {"host": 2, "writeback": 1, "spend": 1, "scheduler": 1} + assert calls == {"host": 2, "admit": 3, "writeback": 1, "spend": 1, "scheduler": 1} + + +def test_invalid_host_result_reinvocation_reenters_admission(tmp_path: Path) -> None: + plan = _plan() + calls = {"host": 0, "admit": 0, "writeback": 0, "spend": 0, "scheduler": 0} + writeback, spend, scheduler = _callbacks(calls) + + def host(_request: object) -> dict[str, object]: + calls["host"] += 1 + return {"invalid": True} if calls["host"] == 1 else _host_result(plan) + + def admit(identity: dict[str, object]) -> dict[str, object]: + calls["admit"] += 1 + assert identity["attempt"] == (1 if calls["admit"] == 1 else 2) + return {"admitted": calls["admit"] != 2, "reason": "minimum_interval_wait"} + + common = dict(host_runner=host, admit_start=admit, project=tmp_path, + runtime_root=tmp_path / "runtime", goal_id="fixture-goal", timeout_seconds=5, + execute=True, task_validator=_passing_validator, writeback=writeback, + spend=spend, scheduler=scheduler) + first = run_loopx_turn_once(plan, **common) + waiting = run_loopx_turn_once(plan, retry_failed=True, **common) + recovered = run_loopx_turn_once(plan, retry_failed=True, **common) + + assert first["result_kind"] == "validation_failed" + assert waiting["status"] == "interval_wait" and waiting["effects"]["host_invoked"] is False + assert recovered["status"] == "committed" + assert calls == {"host": 2, "admit": 3, "writeback": 1, "spend": 1, "scheduler": 1} def test_run_once_bounds_provider_capacity_retries_without_spending_quota( From 998c79535711361594a8f141e4809d3111b8f63b Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Wed, 23 Sep 2026 14:54:06 +0800 Subject: [PATCH 3/5] docs: record managed admission boundary in both RFCs Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../rfcs/automatic-execution-admission-v0.md | 32 +++++++++++++------ .../automatic-execution-admission-v0.zh-CN.md | 23 +++++++++---- 2 files changed, 40 insertions(+), 15 deletions(-) diff --git a/docs/architecture/rfcs/automatic-execution-admission-v0.md b/docs/architecture/rfcs/automatic-execution-admission-v0.md index b9d1e2bac2..c809163026 100644 --- a/docs/architecture/rfcs/automatic-execution-admission-v0.md +++ b/docs/architecture/rfcs/automatic-execution-admission-v0.md @@ -4,7 +4,7 @@ - **Delivery maturity:** Partial, proposed implementation; not promoted - **Owners:** Quota, scheduler and host-runtime maintainers - **Created / last normative revision:** 2026-09-23 -- **Implementation baseline:** `23edcb19c` +- **Implementation baseline:** `79241d7ef` - **Language mirror:** [中文版](automatic-execution-admission-v0.zh-CN.md) - **Related contracts:** [roadmap](loopx-overall-roadmap-v0.md), [quota](../../quota-allocation.md), [cadence hint](../../operations/long-task-cadence-policy.md), [session execution modes](agent-session-execution-modes-v0.md) @@ -19,9 +19,11 @@ cannot rewrite it. A quota slot, a timer tick and a model invocation remain different events. The target design requires every controlled new host invocation to pass temporal admission as well as the existing budget, permission, binding and work gates. -No configured interval preserves existing behavior. The first milestone changes Codex App schedule recommendations, including reset -and backoff. It does not yet enforce launch admission. M2 will require every new -host invocation, retry and continuation to pass admission; cached settlement will not. +No configured interval preserves existing behavior. M1 changes Codex App schedule +recommendations, including reset and backoff. M2 adds pre-host admission to managed +`turn run-once`; App timer and other launchers remain separate qualification work. +Every new managed host invocation and failed-result retry passes admission; cached +settlement does not. This RFC does not authorize changes to existing automations, model selection, quota allocation, remote services or public publishing. @@ -108,7 +110,7 @@ protection against clock manipulation. | Host path | Required contract | Initial boundary | | --- | --- | --- | -| Managed `turn run-once` | Atomic admission immediately before a new host attempt | Planned for M2, including failed-result recovery | +| Managed `turn run-once` | Atomic admission before a new host attempt, including failed-result recovery | M2 candidate; isolated CLI and concurrency tests, no external host promotion | | Local legacy scheduler / external launchers | Route launches through admitted Turn or implement the same owner call | Not yet qualified; do not advertise enforcement | | Codex App automation | Apply floor-compatible timer, read actual schedule, ACK only matching facts | M1 schedule recommendation floor; hook coverage not qualified | | Attached interactive/manual session | Explicit manual intent; existing authority gates remain | Caller records reason; automatic continuation cannot masquerade as manual | @@ -209,10 +211,22 @@ product journey. The broader product goal remains open until M2/M3 acceptance. ## Appendix: implementation ledger -Baseline audit at `23edcb19c`: no durable owner minimum interval. M1 is an -unmerged candidate; M2/M3/M4 and live App qualification remain unverified. -Tests and PR validation must distinguish deterministic evidence from host -promotion. No existing automation is activated or rebound by this proposal. +Baseline audit at `23edcb19c`: no durable owner minimum interval. M1 merged in +PR #4921. The M2 candidate reserves a managed Turn start in the same quota +policy file and lock before host invocation. Denial returns the next eligible +time without a host call, writeback or quota spend; a failed host consumes its +start, and settlement replay skips admission. Goal floors apply per agent; +automation floors require an explicit stable `--automation-id`. A manual start +requires `--manual-interval-bypass-reason`, records a start, and bypasses only +the interval. The local CLI remains a same-UID trust boundary. + +The App timer-to-hook path, non-Turn launchers, packaged settings UI and live +model-host promotion remain unqualified. No existing automation is activated +or rebound by this proposal. +M1 policy files are read as v1 and upgraded in place to v2 on the first +configuration write or admitted start. The path stays stable; older binaries +reject the v2 schema rather than silently discarding start records. Pause the +launcher before downgrade. ### Hook research — 2026-09-23 diff --git a/docs/architecture/rfcs/automatic-execution-admission-v0.zh-CN.md b/docs/architecture/rfcs/automatic-execution-admission-v0.zh-CN.md index ca8dd60f64..76726cfa69 100644 --- a/docs/architecture/rfcs/automatic-execution-admission-v0.zh-CN.md +++ b/docs/architecture/rfcs/automatic-execution-admission-v0.zh-CN.md @@ -4,7 +4,7 @@ - **交付成熟度:** Partial,候选实现,尚未推广 - **维护边界:** quota、scheduler、host runtime - **创建 / 规范修订:** 2026-09-23 -- **实现基线:** `23edcb19c` +- **实现基线:** `79241d7ef` - **语言镜像:** [English](automatic-execution-admission-v0.md) - **相关契约:** [路线图](loopx-overall-roadmap-v0.zh-CN.md)、[quota](../../quota-allocation.md)、[节奏提示](../../operations/long-task-cadence-policy.md)、[执行模式](agent-session-execution-modes-v0.md) @@ -16,8 +16,10 @@ scheduler 退避消费这一约束,不能改写它。配额槽、定时器唤醒和模型调用是三个事件。 目标设计要求受控执行器启动每次新的 host 调用前,同时满足时间准入、预算、权限、绑定和工作门禁。 -未配置时保持原行为。M1 优先修改 Codex App 的调度建议、重置与退避,不宣称已经强制执行启动准入。 -M2 再要求每次新 host 调用、重试与续跑重新准入;缓存结果和结算不重新消耗准入。本 RFC 不授权修改现有自动化、模型选择、配额分配、远程服务或公开发布。 +未配置时保持原行为。M1 修改 Codex App 的调度建议、重置与退避。M2 为 managed +`turn run-once` 增加 host 启动前准入;App 定时器和其他 launcher 仍需单独验收。 +每次新的 managed host 调用及失败结果重试重新准入;缓存结果和结算不重复消耗准入。 +本 RFC 不授权修改现有自动化、模型选择、配额分配、远程服务或公开发布。 ## 2. 问题与不变量 @@ -76,7 +78,7 @@ App、Turn、前端、Lark 不得另存一套策略。共享 authority provider | 入口 | 必须履行的契约 | 本阶段边界 | | --- | --- | --- | -| Managed `turn run-once` | 每次新 host 尝试前原子准入 | M2 计划,含错误结果后的恢复 | +| Managed `turn run-once` | 每次新 host 尝试及失败重试前原子准入 | M2 候选;已做隔离 CLI 与并发测试,未推广外部宿主 | | 旧 local scheduler / 外部 launcher | 经受控 Turn 启动,或调用相同准入 owner | 尚未验收,不宣传为已强制执行 | | Codex App automation | 应用满足下限的定时器,回读真实值,事实匹配后 ACK | M1 调度建议下限;hook 覆盖范围未验收 | | 附着式交互 / 手动会话 | 显式手动意图;其余门禁保留 | 记录原因;自动续跑不能冒充手动 | @@ -155,8 +157,17 @@ M1 对 App 调度管理有独立价值,但不代表多宿主产品旅程完成 ## 附录:实现记录 -`23edcb19c` 基线没有持久用户最短间隔。M1 是未合并候选;M2/M3/M4 与真实 App 验收尚未完成。 -测试和 PR 必须区分确定性验证与宿主推广。本提案不激活、不改绑任何已有自动化。 +`23edcb19c` 基线没有持久用户最短间隔。M1 已由 PR #4921 合并。M2 候选在 +同一 quota 策略文件与锁中预留 managed Turn 的启动位,再调用 host。拒绝时返回下次 +可执行时间,不调用 host、不写回、不花 quota;失败的 host 消耗已预留间隔,结算回放 +跳过准入。Goal 下限按 agent 生效;automation 下限要求显式稳定的 `--automation-id`。 +手动启动要求 `--manual-interval-bypass-reason`,记录这次启动,只绕过时间下限。 +本地 CLI 仍以相同 OS 用户为信任边界。 + +App 定时器到 hook、非 Turn launcher、打包设置界面及真实模型宿主推广仍未验收。 +本提案不激活、不改绑任何已有自动化。 +M1 的 v1 策略文件在首次配置写入或获准启动时原地升级为 v2,文件路径保持不变。 +旧版程序会拒绝 v2 schema,避免静默丢弃启动记录;降级前必须暂停 launcher。 ### Hook 调研 — 2026-09-23 From e77aaed0e9281a82b4ca86a0e2f4008939881786 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 12:44:26 +0800 Subject: [PATCH 4/5] fix(authority): make managed cadence starts recoverable across a crash A committed reservation could precede the first durable Turn journal attempt, so the same Turn asked again with the same request identity and was rejected as a duplicate forever, even past the owner floor and with a manual reason. Managed starts are now two-phase in the same cadence store: admission writes a reserved start, and the Turn executor confirms it only after the host attempt is durable in the journal. The same identity resumes a reserved start once the floor is reached, while a confirmed start stays fail-closed. A record without the phase field is read as an attempted start, so older or hand-edited stores fail closed. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/turn.py | 129 ++++++++++--- .../control_plane/effect_runtime_handlers.ts | 3 +- .../control_plane/quota/automation_cadence.ts | 62 +++++- loopx/control_plane/turn_driver/executor.py | 11 ++ .../automation_cadence.test.ts | 43 ++++- tests/test_loopx_turn_executor.py | 179 ++++++++++++++++++ 6 files changed, 388 insertions(+), 39 deletions(-) diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index 5dbfc65143..8368798a4c 100644 --- a/loopx/cli_commands/turn.py +++ b/loopx/cli_commands/turn.py @@ -7,7 +7,7 @@ import time from collections.abc import Callable, Mapping from pathlib import Path -from typing import Any +from typing import Any, NamedTuple from ..cli_rollout import append_cli_rollout_event from ..capabilities.explore.composition_frontier import ( @@ -91,6 +91,86 @@ FormatSelector = Callable[..., str] +class ManagedCadenceStart(NamedTuple): + """Owner-cadence callbacks for one managed Turn start. + + `admit` reserves (or resumes) the interval slot before any host or journal + attempt; `confirm` marks that reservation as a real host attempt once the + Turn journal is durable. Keeping them separate means a crash in between + leaves a resumable reservation rather than a permanently rejected Turn. + """ + + admit: Callable[[Mapping[str, Any]], dict[str, Any]] + confirm: Callable[[], None] + + +def managed_cadence_start( + *, + runtime_root: Path, + goal_id: str, + agent_id: str | None, + automation_id: str | None, + manual_reason: str | None, + on_admitted: Callable[[], None] | None = None, +) -> ManagedCadenceStart: + """Bind one managed Turn start to the TypeScript owner-cadence store.""" + + admitted_request: dict[str, Any] = {} + + def admit(identity: Mapping[str, Any]) -> dict[str, Any]: + now_ms = time.time_ns() // 1_000_000 + request_id = f"{identity['turn_key']}:{identity['attempt']}" + admission = effect_runtime_result( + "quota.automation_cadence.admit", + { + "runtime_root": str(runtime_root), + "goal_id": goal_id, + "agent_id": agent_id, + "automation_id": automation_id, + "request_id": request_id, + "trigger_at_ms": now_ms, + "now_ms": now_ms, + "manual_reason": manual_reason, + }, + retry_safe=False, + ) + admitted_request.clear() + if admission.get("admitted") is True: + admitted_request["request_id"] = request_id + admitted_request["reserved"] = admission.get("reserved") is True + if on_admitted is not None: + on_admitted() + return { + key: admission.get(key) + for key in ( + "admitted", "reserved", "resumed", "reason", "next_eligible_at_ms", + "min_interval_minutes", "pre_model_admission", + ) + } + + def confirm() -> None: + request_id = admitted_request.get("request_id") + if admitted_request.get("reserved") is not True or not isinstance(request_id, str): + return + confirmation = effect_runtime_result( + "quota.automation_cadence.confirm_start", + { + "runtime_root": str(runtime_root), + "goal_id": goal_id, + "agent_id": agent_id, + "automation_id": automation_id, + "request_id": request_id, + }, + retry_safe=True, + ) + if confirmation.get("confirmed") is not True: + raise ValueError( + "managed Turn start could not be confirmed against the owner cadence store" + ) + + return ManagedCadenceStart(admit=admit, confirm=confirm) + + def handle_turn_command( @@ -1043,36 +1123,22 @@ def post_settlement_reward_memory( settlement_evidence=settlement_evidence, ) - def admit_managed_start(identity: Mapping[str, Any]) -> dict[str, Any]: - now_ms = time.time_ns() // 1_000_000 - admitted = effect_runtime_result( - "quota.automation_cadence.admit", - { - "runtime_root": str(runtime_root), - "goal_id": args.goal_id, - "agent_id": args.agent_id, - "automation_id": args.automation_id, - "request_id": f"{identity['turn_key']}:{identity['attempt']}", - "trigger_at_ms": now_ms, - "now_ms": now_ms, - "manual_reason": args.manual_interval_bypass_reason, - }, - retry_safe=False, + def on_managed_start_admitted() -> None: + ensure_turn_heartbeat_settlement_receipt( + runtime_root, + settlement_identity, + semantic_replan_guard_scoped=replan_guard_scoped, + semantic_replan_obligation_id=replan_obligation_id, ) - if admitted.get("admitted") is True: - ensure_turn_heartbeat_settlement_receipt( - runtime_root, - settlement_identity, - semantic_replan_guard_scoped=replan_guard_scoped, - semantic_replan_obligation_id=replan_obligation_id, - ) - return { - key: admitted.get(key) - for key in ( - "admitted", "reserved", "reason", "next_eligible_at_ms", - "min_interval_minutes", "pre_model_admission", - ) - } + + managed_cadence = managed_cadence_start( + runtime_root=runtime_root, + goal_id=args.goal_id, + agent_id=args.agent_id, + automation_id=args.automation_id, + manual_reason=args.manual_interval_bypass_reason, + on_admitted=on_managed_start_admitted, + ) payload = run_loopx_turn_once( payload, @@ -1102,7 +1168,8 @@ def admit_managed_start(identity: Mapping[str, Any]) -> dict[str, Any]: if args.execute and args.host == "codex-cli" else None ), - admit_start=admit_managed_start if args.execute else None, + admit_start=managed_cadence.admit if args.execute else None, + confirm_start=managed_cadence.confirm if args.execute else None, ) else: raise ValueError("turn requires the `plan` or `run-once` subcommand") diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 1faea5a3ad..53d920e0d9 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -1,4 +1,4 @@ -import {admitAutomationStart, manageAutomationCadence, projectCadenceSchedule} from "./quota/automation_cadence.ts"; +import {admitAutomationStart, confirmAutomationStart, manageAutomationCadence, projectCadenceSchedule} from "./quota/automation_cadence.ts"; import {manageLocalAuthorityArchive} from "./coordination/local_authority_archive.ts"; import {selectPeriodicReportProgress, selectPeriodicReportApprovalRetry} from "./capabilities/periodic_report_progress.ts"; import {planIssueFixMonitorReconciliation} from "./capabilities/issue_fix_monitor_reconciliation.ts"; @@ -480,6 +480,7 @@ export function createEffectRuntimeHandlers( ["scheduler.state_transition.evaluate", evaluateSchedulerStateTransition], ["quota.automation_cadence.manage", manageAutomationCadence], ["quota.automation_cadence.admit", admitAutomationStart], + ["quota.automation_cadence.confirm_start", confirmAutomationStart], ["quota.automation_cadence.schedule", projectCadenceSchedule], ["scheduler.state.evaluate", evaluateSchedulerStateOperation], ["scheduler.state.load", loadSchedulerState], diff --git a/loopx/control_plane/quota/automation_cadence.ts b/loopx/control_plane/quota/automation_cadence.ts index 79d6c848e2..096d5b0cfc 100644 --- a/loopx/control_plane/quota/automation_cadence.ts +++ b/loopx/control_plane/quota/automation_cadence.ts @@ -12,7 +12,13 @@ const SCHEMA = "automation_cadence_store_v2"; const RESULT = "automation_cadence_result_v1"; type Scope = {agent_id: string | null; automation_id: string | null}; type Rule = Scope & {min_interval_minutes: number; revision: number; owner_reference: string}; -type Start = Scope & {started_at_ms: number; trigger_at_ms: number; request_id: string; manual_reason: string | null}; +/** `reserved` is a start whose host may still be unstarted; `started` is a start + * whose host attempt is already durable in the Turn journal. Only the former can + * be resumed by the same request identity, and a record without this field is + * read as `started` so an older or hand-edited store fails closed. */ +type StartState = "reserved" | "started"; +type Start = Scope & {started_at_ms: number; trigger_at_ms: number; request_id: string; + manual_reason: string | null; state: StartState}; type Store = {schema_version: typeof SCHEMA; goal_id: string; revision: number; rules: Rule[]; starts: Start[]}; const fail = (message: string): never => {throw new EffectRuntimeRequestError(message, "automation_cadence_invalid");}; function text(value: unknown, name: string): string { @@ -62,7 +68,9 @@ function decode(value: unknown, goal: string): Store { const trigger_at_ms = integer(r.trigger_at_ms, "trigger_at_ms"); if (trigger_at_ms > started_at_ms) fail("cadence start trigger follows its start"); return {...s, started_at_ms, trigger_at_ms, request_id: text(r.request_id, "request_id"), - manual_reason: r.manual_reason == null ? null : text(r.manual_reason, "manual_reason")}; + manual_reason: r.manual_reason == null ? null : text(r.manual_reason, "manual_reason"), + state: r.state === undefined ? "started" + : requireStringLiteral(r.state, ["reserved", "started"] as const, "start state")}; }); return {schema_version: SCHEMA, goal_id: goal, revision, rules, starts}; } @@ -99,7 +107,13 @@ function projection(store: Store, s: Scope, nowMs = Date.now()): JsonObject { }; } -/** Reserve a new managed-host start under the same lock as policy changes. */ +/** Reserve a managed-host start under the same lock as policy changes. + * + * A reserved start whose host attempt never became durable may be resumed by the + * same request identity, so a crash between reservation and the first journal + * attempt cannot strand the Turn. A start whose host attempt is already durable + * stays fail-closed, and the caller must confirm the reservation once the + * attempt is recorded. */ export async function admitAutomationStart(p: JsonObject): Promise { const goal = text(p.goal_id, "goal_id"), s = scope(p); if (!s.agent_id) fail("admission requires agent_id"); @@ -119,9 +133,17 @@ export async function admitAutomationStart(p: JsonObject): Promise { .map(rule => rule.automation_id)); const prior = store.starts.filter(start => start.agent_id === s.agent_id && activeScopes.has(start.automation_id)); - if (prior.some(start => start.request_id === request || trigger < start.trigger_at_ms)) { + const held = prior.find(start => start.request_id === request); + if (prior.some(start => trigger < start.trigger_at_ms) || (held !== undefined && held.state === "started")) { return {...current, admitted: false, reserved: false, reason: "duplicate_or_stale_trigger"}; } + if (held !== undefined) { + // Same identity, no durable host attempt: keep the original interval anchor. + return manual === null && current.eligible_now !== true + ? {...current, admitted: false, reserved: false, reason: "minimum_interval_wait"} + : {...current, admitted: true, reserved: true, resumed: true, + reason: "resumed_unstarted_reservation"}; + } if (manual === null && current.eligible_now !== true) { return {...current, admitted: false, reserved: false, reason: "minimum_interval_wait"}; } @@ -130,13 +152,43 @@ export async function admitAutomationStart(p: JsonObject): Promise { for (const item of scopes) { store.starts = store.starts.filter(start => key(start) !== key(item)); store.starts.push({...item, started_at_ms: now, trigger_at_ms: trigger, request_id: request, - manual_reason: manual}); + manual_reason: manual, state: "reserved"}); } await atomicWriteJson(path, store); return {...projection(store, s, now), admitted: true, reserved: true, reason: manual === null ? "admitted" : "explicit_manual_interval_bypass"}; }); } + +/** Mark a reservation as an attempted host start once the Turn journal is durable. + * + * Confirmation is idempotent and never creates a store, so the default-off path + * stays free of new files. Until it lands, the record stays resumable; after it + * lands, the same request identity is rejected fail-closed. */ +export async function confirmAutomationStart(p: JsonObject): Promise { + const goal = text(p.goal_id, "goal_id"), s = scope(p); + if (!s.agent_id) fail("confirmation requires agent_id"); + const path = cadenceStorePath(requireNonEmptyString(p.runtime_root, "runtime_root"), goal); + const request = text(p.request_id, "request_id"); + const record = async (store: Store): Promise => { + const matches = store.starts.filter(start => start.agent_id === s.agent_id && start.request_id === request); + const pending = matches.some(start => start.state === "reserved"); + if (pending) { + store.starts = store.starts.map(start => start.agent_id === s.agent_id && start.request_id === request + ? {...start, state: "started"} : start); + await atomicWriteJson(path, store); + } + return {schema_version: RESULT, ok: true, goal_id: goal, ...s, request_id: request, confirmed: true, + reason: pending ? "start_confirmed" : "already_confirmed"}; + }; + // An absent or empty reservation must not create a store or lock file. + const snapshot = await load(path, goal); + if (!snapshot.starts.some(start => start.agent_id === s.agent_id && start.request_id === request)) { + return {schema_version: RESULT, ok: true, goal_id: goal, ...s, request_id: request, confirmed: false, + reason: "reservation_missing"}; + } + return withFileMutationLock(path, async () => record(await load(path, goal))); +} /** Pure typed calculation shared by scheduler adapters; no policy mutation. */ export function projectCadenceProgression(p: JsonObject): JsonObject { const floor = minutes(p.min_interval_minutes); diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index 54e20a23d4..61ada9ba7a 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -754,6 +754,7 @@ def _host_result_stage( journal: dict[str, Any], journal_path: Path, effects: dict[str, bool], + confirm_start: Callable[[], None] | None = None, ) -> tuple[dict[str, Any] | None, list[str], dict[str, Any] | None]: completed_phases = list(journal.get("completed_phases") or []) result = ( @@ -764,6 +765,10 @@ def _host_result_stage( if "typed_result" not in completed_phases: journal["host_attempt_count"] = int(journal.get("host_attempt_count") or 0) + 1 _write_journal(journal_path, journal) + # The attempt is durable now, so a later restart must not resume this + # reservation. Confirmation failure stops before the host starts. + if confirm_start is not None: + confirm_start() host_observation = ( _run_host_runner(request, runner=host_runner) if host_runner is not None @@ -1224,6 +1229,7 @@ def run_loopx_turn_once( scheduler: Scheduler | None = None, post_settlement: PostSettlement | None = None, admit_start: Callable[[Mapping[str, Any]], dict[str, Any]] | None = None, + confirm_start: Callable[[], None] | None = None, ) -> dict[str, Any]: if host_runner is not None and host_argv is not None: raise ValueError("run-once accepts either host_argv or host_runner, not both") @@ -1426,6 +1432,11 @@ def finish_recovery(payload: dict[str, Any]) -> dict[str, Any]: journal=journal, journal_path=journal_path, effects=effects, + confirm_start=( + confirm_start + if admission is not None and admission.get("reserved") is True + else None + ), ) if terminal is not None: return finish_recovery(terminal) diff --git a/tests/control_plane_ts/automation_cadence.test.ts b/tests/control_plane_ts/automation_cadence.test.ts index 4236ce2a14..9b9422f8f7 100644 --- a/tests/control_plane_ts/automation_cadence.test.ts +++ b/tests/control_plane_ts/automation_cadence.test.ts @@ -4,7 +4,7 @@ import {mkdtemp, mkdir, rm, readFile, writeFile, readdir} from "node:fs/promises import {tmpdir} from "node:os"; import {dirname, join} from "node:path"; import {execFileSync} from "node:child_process"; -import {admitAutomationStart as admit, cadenceStorePath, manageAutomationCadence as manage, projectCadenceProgression as progression} from "../../loopx/control_plane/quota/automation_cadence.ts"; +import {admitAutomationStart as admit, cadenceStorePath, confirmAutomationStart as confirm, manageAutomationCadence as manage, projectCadenceProgression as progression} from "../../loopx/control_plane/quota/automation_cadence.ts"; test("owner floor inherits without changing another agent; reductions and concurrent writes require authority", async () => { const root = await mkdtemp(join(tmpdir(), "cadence-")); @@ -56,7 +56,10 @@ test("managed starts are atomic, durable, per agent, and due at the exact interv assert.equal(concurrent.filter(r => r.admitted === false).length, 1); assert.equal((await start("early", 1000 + 86_400_000 - 1)).reason, "minimum_interval_wait"); assert.equal((await start("due", 1000 + 86_400_000)).admitted, true); - assert.equal((await start("due", 1000 + 2 * 86_400_000)).reason, "duplicate_or_stale_trigger"); + // An unconfirmed reservation resumes the same identity; a confirmed start is fail-closed. + assert.equal((await start("due", 1000 + 2 * 86_400_000)).reason, "resumed_unstarted_reservation"); + await confirm({...base, request_id: "due"}); + assert.equal((await start("due", 1000 + 3 * 86_400_000)).reason, "duplicate_or_stale_trigger"); assert.equal((await start("stale", 1000 + 86_400_000, {trigger_at_ms: 999})).admitted, false); assert.equal((await admit({...base, agent_id: "b", request_id: "other-agent", now_ms: 1001, trigger_at_ms: 1001})).admitted, true); const manual = await start("manual", 1000 + 86_400_001, {manual_reason: "owner-request"}); @@ -109,6 +112,41 @@ test("a restarted process observes the durable start before admitting another", } finally {await rm(root, {recursive: true, force: true});} }); +test("an unconfirmed reservation resumes while a confirmed start fails closed", async () => { + const root = await mkdtemp(join(tmpdir(), "cadence-resume-")); + const base = {runtime_root: root, goal_id: "fixture", agent_id: "a", automation_id: null}; + const start = (request_id: string, now_ms: number, extra = {}) => admit({...base, request_id, + now_ms, trigger_at_ms: now_ms, ...extra}); + const path = cadenceStorePath(root, "fixture"); + try { + assert.equal((await confirm({...base, request_id: "absent"})).reason, "reservation_missing"); + assert.deepEqual(await readdir(root), []); // neither phase creates files with no policy + await manage({...base, operation: "configure", expected_revision: 0, min_interval_minutes: 60, + owner_reference: "owner-request", execute: true}); + assert.equal((await start("turn:1", 1000)).reserved, true); + assert.equal((await start("turn:1", 1000 + 3_600_000 - 1)).reason, "minimum_interval_wait"); + const resumed = await start("turn:1", 1000 + 3_600_000); + assert.equal(resumed.admitted, true); + assert.equal(resumed.resumed, true); + assert.equal(resumed.reason, "resumed_unstarted_reservation"); + const reserved = JSON.parse(await readFile(path, "utf8")).starts[0]; + assert.equal(reserved.state, "reserved"); + assert.equal(reserved.started_at_ms, 1000); // the owner floor anchor never moves + assert.equal((await confirm({...base, request_id: "turn:1"})).reason, "start_confirmed"); + assert.equal((await confirm({...base, request_id: "turn:1"})).reason, "already_confirmed"); + assert.equal(JSON.parse(await readFile(path, "utf8")).starts[0].state, "started"); + // A confirmed host attempt is fail-closed for the same identity, bypass or not. + assert.equal((await start("turn:1", 1000 + 2 * 3_600_000)).reason, "duplicate_or_stale_trigger"); + assert.equal((await start("turn:1", 1000 + 2 * 3_600_000, + {manual_reason: "owner-request"})).reason, "duplicate_or_stale_trigger"); + // A record written without the phase field is read as an attempted start. + const store = JSON.parse(await readFile(path, "utf8")); + delete store.starts[0].state; + await writeFile(path, JSON.stringify(store)); + assert.equal((await start("turn:1", 1000 + 3 * 3_600_000)).reason, "duplicate_or_stale_trigger"); + } finally {await rm(root, {recursive: true, force: true});} +}); + test("an M1 policy file upgrades in place without losing its configured floor", async () => { const root = await mkdtemp(join(tmpdir(), "cadence-migration-")); try { @@ -124,6 +162,7 @@ test("an M1 policy file upgrades in place without losing its configured floor", assert.equal(saved.schema_version, "automation_cadence_store_v2"); assert.equal(saved.rules[0].min_interval_minutes, 60); assert.equal(saved.starts.length, 1); + assert.equal(saved.starts[0].state, "reserved"); } finally {await rm(root, {recursive: true, force: true});} }); diff --git a/tests/test_loopx_turn_executor.py b/tests/test_loopx_turn_executor.py index 094e2b44f4..daa273ea5c 100644 --- a/tests/test_loopx_turn_executor.py +++ b/tests/test_loopx_turn_executor.py @@ -7,6 +7,9 @@ import pytest +from loopx.cli_commands import turn as turn_command +from loopx.cli_commands.turn import ManagedCadenceStart, managed_cadence_start +from loopx.control_plane.effect_runtime import effect_runtime_result from loopx.control_plane.turn_driver import executor as turn_executor from loopx.control_plane.turn_driver import ( LOOPX_TURN_RESULT_SCHEMA_VERSION, @@ -945,6 +948,182 @@ def allow(_identity: object) -> dict[str, object]: assert calls["host"] == 1 +def _configure_managed_floor(runtime_root: Path, *, minutes_value: int) -> None: + """Write the owner floor through the shipped TypeScript cadence store.""" + + configured = effect_runtime_result( + "quota.automation_cadence.manage", + { + "runtime_root": str(runtime_root), + "goal_id": "fixture-goal", + "agent_id": None, + "automation_id": None, + "operation": "configure", + "expected_revision": 0, + "min_interval_minutes": minutes_value, + "owner_reference": "fixture-owner", + "execute": True, + }, + ) + assert configured["min_interval_minutes"] == minutes_value, configured + + +def _cadence_starts(runtime_root: Path) -> list[dict[str, object]]: + """Read the real TypeScript cadence store that admits managed starts.""" + + stores = [] + for path in sorted(runtime_root.rglob("*.json")): + if path.name.endswith(".lock.holder.json"): + continue + payload = json.loads(path.read_text(encoding="utf-8")) + if ( + isinstance(payload, dict) + and str(payload.get("schema_version", "")).startswith("automation_cadence_store") + and payload.get("goal_id") == "fixture-goal" + ): + stores.append(payload) + assert len(stores) == 1, stores + return list(stores[0]["starts"]) + + +class _FrozenTurnClock: + """Freeze the managed-start clock so a test can cross the owner floor.""" + + def __init__(self, now_ms: int) -> None: + self._now_ns = now_ms * 1_000_000 + + def time_ns(self) -> int: + return self._now_ns + + +def _managed_cadence(runtime_root: Path) -> ManagedCadenceStart: + return managed_cadence_start( + runtime_root=runtime_root, + goal_id="fixture-goal", + agent_id="codex-fixture", + automation_id=None, + manual_reason=None, + ) + + +def _managed_start_fixture( + tmp_path: Path, +) -> tuple[dict[str, object], Path, dict[str, int], dict[str, object]]: + plan = _plan() + runtime_root = tmp_path / "runtime" + calls = {"host": 0, "writeback": 0, "spend": 0, "scheduler": 0} + writeback, spend, scheduler = _callbacks(calls) + + def host(_request: object) -> dict[str, object]: + calls["host"] += 1 + return _host_result(plan) + + _configure_managed_floor(runtime_root, minutes_value=1) + return ( + plan, + runtime_root, + calls, + { + "host_runner": host, + "project": tmp_path, + "runtime_root": runtime_root, + "goal_id": "fixture-goal", + "timeout_seconds": 5, + "execute": True, + "task_validator": _passing_validator, + "writeback": writeback, + "spend": spend, + "scheduler": scheduler, + }, + ) + + +def test_reserved_managed_start_recovers_after_death_before_the_first_journal_write( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + """A crash between the cadence reservation and the journal must not strand a Turn. + + The reservation is written by the real TypeScript cadence store and the + process is interrupted before any Turn journal write, so the restart asks + with the same request identity. + """ + + plan, runtime_root, calls, common = _managed_start_fixture(tmp_path) + turn_key = str(plan["transaction"]["turn_key"]) # type: ignore[index] + + def die_after_reservation(identity: Mapping[str, object]) -> dict[str, object]: + reservation = _managed_cadence(runtime_root).admit(identity) + assert reservation["admitted"] is True and reservation["reserved"] is True + raise RuntimeError("simulated process death after the cadence reservation") + + with pytest.raises(RuntimeError, match="simulated process death"): + run_loopx_turn_once(plan, admit_start=die_after_reservation, **common) + + assert [(row["state"], row["request_id"]) for row in _cadence_starts(runtime_root)] == [ + ("reserved", f"{turn_key}:1") + ] + assert not list((runtime_root / "goals" / "fixture-goal" / "turns").glob("*.json")) + assert calls == {"host": 0, "writeback": 0, "spend": 0, "scheduler": 0} + + started_at_ms = int(_cadence_starts(runtime_root)[0]["started_at_ms"]) + monkeypatch.setattr(turn_command, "time", _FrozenTurnClock(started_at_ms + 120_000)) + restart = _managed_cadence(runtime_root) + recovered = run_loopx_turn_once( + plan, admit_start=restart.admit, confirm_start=restart.confirm, **common + ) + + assert recovered["status"] == "committed", recovered + assert recovered["admission"]["resumed"] is True + assert calls == {"host": 1, "writeback": 1, "spend": 1, "scheduler": 1} + assert [(row["state"], row["request_id"]) for row in _cadence_starts(runtime_root)] == [ + ("started", f"{turn_key}:1") + ] + + +def test_reserved_managed_start_recovers_after_death_before_the_attempt_record( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + """The same recovery holds when the journal exists without an attempt.""" + + plan, runtime_root, calls, common = _managed_start_fixture(tmp_path) + turn_key = str(plan["transaction"]["turn_key"]) # type: ignore[index] + cadence = _managed_cadence(runtime_root) + journal_writes = turn_executor._write_journal + + def die_before_attempt_record(path: Path, journal: dict[str, object]) -> None: + if "host_attempt_count" in journal: + raise RuntimeError("simulated process death before the attempt record") + journal_writes(path, journal) + + monkeypatch.setattr(turn_executor, "_write_journal", die_before_attempt_record) + with pytest.raises(RuntimeError, match="simulated process death"): + run_loopx_turn_once( + plan, admit_start=cadence.admit, confirm_start=cadence.confirm, **common + ) + monkeypatch.setattr(turn_executor, "_write_journal", journal_writes) + + journal = _journal(runtime_root) + assert "host_attempt_count" not in journal + assert journal["admission"]["reserved"] is True + assert [(row["state"], row["request_id"]) for row in _cadence_starts(runtime_root)] == [ + ("reserved", f"{turn_key}:1") + ] + assert calls == {"host": 0, "writeback": 0, "spend": 0, "scheduler": 0} + + started_at_ms = int(_cadence_starts(runtime_root)[0]["started_at_ms"]) + monkeypatch.setattr(turn_command, "time", _FrozenTurnClock(started_at_ms + 120_000)) + restart = _managed_cadence(runtime_root) + recovered = run_loopx_turn_once( + plan, admit_start=restart.admit, confirm_start=restart.confirm, **common + ) + + assert recovered["status"] == "committed", recovered + assert calls == {"host": 1, "writeback": 1, "spend": 1, "scheduler": 1} + assert [(row["state"], row["request_id"]) for row in _cadence_starts(runtime_root)] == [ + ("started", f"{turn_key}:1") + ] + + def test_run_once_rejects_oversized_built_in_host_result(tmp_path: Path) -> None: plan = _plan() calls = {"writeback": 0, "spend": 0, "scheduler": 0} From f121ef6af0cc52ea6b103d6ab476cf4a8c3485bf Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Thu, 24 Sep 2026 12:44:26 +0800 Subject: [PATCH 5/5] docs(rfc): record the two-phase managed start handoff Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../rfcs/automatic-execution-admission-v0.md | 10 +++++++++- .../rfcs/automatic-execution-admission-v0.zh-CN.md | 7 ++++++- 2 files changed, 15 insertions(+), 2 deletions(-) diff --git a/docs/architecture/rfcs/automatic-execution-admission-v0.md b/docs/architecture/rfcs/automatic-execution-admission-v0.md index c809163026..a3e5c19190 100644 --- a/docs/architecture/rfcs/automatic-execution-admission-v0.md +++ b/docs/architecture/rfcs/automatic-execution-admission-v0.md @@ -110,7 +110,7 @@ protection against clock manipulation. | Host path | Required contract | Initial boundary | | --- | --- | --- | -| Managed `turn run-once` | Atomic admission before a new host attempt, including failed-result recovery | M2 candidate; isolated CLI and concurrency tests, no external host promotion | +| Managed `turn run-once` | Atomic admission before a new host attempt, including failed-result recovery | M2 candidate; isolated CLI, concurrency and crash-recovery tests, no external host promotion | | Local legacy scheduler / external launchers | Route launches through admitted Turn or implement the same owner call | Not yet qualified; do not advertise enforcement | | Codex App automation | Apply floor-compatible timer, read actual schedule, ACK only matching facts | M1 schedule recommendation floor; hook coverage not qualified | | Attached interactive/manual session | Explicit manual intent; existing authority gates remain | Caller records reason; automatic continuation cannot masquerade as manual | @@ -223,6 +223,14 @@ the interval. The local CLI remains a same-UID trust boundary. The App timer-to-hook path, non-Turn launchers, packaged settings UI and live model-host promotion remain unqualified. No existing automation is activated or rebound by this proposal. +A managed start is two-phase in the same store: admission reserves the interval +slot, and the Turn executor confirms that reservation only after the host +attempt is durable in its journal. A crash between the two leaves the +reservation resumable by the same Turn identity once the floor is reached, so a +reserved-but-unstarted start never strands a Turn; a confirmed start stays +fail-closed for the same identity, and an explicit manual reason cannot bypass +that. A store record written without the phase field is read as an attempted +start, so an older or hand-edited file fails closed rather than resuming. M1 policy files are read as v1 and upgraded in place to v2 on the first configuration write or admitted start. The path stays stable; older binaries reject the v2 schema rather than silently discarding start records. Pause the diff --git a/docs/architecture/rfcs/automatic-execution-admission-v0.zh-CN.md b/docs/architecture/rfcs/automatic-execution-admission-v0.zh-CN.md index 76726cfa69..bec08f6f4a 100644 --- a/docs/architecture/rfcs/automatic-execution-admission-v0.zh-CN.md +++ b/docs/architecture/rfcs/automatic-execution-admission-v0.zh-CN.md @@ -78,7 +78,7 @@ App、Turn、前端、Lark 不得另存一套策略。共享 authority provider | 入口 | 必须履行的契约 | 本阶段边界 | | --- | --- | --- | -| Managed `turn run-once` | 每次新 host 尝试及失败重试前原子准入 | M2 候选;已做隔离 CLI 与并发测试,未推广外部宿主 | +| Managed `turn run-once` | 每次新 host 尝试及失败重试前原子准入 | M2 候选;已做隔离 CLI、并发与跨边界崩溃恢复测试,未推广外部宿主 | | 旧 local scheduler / 外部 launcher | 经受控 Turn 启动,或调用相同准入 owner | 尚未验收,不宣传为已强制执行 | | Codex App automation | 应用满足下限的定时器,回读真实值,事实匹配后 ACK | M1 调度建议下限;hook 覆盖范围未验收 | | 附着式交互 / 手动会话 | 显式手动意图;其余门禁保留 | 记录原因;自动续跑不能冒充手动 | @@ -166,6 +166,11 @@ M1 对 App 调度管理有独立价值,但不代表多宿主产品旅程完成 App 定时器到 hook、非 Turn launcher、打包设置界面及真实模型宿主推广仍未验收。 本提案不激活、不改绑任何已有自动化。 +managed 启动在同一 store 内分两步:准入预留间隔位,Turn executor 只在该 host 尝试 +已写入 Turn journal 之后确认该预留。两步之间进程退出时,同一 Turn 身份在满足时间 +下限后仍可恢复,因此"已预留但未启动"不会永久卡住 Turn;已确认的启动对同一身份保持 +fail-closed,显式手动理由也无法绕过。缺少阶段字段的旧记录按"已尝试启动"读取, +旧版或手工改写的文件因此 fail-closed,而不是被当作可恢复预留。 M1 的 v1 策略文件在首次配置写入或获准启动时原地升级为 v2,文件路径保持不变。 旧版程序会拒绝 v2 schema,避免静默丢弃启动记录;降级前必须暂停 launcher。