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