diff --git a/loopx/chat_status_api.py b/loopx/chat_status_api.py index 273dfbf031..68b2f83e49 100644 --- a/loopx/chat_status_api.py +++ b/loopx/chat_status_api.py @@ -11,8 +11,14 @@ from .chat import redact_local_paths from .chat_goal_subagent_api import goal_subagent_configuration_enabled from .chat_workspace_directory import workspace_goal_directory -from .codex_app_thread_activity import codex_thread_observers -from .control_plane.agents.host_thread_activity import attach_host_thread_activity +from .codex_app_thread_activity import ( + codex_delivery_expectations, + codex_thread_observers, +) +from .control_plane.agents.host_thread_activity import ( + attach_host_delivery_windows, + attach_host_thread_activity, +) from .control_plane.effect_runtime import ( EffectRuntimePermanentIOError, EffectRuntimeRemoteError, @@ -197,6 +203,10 @@ def _status(self, *, delivery_review: bool = False) -> None: attach_host_thread_activity( projection, observers=codex_thread_observers() ) + attach_host_delivery_windows( + projection, + expectations=codex_delivery_expectations, + ) if delivery_review: if projection.get("ok") is not True: self._send_error("Delivery review sources are unavailable.", status=503) diff --git a/loopx/codex_app_thread_activity.py b/loopx/codex_app_thread_activity.py index 4502617135..746783dde5 100644 --- a/loopx/codex_app_thread_activity.py +++ b/loopx/codex_app_thread_activity.py @@ -13,11 +13,14 @@ import os import re import sqlite3 +import tomllib from collections.abc import Iterable, Mapping from pathlib import Path from typing import Any from .control_plane.agents.host_thread_activity import ( + HostDeliveryExpectation, + HostDeliveryScope, HostThreadActivity, HostThreadObserver, HostThreadState, @@ -26,6 +29,10 @@ # Remote surfaces such as codex-app-ssh keep their store on another machine. CODEX_LOCAL_STORE_SURFACES = frozenset({"codex-app", "codex-cli-tui", "codex-ide-plugin"}) + +# An installed app automation drives the app surface; the home's other local +# surfaces keep their own loop, so the automation cadence is not theirs to carry. +CODEX_APP_BINDING_SURFACE = "codex-app" CODEX_HOMES_ENV = "LOOPX_CODEX_HOMES" _STATE_DB_RE = re.compile(r"^state_(\d+)\.sqlite$") @@ -37,6 +44,11 @@ _TAIL_LIMIT_BYTES = 8 * 1024 * 1024 _ROLLOUT_CACHE_LIMIT = 256 +# The canonical rrule parser reads INTERVAL with JavaScript's integer +# conversion, which reports no value outside this range. +_MAX_SAFE_INTEGER = 2**53 - 1 +_MAX_SAFE_INTEGER_DIGITS = len(str(_MAX_SAFE_INTEGER)) + _rollout_cache: dict[tuple[str, int, int], HostThreadActivity] = {} @@ -293,3 +305,137 @@ def observe(thread_ids: Iterable[str]) -> dict[str, HostThreadActivity]: return observe_codex_threads(thread_ids, homes=homes) return {surface: observe for surface in CODEX_LOCAL_STORE_SURFACES} + + +def _automation_rrule_interval_minutes(rrule: Any) -> int | None: + """Read the interval a minutely automation rrule asks for. + + The canonical parse is ``schedulerRruleIntervalMinutes`` in the TypeScript + scheduler (``scheduler/state_store.ts``). Its Python twin + (``scheduler_rrule_interval_minutes``) is an effect-runtime call, so reading + an installed automation through it would give the status route a scheduler + runtime dependency. This mirrors the canonical contract rather than keeping + a second, looser dialect: the same normalization, the same first-``=`` + split, an exact ``FREQ=MINUTELY``, and a positive integer ``INTERVAL``. An + rrule this adapter cannot read reports no interval, never one. + """ + + text = str(rrule or "").strip() + text = re.sub(r"\s+", " ", text) + if text.upper().startswith("RRULE:"): + text = text[6:].strip() + parts: dict[str, str] = {} + for part in text.split(";"): + separator = part.find("=") + if separator < 0: + continue + parts[part[:separator].strip().upper()] = part[separator + 1 :].strip() + if parts.get("FREQ", "").upper() != "MINUTELY": + return None + # ``INTERVAL`` is read the way the canonical parser reads it: absent, + # non-numeric and out-of-range all report no interval, so an unreadable + # rrule can never be measured as if it fired every minute. + raw_interval = parts.get("INTERVAL", "") + if re.fullmatch(r"[+-]?[0-9]+", raw_interval) is None: + return None + # Rule the value out by its significant digits before asking the interpreter + # to build it. Python refuses an integer of more than + # ``sys.get_int_max_str_digits()`` digits, and an unreadable manifest must + # never escape this provider: the status route turns a local projection + # failure into one 500 for the whole workspace, so a single oversized + # manifest would hide every healthy lane. The canonical parser reports no + # value for anything this large either. + negative = raw_interval.startswith("-") + digits = raw_interval.lstrip("+-").lstrip("0") + if len(digits) > _MAX_SAFE_INTEGER_DIGITS: + return None + # Convert the significant digits only: the interpreter's limit counts + # leading zeros too, so the length check above must bound what is converted. + interval = int(("-" if negative else "") + digits) if digits else 0 + if interval > _MAX_SAFE_INTEGER or interval < -_MAX_SAFE_INTEGER: + return None + return interval if interval > 0 else None + + +def _installed_automation_intervals( + homes: list[Path] | None = None, +) -> dict[tuple[str, str, str], int]: + """Map each installed heartbeat automation to the lane it serves. + + A lane is identified the way the canonical resolver identifies it + (``loopx.upgrade.resolve_codex_app_automation_rrule``): Goal, agent and the + bound ``target_thread_id``. Only a ``kind = "heartbeat"`` manifest carries a + heartbeat cadence, and an automation that cannot be placed on one lane -- or + two that claim the same lane -- is omitted, so the window stays ``unknown`` + instead of being measured against another lane's cadence. + """ + + from .upgrade import infer_agent_id_from_prompt, infer_goal_id_from_prompt + + intervals: dict[tuple[str, str, str], int] = {} + ambiguous: set[tuple[str, str, str]] = set() + for home in codex_homes() if homes is None else homes: + for path in sorted((home / "automations").glob("*/automation.toml")): + try: + item = tomllib.loads(path.read_text(encoding="utf-8")) + except (OSError, UnicodeDecodeError, tomllib.TOMLDecodeError): + continue + if str(item.get("status") or "").upper() != "ACTIVE": + continue + # Another kind installed for the same Goal is a different loop, not + # this lane's heartbeat. + if str(item.get("kind") or "").strip().lower() != "heartbeat": + continue + prompt = item.get("prompt") + if not isinstance(prompt, str) or not prompt.strip(): + continue + goal_id = infer_goal_id_from_prompt(prompt) + agent_id = infer_agent_id_from_prompt(prompt) + target_thread_id = str(item.get("target_thread_id") or "").strip() + interval = _automation_rrule_interval_minutes(item.get("rrule")) + if not goal_id or not agent_id or not target_thread_id or interval is None: + continue + key = (goal_id, agent_id, target_thread_id) + if key in intervals: + # Two automations claim one lane: neither is that lane's cadence. + del intervals[key] + ambiguous.add(key) + continue + if key not in ambiguous: + intervals[key] = interval + return intervals + + +def codex_delivery_expectations( + scopes: Iterable[HostDeliveryScope], + *, + homes: list[Path] | None = None, +) -> dict[HostDeliveryScope, HostDeliveryExpectation]: + """Resolve each bound app lane's cadence from its installed automation. + + A lane is placed by its canonical binding, so an automation installed for + another thread of the same Goal is not read as this lane's. A lane the + installed automations cannot place is omitted, so its window is reported as + ``unknown`` rather than measured against a guessed cadence. + """ + + wanted = list(scopes) + if not any(scope.host_surface == CODEX_APP_BINDING_SURFACE for scope in wanted): + return {} + intervals = _installed_automation_intervals(homes) + expectations: dict[HostDeliveryScope, HostDeliveryExpectation] = {} + for scope in wanted: + if scope.host_surface != CODEX_APP_BINDING_SURFACE: + continue + if not scope.thread_id: + # Without the binding, this lane cannot be told apart from another + # lane of the same Goal and agent. + continue + interval = intervals.get((scope.goal_id, scope.agent_id, scope.thread_id)) + if interval is None: + continue + expectations[scope] = HostDeliveryExpectation( + expected_interval_minutes=interval, + source="codex_app_automation_rrule", + ) + return expectations diff --git a/loopx/control_plane/agents/host_thread_activity.py b/loopx/control_plane/agents/host_thread_activity.py index 4a579c69a1..5a9b62f1cc 100644 --- a/loopx/control_plane/agents/host_thread_activity.py +++ b/loopx/control_plane/agents/host_thread_activity.py @@ -16,7 +16,9 @@ from typing import Any HOST_THREAD_ACTIVITY_SCHEMA_VERSION = "loopx_host_thread_activity_v0" +HOST_DELIVERY_WINDOW_SCHEMA_VERSION = "loopx_host_delivery_window_v0" MAX_OBSERVED_THREADS_PER_GOAL = 32 +DEFAULT_DELIVERY_WINDOW_TOLERANCE = 2 class HostThreadState(str, Enum): @@ -138,3 +140,279 @@ def attach_host_thread_activity( ).value, "threads": threads, } + + +class HostDeliveryWindowState(str, Enum): + # The lane produced activity inside the cadence it is expected to keep. + FRESH = "fresh" + # No activity was observed inside the expected window. This names the + # symptom only; it does not by itself name a cause. + STALE = "stale" + # No observation of this lane was available at all. + MISSING = "missing" + # The window could not be computed, so it is not reported as fresh. + UNKNOWN = "unknown" + + +class HostDeliveryWindowUnknownReason(str, Enum): + NO_EXPECTATION = "no_expectation" + NO_OBSERVATION = "no_observation" + UNPARSEABLE_OBSERVATION = "unparseable_observation" + + +_REASONED_DELIVERY_WINDOW_STATES = frozenset( + {HostDeliveryWindowState.MISSING, HostDeliveryWindowState.UNKNOWN} +) + + +@dataclass(frozen=True) +class HostDeliveryExpectation: + """The cadence one bound lane is expected to produce activity at. + + ``source`` names where the cadence came from, so a reader can tell an agreed + loopX cadence from a host-reported one without guessing. + """ + + expected_interval_minutes: int | None + source: str + + def __post_init__(self) -> None: + if ( + self.expected_interval_minutes is not None + and self.expected_interval_minutes <= 0 + ): + raise ValueError("expected_interval_minutes must be positive or null") + if not self.source.strip(): + raise ValueError("source is required for a delivery expectation") + + +@dataclass(frozen=True) +class HostDeliveryScope: + """One bound lane a delivery window can be reported for. + + ``thread_id`` is the canonical binding the lane was observed through, so a + provider can tell this lane's own automation from another lane's. It is + carried only to resolve the expectation: the window payload never names it. + """ + + goal_id: str + agent_id: str + host_surface: str + thread_id: str | None = None + + +@dataclass(frozen=True) +class HostDeliveryWindow: + state: HostDeliveryWindowState + reason: HostDeliveryWindowUnknownReason | None = None + last_observed_at: str | None = None + age_seconds: int | None = None + age_hours: float | None = None + expected_interval_minutes: int | None = None + window_minutes: int | None = None + tolerance: int | None = None + source: str | None = None + + def __post_init__(self) -> None: + reasoned = self.state in _REASONED_DELIVERY_WINDOW_STATES + if reasoned != (self.reason is not None): + raise ValueError( + "a missing or unknown delivery window requires exactly one reason" + ) + + @classmethod + def unknown( + cls, + reason: HostDeliveryWindowUnknownReason, + *, + state: HostDeliveryWindowState = HostDeliveryWindowState.UNKNOWN, + source: str | None = None, + ) -> HostDeliveryWindow: + return cls(state=state, reason=reason, source=source) + + def to_payload(self) -> dict[str, Any]: + payload: dict[str, Any] = { + "schema_version": HOST_DELIVERY_WINDOW_SCHEMA_VERSION, + "state": self.state.value, + "reason": self.reason.value if self.reason else None, + "last_observed_at": self.last_observed_at, + "age_seconds": self.age_seconds, + "age_hours": self.age_hours, + "expected_interval_minutes": self.expected_interval_minutes, + "window_minutes": self.window_minutes, + "tolerance": self.tolerance, + "source": self.source, + } + return {key: value for key, value in payload.items() if value is not None} + + +def _parse_observation_time(value: str | None) -> datetime | None: + text = str(value or "").strip() + if not text: + return None + try: + parsed = datetime.fromisoformat(text.replace("Z", "+00:00")) + except ValueError: + return None + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed.astimezone(timezone.utc) + + +def build_host_delivery_window( + activity: HostThreadActivity | None, + expectation: HostDeliveryExpectation | None, + *, + now: datetime | None = None, + tolerance: int = DEFAULT_DELIVERY_WINDOW_TOLERANCE, +) -> HostDeliveryWindow: + """Compare one observed lane against the cadence it is expected to keep. + + Absence of evidence is never reported as healthy: an unobserved lane is + ``missing`` and an incomputable window is ``unknown``. Neither is ``fresh``. + """ + + if tolerance < 1: + raise ValueError("tolerance must be at least one interval") + if expectation is None or not expectation.expected_interval_minutes: + return HostDeliveryWindow.unknown(HostDeliveryWindowUnknownReason.NO_EXPECTATION) + if activity is None or activity.state is HostThreadState.UNKNOWN: + return HostDeliveryWindow.unknown( + HostDeliveryWindowUnknownReason.NO_OBSERVATION, + state=HostDeliveryWindowState.MISSING, + source=expectation.source, + ) + observed = _parse_observation_time(activity.last_event_at) + if observed is None: + return HostDeliveryWindow.unknown( + HostDeliveryWindowUnknownReason.UNPARSEABLE_OBSERVATION, + source=expectation.source, + ) + window_minutes = expectation.expected_interval_minutes * tolerance + reference = now or datetime.now(timezone.utc) + age_seconds = max(0, int((reference - observed).total_seconds())) + return HostDeliveryWindow( + state=( + HostDeliveryWindowState.FRESH + if age_seconds <= window_minutes * 60 + else HostDeliveryWindowState.STALE + ), + last_observed_at=activity.last_event_at, + age_seconds=age_seconds, + age_hours=round(age_seconds / 3600, 2), + expected_interval_minutes=expectation.expected_interval_minutes, + window_minutes=window_minutes, + tolerance=tolerance, + source=expectation.source, + ) + + +# Maps bound lanes to the cadence each is expected to keep; lanes it cannot +# resolve are omitted rather than guessed. +HostDeliveryExpectationProvider = Callable[ + [Iterable[HostDeliveryScope]], Mapping[HostDeliveryScope, HostDeliveryExpectation] +] + + +def _bound_thread_ids(goal: Mapping[str, Any]) -> dict[tuple[str, str], str | None]: + """Map each ``(agent, surface)`` lane of a Goal to its canonical thread. + + A lane the Goal binds to more than one thread has no single identity, so it + maps to ``None`` and its window stays ``unknown``. + """ + + resolved: dict[tuple[str, str], str | None] = {} + for agent_id, host_surface, thread_id in _goal_bindings(goal): + key = (agent_id, host_surface) + resolved[key] = None if key in resolved else thread_id + return resolved + + +def _scope_from_row( + goal_id: str, row: Mapping[str, Any], bound: Mapping[tuple[str, str], str | None] +) -> HostDeliveryScope | None: + agent_id = str(row.get("agent_id") or "").strip() + host_surface = str(row.get("host_surface") or "").strip() + if not goal_id or not agent_id or not host_surface: + return None + return HostDeliveryScope( + goal_id=goal_id, + agent_id=agent_id, + host_surface=host_surface, + thread_id=bound.get((agent_id, host_surface)), + ) + + +def attach_host_delivery_windows( + status_payload: dict[str, Any], + *, + expectations: HostDeliveryExpectationProvider, + tolerance: int = DEFAULT_DELIVERY_WINDOW_TOLERANCE, + now: datetime | None = None, +) -> None: + """Add ``delivery_window`` to each observed lane of ``host_thread_activity``. + + This reads the projection that ``attach_host_thread_activity`` already + produced, so no host record is read twice. Each lane is resolved through the + canonical binding the Goal records, so a cadence installed for another + thread of the same Goal is not read as this lane's. A lane the expectation + provider cannot resolve is reported as ``unknown``, never as healthy. + """ + + run_history = status_payload.get("run_history") + goals = run_history.get("goals") if isinstance(run_history, Mapping) else None + if not isinstance(goals, list): + return + rows_by_goal: list[ + tuple[str, list[dict[str, Any]], dict[tuple[str, str], str | None]] + ] = [] + scopes: list[HostDeliveryScope] = [] + for goal in goals: + if not isinstance(goal, dict): + continue + activity = goal.get("host_thread_activity") + threads = activity.get("threads") if isinstance(activity, Mapping) else None + if not isinstance(threads, list): + continue + rows = [row for row in threads if isinstance(row, dict)] + if not rows: + continue + goal_id = str(goal.get("id") or "").strip() + bound = _bound_thread_ids(goal) + rows_by_goal.append((goal_id, rows, bound)) + for row in rows: + scope = _scope_from_row(goal_id, row, bound) + if scope is not None: + scopes.append(scope) + resolved = expectations(scopes) if scopes else {} + for goal_id, rows, bound in rows_by_goal: + for row in rows: + scope = _scope_from_row(goal_id, row, bound) + row["delivery_window"] = build_host_delivery_window( + _activity_from_row(row), + resolved.get(scope) if scope is not None else None, + now=now, + tolerance=tolerance, + ).to_payload() + + +def _activity_from_row(row: Mapping[str, Any]) -> HostThreadActivity: + state_value = str(row.get("state") or "") + try: + state = HostThreadState(state_value) + except ValueError: + return HostThreadActivity.unknown(HostThreadUnknownReason.RECORD_UNRECOGNIZED) + reason_value = str(row.get("reason") or "") + reason: HostThreadUnknownReason | None = None + if reason_value: + try: + reason = HostThreadUnknownReason(reason_value) + except ValueError: + reason = HostThreadUnknownReason.RECORD_UNRECOGNIZED + return HostThreadActivity( + state=state, + reason=reason, + turn_started_at=row.get("turn_started_at"), + last_turn_ended_at=row.get("last_turn_ended_at"), + last_event_at=row.get("last_event_at"), + ) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 78f7daceb1..f7b506e010 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -527,7 +527,7 @@ }, { "site": "loopx/chat_status_api.py::.ChatStatusRequestMixin._status::codec_read:load_registry#1", - "line": 152, + "line": 158, "column": 28, "kind": "codec_read", "api": "load_registry", @@ -535,7 +535,7 @@ }, { "site": "loopx/chat_status_api.py::.ChatStatusRequestMixin._status::codec_read:load_registry#2", - "line": 179, + "line": 185, "column": 21, "kind": "codec_read", "api": "load_registry", diff --git a/tests/test_host_thread_activity.py b/tests/test_host_thread_activity.py index 2c112b6a38..9f75eae026 100644 --- a/tests/test_host_thread_activity.py +++ b/tests/test_host_thread_activity.py @@ -2,6 +2,7 @@ import json import sqlite3 +from datetime import datetime, timedelta, timezone from pathlib import Path from types import SimpleNamespace from typing import Any @@ -10,12 +11,23 @@ import loopx.chat_status_api as chat_status_api import loopx.codex_app_thread_activity as codex_activity -from loopx.codex_app_thread_activity import codex_homes, codex_thread_observers, observe_codex_threads +from loopx.codex_app_thread_activity import ( + _automation_rrule_interval_minutes, + codex_delivery_expectations, + codex_homes, + codex_thread_observers, + observe_codex_threads, +) from loopx.control_plane.agents.host_thread_activity import ( + HostDeliveryExpectation, + HostDeliveryScope, + HostDeliveryWindowState, HostThreadActivity, HostThreadState, HostThreadUnknownReason, + attach_host_delivery_windows, attach_host_thread_activity, + build_host_delivery_window, ) T0 = "2026-09-25T10:00:00.000Z" @@ -356,9 +368,583 @@ def _send_error(self, message: str, **kwargs: Any) -> None: Handler()._status() threads = sent[0]["run_history"]["goals"][0]["host_thread_activity"]["threads"] - assert threads == [{"agent_id": "a", "host_surface": "codex-app", "state": "turn_open", "turn_started_at": T0, "last_event_at": T1}] + assert threads == [ + { + "agent_id": "a", + "host_surface": "codex-app", + "state": "turn_open", + "turn_started_at": T0, + "last_event_at": T1, + # This home installs no automation, so the expected cadence is + # unknown rather than assumed healthy. + "delivery_window": { + "schema_version": "loopx_host_delivery_window_v0", + "state": "unknown", + "reason": "no_expectation", + }, + } + ] def test_remote_codex_surfaces_have_no_local_observer() -> None: assert "codex-app-ssh" not in codex_thread_observers() assert {"codex-app", "codex-cli-tui", "codex-ide-plugin"} <= set(codex_thread_observers()) + + +def test_automation_rrule_interval_reads_only_minutely_intervals() -> None: + assert _automation_rrule_interval_minutes("FREQ=MINUTELY;INTERVAL=30") == 30 + assert _automation_rrule_interval_minutes("freq=minutely;interval=5") == 5 + assert _automation_rrule_interval_minutes("RRULE:FREQ=MINUTELY;INTERVAL=2") == 2 + # Field order, spacing and duplicates follow the canonical parser. + assert _automation_rrule_interval_minutes("INTERVAL=3;FREQ=MINUTELY") == 3 + assert _automation_rrule_interval_minutes("FREQ=MINUTELY; INTERVAL = 7") == 7 + assert _automation_rrule_interval_minutes("FREQ=MINUTELY;INTERVAL=5;INTERVAL=9") == 9 + assert _automation_rrule_interval_minutes( + f"FREQ=MINUTELY;INTERVAL={2**53 - 1}" + ) == 2**53 - 1 + # Leading zeros are not part of the value, and must not be converted. + assert _automation_rrule_interval_minutes( + "FREQ=MINUTELY;INTERVAL=" + "0" * 5000 + "5" + ) == 5 + # Anything this adapter cannot read reports no interval rather than one. + assert _automation_rrule_interval_minutes("FREQ=HOURLY;INTERVAL=2") is None + assert _automation_rrule_interval_minutes("FREQ=MINUTELY;INTERVAL=0") is None + assert _automation_rrule_interval_minutes(None) is None + + +@pytest.mark.parametrize( + "rrule", + [ + # An omitted INTERVAL is unreadable, not a one-minute cadence. + "FREQ=MINUTELY", + "FREQ=MINUTELY;INTERVAL=-3", + "FREQ=MINUTELY;INTERVAL=bad", + "FREQ=MINUTELY;INTERVAL=", + "FREQ=MINUTELY;INTERVAL=1.5", + # The frequency must match exactly, not as a substring. + "FREQ=MINUTELYISH;INTERVAL=7", + "MINUTELY;INTERVAL=7", + "", + ], +) +def test_an_unreadable_automation_rrule_never_reports_a_cadence(rrule: str) -> None: + assert _automation_rrule_interval_minutes(rrule) is None + + +@pytest.mark.parametrize( + "interval", + [ + "9" * 16, # one digit past the safe integer range + str(2**53), + "1" * 4301, # past the interpreter's own decimal digit limit + "0" * 5000, # leading zeros the interpreter still counts as digits + "-" + "1" * 4301, + ], + ids=["16-nines", "safe-plus-one", "4301-digits", "5000-zeros", "negative-4301"], +) +def test_an_out_of_range_interval_is_reported_without_converting_it( + interval: str, +) -> None: + """The provider must not ask the interpreter for a value it refuses to build. + + ``int()`` raises above ``sys.get_int_max_str_digits()``, and the status route + turns any local projection failure into one 500 for the whole workspace. + """ + + assert _automation_rrule_interval_minutes( + f"FREQ=MINUTELY;INTERVAL={interval}" + ) is None + + +def _install_automation( + home: CodexHome, + *, + automation_id: str, + goal_id: str, + agent_id: str, + rrule: str, + status: str = "ACTIVE", + kind: str = "heartbeat", + target_thread_id: str = "t-open", +) -> None: + path = home.root / "automations" / automation_id / "automation.toml" + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text( + f'version = 1\nid = "{automation_id}"\nkind = "{kind}"\n' + f'status = "{status}"\nrrule = "{rrule}"\ntarget_thread_id = "{target_thread_id}"\n' + f'prompt = "Advance `{goal_id}` from active state. Agent: `{agent_id}`."\n', + encoding="utf-8", + ) + + +def test_delivery_expectations_resolve_only_installed_active_app_automations( + tmp_path: Path, +) -> None: + home = CodexHome(tmp_path / ".codex") + _install_automation( + home, automation_id="a", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=30", + ) + _install_automation( + home, automation_id="b", goal_id="bound", agent_id="b", + rrule="FREQ=MINUTELY;INTERVAL=30", status="PAUSED", + ) + _install_automation( + home, automation_id="c", goal_id="bound", agent_id="c", + rrule="FREQ=HOURLY;INTERVAL=2", + ) + + lane = HostDeliveryScope( + goal_id="bound", agent_id="a", host_surface="codex-app", thread_id="t-open" + ) + resolved = codex_delivery_expectations( + [ + lane, + HostDeliveryScope(goal_id="bound", agent_id="b", host_surface="codex-app", thread_id="t-open"), + HostDeliveryScope(goal_id="bound", agent_id="c", host_surface="codex-app", thread_id="t-open"), + # The app automation is not the loop that drives another surface. + HostDeliveryScope(goal_id="bound", agent_id="a", host_surface="codex-cli-tui", thread_id="t-open"), + # Without the canonical binding a lane has no identity to place. + HostDeliveryScope(goal_id="bound", agent_id="a", host_surface="codex-app"), + ], + homes=[home.root], + ) + + assert resolved == { + lane: HostDeliveryExpectation( + expected_interval_minutes=30, source="codex_app_automation_rrule" + ) + } + + +def test_delivery_expectations_ignore_another_thread_or_kind_of_the_same_goal( + tmp_path: Path, +) -> None: + home = CodexHome(tmp_path / ".codex") + _install_automation( + home, automation_id="lane", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=5", target_thread_id="t-lane", + ) + # Same Goal and agent, but another bound thread. + _install_automation( + home, automation_id="other-thread", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=60", target_thread_id="t-other", + ) + # Same Goal, agent and thread, but a different kind of automation. + _install_automation( + home, automation_id="cron", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=7", kind="cron", target_thread_id="t-lane", + ) + + lane = HostDeliveryScope( + goal_id="bound", agent_id="a", host_surface="codex-app", thread_id="t-lane" + ) + assert codex_delivery_expectations([lane], homes=[home.root]) == { + lane: HostDeliveryExpectation( + expected_interval_minutes=5, source="codex_app_automation_rrule" + ) + } + # A thread with no automation of its own inherits neither neighbour's cadence. + assert ( + codex_delivery_expectations( + [ + HostDeliveryScope( + goal_id="bound", agent_id="a", host_surface="codex-app", thread_id="t-unbound" + ) + ], + homes=[home.root], + ) + == {} + ) + + +def test_two_automations_claiming_one_lane_report_no_cadence(tmp_path: Path) -> None: + home = CodexHome(tmp_path / ".codex") + for automation_id, interval in (("first", 5), ("second", 9)): + _install_automation( + home, automation_id=automation_id, goal_id="bound", agent_id="a", + rrule=f"FREQ=MINUTELY;INTERVAL={interval}", target_thread_id="t-lane", + ) + lane = HostDeliveryScope( + goal_id="bound", agent_id="a", host_surface="codex-app", thread_id="t-lane" + ) + assert codex_delivery_expectations([lane], homes=[home.root]) == {} + + +def _attached_delivery_window( + payload: dict[str, Any], *, expectations: Any, now: datetime +) -> dict[str, Any]: + attach_host_thread_activity(payload, observers=codex_thread_observers()) + attach_host_delivery_windows(payload, expectations=expectations, now=now) + threads = payload["run_history"]["goals"][0]["host_thread_activity"]["threads"] + return threads[0]["delivery_window"] + + +def test_status_route_measures_a_lane_against_its_own_automation( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + home = CodexHome(tmp_path / ".codex") + home.thread("t-open", [_event(T2, "task_complete")]) + _install_automation( + home, automation_id="lane", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=5", target_thread_id="t-open", + ) + monkeypatch.setenv("LOOPX_CODEX_HOMES", str(home.root)) + payload = _status({"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-open"}) + + window = _attached_delivery_window( + payload, + expectations=codex_delivery_expectations, + now=datetime(2026, 9, 25, 10, 7, tzinfo=timezone.utc), + ) + + assert window["state"] == "fresh" + assert window["expected_interval_minutes"] == 5 + assert window["window_minutes"] == 10 + assert window["age_seconds"] == 60 + assert window["source"] == "codex_app_automation_rrule" + + +def test_status_route_does_not_measure_a_lane_against_another_threads_automation( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + home = CodexHome(tmp_path / ".codex") + # Recent manual activity on the lane's own thread, and one automation + # installed for another thread of the same Goal with a wide cadence. + home.thread("t-open", [_event(T2, "task_complete")]) + _install_automation( + home, automation_id="other-thread", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=60", target_thread_id="t-other", + ) + monkeypatch.setenv("LOOPX_CODEX_HOMES", str(home.root)) + payload = _status({"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-open"}) + + window = _attached_delivery_window( + payload, + expectations=codex_delivery_expectations, + now=datetime(2026, 9, 25, 10, 7, tzinfo=timezone.utc), + ) + + assert window == { + "schema_version": "loopx_host_delivery_window_v0", + "state": "unknown", + "reason": "no_expectation", + } + + +@pytest.mark.parametrize("rrule", ["FREQ=MINUTELY;INTERVAL=-3", "FREQ=MINUTELY", "FREQ=MINUTELYISH;INTERVAL=1"]) +def test_status_route_never_reports_fresh_from_an_unreadable_cadence( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, rrule: str +) -> None: + home = CodexHome(tmp_path / ".codex") + home.thread("t-open", [_event(T2, "task_complete")]) + _install_automation( + home, automation_id="lane", goal_id="bound", agent_id="a", + rrule=rrule, target_thread_id="t-open", + ) + monkeypatch.setenv("LOOPX_CODEX_HOMES", str(home.root)) + payload = _status({"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-open"}) + + window = _attached_delivery_window( + payload, + expectations=codex_delivery_expectations, + now=datetime(2026, 9, 25, 10, 7, tzinfo=timezone.utc), + ) + + assert window["state"] == "unknown" + assert window["reason"] == "no_expectation" + + +def _just_now() -> str: + """A host timestamp seconds old, in the rollout's own format. + + The status route reads the wall clock, so a healthy lane can only be + measured as fresh against activity this recent. + """ + + stamp = datetime.now(timezone.utc) - timedelta(seconds=30) + return f"{stamp:%Y-%m-%dT%H:%M:%S}.{stamp.microsecond // 1000:03d}Z" + + +def _route_response( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, payload: dict[str, Any] +) -> dict[str, Any]: + """Run the real status route once; any error response fails the test.""" + + monkeypatch.setattr(chat_status_api, "collect_status", lambda **_kwargs: payload) + sent: list[dict[str, Any]] = [] + + class Handler(chat_status_api.ChatStatusRequestMixin): + path = "/status.json" + server = SimpleNamespace( + selected_goal_id=None, + registry_path=tmp_path / "registry.json", + runtime_root_override=None, + scan_roots=[], + runtime_root=tmp_path / "runtime", + limit=10, + goal_subagent_configuration_enabled=False, + ) + + def _send_json(self, payload: dict[str, Any], *, status: int = 200) -> None: + sent.append(payload) + + def _send_error(self, message: str, **kwargs: Any) -> None: + raise AssertionError(f"the status route errored: {message} {kwargs}") + + Handler()._status() + assert len(sent) == 1, "the route must answer once with the workspace projection" + return sent[0] + + +def _lane_window(payload: dict[str, Any]) -> dict[str, Any]: + threads = payload["run_history"]["goals"][0]["host_thread_activity"]["threads"] + return threads[0]["delivery_window"] + + +@pytest.mark.parametrize("oversized_lane", ["own", "unrelated"]) +def test_an_oversized_interval_cannot_take_down_the_workspace_status( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, oversized_lane: str +) -> None: + """One unreadable manifest must not cost the whole workspace its readback.""" + + home = CodexHome(tmp_path / ".codex") + home.thread("t-open", [_event(_just_now(), "task_complete")]) + oversized = "FREQ=MINUTELY;INTERVAL=" + "1" * 4301 + if oversized_lane == "own": + _install_automation( + home, automation_id="lane", goal_id="bound", agent_id="a", + rrule=oversized, target_thread_id="t-open", + ) + else: + _install_automation( + home, automation_id="lane", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=5", target_thread_id="t-open", + ) + _install_automation( + home, automation_id="other", goal_id="unrelated", agent_id="z", + rrule=oversized, target_thread_id="t-other", + ) + monkeypatch.setenv("LOOPX_CODEX_HOMES", str(home.root)) + + window = _lane_window( + _route_response( + tmp_path, + monkeypatch, + _status({"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-open"}), + ) + ) + + if oversized_lane == "own": + # Unreadable is unknown; it is never a guessed cadence. + assert window == { + "schema_version": "loopx_host_delivery_window_v0", + "state": "unknown", + "reason": "no_expectation", + } + else: + # The unrelated manifest is skipped, so the healthy lane keeps its window. + assert window["state"] == "fresh" + assert window["expected_interval_minutes"] == 5 + + +def test_correcting_an_oversized_interval_restores_the_lane_window( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + home = CodexHome(tmp_path / ".codex") + home.thread("t-open", [_event(_just_now(), "task_complete")]) + _install_automation( + home, automation_id="lane", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=" + "1" * 4301, target_thread_id="t-open", + ) + monkeypatch.setenv("LOOPX_CODEX_HOMES", str(home.root)) + payload = _status({"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-open"}) + + assert _lane_window(_route_response(tmp_path, monkeypatch, payload))["state"] == "unknown" + + _install_automation( + home, automation_id="lane", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=5", target_thread_id="t-open", + ) + + assert _lane_window(_route_response(tmp_path, monkeypatch, payload))["state"] == "fresh" + + +def _offered_scopes(payload: dict[str, Any]) -> list[HostDeliveryScope]: + seen: list[HostDeliveryScope] = [] + + def expectations(scopes: Any) -> dict[HostDeliveryScope, HostDeliveryExpectation]: + seen.extend(scopes) + return {} + + attach_host_delivery_windows(payload, expectations=expectations) + return seen + + +def test_delivery_scope_carries_the_canonical_binding_without_serializing_it() -> None: + payload = _status({"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-open"}) + attach_host_thread_activity( + payload, + observers={"codex-app": lambda _ids: { + "t-open": HostThreadActivity(state=HostThreadState.IDLE, last_event_at=T1), + }}, + ) + + assert [ + (scope.goal_id, scope.agent_id, scope.host_surface, scope.thread_id) + for scope in _offered_scopes(payload) + ] == [("bound", "a", "codex-app", "t-open")] + # The binding stays in coordination; the projection never names the thread. + assert "t-open" not in json.dumps( + payload["run_history"]["goals"][0]["host_thread_activity"] + ) + + +def test_a_lane_bound_to_two_threads_is_offered_without_an_identity() -> None: + payload = _status( + {"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-one"}, + {"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-two"}, + ) + attach_host_thread_activity( + payload, + observers={"codex-app": lambda _ids: { + "t-one": HostThreadActivity(state=HostThreadState.IDLE, last_event_at=T1), + "t-two": HostThreadActivity(state=HostThreadState.IDLE, last_event_at=T1), + }}, + ) + + assert [scope.thread_id for scope in _offered_scopes(payload)] == [None, None] + assert all( + row["delivery_window"] == { + "schema_version": "loopx_host_delivery_window_v0", + "state": "unknown", + "reason": "no_expectation", + } + for row in payload["run_history"]["goals"][0]["host_thread_activity"]["threads"] + ) + + +def test_delivery_expectations_read_no_home_without_an_app_lane(tmp_path: Path) -> None: + absent = tmp_path / "absent-home" + + assert ( + codex_delivery_expectations( + [HostDeliveryScope(goal_id="bound", agent_id="a", host_surface="codex-cli-tui")], + homes=[absent], + ) + == {} + ) + assert not absent.exists() + + +def test_delivery_window_never_reports_absence_as_fresh() -> None: + expected = HostDeliveryExpectation(expected_interval_minutes=5, source="test") + observed = HostThreadActivity(state=HostThreadState.IDLE, last_event_at=T1) + now = datetime(2026, 9, 25, 10, 20, tzinfo=timezone.utc) + + assert build_host_delivery_window(None, expected, now=now).to_payload() == { + "schema_version": "loopx_host_delivery_window_v0", + "state": "missing", + "reason": "no_observation", + "source": "test", + } + assert build_host_delivery_window(observed, None, now=now).to_payload() == { + "schema_version": "loopx_host_delivery_window_v0", + "state": "unknown", + "reason": "no_expectation", + } + unreadable = HostThreadActivity(state=HostThreadState.IDLE, last_event_at="not-a-time") + assert ( + build_host_delivery_window(unreadable, expected, now=now).state + is HostDeliveryWindowState.UNKNOWN + ) + unobserved = HostThreadActivity.unknown(HostThreadUnknownReason.STORE_UNAVAILABLE) + assert ( + build_host_delivery_window(unobserved, expected, now=now).state + is HostDeliveryWindowState.MISSING + ) + + +def test_delivery_window_separates_fresh_from_stale_by_the_expected_cadence() -> None: + expected = HostDeliveryExpectation( + expected_interval_minutes=5, source="codex_app_automation_rrule" + ) + observed = HostThreadActivity(state=HostThreadState.IDLE, last_event_at=T1) + + fresh = build_host_delivery_window( + observed, expected, now=datetime(2026, 9, 25, 10, 11, tzinfo=timezone.utc) + ) + assert fresh.to_payload() == { + "schema_version": "loopx_host_delivery_window_v0", + "state": "fresh", + "last_observed_at": T1, + "age_seconds": 360, + "age_hours": 0.1, + "expected_interval_minutes": 5, + "window_minutes": 10, + "tolerance": 2, + "source": "codex_app_automation_rrule", + } + + stale = build_host_delivery_window( + observed, expected, now=datetime(2026, 9, 25, 10, 31, tzinfo=timezone.utc) + ) + assert stale.state is HostDeliveryWindowState.STALE + assert stale.age_seconds == 1560 + assert stale.last_observed_at == T1 + + # A later observation moves the window forward instead of the clock. + later = HostThreadActivity(state=HostThreadState.IDLE, last_event_at=T2) + assert ( + build_host_delivery_window( + later, expected, now=datetime(2026, 9, 25, 10, 11, tzinfo=timezone.utc) + ).state + is HostDeliveryWindowState.FRESH + ) + + +def test_attach_reports_the_window_from_an_installed_automation(tmp_path: Path) -> None: + home = CodexHome(tmp_path / ".codex") + home.thread("t-open", [_event(T0, "task_started"), _item(T1)]) + _install_automation( + home, automation_id="a", goal_id="bound", agent_id="a", + rrule="FREQ=MINUTELY;INTERVAL=5", + ) + + payload = _status({"agent_id": "a", "host_surface": "codex-app", "thread_id": "t-open"}) + reads: list[list[str]] = [] + + def observe(thread_ids: Any) -> dict[str, HostThreadActivity]: + reads.append(list(thread_ids)) + return observe_codex_threads(thread_ids, homes=[home.root]) + + attach_host_thread_activity(payload, observers={"codex-app": observe}) + attach_host_delivery_windows( + payload, + expectations=lambda scopes: codex_delivery_expectations(scopes, homes=[home.root]), + now=datetime(2026, 9, 25, 10, 31, tzinfo=timezone.utc), + ) + + observation = payload["run_history"]["goals"][0]["host_thread_activity"] + assert observation["threads"][0]["delivery_window"] == { + "schema_version": "loopx_host_delivery_window_v0", + "state": "stale", + "last_observed_at": T1, + "age_seconds": 1560, + "age_hours": 0.43, + "expected_interval_minutes": 5, + "window_minutes": 10, + "tolerance": 2, + "source": "codex_app_automation_rrule", + } + # The window reads the projection that already observed the host. + assert reads == [["t-open"]] + assert "t-open" not in json.dumps(observation) + + unattached = _status() + expected = json.loads(json.dumps(unattached)) + attach_host_delivery_windows( + unattached, + expectations=lambda scopes: pytest.fail("no observed lane may ask for a cadence"), + ) + assert unattached == expected