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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 22 additions & 4 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,9 @@ src/panopticon/
# checkout, mounted rw at /workspace; point origin at the forge, then init any
# submodules — in that order, since relative .gitmodules URLs resolve against
# origin); spawner.py = the spawn loop (claim an unclaimed task → spawn its
# container; prefills claude's input box with the task memo on a first spawn);
# container; prefills claude's input box with the task memo on a first spawn;
# `pause` stops a snoozed task's container + releases its claim, and the same
# snooze gate keeps spawn/reconcile/heal off it until the deadline lapses);
# prefill.py = the detached input-box prefill
# poller (mirrors cloude-cade: pipe-pane watch for ESC[?2004h → paste-buffer the
# description, unsent); daemon.py = the provision-only pull loop;
Expand Down Expand Up @@ -195,6 +197,9 @@ on every PR (the same commands the Makefile wraps).
- `tests/test_discovery.py` — workflow discovery (Slice 8): the built-in package + an optional
path are scanned for `Workflow` subclasses; a dropped-in module registers with no core change;
underscored/non-workflow files are ignored; duplicate names are rejected.
- `tests/core/test_snooze.py` — the snooze predicate both clock-owners share: an active vs lapsed
deadline, the sticky sentinel, naive datetimes, and an unreadable fact reading *inactive* rather
than raising (a bad recorded value must never stall a spawn).
- `tests/test_git.py` — local git ops: unit tests pin the emitted `git` commands and slug-gating
for `GitWorktrees` and the per-task-clone ops `GitClones` (clone/branch/set-origin, ADR 0011);
a `skipif` integration test creates a real worktree.
Expand All @@ -220,10 +225,15 @@ on every PR (the same commands the Makefile wraps).
`exit_reason` — OOM kill/exit code — including for an already-`down` task), `heal` **self-heal**
(a claimed-by-us non-terminal task whose tmux session is gone → respawn via the idempotent spawn
path; skips healthy/unclaimed/terminal tasks; the crash-loop cap — surfaced once as `failed`
with a press-R detail — + survivor-window budget reset), and the `spawnable_tasks` filter; an
integration test claims + spawns against the real task service over REST (fake git/runner).
with a press-R detail — + survivor-window budget reset), `pause` **snooze → stop** (a snoozed
task's container is stopped and its claim released — even mid-turn; no-op when unclaimed/terminal/
shell/awake/already stopped — and the matching gates that keep `spawn_one`/`reconcile`/`heal` off a
paused task within the same pass), and the `spawnable_tasks` filter; integration tests claim +
spawn against the real task service over REST (fake git/runner) and drive the whole snooze cycle
(stopped + `queued` → the deadline lapses on an injected clock → claimed + respawned).
- `tests/test_host.py` — the unified per-host daemon (ADR 0008/0011): a unit test isolates a
failing task and another pins that each pass also `heal`s every task; an integration test drives
failing task, another pins that each pass also `heal`s every task, and another that it `pause`s
every task **before** spawning it (so nothing re-spawns what a snooze just stopped); an integration test drives
spawn → set slug → provision against the real task service over REST (claimed + spawned, then
branched, no re-spawn).
- `tests/test_daemon.py` — the observe-and-provision loop + its launch: unit tests drive
Expand Down Expand Up @@ -326,6 +336,14 @@ on every PR (the same commands the Makefile wraps).
in — but, being a transition, a free move runs through an agent skill (the user directs the
agent), not the dashboard. `force_transition` is the engine primitive (e.g. going back to
coding is just `set_state(ITERATING)` — not a named operation).
- **Snooze** — `Task.snoozed_until`, an operator-owned "not now" deadline (dashboard `e` = 12h,
`E` = the sticky sentinel, a second `e` clears it). The task service stores it **verbatim** and
never compares it to a clock (the determinism invariant); the arithmetic is `core/snooze.py`
(`is_snoozed(until, now)` — pure, `now` passed in), shared by the two callers that own a clock:
the dashboard mutes the row, and the **session service stops the container** (`Spawner.pause` →
stop + release the claim, so the task reads `queued`). Waking needs no separate path — a lapsed
or cleared deadline leaves an unclaimed, un-snoozed task, which `spawn_one` claims and respawns
with its per-task clone and CLI session history intact.
- **Turn-flip / blocked** — the live `Task.turn` flips *within* a state via
`PUT /tasks/{id}/turn` (the agnostic agent↔user ball tracking). The **contract** for the
in-container hooks: the agent's stop hook sets `turn=user` (**unless a background task — a
Expand Down
49 changes: 49 additions & 0 deletions src/panopticon/core/snooze.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
"""The operator snooze predicate — pure arithmetic over a recorded deadline (no clock read).

``Task.snoozed_until`` is a recorded fact: the task service stores the operator's deadline verbatim
and never compares it to a clock (the determinism invariant). Deciding whether a deadline is *active*
is therefore the caller's job, and two callers need the same answer:

- the **dashboard**, to mute a snoozed row (its display clock), and
- the **session service**, to stop a snoozed task's container (its host clock, ADR 0008).

So the arithmetic lives here, in ``core``, with ``now`` passed **in** — LLM-free, I/O-free, and
clock-free like the rest of the package.
"""

