This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/agent-team/agent_team/draft_pr_monitor.py
Adam Moussa cb84629c2b feat(agent-team): P3 Phases A/B/E — safety tooling, wiring, docs
Phase A (safety):
- scripts/p3_rollback.sh (+test): restore all privileged P3 surfaces from a
  recorded baseline; --dry-run default, --apply gated. Correct App-uninstall
  (App JWT) model; per-task env-reviewer restore by numeric id; real
  protection post-restore assert (normalize reads argv, fails loud, divergent
  state exits non-zero — regression-tested). KNOWN-LIMITATIONS header flags the
  branch-protection GET->PUT transform + live-validation for the C1 gate.
- scripts/assert_no_write_token.py (+test): box/CI audit that no write token
  (incl. ghu_/ghr_ prefixes + App PEM) lives on the box.
- draft_pr_monitor.py (+test): runaway (>3/15min) + stale (7d) draft-PR sweep,
  wired into tick() and bound a read-only provider in serve.

Phase B (wiring): systemd EnvironmentFile P3 vars + verification; new-draft-PR
lifecycle notice.

Phase E (docs): P3-LIVE-FLIP-PLAN/README/ci-README reflect CI-live-since-6/22 +
box-integration; runbook consolidated (rollback Incident 7 + box-env wiring);
removed a stray duplicate runbook.

Suite: 1360 passed, ruff clean. Branch only; not merged/deployed.
REMAINING HUMAN GATES: C1 /sh-security-review + GPT-4.1 cross-review on the
enabled workflow + rollback script; D box deploy + smoke + merge.
2026-06-23 19:52:04 -04:00

422 lines
16 KiB
Python

