From e4658d8e1cf5eff07b57fc8ca7eee75df319e0e2 Mon Sep 17 00:00:00 2001 From: GZY-SUPER-HACKER <162807803+GZY-SUPER-HACKER@users.noreply.github.com> Date: Sat, 10 Oct 2026 21:11:58 +0800 Subject: [PATCH 1/3] feat(scheduler): report a host delivery window for bound app lanes A bound Codex App lane can stop receiving scheduled turns while its goal keeps reading as runnable. The installed automation and its ACK are configuration evidence, and every receipt check loopX holds is keyed on a turn identity the host produces, so none of them answer "has this lane moved lately". Add a read-only delivery window to the existing status projection. It compares a lane's own last observed host activity against the cadence its installed automation will actually fire at, and reports fresh / stale / missing / unknown with the observation age. Absence of evidence is never fresh: an unobserved lane is missing, and a lane whose cadence cannot be read is unknown. The comparison is host-agnostic; only the expectation provider is Codex-specific, and it reads the automation manifest rather than the scheduler runtime so that a status request keeps no dependency on that runtime. Related to #3927. Does not close it and does not name a root cause. Signed-off-by: GZY-SUPER-HACKER <162807803+GZY-SUPER-HACKER@users.noreply.github.com> --- loopx/chat_status_api.py | 14 +- loopx/codex_app_thread_activity.py | 84 ++++++ .../agents/host_thread_activity.py | 250 ++++++++++++++++++ tests/test_host_thread_activity.py | 223 +++++++++++++++- 4 files changed, 567 insertions(+), 4 deletions(-) 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..2a3f802a8c 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$") @@ -293,3 +300,80 @@ 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 lives in the TypeScript scheduler + (``scheduler/state_store.ts``); this reads the same ``INTERVAL=`` field from + the installed automation so that a status projection needs no scheduler + runtime. An rrule this adapter cannot read reports no interval, never one. + """ + + text = str(rrule or "").strip().upper() + if "FREQ=MINUTELY" not in text: + return None + match = re.search(r"\bINTERVAL=(\d+)\b", text) + if match is None: + # A minutely rrule without INTERVAL fires every minute. + return 1 + interval = int(match.group(1)) + return interval if interval > 0 else None + + +def _installed_automation_intervals( + homes: list[Path] | None = None, +) -> dict[tuple[str, str], int]: + """Map each installed app automation to the interval it will fire at.""" + + from .upgrade import infer_agent_id_from_prompt, infer_goal_id_from_prompt + + intervals: dict[tuple[str, str], int] = {} + 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 + 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) + interval = _automation_rrule_interval_minutes(item.get("rrule")) + if not goal_id or not agent_id or interval is None: + continue + intervals[(goal_id, agent_id)] = 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 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 + interval = intervals.get((scope.goal_id, scope.agent_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..b7fc3ea5a3 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,251 @@ 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.""" + + goal_id: str + agent_id: str + host_surface: str + + +@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 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. 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[dict[str, Any], list[dict[str, Any]]]] = [] + 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() + rows_by_goal.append((goal, rows)) + for row in rows: + agent_id = str(row.get("agent_id") or "").strip() + host_surface = str(row.get("host_surface") or "").strip() + if goal_id and agent_id and host_surface: + scopes.append( + HostDeliveryScope( + goal_id=goal_id, + agent_id=agent_id, + host_surface=host_surface, + ) + ) + resolved = expectations(scopes) if scopes else {} + for goal, rows in rows_by_goal: + goal_id = str(goal.get("id") or "").strip() + for row in rows: + scope = HostDeliveryScope( + goal_id=goal_id, + agent_id=str(row.get("agent_id") or "").strip(), + host_surface=str(row.get("host_surface") or "").strip(), + ) + row["delivery_window"] = build_host_delivery_window( + _activity_from_row(row), + resolved.get(scope), + 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/tests/test_host_thread_activity.py b/tests/test_host_thread_activity.py index 2c112b6a38..8b8a16a349 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, 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,216 @@ 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("FREQ=MINUTELY") == 1 + # 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 + + +def _install_automation( + home: CodexHome, + *, + automation_id: str, + goal_id: str, + agent_id: str, + rrule: str, + status: str = "ACTIVE", +) -> 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 = "heartbeat"\n' + f'status = "{status}"\nrrule = "{rrule}"\ntarget_thread_id = "t-open"\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", + ) + + resolved = codex_delivery_expectations( + [ + HostDeliveryScope(goal_id="bound", agent_id="a", host_surface="codex-app"), + HostDeliveryScope(goal_id="bound", agent_id="b", host_surface="codex-app"), + HostDeliveryScope(goal_id="bound", agent_id="c", host_surface="codex-app"), + # The app automation is not the loop that drives another surface. + HostDeliveryScope(goal_id="bound", agent_id="a", host_surface="codex-cli-tui"), + ], + homes=[home.root], + ) + + assert resolved == { + HostDeliveryScope( + goal_id="bound", agent_id="a", host_surface="codex-app" + ): HostDeliveryExpectation( + expected_interval_minutes=30, source="codex_app_automation_rrule" + ) + } + + +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 From 19b92f1008eb87869701f3b5237e0b13c84768f2 Mon Sep 17 00:00:00 2001 From: GZY-SUPER-HACKER <162807803+GZY-SUPER-HACKER@users.noreply.github.com> Date: Sat, 10 Oct 2026 22:55:12 +0800 Subject: [PATCH 2/3] fix(scheduler): keep each delivery window on its own bound lane Address the review on the delivery window. A lane was matched to an installed automation by inferred Goal/agent text alone, so an automation bound to another thread -- or one that is not a heartbeat at all -- could be read as this lane's cadence, and the window then reported fresh from an unrelated source. Place a lane by its canonical binding instead: Goal, agent and the bound target_thread_id, the way resolve_codex_app_automation_rrule places one, and require kind = "heartbeat". Two automations claiming one lane resolve to neither. A lane the Goal binds to more than one thread has no single identity and is offered without one, so its window stays unknown. The binding is joined from the Goal's coordination.thread_agent_bindings, and the thread id still never reaches the window payload. Read INTERVAL through the canonical rrule contract rather than a second, looser dialect of it, so an rrule this adapter cannot read reports no interval instead of a guessed one minute. Regenerate the project registry I/O manifest for the two load_registry sites the added status imports moved. Related to #3927. Does not close it and does not name a root cause. Signed-off-by: GZY-SUPER-HACKER <162807803+GZY-SUPER-HACKER@users.noreply.github.com> --- loopx/codex_app_thread_activity.py | 86 +++++-- .../agents/host_thread_activity.py | 74 ++++-- .../project_registry_io_manifest_v1.json | 4 +- tests/test_host_thread_activity.py | 240 +++++++++++++++++- 4 files changed, 350 insertions(+), 54 deletions(-) diff --git a/loopx/codex_app_thread_activity.py b/loopx/codex_app_thread_activity.py index 2a3f802a8c..053ae322e5 100644 --- a/loopx/codex_app_thread_activity.py +++ b/loopx/codex_app_thread_activity.py @@ -44,6 +44,10 @@ _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 + _rollout_cache: dict[tuple[str, int, int], HostThreadActivity] = {} @@ -305,31 +309,57 @@ def observe(thread_ids: Iterable[str]) -> dict[str, HostThreadActivity]: def _automation_rrule_interval_minutes(rrule: Any) -> int | None: """Read the interval a minutely automation rrule asks for. - The canonical parse lives in the TypeScript scheduler - (``scheduler/state_store.ts``); this reads the same ``INTERVAL=`` field from - the installed automation so that a status projection needs no scheduler - runtime. An rrule this adapter cannot read reports no interval, never one. + 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().upper() - if "FREQ=MINUTELY" not in text: + 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 + interval = int(raw_interval) + if interval > _MAX_SAFE_INTEGER or interval < -_MAX_SAFE_INTEGER: return None - match = re.search(r"\bINTERVAL=(\d+)\b", text) - if match is None: - # A minutely rrule without INTERVAL fires every minute. - return 1 - interval = int(match.group(1)) return interval if interval > 0 else None def _installed_automation_intervals( homes: list[Path] | None = None, -) -> dict[tuple[str, str], int]: - """Map each installed app automation to the interval it will fire at.""" +) -> 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], int] = {} + 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: @@ -338,15 +368,27 @@ def _installed_automation_intervals( 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 interval is None: + if not goal_id or not agent_id or not target_thread_id or interval is None: continue - intervals[(goal_id, agent_id)] = interval + 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 @@ -357,8 +399,10 @@ def codex_delivery_expectations( ) -> dict[HostDeliveryScope, HostDeliveryExpectation]: """Resolve each bound app lane's cadence from its installed automation. - A lane the installed automations cannot place is omitted, so its window is - reported as ``unknown`` rather than measured against a guessed cadence. + 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) @@ -369,7 +413,11 @@ def codex_delivery_expectations( for scope in wanted: if scope.host_surface != CODEX_APP_BINDING_SURFACE: continue - interval = intervals.get((scope.goal_id, scope.agent_id)) + 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( diff --git a/loopx/control_plane/agents/host_thread_activity.py b/loopx/control_plane/agents/host_thread_activity.py index b7fc3ea5a3..5a9b62f1cc 100644 --- a/loopx/control_plane/agents/host_thread_activity.py +++ b/loopx/control_plane/agents/host_thread_activity.py @@ -188,11 +188,17 @@ def __post_init__(self) -> None: @dataclass(frozen=True) class HostDeliveryScope: - """One bound lane a delivery window can be reported for.""" + """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) @@ -308,6 +314,35 @@ def build_host_delivery_window( ] +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], *, @@ -318,15 +353,19 @@ def attach_host_delivery_windows( """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. A lane the expectation provider - cannot resolve is reported as ``unknown``, never as healthy. + 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[dict[str, Any], list[dict[str, Any]]]] = [] + 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): @@ -339,30 +378,19 @@ def attach_host_delivery_windows( if not rows: continue goal_id = str(goal.get("id") or "").strip() - rows_by_goal.append((goal, rows)) + bound = _bound_thread_ids(goal) + rows_by_goal.append((goal_id, rows, bound)) for row in rows: - agent_id = str(row.get("agent_id") or "").strip() - host_surface = str(row.get("host_surface") or "").strip() - if goal_id and agent_id and host_surface: - scopes.append( - HostDeliveryScope( - goal_id=goal_id, - agent_id=agent_id, - host_surface=host_surface, - ) - ) + scope = _scope_from_row(goal_id, row, bound) + if scope is not None: + scopes.append(scope) resolved = expectations(scopes) if scopes else {} - for goal, rows in rows_by_goal: - goal_id = str(goal.get("id") or "").strip() + for goal_id, rows, bound in rows_by_goal: for row in rows: - scope = HostDeliveryScope( - goal_id=goal_id, - agent_id=str(row.get("agent_id") or "").strip(), - host_surface=str(row.get("host_surface") or "").strip(), - ) + scope = _scope_from_row(goal_id, row, bound) row["delivery_window"] = build_host_delivery_window( _activity_from_row(row), - resolved.get(scope), + resolved.get(scope) if scope is not None else None, now=now, tolerance=tolerance, ).to_payload() 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 8b8a16a349..ea9f05d61b 100644 --- a/tests/test_host_thread_activity.py +++ b/tests/test_host_thread_activity.py @@ -394,13 +394,36 @@ def test_remote_codex_surfaces_have_no_local_observer() -> None: 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("FREQ=MINUTELY") == 1 + 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 # 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 + + def _install_automation( home: CodexHome, *, @@ -409,12 +432,14 @@ def _install_automation( 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 = "heartbeat"\n' - f'status = "{status}"\nrrule = "{rrule}"\ntarget_thread_id = "t-open"\n' + 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", ) @@ -437,26 +462,221 @@ def test_delivery_expectations_resolve_only_installed_active_app_automations( 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( [ - HostDeliveryScope(goal_id="bound", agent_id="a", host_surface="codex-app"), - HostDeliveryScope(goal_id="bound", agent_id="b", host_surface="codex-app"), - HostDeliveryScope(goal_id="bound", agent_id="c", host_surface="codex-app"), + 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"), + 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 == { - HostDeliveryScope( - goal_id="bound", agent_id="a", host_surface="codex-app" - ): HostDeliveryExpectation( + 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 _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" From b6c51de241f182235776a653d0b99c41b652d263 Mon Sep 17 00:00:00 2001 From: GZY-SUPER-HACKER <162807803+GZY-SUPER-HACKER@users.noreply.github.com> Date: Sat, 10 Oct 2026 23:37:18 +0800 Subject: [PATCH 3/3] fix(codex-app): bound the interval conversion before the interpreter refuses it Address the second review on the delivery window. INTERVAL was converted with int() ahead of the safe-integer guard, so a manifest whose interval runs past Python's default decimal digit limit raised ValueError inside _installed_automation_intervals. That call sits outside the per-manifest TOML boundary, and the status route answers any local projection failure with one 500 for the whole workspace, so a single oversized manifest -- in this lane or in an unrelated Goal -- cost every healthy lane its readback. Rule the value out by its significant digits before converting it. Leading zeros are stripped first because the interpreter counts them too, so the length check now bounds exactly the digits int() is given. An out-of-range or unreadable interval still reports no expectation, so its lane reads unknown. Related to #3927. Does not close it and does not name a root cause. Signed-off-by: GZY-SUPER-HACKER <162807803+GZY-SUPER-HACKER@users.noreply.github.com> --- loopx/codex_app_thread_activity.py | 16 +++- tests/test_host_thread_activity.py | 149 ++++++++++++++++++++++++++++- 2 files changed, 163 insertions(+), 2 deletions(-) diff --git a/loopx/codex_app_thread_activity.py b/loopx/codex_app_thread_activity.py index 053ae322e5..746783dde5 100644 --- a/loopx/codex_app_thread_activity.py +++ b/loopx/codex_app_thread_activity.py @@ -47,6 +47,7 @@ # 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] = {} @@ -337,7 +338,20 @@ def _automation_rrule_interval_minutes(rrule: Any) -> int | None: raw_interval = parts.get("INTERVAL", "") if re.fullmatch(r"[+-]?[0-9]+", raw_interval) is None: return None - interval = int(raw_interval) + # 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 diff --git a/tests/test_host_thread_activity.py b/tests/test_host_thread_activity.py index ea9f05d61b..9f75eae026 100644 --- a/tests/test_host_thread_activity.py +++ b/tests/test_host_thread_activity.py @@ -2,7 +2,7 @@ import json import sqlite3 -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from pathlib import Path from types import SimpleNamespace from typing import Any @@ -399,6 +399,13 @@ def test_automation_rrule_interval_reads_only_minutely_intervals() -> None: 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 @@ -424,6 +431,31 @@ def test_an_unreadable_automation_rrule_never_reports_a_cadence(rrule: str) -> N 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, *, @@ -623,6 +655,121 @@ def test_status_route_never_reports_fresh_from_an_unreadable_cadence( 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] = []