diff --git a/docs/development/control-plane-course/07-host-scheduler-and-heartbeat.md b/docs/development/control-plane-course/07-host-scheduler-and-heartbeat.md index 2c3d7cba1..ae8c0f90b 100644 --- a/docs/development/control-plane-course/07-host-scheduler-and-heartbeat.md +++ b/docs/development/control-plane-course/07-host-scheduler-and-heartbeat.md @@ -210,6 +210,30 @@ CLI acknowledged exact applied state Scheduler ACK 本身不构成 delivery,不 spend。 +### 兼容投影不能改变状态权威 / Compatibility preserves state authority + +`app_automation` 与 `codex_app` 是同一 proposal 的宿主投影,不是两个 scheduler +writer。读取优先采用 common state;仅有旧 Codex state 时,保留它经 TS 验证的 +`state_key`。两份投影的 backoff、host facts 与 ACK/failure 命令都必须携带同一个键。 +否则 common state 的当前 identity/index 会被误写到旧键,触发真实的初始档位冲突; +即使 CAS digest 相等也不代表 proposal 的状态归属正确。 + +`app_automation` and `codex_app` project one proposal, not two scheduler writers. +Reads prefer common state; a legacy-only Codex installation retains its +TS-validated legacy key. Both projections must carry that same key through +backoff, host facts and ACK/failure commands. Equal CAS digests do not prove that +a proposal belongs to the chosen state scope. + +兼容 Python API 与手工 Codex CLI 未显式传 `state_key` 时,从当前 packet 取得实际键; +显式键仍须匹配,不得静默改写。Trae 仍只接受 common key。这不是隐式迁移:不能把 +旧状态的非零档位当作空 common state 的首次 ACK,也不能放松 TS 的 reset/CAS 校验。 + +When Python compatibility APIs or manual Codex CLI calls omit `state_key`, they +use the current packet's key. Explicit keys must still match; Trae remains +common-key-only. This is not an implicit migration: a nonzero legacy stage must +not become a first ACK into missing common state, and TS reset/CAS checks remain +unchanged. + ### Proposal、Host Effect 与 Durable Receipt Scheduler 交互包含三个时间点,不能压成一个 `RRULE matches`: diff --git a/loopx/cli_commands/quota_context.py b/loopx/cli_commands/quota_context.py index df18f3e14..e54a098c6 100644 --- a/loopx/cli_commands/quota_context.py +++ b/loopx/cli_commands/quota_context.py @@ -26,7 +26,6 @@ ) from ..control_plane.scheduler.state import ( APP_AUTOMATION_STATEFUL_BACKOFF_STATE_KEY, - CODEX_APP_STATEFUL_BACKOFF_STATE_KEY, ) from ..status import AUTONOMOUS_REPLAN_PERIODIC_LOOKBACK, collect_status from ..turn_identity import mint_turn_instance_id, normalize_turn_instance_id @@ -232,7 +231,7 @@ def validate_quota_command_context_request( default_state_key = ( APP_AUTOMATION_STATEFUL_BACKOFF_STATE_KEY if selected_surface == HostSurface.TRAE_APP.value - else CODEX_APP_STATEFUL_BACKOFF_STATE_KEY + else None ) if ( selected_surface == HostSurface.TRAE_APP.value diff --git a/loopx/control_plane/quota/scheduler_ack.py b/loopx/control_plane/quota/scheduler_ack.py index 03d22dffb..744c4f19a 100644 --- a/loopx/control_plane/quota/scheduler_ack.py +++ b/loopx/control_plane/quota/scheduler_ack.py @@ -25,7 +25,7 @@ def _scheduler_packet( before: dict[str, Any], *, surface: str, - state_key: str, + state_key: str | None, ) -> tuple[dict[str, Any], dict[str, Any], dict[str, Any]]: scheduler_hint = ( before.get("scheduler_hint") @@ -36,7 +36,10 @@ def _scheduler_packet( packet_key = "app_automation" elif ( surface == CODEX_APP_SURFACE - and state_key == APP_AUTOMATION_STATEFUL_BACKOFF_STATE_KEY + and ( + state_key == APP_AUTOMATION_STATEFUL_BACKOFF_STATE_KEY + or (state_key is None and isinstance(scheduler_hint.get("app_automation"), dict)) + ) ): packet_key = "app_automation" else: @@ -54,6 +57,15 @@ def _scheduler_packet( return scheduler_hint, surface_packet, stateful_backoff +def _followup_state_key( + before: dict[str, Any], *, surface: str, state_key: str | None, +) -> str: + if state_key is not None: + return state_key + _, _, backoff = _scheduler_packet(before, surface=surface, state_key=None) + return str(backoff.get("state_key") or CODEX_APP_STATEFUL_BACKOFF_STATE_KEY) + + def _current_hint_identity( before: dict[str, Any], *, @@ -277,7 +289,7 @@ def record_quota_scheduler_ack_for_decision( agent_id: str | None, execute: bool = False, surface: str = CODEX_APP_SURFACE, - state_key: str = CODEX_APP_STATEFUL_BACKOFF_STATE_KEY, + state_key: str | None = None, applied_rrule: str | None = None, reset_token: str | None = None, identity_signature: str | None = None, @@ -286,6 +298,7 @@ def record_quota_scheduler_ack_for_decision( use_current_hint: bool = False, host_match_observed: bool = False, ) -> dict[str, Any]: + state_key = _followup_state_key(before, surface=surface, state_key=state_key) safe_agent_id = normalize_todo_claimed_by(agent_id) if host_match_observed and ( not str(applied_rrule or "").strip() @@ -374,12 +387,13 @@ def record_quota_scheduler_failure_for_decision( agent_id: str | None, execute: bool = False, surface: str = CODEX_APP_SURFACE, - state_key: str = CODEX_APP_STATEFUL_BACKOFF_STATE_KEY, + state_key: str | None = None, failed_rrule: str | None = None, observed_host_rrule: str | None = None, failure_kind: str = "host_tool_failure", generated_at: str | None = None, ) -> dict[str, Any]: + state_key = _followup_state_key(before, surface=surface, state_key=state_key) safe_agent_id = normalize_todo_claimed_by(agent_id) target_rrule = normalize_scheduler_rrule(failed_rrule) try: diff --git a/loopx/control_plane/scheduler/app_automation_compat.py b/loopx/control_plane/scheduler/app_automation_compat.py index cfa802673..0795ee07d 100644 --- a/loopx/control_plane/scheduler/app_automation_compat.py +++ b/loopx/control_plane/scheduler/app_automation_compat.py @@ -27,7 +27,7 @@ def build_codex_app_compatibility_projection( build_failure_hint: Callable[..., dict[str, Any]], build_fallback_hint: Callable[..., dict[str, Any]], ) -> dict[str, Any]: - """Translate the canonical App packet to the exact legacy Codex shape.""" + """Translate the App packet without changing its durable authority scope.""" legacy = copy.deepcopy(app_automation) legacy["applicability"] = "applicable" @@ -37,14 +37,18 @@ def build_codex_app_compatibility_projection( else None ) backoff = legacy.get("stateful_backoff") + state_key = ( + backoff["state_key"] + if isinstance(backoff, dict) + else CODEX_APP_STATEFUL_BACKOFF_STATE_KEY + ) if isinstance(backoff, dict): backoff["schema_version"] = CODEX_APP_STATEFUL_BACKOFF_SCHEMA_VERSION - backoff["state_key"] = CODEX_APP_STATEFUL_BACKOFF_STATE_KEY legacy_facts = ( { **scheduler_host_facts, "surface": CODEX_APP_SURFACE, - "state_key": CODEX_APP_STATEFUL_BACKOFF_STATE_KEY, + "state_key": state_key, } if isinstance(scheduler_host_facts, Mapping) else None @@ -64,6 +68,7 @@ def build_codex_app_compatibility_projection( else None ) legacy["failure_hint"] = build_failure_hint( + state_key=state_key, goal_id=goal_id, agent_id=agent_id, failed_rrule=legacy.get("recommended_rrule"), @@ -86,6 +91,7 @@ def build_codex_app_compatibility_projection( else {} ) legacy["ack_hint"] = build_ack_hint( + state_key=state_key, goal_id=goal_id, agent_id=agent_id, applied_rrule=canonical_args.get("applied_rrule"), diff --git a/loopx/control_plane/scheduler/scheduler_hint.py b/loopx/control_plane/scheduler/scheduler_hint.py index fe1477a76..390c74ca9 100644 --- a/loopx/control_plane/scheduler/scheduler_hint.py +++ b/loopx/control_plane/scheduler/scheduler_hint.py @@ -686,6 +686,12 @@ def build( if context is not None and context.app_automation_applicable else CODEX_APP_SURFACE ) + # The TS store has already validated this scope. Preserve it through + # both host projections, including a legacy-only Codex installation. + app_state_key = ( + (self.codex_app_scheduler_state or {}).get("state_key") + or APP_AUTOMATION_STATEFUL_BACKOFF_STATE_KEY + ) cadence = self._cadence_projections(codex_interval, codex_max, multiplier, cadence_progression_override) local_cadence_progression, app_cadence_progression = cadence["local"], cadence["app"] app_host_max, codex_max, floor = cadence["app_max"], cadence["local_max"], cadence["floor"] @@ -880,7 +886,7 @@ def build( ), "stateful_backoff": { "schema_version": APP_AUTOMATION_STATEFUL_BACKOFF_SCHEMA_VERSION, - "state_key": APP_AUTOMATION_STATEFUL_BACKOFF_STATE_KEY, + "state_key": app_state_key, "identity_signature": identity_signature, "reset_token": reset_token, "progression_index": current_index, @@ -921,7 +927,7 @@ def build( "goal_id": str(goal_id), "agent_id": str(agent_id), "surface": app_surface, - "state_key": APP_AUTOMATION_STATEFUL_BACKOFF_STATE_KEY, + "state_key": app_state_key, "reset_token": reset_token, "identity_signature": identity_signature, "progression_index": current_index, @@ -940,6 +946,7 @@ def build( app_automation["recommended_rrule"] = current_rrule if goal_id and agent_id: app_automation["failure_hint"] = build_app_automation_scheduler_failure_hint( + state_key=app_state_key, goal_id=goal_id, agent_id=agent_id, failed_rrule=current_rrule, @@ -965,6 +972,7 @@ def build( ) if ack_needed and goal_id and agent_id: app_automation["ack_hint"] = build_app_automation_scheduler_ack_hint( + state_key=app_state_key, goal_id=goal_id, agent_id=agent_id, applied_rrule=current_rrule, diff --git a/loopx/control_plane/scheduler/state.py b/loopx/control_plane/scheduler/state.py index ab2db5f74..5c669b295 100644 --- a/loopx/control_plane/scheduler/state.py +++ b/loopx/control_plane/scheduler/state.py @@ -319,8 +319,9 @@ def load_app_automation_scheduler_state( ) if current is not None or surface != CODEX_APP_SURFACE: return current - # Codex alone reads its pre-app_automation key so the next successful ACK - # can rewrite the cadence state under the provider-neutral contract. + # Codex alone reads its pre-app_automation key. Follow-ups retain that + # validated scope; copying its nonzero progression into a missing common + # state would incorrectly turn an acknowledged continuation into a reset. return load_scheduler_state( runtime_root, goal_id=goal_id, diff --git a/loopx/quota.py b/loopx/quota.py index e56f59ede..0dffa5ff2 100644 --- a/loopx/quota.py +++ b/loopx/quota.py @@ -93,7 +93,6 @@ SchedulerExecutionContextResolution, ) from .control_plane.scheduler.state import ( - CODEX_APP_STATEFUL_BACKOFF_STATE_KEY, CODEX_APP_SURFACE, ) from .control_plane.todos.contract import ( @@ -984,7 +983,7 @@ def record_quota_scheduler_ack( agent_id: str | None = None, available_capabilities: Any = None, surface: str = CODEX_APP_SURFACE, - state_key: str = CODEX_APP_STATEFUL_BACKOFF_STATE_KEY, + state_key: str | None = None, applied_rrule: str | None = None, reset_token: str | None = None, identity_signature: str | None = None, @@ -1023,7 +1022,7 @@ def record_quota_scheduler_ack( agent_id=safe_agent_id, execute=execute, surface=str(surface or CODEX_APP_SURFACE).strip() or CODEX_APP_SURFACE, - state_key=str(state_key or CODEX_APP_STATEFUL_BACKOFF_STATE_KEY).strip(), + state_key=str(state_key).strip() if state_key is not None else None, applied_rrule=applied_rrule, reset_token=reset_token, identity_signature=identity_signature, diff --git a/tests/control_plane/test_scheduler_compat_state_key.py b/tests/control_plane/test_scheduler_compat_state_key.py new file mode 100644 index 000000000..2f4f70062 --- /dev/null +++ b/tests/control_plane/test_scheduler_compat_state_key.py @@ -0,0 +1,250 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +import json +import os +from pathlib import Path +import subprocess + +import pytest + +from examples.control_plane.quota_plan_fixtures import ( + SCOPED_AGENT_ID, + write_cli_fixture, +) +from loopx.control_plane.quota.scheduler_ack import ( + record_quota_scheduler_ack_for_decision, + record_quota_scheduler_failure_for_decision, +) +from loopx.control_plane.scheduler.state import ( + APP_AUTOMATION_STATEFUL_BACKOFF_STATE_KEY as APP_KEY, + LEGACY_CODEX_APP_STATEFUL_BACKOFF_STATE_KEY as LEGACY_KEY, + build_scheduler_state, + load_scheduler_state, + scheduler_state_path, + write_scheduler_state, +) +from loopx.control_plane.testing.canary_harness import run_json_cli_result + +REPO_ROOT = Path(__file__).resolve().parents[2] +GOAL_ID = "needs-operator" + + +@pytest.mark.parametrize("state_layout", ["canonical", "legacy", "coexisting"]) +@pytest.mark.parametrize("operation", ["ack", "host_failure"]) +@pytest.mark.parametrize("caller", ["native_cli", "python_adapter", "python_cli"]) +def test_compat_followup_preserves_the_read_authority_through_public_entry_points( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + state_layout: str, + operation: str, + caller: str, +) -> None: + registry, runtime, project = write_cli_fixture( + tmp_path / "fixture", scoped_agents=True + ) + monkeypatch.setenv("CODEX_HOME", str(tmp_path / "codex-home")) + monkeypatch.delenv("CODEX_THREAD_ID", raising=False) + turn_id = "scheduler-compat-turn" + + def guard() -> dict: + code, result = run_json_cli_result( + "quota", + "should-run", + "--goal-id", + GOAL_ID, + "--agent-id", + SCOPED_AGENT_ID, + "--codex-app", + "--turn-instance-id", + turn_id, + registry_path=registry, + runtime_root=runtime, + cwd=project, + ) + assert code == 0, result + return result + + first = guard()["scheduler_hint"]["app_automation"] + selected_key = LEGACY_KEY if state_layout == "legacy" else APP_KEY + acknowledged = build_scheduler_state( + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + state_key=selected_key, + reset_token=first["stateful_backoff"]["reset_token"], + identity_signature=first["stateful_backoff"]["identity_signature"], + progression_index=0, + progression_minutes=[30, 60], + last_applied_rrule="FREQ=MINUTELY;INTERVAL=30", + updated_at=(datetime.now(timezone.utc) - timedelta(minutes=31)).isoformat(), + ) + write_scheduler_state( + runtime, + acknowledged, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + state_key=selected_key, + ) + old_path = scheduler_state_path( + runtime, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + state_key=LEGACY_KEY, + ) + if state_layout == "coexisting": + obsolete = build_scheduler_state( + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + state_key=LEGACY_KEY, + reset_token="obsolete-reset", + identity_signature="obsolete-identity", + progression_index=0, + progression_minutes=[3, 6, 10], + last_applied_rrule="FREQ=MINUTELY;INTERVAL=3", + updated_at=datetime.now(timezone.utc).isoformat(), + ) + write_scheduler_state( + runtime, + obsolete, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + state_key=LEGACY_KEY, + ) + old_bytes = old_path.read_bytes() + + decision = guard() + app = decision["scheduler_hint"]["app_automation"] + compat = decision["scheduler_hint"]["codex_app"] + assert app["stateful_backoff"]["progression_index"] == 1 + assert app["stateful_backoff"]["state_key"] == selected_key + hint = compat["ack_hint" if operation == "ack" else "failure_hint"] + if caller == "native_cli": + completed = subprocess.run( + [ + str(REPO_ROOT / "scripts" / "loopx"), + "--format", + "json", + *hint["cli_args"], + ], + cwd=project, + check=False, + capture_output=True, + text=True, + env={**os.environ, "LOOPX_PYTHON": str(tmp_path / "python-must-not-run")}, + timeout=30, + ) + assert completed.returncode == 0, completed.stdout + completed.stderr + result = json.loads(completed.stdout) + elif caller == "python_cli": + command = ( + "scheduler-ack-current" if operation == "ack" else "scheduler-fail-current" + ) + operation_args = ( + ["--applied-rrule", "FREQ=MINUTELY;INTERVAL=60"] + if operation == "ack" + else [ + "--failed-rrule", + "FREQ=MINUTELY;INTERVAL=60", + "--app-automation-current-rrule", + "FREQ=MINUTELY;INTERVAL=30", + ] + ) + code, result = run_json_cli_result( + "quota", + command, + "--goal-id", + GOAL_ID, + "--agent-id", + SCOPED_AGENT_ID, + "--codex-app", + "--turn-instance-id", + turn_id, + "--execute", + *operation_args, + registry_path=registry, + runtime_root=runtime, + cwd=project, + ) + assert code == 0, result + elif operation == "ack": + result = record_quota_scheduler_ack_for_decision( + decision, + runtime_root=runtime, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + execute=True, + applied_rrule="FREQ=MINUTELY;INTERVAL=60", + use_current_hint=True, + ) + else: + result = record_quota_scheduler_failure_for_decision( + decision, + runtime_root=runtime, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + execute=True, + failed_rrule="FREQ=MINUTELY;INTERVAL=60", + observed_host_rrule="FREQ=MINUTELY;INTERVAL=30", + ) + assert compat["stateful_backoff"]["state_key"] == selected_key + assert result["ok"] is True, result + assert result["state_key"] == selected_key + assert result["scheduler_commit"]["status"] == "written" + state = load_scheduler_state( + runtime, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + state_key=selected_key, + ) + assert state is not None + assert state["progression_index"] == 1 + if operation == "ack": + assert state["last_applied_rrule"] == "FREQ=MINUTELY;INTERVAL=60" + else: + assert state["host_update_failures"] + if state_layout == "coexisting": + assert old_path.read_bytes() == old_bytes + if state_layout == "legacy": + assert ( + load_scheduler_state( + runtime, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + state_key=APP_KEY, + ) + is None + ) + + +def test_explicit_state_key_is_not_silently_retargeted(tmp_path: Path) -> None: + before = { + "scheduler_hint": { + "codex_app": { + "stateful_backoff": { + "state_key": APP_KEY, + "reset_token": "reset", + "identity_signature": "identity", + "progression_index": 0, + "progression_minutes": [30, 60], + "current_rrule": "FREQ=MINUTELY;INTERVAL=30", + } + } + } + } + result = record_quota_scheduler_ack_for_decision( + before, + runtime_root=tmp_path, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + execute=True, + state_key=LEGACY_KEY, + applied_rrule="FREQ=MINUTELY;INTERVAL=30", + ) + assert result["ok"] is False + assert "state-key does not match" in result["reason"] + assert not scheduler_state_path( + tmp_path, + goal_id=GOAL_ID, + agent_id=SCOPED_AGENT_ID, + state_key=LEGACY_KEY, + ).exists()