Skip to content
Open
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
9 changes: 9 additions & 0 deletions src/panopticon/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -232,6 +232,15 @@ def set_snooze(self, task_id: str, until: str | None) -> JsonObj:
self._json(self._http.put(f"/tasks/{task_id}/snooze", json={"until": until})),
)

def set_waiting_on(self, task_id: str, waiting_on: str | None) -> JsonObj:
"""Record (or clear, with ``None``) why the task is parked on a third party."""
return cast(
JsonObj,
self._json(
self._http.put(f"/tasks/{task_id}/waiting-on", json={"waiting_on": waiting_on})
),
)

def set_paused(self, task_id: str, paused: bool) -> JsonObj:
"""Park the task (reaping its container, keeping its session) or bring it back."""
return cast(
Expand Down
30 changes: 30 additions & 0 deletions src/panopticon/core/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,30 @@ class Actor(str, Enum):
AGENT = "agent"


class WaitingOn(str, Enum):
"""Why a task is parked on a party that is **neither** the user nor the agent.

Deliberately *not* a third :class:`Actor`. ``turn`` is machine-driven — the container's Stop
hook sets it to ``user`` and its UserPromptSubmit hook sets it to ``agent`` on every turn
boundary — so a third turn value would be clobbered the moment the agent did anything. And
``Actor`` is load-bearing in the state machine (``turn_on_enter``, ``advanced_by``,
responsibility gating), where a third party has no meaning: nothing external ever *advances* a
task. This is the separate axis the turn can't carry.

Also distinct from :attr:`Task.blocked`, which is the **agent's** own declaration that it is
stuck and is cleared explicitly. This is **derived** by the session service from the forge and
clears itself when the underlying condition does, so the two never need reconciling.

The point is triage: a task at ``turn=user`` looks actionable, and opening it only to find it's
parked on someone else's review is the cost this removes.
"""

#: The PR is open and the forge says a review is still required. Nobody here can move it.
EXTERNAL_REVIEW = "external-review"
#: Checks are still running. Transient, but not actionable while it lasts.
CI = "ci"


class Status(str, Enum):
"""Resolution status of a single responsibility."""

Expand Down Expand Up @@ -268,6 +292,12 @@ class Task:
#: A deliberate "waiting on something" marker the agent sets; it is **orthogonal to the
#: turn** and survives turn flips (cloude-cade's `:blocked:`), cleared only explicitly.
blocked: bool = False
#: Why this task is parked on a third party (see :class:`WaitingOn`), or ``None`` when it isn't.
#: **Derived**, not declared: the session service reads the forge each pass and records what it
#: finds, so it clears itself when the PR is approved or the checks go green. The control plane
#: never computes it — it has no forge access and stays LLM-free and network-free by design.
#: Orthogonal to ``turn`` and to ``blocked``, and touched by neither.
waiting_on: WaitingOn | None = None
#: A brief, one-line reminder of what the task is, collected when the task is created (shown
#: in the dashboard's task summary) — a human label of *intent*, not a full description (that
#: lives in the task's plan artifact). Distinct from the ``slug`` (a short identifier the
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
"""add task waiting_on

Revision ID: 270851e082e2
Revises: af3f4b8c8271
Create Date: 2026-09-25 01:46:13.146962
"""

from __future__ import annotations

from collections.abc import Sequence

import sqlalchemy as sa
from alembic import op

# revision identifiers, used by Alembic.
revision: str = "270851e082e2"
down_revision: str | None = "af3f4b8c8271"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None


def upgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
with op.batch_alter_table("task", schema=None) as batch_op:
batch_op.add_column(sa.Column("waiting_on", sa.String(), nullable=True))

# ### end Alembic commands ###


def downgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
with op.batch_alter_table("task", schema=None) as batch_op:
batch_op.drop_column("waiting_on")

# ### end Alembic commands ###
8 changes: 8 additions & 0 deletions src/panopticon/sessionservice/host.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
from panopticon.sessionservice.images import ImageBuilder
from panopticon.sessionservice.local_runner import DEFAULT_IMAGE, LocalRunner
from panopticon.sessionservice.provisioner import Provisioner
from panopticon.sessionservice.review_watcher import ReviewWatcher
from panopticon.sessionservice.shell_runner import ShellRunner
from panopticon.sessionservice.spawner import Spawner
from panopticon.sessionservice.stall import (
Expand All @@ -65,6 +66,7 @@ def __init__(
provisioner: Provisioner,
ask_worker: AskWorker | None = None,
*,
review_watcher: ReviewWatcher | None = None,
stall_monitor: StallMonitor | None = None,
sleep: Callable[[float], None] = time.sleep,
interval: float = 2.0,
Expand All @@ -73,6 +75,7 @@ def __init__(
self._spawner = spawner
self._provisioner = provisioner
self._ask_worker = ask_worker
self._review_watcher = review_watcher
self._stall_monitor = stall_monitor
self._sleep = sleep
self._interval = interval
Expand Down Expand Up @@ -113,6 +116,9 @@ def tick(self, tasks: list[JsonObj]) -> None:
# After heal: both self-gate, but reaping first would stop a container that heal
# then sees sessionless. It skips paused tasks, so the order is belt-and-braces.
self._spawner.reap_paused(task)
if self._review_watcher is not None:
# Read-only and throttled: no container work, just the forge → triage marker.
self._review_watcher.observe(task)
if self._stall_monitor is not None:
self._stall_monitor.tick(task)
except Exception: # a transient git/REST/FS error on one task must not stall the others
Expand Down Expand Up @@ -224,6 +230,7 @@ def run_host(
)
provisioner = Provisioner(client, clones_root=tasks_root, git=git, executions=executions)
ask_worker = AskWorker(client, runner, spawner, runner_id=runner_id)
review_watcher = ReviewWatcher(client)
stall_monitor = StallMonitor(
client,
runner,
Expand All @@ -240,6 +247,7 @@ def run_host(
spawner,
provisioner,
ask_worker=ask_worker,
review_watcher=review_watcher,
stall_monitor=stall_monitor,
interval=interval,
sleep=sleep,
Expand Down
217 changes: 217 additions & 0 deletions src/panopticon/sessionservice/review_watcher.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,217 @@
"""Derive ``Task.waiting_on`` from the forge (ADR 0008's observe-and-record shape).

The sibling of :class:`~panopticon.sessionservice.provisioner.Provisioner` and
:class:`~panopticon.sessionservice.ask_worker.AskWorker` for the triage side. The task service
records *why* a task is parked on a third party but cannot work it out — it has no forge access and
stays network-free by design. The session service runs where ``gh`` is authenticated, so it owns
the derivation: each pass it reads the PR of any task that has one and reports what it finds.

**Derived, not declared.** Nothing has to remember to set this and nothing has to remember to clear
it: when the PR is approved or the checks go green, the next pass reports ``None`` and the marker
disappears on its own. That is the whole reason this is a watcher rather than an agent skill — a
skill that forgets to clear leaves a task looking parked forever, which is worse than no marker at
all, because the operator learns to distrust it.

LLM-free. ``run`` is injectable so the emitted ``gh`` commands are unit-testable without a forge.
"""

from __future__ import annotations

import json
import logging
import time
from collections.abc import Callable, Sequence

from panopticon.client import JsonObj, TaskServiceClient
from panopticon.core.git import CommandRunner, _subprocess_run
from panopticon.core.models import WaitingOn
from panopticon.core.state import TERMINAL_LABELS

_log = logging.getLogger(__name__)

#: Seconds before a task's PR is re-read. The host daemon wakes on the change feed, not a timer, so
#: a busy fleet can tick many times a second — without this, each tick would be one `gh` call per
#: task with a PR. Review state changes on human timescales; a minute of staleness costs nothing.
POLL_INTERVAL_SECONDS = 60.0

#: Check states that mean CI hasn't finished. Anything else (SUCCESS, FAILURE, …) has concluded —
#: a *failing* check is not "waiting on CI", it's work for whoever owns the task.
_PENDING_CHECK_STATES = frozenset({"PENDING", "QUEUED", "IN_PROGRESS", "WAITING", "REQUESTED"})


def _pending_from_others(review_requests: object, me: str | None) -> bool:
"""Whether a review is pending from someone other than ``me``.

``reviewRequests`` holds users (``login``) and teams (``name``/``slug``) the review was asked
of. Empty means nobody was asked — a protection rule wanting *a* review, which is ours to give.
A team entry has no login to compare, and is treated as external by construction.
"""
if not isinstance(review_requests, list) or not review_requests:
return False
for req in review_requests:
if not isinstance(req, dict):
continue
login = req.get("login")
if login is None:
return True # a team (or an unfamiliar shape) — we are never a team, so: somebody else
if me is not None and str(login).lower() != me.lower():
return True # a named someone who demonstrably isn't us
# Either the request is ours, or we couldn't establish an identity to compare against. Both
# fall through to "ours": the asymmetry matters, because a wrong "external" *hides* work.
return False


class ReviewWatcher:
"""Reads each task's PR and records why it's parked, or that it no longer is."""

def __init__(
self,
client: TaskServiceClient,
*,
run: CommandRunner = _subprocess_run,
now: Callable[[], float] = time.monotonic,
poll_interval: float = POLL_INTERVAL_SECONDS,
) -> None:
self._client = client
self._run = run
self._now = now
self._poll_interval = poll_interval
#: task id → monotonic time of its last successful read, for the throttle above.
self._last_polled: dict[str, float] = {}
#: The authenticated forge login, resolved once on first use (it can't change under a
#: running daemon). ``False`` records a failed lookup so we don't retry it every pass.
self._me: str | None | bool = None

def observe(self, task: JsonObj) -> WaitingOn | None:
"""Record why ``task`` is parked on a third party, returning what was recorded.

Self-gating, so the host daemon can call it on every task each pass: a task with no PR, a
terminal one, or one polled within :data:`POLL_INTERVAL_SECONDS` is skipped and returns the
value already on the task.

A read failure is **not** treated as "nothing to wait on" — the previously recorded value is
left alone. Reporting ``None`` because ``gh`` was rate-limited would quietly mark a parked
task as actionable, which is the exact error this feature exists to prevent.
"""
task_id = task["id"]
current = task.get("waiting_on")
if not task.get("url") or task["state"] in TERMINAL_LABELS:
return self._as_enum(current)
if task.get("paused"):
return self._as_enum(current) # no container, no triage value — don't spend the call
last = self._last_polled.get(task_id)
if last is not None and self._now() - last < self._poll_interval:
return self._as_enum(current)

pr = self._read_pr(str(task["url"]))
if pr is None:
return self._as_enum(current) # unreadable — keep what we had, don't guess
self._last_polled[task_id] = self._now()

waiting_on = self._derive(pr, self._viewer())
if waiting_on != self._as_enum(current):
self._client.set_waiting_on(task_id, waiting_on.value if waiting_on else None)
_log.info(
"task %s: waiting_on %s → %s",
task_id,
current,
waiting_on.value if waiting_on else None,
)
return waiting_on

@staticmethod
def _as_enum(value: object) -> WaitingOn | None:
"""The task's recorded value as an enum. Tolerant: an unrecognized string (a newer runner
wrote a reason this one doesn't know) reads as ``None`` rather than raising mid-pass."""
if not isinstance(value, str):
return None
try:
return WaitingOn(value)
except ValueError:
return None

def _read_pr(self, url: str) -> JsonObj | None:
"""The PR's review + check state, or ``None`` if it can't be read.

Every failure mode lands here and is swallowed deliberately: `gh` absent, unauthenticated,
rate-limited, the URL not being a PR, a network blip. None of those should stall a host
pass, and none of them are evidence about the PR.
"""
try:
out = self._run(
[
"gh",
"pr",
"view",
url,
"--json",
"state,reviewDecision,reviewRequests,statusCheckRollup",
],
check=False,
)
except Exception:
_log.debug("gh pr view failed for %s", url, exc_info=True)
return None
try:
parsed = json.loads(out)
except (ValueError, TypeError):
return None # `gh` printed an error rather than JSON (not a PR, no auth, …)
return parsed if isinstance(parsed, dict) else None

def _viewer(self) -> str | None:
"""The authenticated forge login, or ``None`` if it can't be determined.

Cached for the daemon's life. On failure we return ``None`` and remember that, which
makes :meth:`_derive` conservative: with no identity to compare against, a pending
review is treated as *ours* rather than guessed to be someone else's."""
if self._me is None:
try:
out = self._run(["gh", "api", "user", "--jq", ".login"], check=False)
self._me = out.strip() or False
except Exception:
_log.debug("gh api user failed; treating reviews as ours", exc_info=True)
self._me = False
return self._me if isinstance(self._me, str) else None

@staticmethod
def _derive(pr: JsonObj, me: str | None) -> WaitingOn | None:
"""Fold the PR's state into a reason, or ``None`` when the ball is ours.

``reviewDecision == "REVIEW_REQUIRED"`` is **not** on its own evidence of an external wait:
on a branch with a protection rule it only means *a* review is required and none has been
given. Every open task PR here reads that way, almost always with nobody requested — which
makes the pending review **ours**. Labelling those external would dim exactly the work the
operator most needs to see, so the question asked here is narrower: is a review pending from
someone *other than us*?

A requested **team** counts as external: somebody on it owes the review, and it isn't
specifically us. With no identity available (``me is None``) we stay conservative and treat
a pending review as ours — under-marking costs a glance, over-marking hides real work.

Precedence is review-before-CI, deliberately. Both mean "not yours", but a required review
sits for days while checks resolve in minutes, so the review is the more useful thing to
show. Three states that look like waiting but aren't:

* ``CHANGES_REQUESTED`` — the reviewer has acted and handed it *back*; that is work.
* ``REVIEW_REQUIRED`` with nobody (or only us) requested — that is work, and it's ours.
* a failing (not pending) check — also work, for the same reason.
"""
if pr.get("state") != "OPEN":
return None # merged or closed — nothing left to wait for
if pr.get("reviewDecision") == "REVIEW_REQUIRED" and _pending_from_others(
pr.get("reviewRequests"), me
):
return WaitingOn.EXTERNAL_REVIEW
rollup = pr.get("statusCheckRollup")
if isinstance(rollup, Sequence) and not isinstance(rollup, str | bytes):
for check in rollup:
if not isinstance(check, dict):
continue
# `status` is the workflow-run field; `state` the commit-status one. A rollup mixes
# both kinds, so a check is pending if *either* says so.
if (
str(check.get("status") or "").upper() in _PENDING_CHECK_STATES
or str(check.get("state") or "").upper() in _PENDING_CHECK_STATES
):
return WaitingOn.CI
return None
Loading
Loading