from __future__ import annotations

from datetime import UTC, datetime

#: The reserved "sticky" deadline: a snooze that never expires, cleared only explicitly. Recorded by
#: the dashboard's `E` hotkey; defined here so the value has exactly one definition.
INDEFINITE_UNTIL = "9999-12-31T23:59:59+00:00"


def snooze_remaining(snoozed_until: str | None, now: datetime) -> float | None:
"""Active seconds remaining on ``snoozed_until`` at ``now``, else ``None``.

``float("inf")`` for :data:`INDEFINITE_UNTIL`; ``None`` when there's no deadline, when it has
already passed, or when it isn't a parseable ISO-8601 string (an unreadable fact is inactive
rather than an error — it must never take down a display or stall a spawn). Naive datetimes on
either side are read as UTC.
"""
if not isinstance(snoozed_until, str):
return None
if snoozed_until == INDEFINITE_UNTIL:
return float("inf")
try:
deadline = datetime.fromisoformat(snoozed_until)
except ValueError:
return None
if deadline.tzinfo is None:
deadline = deadline.replace(tzinfo=UTC)
if now.tzinfo is None:
now = now.replace(tzinfo=UTC)
seconds = (deadline - now).total_seconds()
return seconds if seconds > 0 else None


def is_snoozed(snoozed_until: str | None, now: datetime) -> bool:
"""Whether ``snoozed_until`` is an **active** snooze at ``now`` (see :func:`snooze_remaining`)."""
return snooze_remaining(snoozed_until, now) is not None
14 changes: 10 additions & 4 deletions src/panopticon/sessionservice/host.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,10 +69,15 @@ def __init__(
self._interval = interval

def tick(self, tasks: list[JsonObj]) -> None:
"""One pass over a task snapshot: spawn each spawnable task, provision each slugged one,
publish each pending push, reconcile each claimed one's container-lifecycle status
(down-detection), and heal each orphan (a claimed task whose tmux session is gone →
respawn). All self-gate, so re-running over an unchanged snapshot is a no-op.
"""One pass over a task snapshot: pause each snoozed task (stop its container, release the
claim), spawn each spawnable task, provision each slugged one, publish each pending push,
reconcile each claimed one's container-lifecycle status (down-detection), and heal each orphan
(a claimed task whose tmux session is gone → respawn). All self-gate, so re-running over an
unchanged snapshot is a no-op.

``pause`` runs **first**: it's what frees a snoozed task's resources, and every step after it
skips a snoozed task, so nothing in the same pass undoes it. Waking needs no step of its own —
a task whose deadline has lapsed is unclaimed and un-snoozed, so ``spawn_one`` brings it back.

``publish`` runs **before** ``cleanup``: a task can request its push and reach a terminal
state in the same breath, and cleanup deletes the per-task clone the push reads from.
Expand All @@ -89,6 +94,7 @@ def tick(self, tasks: list[JsonObj]) -> None:
_log.warning("flagging heal failed for task %s", task.get("id"), exc_info=True)
for task in tasks:
try:
self._spawner.pause(task) # snoozed → stop the container, release the claim
self._spawner.spawn_one(task)
self._provisioner.provision(task)
self._publisher.publish(task)
Expand Down
86 changes: 80 additions & 6 deletions src/panopticon/sessionservice/spawner.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,9 @@
that are **unclaimed** and **non-terminal**, **claims** one for this host (the claim is the spawn
gate — exactly one runner owns it; a lost race is a 409 we skip), prepares its writable per-task
clone (`prepare_workspace`), and spawns the container via the runner with the repo's secrets + the
``/workspace`` mount. Provisioning (slug → branch) is the sibling loop (`ProvisionDaemon`); the
``/workspace`` mount. It also **pauses** a task the operator snoozed (stop the container, release the
claim) — the deadline is a recorded fact, and reading it against a clock is the host's job, not the
control plane's. Provisioning (slug → branch) is the sibling loop (`ProvisionDaemon`); the
unified host daemon runs both. LLM-free.
"""

Expand All @@ -17,6 +19,7 @@
import subprocess
import time
from collections.abc import Callable
from datetime import UTC, datetime
from pathlib import Path

import httpx
Expand All @@ -25,11 +28,12 @@
from panopticon.core.dirs import hook_file_path
from panopticon.core.features import require_available_agent_cli
from panopticon.core.models import ContainerStatus, LifecyclePhase, resolve_agent_cli
from panopticon.core.snooze import is_snoozed
from panopticon.core.state import TERMINAL_LABELS
from panopticon.sessionservice.clones import CloneCache
from panopticon.sessionservice.executions import WorkflowExecutions
from panopticon.sessionservice.images import ImageBuilder
from panopticon.sessionservice.local_runner import LocalRunner, base_image
from panopticon.sessionservice.local_runner import LocalRunner, base_image, session_name
from panopticon.sessionservice.shell_runner import ShellRunner
from panopticon.sessionservice.spawn import cleanup_workspace, prepare_workspace

Expand Down Expand Up @@ -104,6 +108,7 @@ def __init__(
rmtree: Callable[[str], None] = shutil.rmtree,
docker_cleanup: Callable[[str], None] | None = None,
now: Callable[[], float] = time.monotonic,
utcnow: Callable[[], datetime] = lambda: datetime.now(UTC),
max_respawns: int = MAX_RESPAWNS,
respawn_reset: float = RESPAWN_RESET_SECONDS,
) -> None:
Expand Down Expand Up @@ -131,6 +136,10 @@ def __init__(
docker_cleanup if docker_cleanup is not None else runner.delete_workspace_contents
)
self._now = now
#: Wall clock for the snooze predicate (:meth:`pause`). Distinct from ``now`` (monotonic, for
#: the crash-loop window): a snooze deadline is an absolute time, so it needs a real clock.
#: Read **here**, on the host — the control plane records the deadline and never compares it.
self._utcnow = utcnow
self._max_respawns = max_respawns
self._respawn_reset = respawn_reset
#: task_id → (respawns in the current burst, monotonic time of the last respawn), the
Expand All @@ -147,9 +156,16 @@ def spawn_one(self, task: JsonObj) -> str | None:
Reports each spawn phase to the task service as it goes (``CLAIMING`` → ``PREPARING`` →
``BUILDING`` → ``STARTING`` → ``AWAITING``) so the dashboard can surface the steps to becoming
live; a step raising is reported as ``FAILED`` (with the error) before re-raising, so the
host daemon's per-task isolation still applies but the failure is visible, not silent."""
host daemon's per-task isolation still applies but the failure is visible, not silent.

An **actively snoozed** task is skipped: the operator muted it, so it must not come up (and a
task :meth:`pause` just released must not be re-spawned in the same pass). Once the deadline
lapses — or the operator clears it — this is also what wakes the task, spawning it through the
full visible lifecycle with its clone and CLI session history intact."""
if task["state"] in TERMINAL_LABELS or task.get("claimed_by"):
return None
if self._is_snoozed(task):
return None # muted by the operator — don't bring it up (see `pause`)
try:
self._client.claim(task["id"], self._runner_id)
except httpx.HTTPStatusError as exc:
Expand Down Expand Up @@ -306,6 +322,49 @@ def _report(self, task_id: str, phase: LifecyclePhase, detail: str | None = None
with contextlib.suppress(httpx.HTTPError):
self._client.report_lifecycle(task_id, self._runner_id, phase.value, detail)

def _is_snoozed(self, task: JsonObj) -> bool:
"""Whether the operator's snooze on ``task`` is **active** right now (this host's clock).

