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
163 changes: 163 additions & 0 deletions src/panopticon/sessionservice/review_watcher.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,163 @@
"""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"})


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] = {}

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)
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,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

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

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. Two states that look like waiting but aren't:

* ``CHANGES_REQUESTED`` — the reviewer has acted and handed it *back*; that is work.
* 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":
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
23 changes: 22 additions & 1 deletion src/panopticon/taskservice/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,16 @@
from pydantic import BaseModel, ConfigDict, Field

from panopticon.core.artifacts import ArtifactError
from panopticon.core.models import Actor, Ask, AskStatus, LifecyclePhase, Repo, Status, Task
from panopticon.core.models import (
Actor,
Ask,
AskStatus,
LifecyclePhase,
Repo,
Status,
Task,
WaitingOn,
)
from panopticon.core.store import AlreadyExists, NotFound, StoreError
from panopticon.core.workflow import IllegalTransition, InvalidWorkflow, ResponsibilitiesNotMet
from panopticon.taskservice.service import (
Expand Down Expand Up @@ -90,6 +99,7 @@ class TaskSummaryOut(BaseModel):
url: str | None
snoozed_until: str | None = None
paused: bool = False # operator parked it: container reaped, session + workspace kept
waiting_on: WaitingOn | None = None # parked on a third party (derived from the forge)
branch: str | None
clone: str | None
claimed_by: str | None
Expand Down Expand Up @@ -131,6 +141,7 @@ class TaskOut(BaseModel):
None # operator-owned attention mute deadline (ISO-8601); None = not snoozed
)
paused: bool = False # operator parked it: container reaped, session + workspace kept
waiting_on: WaitingOn | None = None # parked on a third party (derived from the forge)
branch: str | None
clone: str | None
claimed_by: str | None # the runner that owns this task (the spawn gate), or None
Expand Down Expand Up @@ -349,6 +360,10 @@ class PauseIn(BaseModel):
paused: bool


class WaitingOnIn(BaseModel):
waiting_on: WaitingOn | None


class SortWeightIn(BaseModel):
sort_weight: int

Expand Down Expand Up @@ -742,6 +757,12 @@ async def set_turn(task_id: str, body: TurnIn) -> TaskOut:
async def set_blocked(task_id: str, body: BlockedIn) -> TaskOut:
return _task_out(await service.set_blocked(task_id, body.blocked))

@app.put("/tasks/{task_id}/waiting-on")
async def set_waiting_on(task_id: str, body: WaitingOnIn) -> TaskOut:
"""Record (or clear) why a task is parked on a third party. Reported by the session
service, which reads the forge; the control plane has no forge access of its own."""
return _task_out(await service.set_waiting_on(task_id, body.waiting_on))

@app.put("/tasks/{task_id}/pause")
async def set_paused(task_id: str, body: PauseIn) -> TaskOut:
"""Park (or restore) a task. Records the fact only — the session service reaps the
Expand Down
Loading
Loading