From 6121a4378a0372fdbfd85a79eb68c227cc3e26d4 Mon Sep 17 00:00:00 2001 From: Dimitri Krattiger Date: Fri, 25 Sep 2026 01:52:39 -0600 Subject: [PATCH] feat(task): show when a task is parked on a third party MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A task awaiting someone else's review sits at `turn=user` and renders yellow — indistinguishable from work that is actually yours. The only way to find out was to open it, which is the round trip this removes. Of 36 live tasks, 32 were at `turn=user`; several were waiting on a reviewer. Adds `Task.waiting_on` (`WaitingOn`: external-review | ci), derived by the session service from the forge and rendered dim in the turn cell. NOT a third `Actor`, which was the obvious shape and the wrong one. `turn` is machine-driven — the container's Stop hook sets it to `user` and its UserPromptSubmit hook 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 across `turn_on_enter`, `advanced_by`, and responsibility gating (25 call sites), where a third party has no meaning: nothing external ever *advances* a task. Derived, not declared, and deliberately so. Nothing has to remember to set it 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. An agent skill that forgets to clear leaves a task looking parked forever — worse than no marker, because you learn to distrust it. Also distinct from `blocked`, which is the agent's own "I am stuck" and is cleared explicitly. Different lifecycles, so they stay separate fields; `blocked` still outranks this in the cell, since being stuck needs attention. The derivation refuses to guess. An unreadable PR (gh absent, rate-limited, unauthenticated) leaves the previous value alone rather than reporting None — marking a parked task actionable because a shell command failed is precisely the error this exists to prevent. CHANGES_REQUESTED and failing checks are *work*, not waits, so neither sets the marker. Reads are throttled to once a minute per task: the host wakes on the change feed, not a timer, so an unthrottled watcher would be one `gh` call per task per tick. Co-Authored-By: Claude Opus 5 (1M context) --- src/panopticon/client.py | 9 + src/panopticon/core/models.py | 30 +++ ...260925_270851e082e2_add_task_waiting_on.py | 35 ++++ src/panopticon/sessionservice/host.py | 8 + .../sessionservice/review_watcher.py | 163 ++++++++++++++++ src/panopticon/taskservice/api.py | 23 ++- src/panopticon/taskservice/service.py | 20 ++ .../taskservice/store_sqlalchemy.py | 14 +- src/panopticon/terminal/dashboard.py | 26 +++ tests/core/test_store.py | 11 +- tests/sessionservice/test_host.py | 5 + tests/sessionservice/test_review_watcher.py | 177 ++++++++++++++++++ tests/taskservice/test_api.py | 35 ++++ tests/terminal/test_dashboard.py | 23 +++ 14 files changed, 576 insertions(+), 3 deletions(-) create mode 100644 src/panopticon/migrations/versions/20260925_270851e082e2_add_task_waiting_on.py create mode 100644 src/panopticon/sessionservice/review_watcher.py create mode 100644 tests/sessionservice/test_review_watcher.py diff --git a/src/panopticon/client.py b/src/panopticon/client.py index 1073e8e6..b8644386 100644 --- a/src/panopticon/client.py +++ b/src/panopticon/client.py @@ -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( diff --git a/src/panopticon/core/models.py b/src/panopticon/core/models.py index 2acddff1..744200b6 100644 --- a/src/panopticon/core/models.py +++ b/src/panopticon/core/models.py @@ -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.""" @@ -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 diff --git a/src/panopticon/migrations/versions/20260925_270851e082e2_add_task_waiting_on.py b/src/panopticon/migrations/versions/20260925_270851e082e2_add_task_waiting_on.py new file mode 100644 index 00000000..7f6dd9b8 --- /dev/null +++ b/src/panopticon/migrations/versions/20260925_270851e082e2_add_task_waiting_on.py @@ -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 ### diff --git a/src/panopticon/sessionservice/host.py b/src/panopticon/sessionservice/host.py index 1de75033..e298c939 100644 --- a/src/panopticon/sessionservice/host.py +++ b/src/panopticon/sessionservice/host.py @@ -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 ( @@ -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, @@ -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 @@ -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 @@ -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, @@ -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, diff --git a/src/panopticon/sessionservice/review_watcher.py b/src/panopticon/sessionservice/review_watcher.py new file mode 100644 index 00000000..f147bdc2 --- /dev/null +++ b/src/panopticon/sessionservice/review_watcher.py @@ -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 diff --git a/src/panopticon/taskservice/api.py b/src/panopticon/taskservice/api.py index ac63d010..0efff9bb 100644 --- a/src/panopticon/taskservice/api.py +++ b/src/panopticon/taskservice/api.py @@ -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 ( @@ -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 @@ -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 @@ -349,6 +360,10 @@ class PauseIn(BaseModel): paused: bool +class WaitingOnIn(BaseModel): + waiting_on: WaitingOn | None + + class SortWeightIn(BaseModel): sort_weight: int @@ -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 diff --git a/src/panopticon/taskservice/service.py b/src/panopticon/taskservice/service.py index 721580ce..196f7f9a 100644 --- a/src/panopticon/taskservice/service.py +++ b/src/panopticon/taskservice/service.py @@ -32,6 +32,7 @@ Skill, Status, Task, + WaitingOn, compose_container_status, ) from panopticon.core.provisioning import PROVISION_SKILL @@ -791,6 +792,25 @@ async def set_snooze(self, task_id: str, until: str | None) -> Task: _log.debug("task %s: snoozed_until → %s", task_id, until) return task + async def set_waiting_on(self, task_id: str, waiting_on: WaitingOn | None) -> Task: + """Record why a task is parked on a third party, or clear it (``None``). + + A plain recorded fact, like the url or the snooze: ``state``/``turn``/``blocked`` are left + untouched. The control plane cannot *derive* this — it has no forge access by design — so + the session service observes the PR each pass and reports what it sees, the same + observed-not-pushed shape as provisioning and ask delivery. + + Idempotent by nature: the watcher reports the same value on every pass while the condition + holds, and reports ``None`` the moment it clears, so this never needs a separate clear call. + """ + task = await self.get_task(task_id) + if task.waiting_on == waiting_on: + return task # unchanged — don't churn the change feed the dashboard long-polls + task.waiting_on = waiting_on + await self._save_task(task) + _log.info("task %s: waiting_on → %s", task_id, waiting_on.value if waiting_on else None) + return task + async def set_paused(self, task_id: str, paused: bool) -> Task: """Park a task (or bring it back) so an operator can reclaim its container's memory. diff --git a/src/panopticon/taskservice/store_sqlalchemy.py b/src/panopticon/taskservice/store_sqlalchemy.py index befb25f7..bca17059 100644 --- a/src/panopticon/taskservice/store_sqlalchemy.py +++ b/src/panopticon/taskservice/store_sqlalchemy.py @@ -42,7 +42,15 @@ ) from sqlalchemy.pool import StaticPool -from panopticon.core.models import Actor, HistoryEntry, Repo, Responsibility, Status, Task +from panopticon.core.models import ( + Actor, + HistoryEntry, + Repo, + Responsibility, + Status, + Task, + WaitingOn, +) from panopticon.core.store import ( AlreadyExists, IntegrityError, @@ -131,6 +139,7 @@ class _TaskRow(_Base): url: Mapped[str | None] = mapped_column(default=None) snoozed_until: Mapped[str | None] = mapped_column(default=None) paused: Mapped[bool] = mapped_column(default=False) + waiting_on: Mapped[str | None] = mapped_column(default=None) branch: Mapped[str | None] = mapped_column(default=None) clone: Mapped[str | None] = mapped_column(default=None) claimed_by: Mapped[str | None] = mapped_column(default=None) @@ -163,6 +172,7 @@ def to_domain(self) -> Task: url=self.url, snoozed_until=self.snoozed_until, paused=self.paused, + waiting_on=WaitingOn(self.waiting_on) if self.waiting_on else None, branch=self.branch, clone=self.clone, claimed_by=self.claimed_by, @@ -192,6 +202,7 @@ def from_domain(cls, task: Task) -> _TaskRow: url=task.url, snoozed_until=task.snoozed_until, paused=task.paused, + waiting_on=task.waiting_on.value if task.waiting_on else None, branch=task.branch, clone=task.clone, claimed_by=task.claimed_by, @@ -426,6 +437,7 @@ async def _update_task(self, task: Task, stored: Sequence[HistoryEntry]) -> None row.url = task.url row.snoozed_until = task.snoozed_until row.paused = task.paused + row.waiting_on = task.waiting_on.value if task.waiting_on else None row.branch = task.branch row.clone = task.clone row.claimed_by = task.claimed_by diff --git a/src/panopticon/terminal/dashboard.py b/src/panopticon/terminal/dashboard.py index adf33a83..d5243587 100644 --- a/src/panopticon/terminal/dashboard.py +++ b/src/panopticon/terminal/dashboard.py @@ -444,11 +444,34 @@ def _snooze_refresh_delay(task: JsonObj, now: datetime) -> float | None: return max(0.05, seconds - (ceil(seconds / unit) - 1) * unit) +#: Short cell labels for :class:`~panopticon.core.models.WaitingOn`. Kept terse — this shares a +#: narrow column with the turn — and mapped explicitly rather than prettifying the enum value, so a +#: reason added later has to be given a deliberate label instead of leaking its raw wire form. +_WAITING_ON_LABELS = {"external-review": "ext review", "ci": "ci"} + + +def _waiting_on_label(task: JsonObj) -> str | None: + """The cell label for a task parked on a third party, or ``None``. + + Unknown values (a newer runner reporting a reason this dashboard predates) fall back to the raw + string rather than vanishing: a slightly ugly cell beats silently showing the task as + actionable, which is the failure this whole feature exists to prevent.""" + reason = task.get("waiting_on") + if not isinstance(reason, str) or not reason: + return None + return _WAITING_ON_LABELS.get(reason, reason) + + def _turn_cell(task: JsonObj, now: datetime | None = None) -> Text: if now is not None and (label := _snooze_label(task, now)) is not None: return Text(label, style="dim") if task.get("blocked"): return Text(f"{task['turn']} ⚠", style="red") + # Before the turn color, deliberately: a task parked on someone else's review sits at + # `turn=user` and would otherwise render yellow — indistinguishable from work that is actually + # yours. Dim is the point; it reads as "not yours" at a glance, which is the whole feature. + if (waiting := _waiting_on_label(task)) is not None: + return Text(waiting, style="dim") color = "green" if task["turn"] == "agent" else "yellow" return Text(task["turn"], style=color) @@ -528,7 +551,10 @@ def render_detail(task: JsonObj, *, time_summary: str | None = None) -> str: caller, never computed here — it means shelling out to docker to read the task's session transcripts, which is too expensive to do on every highlight/refresh (see `P` / ``action_profile``).""" + waiting = _waiting_on_label(task) turn = f"{task['turn']}{' (blocked)' if task.get('blocked') else ''}" + if waiting is not None: + turn = f"{turn} — waiting on {waiting}" claim = f" claimed: {task['claimed_by']}" if task.get("claimed_by") else "" lines = [ task.get("slug") or task["id"], diff --git a/tests/core/test_store.py b/tests/core/test_store.py index 43945cbf..1934bb6c 100644 --- a/tests/core/test_store.py +++ b/tests/core/test_store.py @@ -16,7 +16,15 @@ import pytest from sqlalchemy import inspect -from panopticon.core.models import Actor, HistoryEntry, Repo, Responsibility, Status, Task +from panopticon.core.models import ( + Actor, + HistoryEntry, + Repo, + Responsibility, + Status, + Task, + WaitingOn, +) from panopticon.core.store import ( AlreadyExists, IntegrityError, @@ -510,6 +518,7 @@ def _fully_populated_task() -> Task: url="https://github.com/acme/widgets/pull/7", snoozed_until="2026-08-06T03:00:00+00:00", paused=True, + waiting_on=WaitingOn.EXTERNAL_REVIEW, branch="panopticon/fix-the-widget", clone="/clones/t-full", claimed_by="local", diff --git a/tests/sessionservice/test_host.py b/tests/sessionservice/test_host.py index 226dc7ab..d614c5d0 100644 --- a/tests/sessionservice/test_host.py +++ b/tests/sessionservice/test_host.py @@ -220,6 +220,11 @@ def reap_paused(self, task: JsonObj) -> None: pass +class _NoOpReviewWatcher: + def observe(self, task: JsonObj) -> None: + return None + + class _NoOpProvisioner: def provision(self, task: JsonObj) -> None: return None diff --git a/tests/sessionservice/test_review_watcher.py b/tests/sessionservice/test_review_watcher.py new file mode 100644 index 00000000..66c46eb0 --- /dev/null +++ b/tests/sessionservice/test_review_watcher.py @@ -0,0 +1,177 @@ +"""Deriving `Task.waiting_on` from the forge: the gates, the derivation truth table, the throttle, +and the failure mode that matters most (an unreadable PR must not look actionable). Fake command +runner — no `gh`, no network, no LLM.""" + +from __future__ import annotations + +import json + +from panopticon.client import JsonObj +from panopticon.core.models import WaitingOn +from panopticon.sessionservice.review_watcher import ReviewWatcher + + +class _FakeClient: + def __init__(self) -> None: + self.recorded: list[tuple[str, str | None]] = [] + + def set_waiting_on(self, task_id: str, waiting_on: str | None) -> JsonObj: + self.recorded.append((task_id, waiting_on)) + return {"id": task_id} + + +def _runner(payload: object, *, calls: list[list[str]] | None = None): + """A fake `gh` that returns ``payload`` as JSON (or a raw string verbatim, to stand in for an + error message rather than JSON).""" + + def run(args, *, check: bool = True) -> str: + if calls is not None: + calls.append(list(args)) + return payload if isinstance(payload, str) else json.dumps(payload) + + return run + + +def _task(**over: object) -> JsonObj: + base: JsonObj = { + "id": "t1", + "state": "REVIEW", + "url": "https://github.com/acme/widgets/pull/7", + "waiting_on": None, + } + base.update(over) + return base + + +_OPEN_REVIEW_REQUIRED = { + "state": "OPEN", + "reviewDecision": "REVIEW_REQUIRED", + "statusCheckRollup": [], +} + + +# -- derivation --------------------------------------------------------------------- + + +def test_open_pr_awaiting_review_is_external_review() -> None: + client = _FakeClient() + w = ReviewWatcher(client, run=_runner(_OPEN_REVIEW_REQUIRED)) # type: ignore[arg-type] + assert w.observe(_task()) is WaitingOn.EXTERNAL_REVIEW + assert client.recorded == [("t1", "external-review")] + + +def test_pending_checks_are_ci() -> None: + client = _FakeClient() + pr = { + "state": "OPEN", + "reviewDecision": "APPROVED", + "statusCheckRollup": [ + {"status": "COMPLETED", "state": "SUCCESS"}, + {"status": "IN_PROGRESS"}, + ], + } + w = ReviewWatcher(client, run=_runner(pr)) # type: ignore[arg-type] + assert w.observe(_task()) is WaitingOn.CI + assert client.recorded == [("t1", "ci")] + + +def test_review_takes_precedence_over_ci() -> None: + # Both mean "not yours", but a required review outlives a check run — show the longer wait. + client = _FakeClient() + pr = { + "state": "OPEN", + "reviewDecision": "REVIEW_REQUIRED", + "statusCheckRollup": [{"status": "IN_PROGRESS"}], + } + w = ReviewWatcher(client, run=_runner(pr)) # type: ignore[arg-type] + assert w.observe(_task()) is WaitingOn.EXTERNAL_REVIEW + + +def test_changes_requested_is_not_waiting() -> None: + # The reviewer acted and handed it back. That's work, not a wait — the most important + # false-positive to avoid, since it's the case where the ball really is ours. + client = _FakeClient() + pr = {"state": "OPEN", "reviewDecision": "CHANGES_REQUESTED", "statusCheckRollup": []} + w = ReviewWatcher(client, run=_runner(pr)) # type: ignore[arg-type] + assert w.observe(_task()) is None + + +def test_failing_check_is_not_waiting_on_ci() -> None: + # A concluded-but-red check is work for us; only an unfinished one is a wait. + client = _FakeClient() + pr = { + "state": "OPEN", + "reviewDecision": "APPROVED", + "statusCheckRollup": [{"status": "COMPLETED", "conclusion": "FAILURE", "state": "FAILURE"}], + } + w = ReviewWatcher(client, run=_runner(pr)) # type: ignore[arg-type] + assert w.observe(_task()) is None + + +def test_merged_pr_clears_the_marker() -> None: + # Self-clearing is the reason this is derived rather than agent-declared. + client = _FakeClient() + pr = {"state": "MERGED", "reviewDecision": "APPROVED", "statusCheckRollup": []} + w = ReviewWatcher(client, run=_runner(pr)) # type: ignore[arg-type] + assert w.observe(_task(waiting_on="external-review")) is None + assert client.recorded == [("t1", None)] + + +# -- gates + failure modes ---------------------------------------------------------- + + +def test_unreadable_pr_leaves_the_existing_value_alone() -> None: + # The failure that would defeat the feature: a rate-limited `gh` must not silently mark a + # parked task actionable. Keep what we had; report nothing. + client = _FakeClient() + w = ReviewWatcher(client, run=_runner("gh: could not determine base repo")) # type: ignore[arg-type] + assert w.observe(_task(waiting_on="external-review")) is WaitingOn.EXTERNAL_REVIEW + assert client.recorded == [] + + +def test_skips_tasks_with_no_pr_terminal_or_paused() -> None: + calls: list[list[str]] = [] + client = _FakeClient() + w = ReviewWatcher(client, run=_runner(_OPEN_REVIEW_REQUIRED, calls=calls)) # type: ignore[arg-type] + assert w.observe(_task(url=None)) is None + assert w.observe(_task(state="COMPLETE")) is None + assert w.observe(_task(paused=True)) is None + assert calls == [] # no forge call for any of them + assert client.recorded == [] + + +def test_unchanged_value_is_not_re_recorded() -> None: + # The watcher runs every pass; re-posting the same value would churn the change feed the + # dashboard long-polls. + client = _FakeClient() + w = ReviewWatcher(client, run=_runner(_OPEN_REVIEW_REQUIRED)) # type: ignore[arg-type] + assert w.observe(_task(waiting_on="external-review")) is WaitingOn.EXTERNAL_REVIEW + assert client.recorded == [] + + +def test_throttle_limits_forge_reads() -> None: + # The host wakes on the change feed, not a timer, so an unthrottled watcher would be one `gh` + # call per task per tick. + calls: list[list[str]] = [] + clock = [1000.0] + client = _FakeClient() + w = ReviewWatcher( + client, + run=_runner(_OPEN_REVIEW_REQUIRED, calls=calls), # type: ignore[arg-type] + now=lambda: clock[0], + poll_interval=60.0, + ) + w.observe(_task()) + w.observe(_task(waiting_on="external-review")) # same tick — throttled + assert len(calls) == 1 + clock[0] += 61.0 + w.observe(_task(waiting_on="external-review")) # past the interval — reads again + assert len(calls) == 2 + + +def test_emitted_command_is_a_json_pr_read() -> None: + calls: list[list[str]] = [] + w = ReviewWatcher(_FakeClient(), run=_runner(_OPEN_REVIEW_REQUIRED, calls=calls)) # type: ignore[arg-type] + w.observe(_task()) + assert calls[0][:4] == ["gh", "pr", "view", "https://github.com/acme/widgets/pull/7"] + assert "state,reviewDecision,statusCheckRollup" in calls[0] diff --git a/tests/taskservice/test_api.py b/tests/taskservice/test_api.py index b352c601..c27ae6f8 100644 --- a/tests/taskservice/test_api.py +++ b/tests/taskservice/test_api.py @@ -345,6 +345,41 @@ def test_set_turn_and_blocked(client: TestClient) -> None: assert blocked.json()["turn"] == "user" # flip-independent: the block left the turn alone +def test_set_waiting_on_records_and_clears(client: TestClient) -> None: + task_id = _new_task(client) + before = client.get(f"/tasks/{task_id}").json() + assert before["waiting_on"] is None + + r = client.put(f"/tasks/{task_id}/waiting-on", json={"waiting_on": "external-review"}) + assert r.status_code == 200 + body = r.json() + assert body["waiting_on"] == "external-review" + # A recorded fact only — the session service derives it; the lifecycle is untouched. In + # particular `turn` must not move: a task parked on a reviewer is still the user's turn. + assert body["state"] == before["state"] + assert body["turn"] == before["turn"] + assert body["blocked"] == before["blocked"] + + cleared = client.put(f"/tasks/{task_id}/waiting-on", json={"waiting_on": None}) + assert cleared.status_code == 200 and cleared.json()["waiting_on"] is None + + +def test_waiting_on_rejects_an_unknown_reason(client: TestClient) -> None: + # A fixed set, not free text — an unrecognized reason is a bug in the reporter, not a new + # category to silently accept. + task_id = _new_task(client) + r = client.put(f"/tasks/{task_id}/waiting-on", json={"waiting_on": "lunch"}) + assert r.status_code == 422 + + +def test_waiting_on_round_trips_into_the_list_payload(client: TestClient) -> None: + # The dashboard reads it off the list, so it has to survive the store, not just the write. + task_id = _new_task(client) + client.put(f"/tasks/{task_id}/waiting-on", json={"waiting_on": "ci"}) + listed = [t for t in client.get("/tasks").json() if t["id"] == task_id] + assert listed and listed[0]["waiting_on"] == "ci" + + def test_set_paused_toggles_the_flag_without_touching_lifecycle(client: TestClient) -> None: task_id = _new_task(client) before = client.get(f"/tasks/{task_id}").json() diff --git a/tests/terminal/test_dashboard.py b/tests/terminal/test_dashboard.py index 528b16c6..b73d80d9 100644 --- a/tests/terminal/test_dashboard.py +++ b/tests/terminal/test_dashboard.py @@ -614,6 +614,29 @@ def test_snooze_label_inactive_for_past_missing_or_invalid() -> None: assert _snooze_label({"snoozed_until": "not-a-date"}, _NOW) is None +def test_turn_cell_dims_a_task_parked_on_a_third_party() -> None: + # The bug this feature fixes: a task awaiting someone else's review sits at turn=user and + # renders yellow, indistinguishable from work that is actually yours. + plain = dashboard._turn_cell({"turn": "user"}) + parked = dashboard._turn_cell({"turn": "user", "waiting_on": "external-review"}) + assert plain.plain == "user" and plain.style == "yellow" + assert parked.plain == "ext review" and parked.style == "dim" + + +def test_blocked_outranks_waiting_on() -> None: + # `blocked` is the agent saying it is stuck — that needs attention, so it must not be dimmed + # away by a third-party wait. + cell = dashboard._turn_cell({"turn": "user", "blocked": True, "waiting_on": "ci"}) + assert cell.style == "red" and "⚠" in cell.plain + + +def test_unknown_waiting_on_reason_still_dims() -> None: + # A newer runner reporting a reason this dashboard predates must not fall through to the + # actionable-looking turn color; show the raw value rather than hiding the wait. + cell = dashboard._turn_cell({"turn": "user", "waiting_on": "design-signoff"}) + assert cell.plain == "design-signoff" and cell.style == "dim" + + def test_pause_key_is_bound_exactly_once() -> None: keys = [hk.key for hk in dashboard.HOTKEYS] assert keys.count("z") == 1