``snoozed_until`` is a recorded fact the task service never interprets (the determinism
invariant): reading it against a clock is the caller's job, and here the caller is the runner.
The arithmetic is shared with the dashboard (:mod:`panopticon.core.snooze`), so the row you see
muted and the container we stop are decided by the same rule."""
return is_snoozed(task.get("snoozed_until"), self._utcnow())

def pause(self, task: JsonObj) -> bool:
"""Stop the container of a task the operator snoozed, and release its claim. Did we stop one?

Snoozing is "not now": the task is muted on the dashboard, so its container shouldn't sit
there holding CPU, RAM and a model session. We stop it and **release the claim**, which also
clears any reported lifecycle phase, so the task composes ``queued`` — and waking it needs no
separate path: once the deadline lapses (or the operator clears it) :meth:`spawn_one` claims
and spawns it again like any queued task, with its per-task clone and the CLI session history
in its config volume intact, so the agent resumes where it left off.

Stopping is unconditional — even mid-turn. ``/workspace`` is a host bind mount, so the
checkout survives; what's lost is the in-container process (an in-flight tool call), not work
on disk.

Self-gates so calling it on every task each pass is safe: only a task **this** runner claims,
non-terminal, actively snoozed, and with a container or session still up (so a task already
paused isn't released again every pass). **Shell** tasks are skipped — their script runs once,
so killing it would be a cancel, not a pause (the same reason :meth:`heal` skips them)."""
task_id = task["id"]
if task.get("claimed_by") != self._runner_id or task["state"] in TERMINAL_LABELS:
return False
if not self._is_snoozed(task):
return False
if self._executions.is_shell(task.get("workflow")):
return False # a shell script is run once — stopping it is a cancel, not a pause
runner = self._runner_for(task)
if not (runner.is_running(task_id) or runner.has_session(task_id)):
return False # nothing left to stop — already paused
_log.info("task %s: snoozed — stopping its container and releasing the claim", task_id)
runner.stop(session_name(task_id))
self._respawns.pop(task_id, None) # a deliberate stop is not a crash — forget the burst
self._client.release(task_id) # → unclaimed, phase cleared → composes `queued`
return True

def reconcile(self, task: JsonObj) -> None:
"""Reconcile a task this runner claims into the right lifecycle status (down-detection).

Expand All @@ -322,6 +381,8 @@ def reconcile(self, task: JsonObj) -> None:
still running is left to keep coming up."""
if task.get("claimed_by") != self._runner_id:
return # not ours (or unclaimed) — spawn_one handles the unclaimed case
if self._is_snoozed(task):
return # we stopped it on purpose (`pause`) — not a death to explain
status = task.get("container_status")
if status not in _IN_PROGRESS and status != ContainerStatus.DOWN.value:
return # live / failed / queued / disconnected — nothing to reconcile
Expand All @@ -347,6 +408,8 @@ def _is_orphan(self, task: JsonObj) -> bool:
operator cancelling), not a crash to respawn — so re-running it would be wrong."""
if task.get("claimed_by") != self._runner_id or task["state"] in TERMINAL_LABELS:
return False
if self._is_snoozed(task):
return False # deliberately stopped by `pause` — a paused task is not an orphan
if self._executions.is_shell(task.get("workflow")):
return False
return not self._runner.has_session(task["id"])
Expand Down Expand Up @@ -502,12 +565,23 @@ def _compose_image(self, workflow: str, repo: JsonObj, agent_cli: str) -> str:
return self._images.build(workflow, repo["id"], layers, agent_cli=agent_cli, verbose=True)


def spawnable_tasks(client: TaskServiceClient) -> Callable[[], list[JsonObj]]:
"""This host's spawn candidates: unclaimed, non-terminal tasks (the runner claims-then-spawns).
def spawnable_tasks(
client: TaskServiceClient, *, utcnow: Callable[[], datetime] = lambda: datetime.now(UTC)
) -> Callable[[], list[JsonObj]]:
"""This host's spawn candidates: unclaimed, non-terminal, un-snoozed tasks (the runner
claims-then-spawns).

An actively snoozed task is no candidate — the operator muted it, and :meth:`Spawner.pause` stops
the container of one that's already up (the same gate :meth:`Spawner.spawn_one` applies, so the
list and the spawn can't disagree).

For M1 (single host) that's every such task the service knows; scoping to this runner's own
assignments is an M5 refinement.
"""
return lambda: [
t for t in client.list_tasks() if not t["claimed_by"] and t["state"] not in TERMINAL_LABELS
t
for t in client.list_tasks()
if not t["claimed_by"]
and t["state"] not in TERMINAL_LABELS
and not is_snoozed(t.get("snoozed_until"), utcnow())
]
Loading
Loading