Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 12 additions & 2 deletions loopx/chat_status_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Comment thread
loopx-agent marked this conversation as resolved.
attach_host_thread_activity,
)
from .control_plane.effect_runtime import (
EffectRuntimePermanentIOError,
EffectRuntimeRemoteError,
Expand Down Expand Up @@ -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)
Expand Down
146 changes: 146 additions & 0 deletions loopx/codex_app_thread_activity.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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$")
Expand All @@ -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] = {}


Expand Down Expand Up @@ -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
Loading
Loading