"""Draft-PR runaway monitor + stale cleanup sweep (P3 box-integration, A4).
Once the box-side BUILD → DISPATCH → VERIFY path is live (CI apply/verify is
already provisioned), a passing task ends at a **draft PR** (the §3.3.2 PASS
terminus). Two failure modes then need a maintenance sweep, mirroring the shape
of :mod:`agent_team.deadline_timer` and :mod:`agent_team.ci_watcher` (a pure,
restart-safe ``tick()``-driven pass whose side effects are injected callables, so
it is unit-testable with NO network and NO live graph):
* **RUNAWAY** — a dispatch loop (a stuck builder, a re-dispatch storm, or a
fixer that keeps re-opening) can open draft PRs far faster than a human can
review them. If **more than 3 draft PRs are opened within any 15-minute
window**, raise an ALARM to ``#agent-team`` (via the existing lifecycle / ALARM
path). The remediation is **operator-driven**: stop auto-dispatch
(``systemctl stop``). The monitor *surfaces* the condition; it never
self-restarts, self-stops, or auto-closes anything — it has no standing write
authority and must not act on infrastructure.
* **STALE** — a draft PR that has sat **idle for more than 7 days** is surfaced
with a single ``#agent-team`` reminder so it is not silently forgotten. The
monitor **never auto-closes** a stale PR — closing is a human decision; the
reminder is the only action.
Flapping backoff (the load-bearing anti-spam discipline): a RUNAWAY condition
typically persists across many ticks (the offending PRs stay inside the window
for the whole 15 minutes), and a STALE PR stays stale until a human acts. Without
backoff the sweep would re-ALARM / re-remind on *every* tick. So the monitor
carries a small :class:`MonitorMemory` (injected, durable-across-ticks): it
ALARMs at most once per ``alarm_cooldown`` and reminds about a given PR at most
once per ``stale_cooldown``. The memory is passed in (not module-global) so the
coordinator owns its lifetime and a test can assert the backoff deterministically.
Fail-soft / fail-closed discipline (mirroring the sibling sweeps):
* Every side effect (``on_alarm`` / ``on_stale_reminder``) is an injected
callable; the monitor performs no transport I/O of its own.
* A draft PR with an unparseable ``opened_at`` is ignored for the runaway count
(it cannot be placed in the window) but a present ``updated_at`` is still
considered for staleness; a PR with neither parseable timestamp is skipped
rather than crashing the sweep.
* A side effect that raises is isolated per condition (logged, recorded) so one
broken post never aborts the pass — the daemon's ``tick()`` loop keeps running.
* All inputs are read fresh each pass (the injected provider re-enumerates the
open draft PRs); the only retained state is the small backoff memory.
"""
from __future__ import annotations
import logging
from collections.abc import Callable
from dataclasses import dataclass, field
from datetime import datetime, timedelta, timezone
from enum import Enum
__all__ = [
"AlarmFn",
"DEFAULT_ALARM_COOLDOWN",
"DEFAULT_ALARM_THRESHOLD",
"DEFAULT_ALARM_WINDOW",
"DEFAULT_STALE_AFTER",
"DEFAULT_STALE_COOLDOWN",
"DraftPr",
"MonitorAction",
"MonitorMemory",
"MonitorOutcome",
"MonitorReport",
"StaleReminderFn",
"run_draft_pr_monitor",
]
logger = logging.getLogger(__name__)
# Concrete thresholds (A4). RUNAWAY: > 3 draft PRs opened within 15 minutes.
DEFAULT_ALARM_THRESHOLD = 3
DEFAULT_ALARM_WINDOW = timedelta(minutes=15)
# STALE: a draft PR idle (no update) for more than 7 days.
DEFAULT_STALE_AFTER = timedelta(days=7)
# Flapping backoff windows. A runaway condition persists across many ticks while
# the offending PRs stay inside the 15-minute window, and a stale PR stays stale
# until a human acts; these cooldowns stop the sweep re-alarming / re-reminding
# every tick. One ALARM per 15 min, one reminder per PR per day.
DEFAULT_ALARM_COOLDOWN = timedelta(minutes=15)
DEFAULT_STALE_COOLDOWN = timedelta(days=1)
class MonitorAction(Enum):
"""The action the sweep took this pass for one surfaced condition.
``ALARMED`` — a runaway burst was detected and the ALARM was posted.
``ALARM_SUPPRESSED`` — a runaway burst was detected but an ALARM was posted
recently (within ``alarm_cooldown``), so it was suppressed (flapping
backoff). ``STALE_REMINDED`` — a stale PR's reminder was posted.
``STALE_SUPPRESSED`` — a stale PR was found but it was reminded about
recently (within ``stale_cooldown``), so the reminder was suppressed.
``ERRORED`` — a side effect raised; the condition was detected but the post
failed (isolated, the sweep continues).
"""
ALARMED = "alarmed"
ALARM_SUPPRESSED = "alarm_suppressed"
STALE_REMINDED = "stale_reminded"
STALE_SUPPRESSED = "stale_suppressed"
ERRORED = "errored"
@dataclass(frozen=True)
class DraftPr:
"""A minimal read-snapshot of one open draft PR (A4).
Only the fields the monitor needs: identity (``number``) and the two
timestamps the runaway-count / staleness checks key off. ``opened_at`` is
when the draft PR was created (runaway window); ``updated_at`` is the last
activity (staleness). Frozen because it is a snapshot — the monitor never
mutates a PR; it acts only through the injected callables.
"""
number: int
opened_at: str | None = None
updated_at: str | None = None
@dataclass(frozen=True)
class MonitorOutcome:
"""The result of one surfaced condition this pass (A4)."""
action: MonitorAction
# For a runaway ALARM: the count of PRs opened inside the window. For a stale
# outcome: the PR number. ``None`` is never expected but keeps the dataclass
# total for the ERRORED path.
detail: str | None = None
error: str | None = None
@dataclass
class MonitorReport:
"""Aggregate result of one draft-PR monitor pass (mirrors the sibling sweeps).
``outcomes`` is one entry per surfaced condition (a runaway ALARM and/or each
stale reminder). The summary counters let the coordinator log/ALARM without
re-walking the list.
"""
outcomes: list[MonitorOutcome] = field(default_factory=list)
# The number of draft PRs counted as opened inside the runaway window this
# pass, for observability (not every counted PR produces an outcome — only a
# *breach* does).
opened_in_window: int = 0
@property
def alarmed(self) -> int:
"""Runaway ALARMs actually posted this pass."""
return sum(1 for o in self.outcomes if o.action is MonitorAction.ALARMED)
@property
def alarm_suppressed(self) -> int:
"""Runaway breaches detected but suppressed by the alarm cooldown."""
return sum(
1 for o in self.outcomes if o.action is MonitorAction.ALARM_SUPPRESSED
)
@property
def stale_reminded(self) -> int:
"""Stale reminders actually posted this pass."""
return sum(1 for o in self.outcomes if o.action is MonitorAction.STALE_REMINDED)
@property
def stale_suppressed(self) -> int:
"""Stale PRs found but suppressed by the per-PR reminder cooldown."""
return sum(
1 for o in self.outcomes if o.action is MonitorAction.STALE_SUPPRESSED
)
@property
def errored(self) -> int:
"""Conditions whose side effect raised (isolated; sweep continued)."""
return sum(1 for o in self.outcomes if o.action is MonitorAction.ERRORED)
@dataclass
class MonitorMemory:
"""Durable-across-ticks backoff memory (the flapping-backoff seam).
The monitor itself is otherwise pure: it reads the open draft PRs fresh each
pass. This small mutable record is the ONE piece of state that must survive
between ticks so a persistent condition does not re-ALARM / re-remind every
pass. The coordinator owns one instance for the daemon's lifetime; a test
constructs its own to assert the backoff deterministically.
``last_alarm_at`` — when a runaway ALARM was last posted (None = never).
``last_reminded_at`` — per-PR-number, when that PR was last reminded about.
Entries for PRs no longer present are pruned each pass so the map cannot grow
without bound across a long-running daemon.
"""
last_alarm_at: datetime | None = None
last_reminded_at: dict[int, datetime] = field(default_factory=dict)
# Injected side-effect seams. Keeping these as callables means the monitor
# performs no transport I/O of its own (testable, faithful to the deadline_timer
# / ci_watcher shape).
AlarmFn = Callable[[int], None]
"""Called once (per cooldown) when a runaway burst is detected. Receives the
count of draft PRs opened inside the window. The handler posts the ALARM to
``#agent-team`` and surfaces the operator remediation (stop auto-dispatch); it
does NOT self-stop."""
StaleReminderFn = Callable[[DraftPr], None]
"""Called once (per cooldown) per stale draft PR. The handler posts a
``#agent-team`` reminder. It NEVER auto-closes the PR."""
def _utc_now() -> datetime:
"""Return the current UTC time (injectable via ``now`` in the sweep)."""
return datetime.now(timezone.utc)
def _parse_iso(value: str | None) -> datetime | None:
"""Parse an ISO-8601 timestamp to an aware UTC datetime, or ``None``.
A missing / unparseable timestamp yields ``None`` so the caller fails soft
(skips that PR for the affected check) rather than crashing the sweep. A
naive stamp is treated as UTC (the ledger / GitHub API write UTC) so the
aware/naive compares below never raise.
"""
if not value:
return None
try:
parsed = datetime.fromisoformat(value)
except (TypeError, ValueError):
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed
def run_draft_pr_monitor(
draft_prs: list[DraftPr],
*,
on_alarm: AlarmFn,
on_stale_reminder: StaleReminderFn,
memory: MonitorMemory | None = None,
alarm_threshold: int = DEFAULT_ALARM_THRESHOLD,
alarm_window: timedelta = DEFAULT_ALARM_WINDOW,
stale_after: timedelta = DEFAULT_STALE_AFTER,
alarm_cooldown: timedelta = DEFAULT_ALARM_COOLDOWN,
stale_cooldown: timedelta = DEFAULT_STALE_COOLDOWN,
now: datetime | None = None,
) -> MonitorReport:
"""Run one draft-PR runaway + stale sweep (A4).
``draft_prs`` is the set of currently-open draft PRs (the coordinator
supplies them fresh each pass via an injected provider). The sweep:
1. **Runaway.** Counts draft PRs whose ``opened_at`` is within ``now -
alarm_window``. If that count is **strictly greater than**
``alarm_threshold`` (the A4 rule: > 3 within 15 min), call ``on_alarm``
once — but only if a prior ALARM is older than ``alarm_cooldown`` (flapping
backoff); otherwise record ``ALARM_SUPPRESSED``. The remediation is the
operator's (``systemctl stop``); the monitor never self-stops.
2. **Stale.** For each draft PR idle longer than ``stale_after`` (``now -
updated_at > stale_after``), call ``on_stale_reminder`` once — but only if
that PR was last reminded longer ago than ``stale_cooldown`` (per-PR
backoff); otherwise record ``STALE_SUPPRESSED``. The monitor never
auto-closes a stale PR.
Side effects are isolated per condition: an ``on_alarm`` / ``on_stale_reminder``
that raises yields an ``ERRORED`` outcome (the failure is recorded and the
backoff watermark is NOT advanced, so the next pass retries) and the sweep
continues with the rest of the batch.
Restart-safety: inputs are read fresh each pass; the only retained state is
``memory`` (the backoff watermarks). A fresh ``memory`` (e.g. after a reboot)
simply means the first post-reboot breach/stale PR is surfaced again — a
re-notification, never a missed or duplicated *action* (the monitor takes no
infrastructure action).
Returns a :class:`MonitorReport` describing what happened.
"""
current = now or _utc_now()
mem = memory if memory is not None else MonitorMemory()
report = MonitorReport()
present_numbers = {pr.number for pr in draft_prs}
# --- Runaway check ----------------------------------------------------- #
window_start = current - alarm_window
opened_in_window = 0
for pr in draft_prs:
opened_at = _parse_iso(pr.opened_at)
if opened_at is not None and opened_at >= window_start:
opened_in_window += 1
report.opened_in_window = opened_in_window
if opened_in_window > alarm_threshold:
_handle_runaway(
opened_in_window,
on_alarm=on_alarm,
mem=mem,
alarm_cooldown=alarm_cooldown,
now=current,
report=report,
)
# --- Stale check ------------------------------------------------------- #
for pr in draft_prs:
updated_at = _parse_iso(pr.updated_at)
if updated_at is None:
# No parseable last-activity stamp -> cannot assess staleness. Skip
# rather than guess (fail-soft); the runaway count is unaffected.
continue
if current - updated_at <= stale_after:
continue
_handle_stale(
pr,
on_stale_reminder=on_stale_reminder,
mem=mem,
stale_cooldown=stale_cooldown,
now=current,
report=report,
)
# Prune backoff watermarks for PRs no longer open so the memory cannot grow
# without bound over a long-running daemon.
for number in list(mem.last_reminded_at):
if number not in present_numbers:
del mem.last_reminded_at[number]
return report
def _handle_runaway(
opened_in_window: int,
*,
on_alarm: AlarmFn,
mem: MonitorMemory,
alarm_cooldown: timedelta,
now: datetime,
report: MonitorReport,
) -> None:
"""Post (or suppress) a runaway ALARM under the flapping-backoff cooldown."""
last = mem.last_alarm_at
if last is not None and now - last < alarm_cooldown:
logger.debug(
"draft-pr-monitor: runaway breach (%d in window) suppressed by cooldown",
opened_in_window,
)
report.outcomes.append(
MonitorOutcome(
action=MonitorAction.ALARM_SUPPRESSED,
detail=str(opened_in_window),
)
)
return
try:
on_alarm(opened_in_window)
except Exception as exc: # noqa: BLE001 - isolate one side-effect failure
logger.exception(
"draft-pr-monitor: runaway ALARM side effect raised; watermark not "
"advanced (will retry next pass)"
)
report.outcomes.append(
MonitorOutcome(
action=MonitorAction.ERRORED,
detail=str(opened_in_window),
error=f"{type(exc).__name__}: {exc}",
)
)
return
# Advance the watermark ONLY after a successful post so a failed post retries.
mem.last_alarm_at = now
logger.warning(
"draft-pr-monitor: RUNAWAY — %d draft PRs opened within the window; "
"ALARM raised (operator remediation: stop auto-dispatch)",
opened_in_window,
)
report.outcomes.append(
MonitorOutcome(action=MonitorAction.ALARMED, detail=str(opened_in_window))
)
def _handle_stale(
pr: DraftPr,
*,
on_stale_reminder: StaleReminderFn,
mem: MonitorMemory,
stale_cooldown: timedelta,
now: datetime,
report: MonitorReport,
) -> None:
"""Post (or suppress) a stale reminder for one PR under per-PR backoff."""
last = mem.last_reminded_at.get(pr.number)
if last is not None and now - last < stale_cooldown:
report.outcomes.append(
MonitorOutcome(action=MonitorAction.STALE_SUPPRESSED, detail=str(pr.number))
)
return
try:
on_stale_reminder(pr)
except Exception as exc: # noqa: BLE001 - isolate one side-effect failure
logger.exception(
"draft-pr-monitor: stale reminder side effect for PR #%s raised; "
"watermark not advanced (will retry next pass)",
pr.number,
)
report.outcomes.append(
MonitorOutcome(
action=MonitorAction.ERRORED,
detail=str(pr.number),
error=f"{type(exc).__name__}: {exc}",
)
)
return
mem.last_reminded_at[pr.number] = now
report.outcomes.append(
MonitorOutcome(action=MonitorAction.STALE_REMINDED, detail=str(pr.number))
)