← Files Compound EngineeringARCHIVED FILE
skills/ce-babysit-pr/scripts/pr-snapshot
172 KB · Oct 5, 2026 · 18:34 UTC
#!/usr/bin/env python3
"""Deterministic snapshot + state helper for ce-babysit-pr.
The agent (SKILL.md) owns judgment and mutations; this script owns the parts
prose cannot do reliably: a combined fetch of both event streams, atomic state
read/write under a file lock, and dedup keyed on remote truth.
The dedup model is claim -> act -> confirm. `snapshot` never marks an item
handled just because it observed it; an item stays actionable until either the
agent confirms it acted (`mark`) OR remote truth removes it (a resolved thread
drops out of the unresolved fetch). So if a resolve/debug pass crashes or fails
before the agent marks it, the item is still actionable on the next tick.
Subcommands:
snapshot --pr N [--repo O/R] --state-dir DIR [--fetch-file F]
Fetch (or load F), diff against on-disk state, persist the
observed state atomically, emit the actionable set as JSON.
mark --state-dir DIR --invocation-id ID --session-started-at TIME
--invocation-budget-seconds N (--thread ID --disposition open|dispatched
| --comment ID --disposition open|dispatched | --check KEY
| --disposition needs-human --residual-file FILE
| --answer-decision ID --answer-file FILE)
A residual-only needs-human mark atomically freezes every source observation
in FILE under one decision. Sources retain their ordinary dispositions.
Record that the agent acted on an item. A `dispatched` thread is
re-emitted when a later reviewer comment moves its last-comment
identity past the one we acted on (our own reply does not
re-trigger). A non-thread feedback item never drops out of
the fetch on its own, so `--comment` is the only way to silence a
handled one. A new head SHA clears dispatched CI checks.
--fetch-file injects a pre-captured combined snapshot instead of calling gh,
so the diff logic is testable without a live PR.
"""
import argparse
import errno
import hashlib
import json
import os
import re
import signal
import subprocess
import sys
import tempfile
import threading
import time
import uuid
from contextlib import contextmanager
from datetime import datetime, timedelta, timezone
from urllib.parse import quote, urlsplit
IS_WINDOWS = sys.platform == "win32"
if IS_WINDOWS:
import msvcrt
else:
import fcntl
# Conclusions that mean "this check needs attention" (failing states).
FAILING = {"FAILURE", "TIMED_OUT", "CANCELLED", "ACTION_REQUIRED", "STARTUP_FAILURE", "STALE"}
DISPOSITION_OPEN = "open"
DISPOSITION_NEEDS_HUMAN = "needs-human"
DISPOSITION_DISPATCHED = "dispatched"
NEEDS_HUMAN_SOURCE_KINDS = {"thread", "comment", "review", "check", "currency"}
CURRENCY_CLAIMED = "claimed"
CURRENCY_CONFIRMED = "confirmed"
CURRENCY_OUTCOME_MUTATION_OBSERVED = "mutation-observed"
CURRENCY_OUTCOME_PROVEN_NO_MUTATION = "proven-no-mutation"
CURRENCY_OUTCOME_AMBIGUOUS = "ambiguous"
CURRENCY_ATTENTION_CLAIM = "claim"
CURRENCY_ATTENTION_DECIDE = "decide"
CURRENCY_ATTENTION_INSPECT = "inspect"
CURRENCY_ATTENTION_RECONCILE = "reconcile"
MANAGER_CONFIRMED = "confirmed"
MANAGER_ABSENT = "absent"
MANAGER_PROBE_ERROR = "probe-error"
RELATIONSHIP_DEPENDENT = "dependent"
RELATIONSHIP_INDEPENDENT = "independent"
RELATIONSHIP_PROBE_ERROR = "probe-error"
BASE_REF_CURRENT = "current"
BASE_REF_RACE = "race"
BASE_REF_STALE_COMPUTATION = "stale-computation"
# GitHub can leave potentialMergeCommit parented on the old base indefinitely after main moves
# (observed 20+ min on tmchow/pr-stack-test, unrefreshed by polls or a body PATCH). Once the live
# base is stable this long with the same head, the cached CLEAN/DIRTY verdict is accepted with a
# disclosure instead of blocking readiness forever.
STALE_MERGE_COMPUTATION_SECONDS = 600
BASE_REF_MERGEABILITY_PENDING = "mergeability-pending"
BASE_REF_PROBE_ERROR = "probe-error"
BASE_REF_LEGACY_STALE = "stale"
# Quiet window a merge-ready wake ordinarily requires. The agent widens it by re-arming when it
# judges more waiting is warranted; the engine never adjusts the window it was given.
DEFAULT_SETTLE_SECONDS = 300.0
# Quiet window before a comment-only wake. Every non-empty top-level body is a candidate the
# resolver must classify, so waking on the first one spends a full ce-resolve-pr-feedback dispatch
# to conclude "status noise". Only this reason waits; real work never consults the clock.
FEEDBACK_COALESCE_SECONDS = 180
# Trajectory check-history states (persisted, compared by ==).
CHECK_UNKNOWN = "unknown"
CHECK_CLEAR = "clear"
CHECK_FAILING = "failing"
# Trajectory single-stream activity labels.
STREAM_CI = "ci"
STREAM_REVIEW = "review"
# Keep a check's recurrence memory across a transient absence (a one-tick gap from a
# workflow-registration lag or a paths-filtered run), but bound growth: evict entries
# unseen for this many ticks.
CHECK_HISTORY_TTL = 30
# One invocation is a bounded monitoring shift, not the lifetime of the PR. Eight hours covers a
# maximum-length GitHub-hosted Actions job (six hours) plus reaction and settle time. Callers may
# choose another fixed value on the first snapshot; later commands must present the same value.
DEFAULT_INVOCATION_BUDGET_SECONDS = 8 * 60 * 60
# The 8h budget is spent in *active watch-capability time*, not raw wall-clock: a span where the
# whole watch process was suspended (laptop asleep) is excluded. Detection is coarse — an activity
# gap wider than this threshold (well above the 150s poll interval) is charged to dead time, minus
# the threshold itself so ordinary polls/ticks never register. The 3-day backstop stays wall-clock.
DEAD_TIME_THRESHOLD_SECONDS = 15 * 60
DEFAULT_INVOCATION_BACKSTOP_SECONDS = 3 * 24 * 60 * 60
CURRENCY_RETRY_BACKOFF_SECONDS = 30
# --- cross-platform state primitives ------------------------------------------
#
# The watcher needs three POSIX facilities that native Windows Python does not
# have: an advisory file lock (fcntl.flock), a process-identity probe for the
# PID-reuse guard (`ps -o lstart=`), and an unconditional rename. Each Windows
# analog differs in a way that matters here, so each is named rather than
# translated literally. See
# docs/solutions/architecture-patterns/posix-process-supervision-on-native-windows.md.
# Poll spacing while waiting on a Windows lock or a blocked rename.
_WIN_RETRY_SECONDS = 0.05
# Bound only the rename retry. A lock wait is unbounded, matching flock.
_WIN_REPLACE_TIMEOUT_SECONDS = 10.0
if IS_WINDOWS:
import ctypes
from ctypes import wintypes
_kernel32 = ctypes.WinDLL("kernel32", use_last_error=True)
_PROCESS_QUERY_LIMITED_INFORMATION = 0x1000
# Declare argtypes/restype explicitly: without them ctypes truncates HANDLEs
# to 32-bit ints on Win64, so a valid handle silently becomes a bad one.
_kernel32.OpenProcess.argtypes = [wintypes.DWORD, wintypes.BOOL, wintypes.DWORD]
_kernel32.OpenProcess.restype = wintypes.HANDLE
_kernel32.CloseHandle.argtypes = [wintypes.HANDLE]
_kernel32.CloseHandle.restype = wintypes.BOOL
_kernel32.GetProcessTimes.argtypes = [
wintypes.HANDLE, ctypes.POINTER(wintypes.FILETIME),
ctypes.POINTER(wintypes.FILETIME), ctypes.POINTER(wintypes.FILETIME),
ctypes.POINTER(wintypes.FILETIME)]
_kernel32.GetProcessTimes.restype = wintypes.BOOL
_kernel32.QueryFullProcessImageNameW.argtypes = [
wintypes.HANDLE, wintypes.DWORD, wintypes.LPWSTR,
ctypes.POINTER(wintypes.DWORD)]
_kernel32.QueryFullProcessImageNameW.restype = wintypes.BOOL
def _win_process_identity(pid):
"""Creation time plus image path — the Windows analog of `ps -o lstart= -o command=`.
Creation time is the part that makes this a PID-reuse guard: Windows recycles PIDs
aggressively, and a recycled PID always carries a later creation time than the watcher
we recorded. Returns None when the process is gone or unopenable, which the caller
treats as "identity not proven" and declines to act on."""
handle = _kernel32.OpenProcess(_PROCESS_QUERY_LIMITED_INFORMATION, False, pid)
if not handle:
return None
try:
created, exited, kernel, user = (wintypes.FILETIME() for _ in range(4))
if not _kernel32.GetProcessTimes(
handle, ctypes.byref(created), ctypes.byref(exited),
ctypes.byref(kernel), ctypes.byref(user)):
return None
started = (created.dwHighDateTime << 32) | created.dwLowDateTime
if not started:
return None
size = wintypes.DWORD(32768)
buf = ctypes.create_unicode_buffer(size.value)
image = (buf.value if _kernel32.QueryFullProcessImageNameW(
handle, 0, buf, ctypes.byref(size)) else "")
finally:
_kernel32.CloseHandle(handle)
return "{} {}".format(started, image)
def _lock_acquire(fd, exclusive):
"""Block until this process owns the lock file.
msvcrt.locking is the Windows analog of flock, with three differences. It locks a byte range
from the current offset rather than the whole file, so the offset is pinned to 0. It has no
shared mode, so a reader takes the same exclusive lock a writer does — free here, because every
critical section is a small local read/write with the network fetch already completed outside
the lock. And its blocking mode (LK_LOCK) gives up after roughly ten seconds instead of waiting,
so retrying the non-blocking mode is what actually reproduces flock's wait-for-the-holder
behavior. As on POSIX, the OS releases the lock if a holder dies, so a crash cannot wedge this."""
if not IS_WINDOWS:
fcntl.flock(fd, fcntl.LOCK_EX if exclusive else fcntl.LOCK_SH)
return
os.lseek(fd, 0, os.SEEK_SET)
while True:
try:
msvcrt.locking(fd, msvcrt.LK_NBLCK, 1)
return
except OSError as e:
# Retry only genuine contention. msvcrt reports a held range as EACCES; anything
# else (EBADF from a closed descriptor, EINVAL from a bad range) is a defect that
# retrying cannot clear, and swallowing it here would spin this loop forever with
# no output — a silent hang, which for a watcher is indistinguishable from the
# "monitoring is quietly not running" failure this port exists to fix.
if e.errno != errno.EACCES:
raise
time.sleep(_WIN_RETRY_SECONDS)
def _lock_release(fd):
if not IS_WINDOWS:
fcntl.flock(fd, fcntl.LOCK_UN)
return
os.lseek(fd, 0, os.SEEK_SET)
try:
msvcrt.locking(fd, msvcrt.LK_UNLCK, 1)
except OSError:
pass
@contextmanager
def _state_lock(state_dir, exclusive=True):
"""Serialize access to the state dir for the duration of the body."""
os.makedirs(state_dir, exist_ok=True)
# O_RDWR|O_CREAT rather than a truncating open("w"): on Windows, truncating a file another
# process holds a byte-range lock on fails, which would make lock acquisition itself the race.
fd = os.open(os.path.join(state_dir, "lock"), os.O_RDWR | os.O_CREAT, 0o600)
try:
_lock_acquire(fd, exclusive)
try:
yield
finally:
_lock_release(fd)
finally:
os.close(fd)
def _replace_atomic(src, dst):
"""Atomically move `src` onto `dst`.
POSIX rename always succeeds over an open destination. Windows refuses while any other handle
to `dst` is open, so an unlocked best-effort reader — or a virus scanner that opened the file
behind us — would otherwise turn a routine concurrent read into a crashed writer. Retry briefly
rather than treating that transient sharing violation as a failure."""
if not IS_WINDOWS:
os.replace(src, dst)
return
deadline = time.monotonic() + _WIN_REPLACE_TIMEOUT_SECONDS
while True:
try:
os.replace(src, dst)
return
except PermissionError:
if time.monotonic() >= deadline:
raise
time.sleep(_WIN_RETRY_SECONDS)
class _WatchSuperseded(Exception):
"""Internal control flow: this watch lost ownership and must unwind immediately."""
class _InvocationSuperseded(Exception):
"""Internal control flow: this watch belongs to an invocation that was replaced."""
def __init__(self, current_invocation_id):
self.current_invocation_id = current_invocation_id
def _now():
return datetime.now(timezone.utc)
def _iso(dt):
return dt.isoformat()
def _run(cmd):
# gh honors CLICOLOR_FORCE / GH_FORCE_TTY even for --json output, and a host that exports them
# (observed under orca) makes every JSON read unparsable. Pin plain output regardless of caller env.
env = os.environ.copy()
env["NO_COLOR"] = "1"
env.pop("CLICOLOR_FORCE", None)
env.pop("GH_FORCE_TTY", None)
return subprocess.run(cmd, capture_output=True, text=True, encoding="utf-8", env=env)
def _run_git(cmd):
env = os.environ.copy()
env["GIT_TERMINAL_PROMPT"] = "0"
env["GCM_INTERACTIVE"] = "never"
env["GIT_ASKPASS"] = ""
env["SSH_ASKPASS"] = ""
return subprocess.run(
cmd, capture_output=True, text=True, encoding="utf-8", env=env, timeout=30,
)
def _run_checked(cmd, label):
r = _run(cmd)
if r.returncode != 0:
raise SystemExit(f"{label} failed: {r.stderr.strip()}")
return r
def _split_repo(repo):
"""Parse "[HOST/]OWNER/NAME" into (owner, name), or (None, None). A host-qualified ref (gh's
documented `[HOST/]OWNER/REPO` selector) drops the host — the last two segments are what the
GraphQL lookup needs; treating the host as the owner would query a nonexistent repo on GHE."""
if not repo:
return None, None
parts = repo.strip("/").split("/")
if len(parts) >= 2:
return parts[-2], parts[-1]
return None, None
def _resolve_repo_ref(repo, url):
"""Resolve (owner, name, host) from --repo + the PR url, else one `gh repo view` call.
The host is parsed from the url and threaded into every `gh api` call so a GitHub Enterprise
PR queries the right host — without it, `gh api` defaults to github.com and a GHE babysitter
fetches the PR via `gh pr view` but then reads review threads / workflow runs from github.com.
Parsing the url (already fetched by `fetch`) also avoids a redundant `gh repo view` per tick."""
owner, name = _split_repo(repo)
host = None
if url:
# https://<host>/OWNER/NAME/pull/N
parts = url.rstrip("/").split("/")
if len(parts) >= 5 and parts[0].startswith("http"):
host = parts[2]
if not owner:
owner, name = parts[-4], parts[-3]
if not owner:
r = _run(["gh", "repo", "view", "--json", "owner,name"])
if r.returncode == 0:
info = json.loads(r.stdout)
owner, name = info.get("owner", {}).get("login"), info.get("name")
if not owner or not name:
raise SystemExit("could not resolve owner/repo; pass --repo OWNER/REPO")
return owner, name, host
def _host_args(host):
"""`gh api --hostname` selector so GHE calls hit the PR's host, not the default github.com."""
return ["--hostname", host] if host else []
def _valid_oid(value):
return isinstance(value, str) and bool(
re.fullmatch(r"(?:[0-9a-fA-F]{40}|[0-9a-fA-F]{64})", value))
def fetch_pr_merge_identity(pr, owner, name, host=None):
"""Read mergeability and the identities it was computed from in one GraphQL observation."""
query = """
query($owner:String!,$repo:String!,$pr:Int!){
repository(owner:$owner,name:$repo){ pullRequest(number:$pr){
mergeable mergeStateStatus headRefOid baseRefOid baseRefName viewerCanUpdateBranch
baseRef { target { oid } }
headRef { target { ... on Commit { parents(first:2) { nodes { oid } } } } }
potentialMergeCommit { oid parents(first:2) { nodes { oid } } }
} }
}"""
r = _run(["gh", "api", "graphql", *_host_args(host), "-f", f"owner={owner}",
"-f", f"repo={name}", "-F", f"pr={pr}", "-f", f"query={query}"])
if r.returncode != 0:
return None
try:
value = json.loads(r.stdout)["data"]["repository"]["pullRequest"]
except (KeyError, TypeError, ValueError, json.JSONDecodeError):
return None
return value if isinstance(value, dict) else None
def fetch_base_ref(owner, name, ref, merge_identity, host=None):
"""Bind GitHub's mergeability observation to an independently probed current base ref.
`baseRefOid` is retained only as historical PR metadata. Current identity comes from both the
PR's `baseRef.target.oid` and an exact Git-ref read; a usable generated test merge must then
name that base and the observed PR head as its two parents. Some valid private-repository OAuth
sessions return 404 from the REST Git-ref endpoint, so that response falls back to the Git
transport while preserving exact-ref matching and non-interactive failure behavior.
"""
historical_oid = (merge_identity or {}).get("baseRefOid")
head_oid = (merge_identity or {}).get("headRefOid")
mergeable = (merge_identity or {}).get("mergeable")
merge_state_status = (merge_identity or {}).get("mergeStateStatus")
graphql_oid = (((merge_identity or {}).get("baseRef") or {}).get("target") or {}).get("oid")
potential_merge = (merge_identity or {}).get("potentialMergeCommit")
merge_commit_oid = potential_merge.get("oid") if isinstance(potential_merge, dict) else None
merge_parent_oids = (((potential_merge or {}).get("parents") or {}).get("nodes")
if isinstance(potential_merge, dict) else None)
if isinstance(merge_parent_oids, list):
merge_parent_oids = [node.get("oid") if isinstance(node, dict) else None
for node in merge_parent_oids]
else:
merge_parent_oids = []
base = {
"host": host or "github.com",
"repository": f"{owner}/{name}",
"ref": ref,
"oid": None,
"graphql_oid": graphql_oid,
"historical_oid": historical_oid,
"merge_commit_oid": merge_commit_oid,
"merge_parent_oids": merge_parent_oids,
"identity": BASE_REF_PROBE_ERROR,
}
if not ref:
return base
encoded_ref = quote(ref, safe="/")
r = _run([
"gh", "api", *_host_args(host),
f"repos/{owner}/{name}/git/ref/heads/{encoded_ref}",
"--jq", ".object.sha",
])
current_oid = (r.stdout or "").strip()
if r.returncode != 0 and re.search(r"\bHTTP\s+404\b", r.stderr or "", re.IGNORECASE):
exact_ref = f"refs/heads/{ref}"
remote_url = (f"https://{host or 'github.com'}/"
f"{quote(owner, safe='')}/{quote(name, safe='')}.git")
try:
# Reuse the gh session (GH_TOKEN or keyring) as git's credential source so a
# headless host with no configured Git credential helper is not anonymous here.
git_result = _run_git([
"git", "-c", "core.askPass=",
"-c", "credential.helper=",
"-c", "credential.helper=!gh auth git-credential",
"ls-remote", "--exit-code", "--refs",
remote_url, exact_ref,
])
except (subprocess.TimeoutExpired, OSError):
git_result = None
if git_result is not None and git_result.returncode == 0:
matches = []
for line in (git_result.stdout or "").splitlines():
fields = line.split()
if len(fields) == 2 and fields[1] == exact_ref and _valid_oid(fields[0]):
matches.append(fields[0])
if len(matches) == 1:
current_oid = matches[0]
r = git_result
if r.returncode != 0 or not _valid_oid(current_oid):
return base
base["oid"] = current_oid
if not (_valid_oid(graphql_oid) and _valid_oid(head_oid)):
return base
if current_oid.lower() != graphql_oid.lower():
base["identity"] = BASE_REF_RACE
return base
if mergeable == "MERGEABLE" and merge_state_status not in (None, "UNKNOWN", "DIRTY"):
if potential_merge is None:
base["identity"] = BASE_REF_MERGEABILITY_PENDING
return base
if (not _valid_oid(merge_commit_oid) or len(merge_parent_oids) != 2
or not all(_valid_oid(oid) for oid in merge_parent_oids)):
return base
if (merge_parent_oids[0].lower() != current_oid.lower()
or merge_parent_oids[1].lower() != head_oid.lower()):
base["identity"] = BASE_REF_RACE
# Only the base parent lags while head and live base agree: a stale test merge,
# eligible for the bounded degrade in diff().
base["stale_computation_candidate"] = (
merge_parent_oids[1].lower() == head_oid.lower())
return base
elif not (mergeable == "CONFLICTING" and merge_state_status == "DIRTY"):
base["identity"] = BASE_REF_MERGEABILITY_PENDING
return base
base["identity"] = BASE_REF_CURRENT
return base
def _base_ref_blocker(base):
identity = (base or {}).get("identity")
if identity is None:
# Compatibility for pre-change --fetch-file fixtures and persisted snapshots.
identity = (base or {}).get("freshness")
if identity in (BASE_REF_CURRENT, BASE_REF_STALE_COMPUTATION):
return None
if identity in (BASE_REF_RACE, BASE_REF_MERGEABILITY_PENDING):
return identity
if identity == BASE_REF_LEGACY_STALE:
return BASE_REF_LEGACY_STALE
return BASE_REF_PROBE_ERROR
def _stack_schema_unavailable(stderr):
"""Recognize only GraphQL schema errors for the private-preview PullRequest stack fields."""
unavailable_markers = ("doesn't exist on type", "does not exist on type",
"cannot query field", "unknown field")
for line in (stderr or "").splitlines():
lowered = line.lower()
if ("pullrequest" in lowered
and re.search(r"\bstack(?:entry)?\b", lowered)
and any(marker in lowered for marker in unavailable_markers)):
return True
return False
def _fetch_default_branch(owner, name, host=None):
"""Read the default branch without relying on private-preview GraphQL fields."""
result = _run(["gh", "api", *_host_args(host), f"repos/{owner}/{name}",
"--jq", ".default_branch"])
if result.returncode != 0:
return None
branch = result.stdout.strip()
return branch if branch and branch != "null" else None
def _thread_identity(t):
"""The remote-truth identity of a thread's latest state."""
return (t.get("last_comment_id"), t.get("last_comment_at"))
def _prior_thread_activity_identity(t):
"""The last external-review baseline already accepted for this parked thread."""
acted_identity = t.get("acted_identity")
if (t.get("disposition") in (DISPOSITION_DISPATCHED, DISPOSITION_NEEDS_HUMAN)
and isinstance(acted_identity, (list, tuple))
and len(acted_identity) == 2):
return tuple(acted_identity)
return _thread_identity(t)
def _default_pr_chain():
"""Backward-compatible neutral shape for injected/older snapshots with no chain facts."""
return {
"manager_status": MANAGER_ABSENT,
"manager_source": None,
"relationship_status": RELATIONSHIP_INDEPENDENT,
"trunk": None,
"default_branch": None,
"current_branch": None,
"target_position": None,
"target_needs_rebase": None,
"upstack_needs_rebase": [],
"entries": [],
"parent_prs": [],
"dependent_prs": [],
}
def _pr_summary(pr):
if not pr:
return None
return {k: pr.get(k) for k in ("number", "url", "state", "isDraft", "baseRefName", "headRefName")
if pr.get(k) is not None}
def _pr_url_identity(url):
"""Normalize a GitHub PR URL to a repository-scoped identity tuple."""
if not isinstance(url, str):
return None
try:
parsed = urlsplit(url.strip())
parts = [part for part in parsed.path.split("/") if part]
if (parsed.scheme.lower() not in ("http", "https") or not parsed.netloc
or len(parts) < 4 or parts[-2].lower() != "pull"):
return None
number = int(parts[-1])
except (TypeError, ValueError):
return None
return (parsed.netloc.lower(), parts[-4].lower(), parts[-3].lower(), number)
def _chain_from_entries(entries, pr, source, trunk=None, current_branch=None,
manager_id=None, manager_number=None, target_url=None):
"""Normalize an ordered manager-owned chain and locate the requested PR within it."""
target_identity = _pr_url_identity(target_url) if target_url is not None else None
target_index = next((i for i, e in enumerate(entries)
if (e.get("number") == pr or str(e.get("number")) == str(pr))
and (target_url is None
or (target_identity is not None
and _pr_url_identity(e.get("url")) == target_identity))), None)
if target_index is None:
return None
target = entries[target_index]
upstack_stale = [
{k: e.get(k) for k in ("number", "position", "name", "url") if e.get(k) is not None}
for e in entries[target_index + 1:] if e.get("needs_rebase") is True
]
chain = _default_pr_chain()
chain.update({
"manager_status": MANAGER_CONFIRMED,
"manager_source": source,
"manager_id": manager_id,
"manager_number": manager_number,
"relationship_status": RELATIONSHIP_DEPENDENT if len(entries) > 1 else RELATIONSHIP_INDEPENDENT,
"trunk": trunk,
"current_branch": current_branch,
"target_position": target.get("position") or target_index + 1,
"target_needs_rebase": target.get("needs_rebase"),
"upstack_needs_rebase": upstack_stale,
"entries": entries,
"parent_prs": [_pr_summary(entries[target_index - 1])] if target_index > 0 else [],
"dependent_prs": [_pr_summary(e) for e in entries[target_index + 1:]],
})
return chain
def _chain_from_stack_view(raw, pr, target_url):
entries = []
for index, branch in enumerate((raw or {}).get("branches") or []):
remote_pr = branch.get("pr") or {}
entries.append({
"position": index + 1,
"name": branch.get("name"),
"head": branch.get("head"),
"base": branch.get("base"),
"is_current": bool(branch.get("isCurrent")),
"is_merged": bool(branch.get("isMerged")),
"is_queued": bool(branch.get("isQueued")),
"needs_rebase": branch.get("needsRebase"),
"number": remote_pr.get("number"),
"url": remote_pr.get("url"),
"state": remote_pr.get("state"),
"isDraft": remote_pr.get("isDraft"),
})
return _chain_from_entries(entries, pr, "gh-stack", (raw or {}).get("trunk"),
(raw or {}).get("currentBranch"), target_url=target_url)
def _fetch_graphql_stack(pr, owner, name, host=None):
"""Read-only remote membership fallback. A successful null stack is distinct from failure."""
query = """
query($owner:String!,$repo:String!,$pr:Int!){
repository(owner:$owner,name:$repo){
defaultBranchRef{ name }
pullRequest(number:$pr){
stackEntry{ position }
stack{ id number size baseRefName
entries(first:100){ nodes{ position pullRequest{
number url state isDraft baseRefName headRefName headRefOid
} } }
}
} }
}"""
args = ["gh", "api", "graphql", *_host_args(host), "-f", f"owner={owner}", "-f", f"repo={name}",
"-F", f"pr={pr}", "-f", f"query={query}"]
r = _run(args)
if r.returncode != 0:
if _stack_schema_unavailable(r.stderr):
default_branch = _fetch_default_branch(owner, name, host)
if default_branch:
return MANAGER_ABSENT, None, default_branch
return MANAGER_PROBE_ERROR, None, None
try:
repository = json.loads(r.stdout)["data"]["repository"]
default_branch = (repository.get("defaultBranchRef") or {}).get("name")
node = repository["pullRequest"]
stack = node.get("stack")
except (KeyError, TypeError, ValueError, json.JSONDecodeError):
return MANAGER_PROBE_ERROR, None, None
if stack is None:
return MANAGER_ABSENT, None, default_branch
entries = []
for raw in (stack.get("entries") or {}).get("nodes") or []:
remote_pr = raw.get("pullRequest") or {}
entries.append({
"position": raw.get("position"),
"name": remote_pr.get("headRefName"),
"head": remote_pr.get("headRefOid"),
"base_ref_name": remote_pr.get("baseRefName"),
"needs_rebase": None,
**(_pr_summary(remote_pr) or {}),
})
chain = _chain_from_entries(entries, pr, "graphql", stack.get("baseRefName"),
manager_id=stack.get("id"), manager_number=stack.get("number"))
return ((MANAGER_CONFIRMED, chain, default_branch) if chain
else (MANAGER_PROBE_ERROR, None, default_branch))
def _manual_relationships(pr, repo, base_ref_name, head_ref_name, owner, name,
host=None, default_branch=None):
"""Find ordinary parent/children with narrow read-only branch filters."""
repo_ref = repo or f"{owner}/{name}"
if host and repo_ref.count("/") == 1:
repo_ref = f"{host}/{repo_ref}"
fields = "number,url,state,isDraft,baseRefName,headRefName"
parent_result = None
if base_ref_name and base_ref_name != default_branch:
parent_result = _run(["gh", "pr", "list", "--repo", repo_ref, "--state", "all",
"--head", base_ref_name, "--limit", "20", "--json", fields])
dependent_result = _run(["gh", "pr", "list", "--repo", repo_ref, "--state", "open",
"--base", head_ref_name or "", "--limit", "100", "--json", fields])
if ((parent_result is not None and parent_result.returncode != 0)
or dependent_result.returncode != 0):
return RELATIONSHIP_PROBE_ERROR, [], []
try:
parent_candidates = json.loads(parent_result.stdout) or [] if parent_result else []
dependent_candidates = json.loads(dependent_result.stdout) or []
except (TypeError, ValueError, json.JSONDecodeError):
return RELATIONSHIP_PROBE_ERROR, [], []
parents = [_pr_summary(p) for p in parent_candidates
if p.get("number") != pr and p.get("headRefName") == base_ref_name]
dependents = [_pr_summary(p) for p in dependent_candidates
if p.get("number") != pr and p.get("state") == "OPEN"
and p.get("baseRefName") == head_ref_name]
status = RELATIONSHIP_DEPENDENT if parents or dependents else RELATIONSHIP_INDEPENDENT
return status, parents, dependents
def fetch_pr_chain(pr, repo, url, base_ref_name, head_ref_name, owner, name, host=None):
"""Classify manager membership and ordinary dependency relationships without mutation.
The local manager is the fast/rich path, but its output is accepted only when it contains the
requested PR. `gh stack view` has no target argument and may describe a different current stack.
"""
local = _run(["gh", "stack", "view", "--json"])
if local.returncode == 0:
try:
chain = _chain_from_stack_view(json.loads(local.stdout), pr, url)
except (TypeError, ValueError, json.JSONDecodeError):
chain = None
if chain:
return chain
manager_status, chain, default_branch = _fetch_graphql_stack(pr, owner, name, host)
if manager_status == MANAGER_CONFIRMED:
chain["default_branch"] = default_branch
return chain
relationship, parents, dependents = _manual_relationships(
pr, repo, base_ref_name, head_ref_name, owner, name, host, default_branch)
result = _default_pr_chain()
result.update({
"manager_status": manager_status,
"relationship_status": relationship,
"default_branch": default_branch,
"parent_prs": parents,
"dependent_prs": dependents,
})
return result
def fetch(pr, repo, light=False):
"""Fetch both event streams via gh into one combined snapshot dict. `light` skips the probes a
downstack signature never reads (merge identity/base ref, eyes reactions, awaiting-approval)."""
repo_args = ["--repo", repo] if repo else []
view = _run_checked(["gh", "pr", "view", str(pr), *repo_args, "--json",
"state,mergeable,mergeStateStatus,reviewDecision,headRefOid,baseRefOid,baseRefName,headRefName,url,number,isDraft,"
"statusCheckRollup,author,comments,reviews"],
"gh pr view")
v = json.loads(view.stdout)
checks = []
for c in v.get("statusCheckRollup") or []:
# CheckRun entries carry startedAt/completedAt; StatusContext (legacy commit statuses)
# carry createdAt. Preserve the provider's real creation identity for rerun reconciliation.
if c.get("__typename") == "StatusContext":
name = c.get("context") or "status"
state = (c.get("state") or "").upper()
checks.append({
"key": name, "name": name,
"status": "COMPLETED" if state in ("SUCCESS", "FAILURE", "ERROR") else "IN_PROGRESS",
"conclusion": {"ERROR": "FAILURE"}.get(state, state) or None,
"details_url": c.get("targetUrl"),
"created_at": c.get("createdAt"),
"started_at": None,
"completed_at": None,
})
else:
wf = c.get("workflowName")
name = c.get("name") or "check"
checks.append({
"key": f"{wf}/{name}" if wf else name, "name": name,
"status": (c.get("status") or "").upper() or "IN_PROGRESS",
"conclusion": (c.get("conclusion") or None) and c["conclusion"].upper(),
"details_url": c.get("detailsUrl"),
"created_at": None,
"started_at": c.get("startedAt"),
"completed_at": c.get("completedAt"),
})
owner, name, host = _resolve_repo_ref(repo, v.get("url"))
merge_identity = None if light else fetch_pr_merge_identity(v.get("number") or pr, owner, name, host)
merge_observation = merge_identity or v
mergeable = merge_observation.get("mergeable")
merge_state_status = merge_observation.get("mergeStateStatus")
head = v.get("headRefOid")
if light:
base = {"host": host or "github.com", "repository": f"{owner}/{name}",
"ref": v.get("baseRefName"), "oid": None, "identity": BASE_REF_PROBE_ERROR}
else:
base = fetch_base_ref(owner, name, v.get("baseRefName"), merge_observation, host)
if merge_identity is not None:
identity_head = merge_identity.get("headRefOid")
identity_base_ref = merge_identity.get("baseRefName")
if (not isinstance(identity_base_ref, str)
or not (_valid_oid(identity_head) and _valid_oid(head))):
base["identity"] = BASE_REF_PROBE_ERROR
elif (identity_head.lower() != head.lower()
or identity_base_ref != v.get("baseRefName")):
# The checks and review payload came from `v`; do not combine them with mergeability
# observed for a different head or base branch.
base["identity"] = BASE_REF_RACE
pr_chain = fetch_pr_chain(v.get("number") or pr, repo, v.get("url"),
v.get("baseRefName"), v.get("headRefName"), owner, name, host)
currency_probe = {
"mergeable": mergeable,
"merge_state_status": merge_state_status,
"base": base,
"pr_chain": pr_chain,
}
host_branch_update_capability = "unknown"
if (merge_state_status == "BEHIND"
and _branch_currency_observation(currency_probe, head) is not None):
capability = merge_observation.get("viewerCanUpdateBranch")
if isinstance(capability, bool):
host_branch_update_capability = capability
return {
"pr_state": v.get("state"),
"pr_is_draft": v.get("isDraft"),
"mergeable": mergeable,
"merge_state_status": merge_state_status,
"review_decision": v.get("reviewDecision") or None,
"head_sha": head,
"base": base,
"host_branch_update_capability": host_branch_update_capability,
"url": v.get("url"),
"checks": checks,
"threads": fetch_threads(pr, owner, name, host),
# Non-thread feedback: top-level PR comments + review submission bodies. ce-resolve-pr-feedback
# handles these too, so a Changes-Requested review body or an actionable top-level comment with
# NO inline thread must not be invisible to the loop. Content-actionability is entirely the
# resolver's judgment; this deterministic layer keeps every non-empty non-author body.
"feedback": _extract_feedback(v),
"awaiting_approval": None if light else fetch_awaiting_approval(owner, name, head, host),
"head_parents": _head_parents(merge_identity),
"pr_chain": pr_chain,
}
def _body_hash(body):
"""Edit identity for a top-level comment / review body: `gh pr view --json` exposes no
updatedAt, so hash the body. Feeds the external_review_moved activity signal (an edit wakes
the watch and extends the settle clock) and `needs-human` reactivation (a human may answer a
parked question by editing the same comment), but deliberately NOT `dispatched` reactivation:
status bots rewrite their bodies on every push, so edit-keyed reactivation of handled items
re-actionized bot comments forever (#1309)."""
return hashlib.sha1((body or "").encode("utf-8")).hexdigest()[:16]
def _extract_feedback(v):
"""Every non-empty top-level PR comment and review body, whoever wrote it.
`gh pr view --json comments,reviews` returns flat arrays (not GraphQL {nodes}). Content,
identity, and surface are evidence for the resolver to judge, not deterministic exclusions --
and that includes the PR author, whose own top-level "please rename X" is the ordinary way a
human asks an agent-opened PR for a change. Excluding the author here made that request
invisible to the tick, so no resolver pass ever ran for it. Loop prevention is the dispatched
mark, which silences a candidate permanently once it has been classified, not identity."""
out = []
for c in v.get("comments") or []:
if (c.get("body") or "").strip():
out.append({"id": c.get("id"), "kind": "comment",
"author": (c.get("author") or {}).get("login"),
"edit_id": _body_hash(c.get("body"))})
for r in v.get("reviews") or []:
if (r.get("body") or "").strip():
out.append({"id": r.get("id"), "kind": "review",
"author": (r.get("author") or {}).get("login"),
"state": r.get("state"), "edit_id": _body_hash(r.get("body"))})
return out
def _head_parents(merge_identity):
"""Parent OIDs of the PR head from the merge-identity GraphQL read, or None when absent. Feeds
the unrequested base-merge detector: a two-parent head whose second parent is the base tip,
produced without a claimed branch_currency item, is a base-into-head merge nobody asked for."""
nodes = ((((merge_identity or {}).get("headRef") or {}).get("target") or {})
.get("parents") or {}).get("nodes")
if not isinstance(nodes, list):
return None
oids = [node.get("oid") if isinstance(node, dict) else None for node in nodes]
return [str(o).lower() for o in oids if _valid_oid(o)]
def fetch_awaiting_approval(owner, name, head, host=None):
"""Count Actions workflow runs on this head that are awaiting maintainer approval —
the fork-PR security gate. Such a run has created NO check-run yet, so it is invisible
to statusCheckRollup; without this, a fork PR blocked on approval reads as 'all checks ok'.
Best-effort: an API/permission failure returns None so the tick can preserve the
last proven gate state rather than treating an unknown result as a proven clear."""
if not head:
return 0
r = _run(["gh", "api", *_host_args(host),
f"repos/{owner}/{name}/actions/runs?head_sha={head}&per_page=50",
"--jq", '[.workflow_runs[] | select(.status==\"action_required\" or .status==\"waiting\" '
'or .conclusion==\"action_required\")] | length'])
if r.returncode != 0:
return None
try:
return int((r.stdout or "").strip() or 0)
except ValueError:
return None
def fetch_threads(pr, owner, name, host=None):
"""Unresolved review threads with their last-comment identity."""
query = """
query($owner:String!,$repo:String!,$pr:Int!,$cursor:String){
repository(owner:$owner,name:$repo){ pullRequest(number:$pr){
reviewThreads(first:100,after:$cursor){
nodes{ id isResolved path line
root: comments(first:1){ nodes{ url } }
comments(last:100){ nodes{ id createdAt lastEditedAt } } }
pageInfo{ hasNextPage endCursor } } } } }"""
out = []
cursor = None
while True:
args = ["gh", "api", "graphql", *_host_args(host), "-f", f"owner={owner}", "-f", f"repo={name}",
"-F", f"pr={pr}", "-f", f"query={query}"]
if cursor:
args += ["-f", f"cursor={cursor}"]
r = _run_checked(args, "gh api graphql")
data = json.loads(r.stdout)["data"]["repository"]["pullRequest"]["reviewThreads"]
for n in data["nodes"]:
if n.get("isResolved"):
continue
cs = n.get("comments", {}).get("nodes") or []
last = cs[-1] if cs else {}
# `last_comment_at` is the reactivation signal — the MAX edit/create time across every
# comment in the thread, not just the last one, so a reviewer editing an *earlier* comment
# (their original request) after the agent's reply still moves the identity and re-opens it.
# Bounded to the last 100 comments (see the query): review threads are ~never that long; an
# edit to a comment outside that window would be missed (acceptable vs paginating per thread).
edit_at = max((c.get("lastEditedAt") or c.get("createdAt") or "" for c in cs), default="")
root = (n.get("root") or {}).get("nodes") or []
out.append({
"thread_id": n["id"],
"last_comment_id": last.get("id"),
"last_comment_at": edit_at or last.get("lastEditedAt") or last.get("createdAt"),
"url": root[0].get("url") if root else None,
"path": n.get("path"),
"line": n.get("line"),
})
if not data["pageInfo"]["hasNextPage"]:
break
cursor = data["pageInfo"]["endCursor"]
return out
def _empty_state(pr, repo, url, now):
owner, name = _split_repo(repo)
created_at = _iso(now)
return {
"pr": {"owner": owner, "repo": name, "number": pr, "url": url},
"head_sha": None, "tick": 0, "state_created_at": created_at,
"started_at": created_at, "invocation_id": None, "invocation_budget_seconds": None,
"last_activity_at": created_at, "dead_time_seconds": 0.0,
"invocation_backstop_seconds": DEFAULT_INVOCATION_BACKSTOP_SECONDS,
"checks": {}, "threads": {}, "ci_dispatched": {},
"human_decisions": [], "answered_human_decisions": [],
"review_decision": None, "mergeable": None, "merge_state_status": None,
"blocked_external_head_sha": None,
"blocked_external_first_seen_at": None,
"feedback_candidate_changed_at": None, "feedback_candidate_ids": [],
"blocked_external_review_last_activity_at": None,
"pr_chain": _default_pr_chain(),
"base": None,
"branch_currency_state": _empty_branch_currency_state(),
"last_change_at": None, "last_action": None, "stop_reason": None,
"watch_generation": None, "watch_pid": None, "watch_process_identity": None,
"trajectory": _empty_trajectory(),
}
def _empty_branch_currency_state():
return {"current_key": None, "items": {}, "semantic_parks": {}, "head_sha": None}
def _load_branch_currency_state(state):
"""Conservatively migrate durable state written before branch currency existed."""
currency = state.get("branch_currency_state")
if not isinstance(currency, dict):
currency = {}
defaults = _empty_branch_currency_state()
for key, value in defaults.items():
currency.setdefault(key, value)
if not isinstance(currency.get("items"), dict):
currency["items"] = {}
for item in currency["items"].values():
if isinstance(item, dict) and item.get("requires_decision"):
item.setdefault("mutation_requires_answer", True)
if not isinstance(currency.get("semantic_parks"), dict):
currency["semantic_parks"] = {}
state["branch_currency_state"] = currency
return currency
def _reactivate_currency_item(currency, item):
"""Re-actionize one source without erasing attempt or mutation history."""
item["disposition"] = DISPOSITION_OPEN
item.pop("claimed_invocation_id", None)
item.pop("confirmed_invocation_id", None)
item.pop("reconciled_invocation_id", None)
item.pop("inspected_semantic_conflict_fingerprint", None)
item.pop("inspection_result", None)
item.pop("requires_decision", None)
item.pop("pending_decision", None)
item.pop("answered_decision", None)
item["parked_semantic_fingerprints"] = sorted(currency.get("semantic_parks") or {})
def _mergeability_certain(mergeable, merge_state_status, base):
"""GitHub mergeability is usable only when its cached base is proven current."""
if _base_ref_blocker(base) is not None:
return False
if mergeable not in ("MERGEABLE", "CONFLICTING"):
return False
if not merge_state_status or merge_state_status == "UNKNOWN":
return False
if merge_state_status == "DIRTY":
return mergeable == "CONFLICTING"
if merge_state_status == "BEHIND":
return mergeable == "MERGEABLE"
return True
def _normal_base_route(cur):
"""Return the only route U1 may classify: an unmanaged PR rooted on its normal base.
Ordinary dependents are intentionally harmless. An open or unknown parent means this target is
itself in a manual dependency chain and therefore outside this unit's authority.
"""
chain = cur.get("pr_chain") or _default_pr_chain()
base = cur.get("base") or {}
if chain.get("manager_status") != MANAGER_ABSENT:
return None
if chain.get("relationship_status") == RELATIONSHIP_PROBE_ERROR:
return None
if not chain.get("default_branch") or base.get("ref") != chain.get("default_branch"):
return None
parents = chain.get("parent_prs") or []
if any(not isinstance(parent, dict) or parent.get("state") not in ("CLOSED", "MERGED")
for parent in parents):
return None
return "normal-base"
def _branch_currency_observation(cur, head):
status = cur.get("merge_state_status")
certain = _mergeability_certain(cur.get("mergeable"), status, cur.get("base"))
if not certain or status not in ("BEHIND", "DIRTY"):
return None
route = _normal_base_route(cur)
base = cur.get("base") or {}
if (route is None or not head
or not all(base.get(key) for key in ("host", "repository", "ref", "oid"))):
return None
identity = {
"host": base["host"],
"base_repository": base["repository"],
"base_ref": base["ref"],
"base_oid": base["oid"],
"head_sha": head,
"status": status,
"route": route,
}
encoded = json.dumps(identity, sort_keys=True, separators=(",", ":")).encode("utf-8")
return {"key": f"currency:{hashlib.sha256(encoded).hexdigest()}", **identity}
def _claimed_currency_items(currency):
"""Return unresolved mutation claims with the current item first.
Base or expected head movement may be the result of the claimed mutation. It is evidence to
reconcile, not permission to discard the claim and start another attempt.
"""
items = currency.get("items") or {}
current_key = currency.get("current_key")
claimed = []
current = items.get(current_key)
if isinstance(current, dict) and current.get("disposition") == CURRENCY_CLAIMED:
claimed.append((current_key, current))
claimed.extend(
(key, item) for key, item in items.items()
if key != current_key and isinstance(item, dict)
and item.get("disposition") == CURRENCY_CLAIMED
)
return claimed
def _prepare_claimed_currency_item(state, item, cur=None, head=None):
invocation_id = state.get("invocation_id")
already_handled_here = (
item.get("claimed_invocation_id") == invocation_id
or item.get("reconciled_invocation_id") == invocation_id
)
base = (cur or {}).get("base") or {}
async_evidence_moved = (
item.get("recovery_state") in (
CURRENCY_OUTCOME_MUTATION_OBSERVED, CURRENCY_OUTCOME_AMBIGUOUS)
and ((head and head != item.get("head_sha"))
or (base.get("oid") and base.get("oid") != item.get("base_oid")))
)
item["attention"] = (CURRENCY_ATTENTION_RECONCILE
if async_evidence_moved or not already_handled_here else None)
item["reconciliation_only"] = True
return item
def _apply_stale_merge_computation(state, cur, head, now):
"""Bounded degrade for a `race` whose only mismatch is GitHub's stale test merge. Tracks the
first time the same (head, live base) pair was seen in that state; after
STALE_MERGE_COMPUTATION_SECONDS the identity becomes `stale-computation` and the base carries
`merge_computation_stale: true` so the report can disclose it. Any head/base movement resets."""
base = cur.get("base") if isinstance(cur.get("base"), dict) else None
if not base:
state["stale_merge_computation"] = None
return
live = base.get("oid")
candidate = (base.get("identity") == BASE_REF_RACE
and base.get("stale_computation_candidate")
and _valid_oid(live) and _valid_oid(head))
if not candidate:
state["stale_merge_computation"] = None
return
key = f"{str(head).lower()}:{str(live).lower()}"
tracked = state.get("stale_merge_computation")
if not isinstance(tracked, dict) or tracked.get("key") != key:
tracked = {"key": key, "first_seen_at": _iso(now)}
state["stale_merge_computation"] = tracked
if _elapsed(tracked.get("first_seen_at"), now) >= STALE_MERGE_COMPUTATION_SECONDS:
base["identity"] = BASE_REF_STALE_COMPUTATION
base["merge_computation_stale"] = True
def _detect_unrequested_base_merge(state, cur, head, head_changed, now):
"""Head-scoped fact: the current head is a two-parent merge of the base tip that no claimed
branch_currency item produced. Ordinary base movement with CLEAN never yields an item, so a
base-into-head merge on such a PR is a defect the loop must report, not maintenance. Persisted
on the head that introduced it and cleared by the next head change; a legitimate DIRTY repair
(mutation_consumed on a claimed item whose head is the merge's other parent) is excluded. A head
whose parents could not be observed stays pending and is re-evaluated every snapshot. The head
seen on the first-ever snapshot of a state dir is the watch's baseline, not a candidate: the
loop has no evidence about who produced it (a user merging main by hand before invoking is
ordinary), and state persists across re-invocations, so only heads that appear during the
watch are classified."""
# Must run before diff() overwrites state["base"] with the current observation.
prior_base = state.get("base") if isinstance(state.get("base"), dict) else {}
persisted = state.get("unrequested_base_merge")
pending = state.get("unrequested_base_merge_pending")
should_evaluate = head_changed or pending == head
if not should_evaluate:
return persisted if isinstance(persisted, dict) and persisted.get("head") == head else None
state["unrequested_base_merge"] = None # cleared unless proven below
if "head_parents" not in cur:
# Pre-detector fixture or persisted snapshot: no parent evidence exists to wait for.
state["unrequested_base_merge_pending"] = None
return None
parents = cur.get("head_parents")
if parents is None:
state["unrequested_base_merge_pending"] = head
return None
state["unrequested_base_merge_pending"] = None
if not isinstance(parents, list) or len(parents) != 2:
return None
base = cur.get("base") if isinstance(cur.get("base"), dict) else {}
base_oids = {
str(o).lower() for o in (
base.get("oid"), base.get("graphql_oid"), base.get("historical_oid"),
prior_base.get("oid"), prior_base.get("graphql_oid"))
if _valid_oid(o)
}
base_parent = next((p for p in parents if p in base_oids), None)
if base_parent is None:
return None
# Only a consumed mutation that produced THIS merge exempts it: the claimed item must name both
# parents — its head is the merge's other parent and its base is the base parent. A confirmed
# item for an earlier head, or for an older base, is not evidence.
other_parents = {p for p in parents if p != base_parent}
for item in (_load_branch_currency_state(state).get("items") or {}).values():
if (isinstance(item, dict) and item.get("mutation_consumed")
and str(item.get("head_sha") or "").lower() in other_parents
and str(item.get("base_oid") or "").lower() == base_parent):
return None
fact = {"head": head, "base_parent": base_parent, "observed_at": _iso(now)}
state["unrequested_base_merge"] = fact
return fact
def _prepare_open_currency_item(state, currency, key, item, now):
parks = currency.get("semantic_parks") or {}
answer_required = _currency_mutation_requires_answer(state, key, item)
inspection_required = (item.get("status") == "DIRTY"
and bool(parks)
and item.get("inspection_result") != "changed"
and not _answered_human_source_matches(
state, "currency", key, _currency_source_observation(item)))
item["inspection_required"] = inspection_required
retry_wait_seconds = 0
retry_not_before = item.get("retry_not_before")
if retry_not_before:
try:
retry_wait_seconds = max(
0, int((_parse_iso8601(retry_not_before) - now).total_seconds()))
except (ValueError, TypeError):
retry_wait_seconds = 0
item["retry_wait_seconds"] = retry_wait_seconds
if answer_required:
item["attention"] = CURRENCY_ATTENTION_DECIDE
elif inspection_required:
item["attention"] = CURRENCY_ATTENTION_INSPECT
elif retry_wait_seconds > 0:
item["attention"] = None
else:
item["attention"] = CURRENCY_ATTENTION_CLAIM
item["reconciliation_only"] = False
def _apply_branch_currency(state, cur, head, head_changed, now):
currency = _load_branch_currency_state(state)
carried_claims = _claimed_currency_items(currency)
carried_semantic_parks = dict(currency.get("semantic_parks") or {})
if head_changed or (currency.get("head_sha") not in (None, head)):
currency = _empty_branch_currency_state()
currency["items"] = dict(carried_claims)
currency["semantic_parks"] = carried_semantic_parks
state["branch_currency_state"] = currency
currency["head_sha"] = head
# An unresolved claim owns the branch-currency lane until it is explicitly reconciled. A base
# move or expected head move may be the mutation's result, so a new observation must not replace
# the old exact key or make the PR look clear in the meantime.
claimed_items = _claimed_currency_items(currency)
if claimed_items:
key, item = claimed_items[0]
currency["current_key"] = key
return _prepare_claimed_currency_item(state, item, cur, head)
observation = _branch_currency_observation(cur, head)
if observation is None:
currency["current_key"] = None
return None
key = observation["key"]
items = currency["items"]
item = items.get(key)
if not isinstance(item, dict):
item = {**observation, "disposition": DISPOSITION_OPEN}
else:
# Refresh non-identity facts without changing an invocation-fenced disposition.
item.update(observation)
item["host_branch_update_capability"] = cur.get(
"host_branch_update_capability", "unknown")
parks = currency.get("semantic_parks") or {}
item["parked_semantic_fingerprints"] = sorted(parks)
item.setdefault("retry_count", 0)
item.setdefault("mutation_consumed", False)
item.setdefault("mutation_requires_answer", False)
disposition = item.get("disposition", DISPOSITION_OPEN)
if disposition == DISPOSITION_OPEN:
_prepare_open_currency_item(state, currency, key, item, now)
elif disposition == CURRENCY_CLAIMED:
_prepare_claimed_currency_item(state, item, cur, head)
else:
item["attention"] = None
item["reconciliation_only"] = False
items[key] = item
currency["current_key"] = key
# Old observations are useful only when they carry an unresolved mutation claim.
retained = {item_key: value for item_key, value in items.items()
if item_key == key
or value.get("disposition") == CURRENCY_CLAIMED}
currency["items"] = retained
return item
def _empty_trajectory():
"""Deterministic cross-tick facts babysit hands the leaves (facts, not judgment):
babysit ships the trajectory, the leaf decides whether it means non-convergence."""
return {
"check_history": {}, # check key -> {state, last_head, recur, seen_tick}
"seen_threads": {}, # unresolved thread id -> first-seen tick
"unresolved_series": [], # unresolved-thread count per tick (last 6)
"stream_series": [], # single-stream activity per tick (last 8)
"problem_keys": [], # last tick's failing checks + non-parked threads (progress detection)
"min_open_problems": None, # lowest total open-problem count seen
"heads_since_progress": 0, # head changes since progress (a new low OR something cleared)
"last_head": None, # head as of the last AGENT tick — hsp counts moves between ticks,
# NOT poll-observed head moves (state["head_sha"] advances on polls)
"invariant_heads": {}, # resolver-supplied invariant_key -> unique heads it was marked on
}
def _load_trajectory(state):
"""Load the trajectory, tolerating a partial or non-dict value from an older on-disk
state.json (persisted in /tmp across script versions) — backfill missing keys so a new
field never KeyErrors an old state, and a null/garbage value never crashes."""
tj = state.get("trajectory")
if not isinstance(tj, dict):
tj = {}
for key, value in _empty_trajectory().items():
tj.setdefault(key, value)
state["trajectory"] = tj
return tj
_INVARIANT_KEY_RE = re.compile(r"^[A-Za-z0-9._:-]{1,120}$")
def _record_invariant_round(state, key, head):
"""Count a resolver-supplied invariant against unique heads. Never infer keys."""
if not key:
return
if not _INVARIANT_KEY_RE.match(key):
raise SystemExit("--invariant-key must be 1-120 chars of A-Za-z0-9._:-")
if not head:
return
tj = _load_trajectory(state)
heads = tj.setdefault("invariant_heads", {})
seen = heads.setdefault(key, [])
if not isinstance(seen, list):
seen = []
heads[key] = seen
if head not in seen:
seen.append(head)
def _invariant_rounds_view(tj):
"""Public trajectory view: resolver keys counted per unique head."""
items = []
for key, heads in sorted((tj.get("invariant_heads") or {}).items()):
if isinstance(heads, list):
items.append({"key": key, "rounds": len(heads)})
return items
def _push_bounded(lst, item, cap):
"""Append to a sliding window that keeps only the last `cap` items."""
lst.append(item)
del lst[:-cap]
def _stream_alternations(series):
"""Count flips between ci-active and review-active ticks — the cross-stream churn signal."""
flips = 0
prev = None
for s in series:
if prev is not None and s != prev:
flips += 1
prev = s
return flips
def _trend(series):
if len(series) < 3:
return "flat"
if series[-1] > series[0]:
return "rising"
if series[-1] < series[0]:
return "falling"
return "flat"
def _record_check_history(state, head, new_checks):
"""Record CI fail->clear->fail recurrence transitions. Runs on EVERY snapshot — agent ticks AND
watch polls — so a CLEAR (or FAIL) observed only between agent ticks is not lost, which would
otherwise make the ping-pong recurrence trigger under-fire under the self-sustaining watch.
Idempotent per transition: once a check is FAILING, re-observing FAIL does not re-increment."""
tj = _load_trajectory(state)
tick = state.get("tick", 0)
hist = tj["check_history"]
for key, c in new_checks.items():
h = hist.setdefault(key, {"state": CHECK_UNKNOWN, "last_head": None, "recur": 0, "seen_tick": tick})
h["seen_tick"] = tick
if c["conclusion"] in FAILING:
if h["state"] == CHECK_CLEAR and h["last_head"] != head: # fail after a clear on a new head
h["recur"] += 1
h["state"] = CHECK_FAILING
h["last_head"] = head
elif c["status"] == "COMPLETED": # observed non-failing terminal — a genuine clear
h["state"] = CHECK_CLEAR
h["last_head"] = head
# IN_PROGRESS/QUEUED: leave prior state untouched (not yet a clear)
# Evict entries unseen for TTL agent-ticks (polls don't advance the tick) — bounds growth without
# erasing a check that was briefly absent (a one-tick gap must not lose real recurrence history).
tj["check_history"] = {k: v for k, v in hist.items() if tick - v.get("seen_tick", tick) <= CHECK_HISTORY_TTL}
def _update_trajectory(state, head, new_checks, new_threads, new_feedback, actionable,
decision_coverage=None):
"""Maintain and emit the deterministic trajectory. Coarse by design: check-name-level
recurrence, backlog trend, cross-stream alternation, no-progress heads. Fine, invariant-
level judgment (log signatures, nit root-clustering) is the leaf's job, not this script's.
`actionable` is the `{ci, threads, comments}` set diff() also returns. Non-thread feedback
(top-level comments + review bodies) counts as review-stream activity and as an open problem
for the stall signal, but the thread-named backlog fields stay scoped to inline threads."""
tj = _load_trajectory(state)
# --- CI recurrence (fail->clear->fail on a *different* head): recorded on EVERY observation
# (agent ticks AND watch polls, via _record_check_history in diff) so a clear seen only between
# agent ticks is not lost. Here we just READ the accumulated history for the trigger. recur_max
# reflects only checks present this tick, so a stale key can't keep it elevated. ---
hist = tj["check_history"]
recur_max = max((hist[k]["recur"] for k in new_checks if k in hist), default=0)
recurring = [{"key": k, "recur": hist[k]["recur"]} for k in new_checks if k in hist and hist[k]["recur"] > 0]
# --- Review-thread backlog: trend of unresolved-thread count + genuinely new threads this
# tick. Scoped to inline threads (non-thread feedback is captured in the total-problem stall
# signal below), so these thread-named fields stay accurate. ---
seen = tj.get("seen_threads", {})
new_arrivals = [tid for tid in new_threads if tid not in seen]
tj["seen_threads"] = {tid: seen.get(tid, state.get("tick", 0)) for tid in new_threads}
_push_bounded(tj["unresolved_series"], len(new_threads), 6)
# --- Cross-stream churn: alternation between ci-only and review-only active ticks.
# Review is active when either threads OR non-thread feedback is actionable. ---
review_active = bool(actionable["threads"] or actionable.get("comments"))
if actionable["ci"] and not review_active:
active = STREAM_CI
elif review_active and not actionable["ci"]:
active = STREAM_REVIEW
else:
active = None # both or neither — not a single-stream tick, don't record
if active:
_push_bounded(tj["stream_series"], active, 8)
# --- No-progress heads: measured from TOTAL open problems (failing checks + non-parked
# unresolved threads + non-parked non-thread feedback), NOT the post-claim `actionable` set —
# marking items dispatched shrinks
# actionable and would fake progress. Reset the counter when the head moves and either the
# total set a new low OR a previously-failing item cleared: progressive migration (A cleared
# while B appears) is progress, not a stall, so it must not accrue heads_since_progress. ---
# Only genuinely-OPEN items are unresolved work: a `dispatched` item is handled (a top-level
# comment never drops out of the fetch, so counting it would keep heads_since_progress climbing
# forever and falsely trip non-convergence on unrelated later work), and `needs-human` is parked.
decision_coverage = decision_coverage or set()
problem_keys = {f"c:{k}" for k, c in new_checks.items()
if c["conclusion"] in FAILING and ("check", k) not in decision_coverage}
problem_keys |= {f"t:{tid}" for tid, t in new_threads.items()
if t.get("disposition") == DISPOSITION_OPEN
and ("thread", tid) not in decision_coverage}
problem_keys |= {f"m:{fid}" for fid, f in new_feedback.items()
if f.get("disposition") == DISPOSITION_OPEN
and (f.get("kind", "comment"), fid) not in decision_coverage}
cleared_something = bool(set(tj.get("problem_keys", [])) - problem_keys)
open_problems = len(problem_keys)
minp = tj.get("min_open_problems")
new_low = minp is None or open_problems < minp
if new_low:
tj["min_open_problems"] = open_problems
# heads_since_progress counts head moves BETWEEN AGENT TICKS (tj["last_head"]), not poll-observed
# moves — a watch poll advances state["head_sha"], so a plain head_changed would read False at the
# agent's tick and starve this stall signal under the default self-sustaining watch.
traj_head_moved = tj.get("last_head") is not None and head != tj.get("last_head")
if new_low or cleared_something:
tj["heads_since_progress"] = 0
elif traj_head_moved:
tj["heads_since_progress"] = tj.get("heads_since_progress", 0) + 1
tj["problem_keys"] = sorted(problem_keys)
tj["last_head"] = head
return {
"recurring_checks": recurring,
"check_recur_max": recur_max,
"unresolved_threads": len(new_threads),
"unresolved_series": list(tj["unresolved_series"]),
"unresolved_trend": _trend(tj["unresolved_series"]),
"new_threads_this_tick": len(new_arrivals),
"stream_alternations": _stream_alternations(tj["stream_series"]),
"heads_since_progress": tj["heads_since_progress"],
"invariant_rounds": _invariant_rounds_view(tj),
}
def _apply_dispositions(items, id_key, prior, identity_fn=None,
reactivate_dispositions=(DISPOSITION_DISPATCHED,)):
"""Claim->act->confirm dedup for a review stream: an item stays actionable until `mark`
records a non-open disposition (persisted in prior[id]['disposition']). Returns
(persisted_by_id, actionable_list, 0).
When identity_fn is given, an item whose disposition is in `reactivate_dispositions` is
*reactivated* — set back to open and re-actionized — once its last-comment / edit identity moves
past the one we acted on. That acted-on identity is captured lazily on the first tick after our
own reply, so the reply becomes the baseline and later reviewer activity reopens the item. This
stops a dispatched-but-unresolved thread with fresh activity from being silenced.
The feedback stream excludes DISPATCHED from `reactivate_dispositions` deliberately: status
bots rewrite their comment bodies on every push, so edit-keyed reactivation re-actionized
handled bot comments forever and merge-ready could never fire (#1309). A brand-new comment is
a new id and is always actionable. Human decisions are a separate coverage layer below; source
dispositions remain ordinary open/dispatched facts."""
persisted, actionable = {}, []
for it in items:
iid = it.get(id_key)
if not iid:
continue
pri = prior.get(iid, {})
disposition = pri.get("disposition", DISPOSITION_OPEN)
acted_identity = pri.get("acted_identity")
if disposition in reactivate_dispositions and identity_fn is not None:
current_identity = identity_fn(it)
if acted_identity is None:
acted_identity = current_identity # first post-action observation: adopt as baseline
elif current_identity != acted_identity:
disposition = DISPOSITION_OPEN # a later human/reviewer reply past our baseline -> reactivate
acted_identity = None
if disposition == DISPOSITION_OPEN:
actionable.append(it)
rec = {**it, "disposition": disposition}
if acted_identity is not None:
rec["acted_identity"] = acted_identity
persisted[iid] = rec
return persisted, actionable, 0
def _needs_human_residual_error(residual):
"""Return why a residual is malformed, or None when the full contract holds."""
if not isinstance(residual, dict) or residual.get("type") != "needs-human":
return "must be one typed needs-human object"
sources = residual.get("sources")
if not isinstance(sources, list) or not sources:
return "sources must be a non-empty list"
seen_sources = set()
for source in sources:
if not isinstance(source, dict):
return "every source must be an object"
source_id, kind = source.get("id"), source.get("kind")
if not isinstance(source_id, str) or not source_id.strip():
return "every source id must be a non-empty string"
if kind not in NEEDS_HUMAN_SOURCE_KINDS:
return "every source kind must be thread, comment, review, check, or currency"
identity = (kind, source_id)
if identity in seen_sources:
return "sources must be unique"
seen_sources.add(identity)
context = residual.get("decision_context")
if not isinstance(context, dict):
return "decision_context must be an object"
for field in ("quoted_feedback", "investigation", "decision_reason"):
if not isinstance(context.get(field), str) or not context[field].strip():
return f"decision_context.{field} must be a non-empty string"
options = context.get("options")
if not isinstance(options, list) or not options:
return "decision_context.options must be a non-empty list"
for option in options:
if not isinstance(option, dict):
return "every decision_context option must be an object"
for field in ("option", "tradeoff"):
if not isinstance(option.get(field), str) or not option[field].strip():
return f"every decision_context option.{field} must be a non-empty string"
if "recommendation" not in context:
return "decision_context.recommendation is required"
recommendation = context.get("recommendation")
if recommendation is not None and (
not isinstance(recommendation, str) or not recommendation.strip()):
return "decision_context.recommendation must be null or a non-empty string"
thread_urls = residual.get("thread_urls")
if not isinstance(thread_urls, list) or any(
not isinstance(url, str) or not url.strip() for url in thread_urls):
return "thread_urls must be a list of non-empty strings"
thread_sources = sum(source["kind"] == "thread" for source in sources)
if len(set(thread_urls)) != len(thread_urls) or len(thread_urls) != thread_sources:
return "thread_urls must contain exactly one unique URL for every thread source"
return None
def _residual_identity(residual):
return json.dumps(residual, sort_keys=True, separators=(",", ":"))
def _source_identity(source):
return source.get("kind"), source.get("id")
def _check_observation_identity(check):
if not isinstance(check, dict):
return None
return {
"head_sha": check.get("head_sha"),
"conclusion": check.get("conclusion"),
"details_url": check.get("details_url"),
"created_at": check.get("created_at"),
"started_at": check.get("started_at"),
"completed_at": check.get("completed_at"),
}
def _currency_source_observation(item):
if not isinstance(item, dict):
return None
observation = {
"key": item.get("key"),
"status": item.get("status"),
"semantic_conflict_fingerprint": item.get("semantic_conflict_fingerprint"),
"recovery_state": item.get("recovery_state"),
"mutation_requires_answer": bool(item.get("mutation_requires_answer", False)),
}
if item.get("status") == "BEHIND" and item.get("recovery_state") is None:
observation["host_branch_update_capability"] = item.get(
"host_branch_update_capability")
return observation
def _current_source_observation(state, kind, source_id, threads, feedback, checks):
if kind == "thread":
record = threads.get(source_id)
return ([record.get("last_comment_id"), record.get("last_comment_at")]
if isinstance(record, dict) else None)
if kind in ("comment", "review"):
record = feedback.get(source_id)
if not isinstance(record, dict) or record.get("kind") != kind:
return None
return [record.get("edit_id")]
if kind == "check":
record = checks.get(source_id)
if not isinstance(record, dict) or record.get("conclusion") not in FAILING:
return None
return _check_observation_identity(record)
if kind == "currency":
currency = _load_branch_currency_state(state)
if currency.get("current_key") != source_id:
return None
return _currency_source_observation(currency.get("items", {}).get(source_id))
return None
def _decision_id(residual, sources):
payload = json.dumps(
{"residual": residual, "sources": sources}, sort_keys=True, separators=(",", ":"))
return f"decision:{hashlib.sha256(payload.encode('utf-8')).hexdigest()}"
def _decision_record(residual, sources, created_at=None):
return {
"id": _decision_id(residual, sources),
"residual": json.loads(json.dumps(residual)),
"sources": json.loads(json.dumps(sources)),
"created_at": created_at,
}
def _install_human_decisions(state, decisions):
"""Install canonical decisions while retaining currency mutation safety independently."""
currency = _load_branch_currency_state(state)
items = currency.get("items", {})
for decision in decisions:
if not isinstance(decision, dict):
continue
sources = decision.get("sources")
if not isinstance(sources, list):
continue
for source in sources:
if not isinstance(source, dict) or source.get("kind") != "currency":
continue
item = items.get(source.get("id"))
if isinstance(item, dict):
item["mutation_requires_answer"] = True
observation = source.get("observation")
if isinstance(observation, dict):
observation["mutation_requires_answer"] = True
if isinstance(decision.get("residual"), dict):
decision["id"] = _decision_id(decision["residual"], sources)
state["human_decisions"] = decisions
def _decision_source_matches(source, kind, source_id, observation):
return (isinstance(source, dict) and source.get("kind") == kind
and source.get("id") == source_id and source.get("observation") == observation)
def _answered_human_source_matches(state, kind, source_id, observation):
return any(
_decision_source_matches(source, kind, source_id, observation)
for decision in state.get("answered_human_decisions") or []
if isinstance(decision, dict)
for source in decision.get("sources") or []
)
def _currency_mutation_requires_answer(state, source_id, item):
observation = _currency_source_observation(item)
return (bool(item.get("mutation_requires_answer"))
and not _answered_human_source_matches(
state, "currency", source_id, observation))
def _source_is_actionable(state, kind, source_id, threads, feedback, checks):
if kind == "thread":
record = threads.get(source_id)
return isinstance(record, dict) and record.get("disposition") == DISPOSITION_OPEN
if kind in ("comment", "review"):
record = feedback.get(source_id)
return (isinstance(record, dict) and record.get("kind") == kind
and record.get("disposition") == DISPOSITION_OPEN)
if kind == "check":
record = checks.get(source_id)
head = state.get("head_sha")
return (isinstance(record, dict) and record.get("conclusion") in FAILING
and source_id not in state.get("ci_dispatched", {}).get(head, []))
if kind == "currency":
currency = _load_branch_currency_state(state)
record = currency.get("items", {}).get(source_id)
return (currency.get("current_key") == source_id and isinstance(record, dict)
and record.get("disposition") in (DISPOSITION_OPEN, CURRENCY_CLAIMED))
return False
def _source_can_enter_human_decision(state, kind, source_id, threads, feedback, checks):
"""Accept exact current evidence even when ordinary dispatch already suppresses it."""
if _current_source_observation(state, kind, source_id, threads, feedback, checks) is None:
return False
if kind != "currency":
return True
currency = _load_branch_currency_state(state)
record = currency.get("items", {}).get(source_id)
return (currency.get("current_key") == source_id and isinstance(record, dict)
and record.get("disposition") in (DISPOSITION_OPEN, CURRENCY_CLAIMED))
def _release_ordinary_suppression_for_decision(state, sources, threads, feedback):
"""Make canonical decision coverage the sole suppression owner for review and CI."""
head = state.get("head_sha")
dispatched_checks = state.get("ci_dispatched", {}).get(head)
for source in sources:
kind, source_id = source.get("kind"), source.get("id")
record = (threads.get(source_id) if kind == "thread"
else feedback.get(source_id) if kind in ("comment", "review") else None)
if isinstance(record, dict):
record["disposition"] = DISPOSITION_OPEN
record.pop("acted_identity", None)
elif kind == "check" and isinstance(dispatched_checks, list):
state["ci_dispatched"][head] = [
key for key in state["ci_dispatched"][head] if key != source_id]
def _migrate_legacy_check_observation(state, source_id, legacy_observations, checks):
legacy = legacy_observations.get(source_id)
current = _current_source_observation(state, "check", source_id, {}, {}, checks)
if (not isinstance(legacy, dict) or not isinstance(current, dict)
or current.get("head_sha") != state.get("head_sha")
or any(current.get(field) != value for field, value in legacy.items())):
return None
return current
def _retire_legacy_check_suppression(state, residual):
if not isinstance(residual, dict) or not isinstance(residual.get("sources"), list):
return
head = state.get("head_sha")
dispatched = state.get("ci_dispatched", {}).get(head)
if not isinstance(dispatched, list):
return
legacy_check_ids = {
source.get("id") for source in residual["sources"]
if isinstance(source, dict) and source.get("kind") == "check"
}
state["ci_dispatched"][head] = [
key for key in dispatched if key not in legacy_check_ids]
def _legacy_decision_records(state, threads, feedback, checks, now):
"""Replace legacy suppression with exact current decisions; incomplete evidence fails open."""
candidates = list(state.get("needs_human_residuals") or [])
candidates.extend(
record.get("needs_human_residual")
for record in list(threads.values()) + list(feedback.values())
if isinstance(record.get("needs_human_residual"), dict)
)
legacy_check_observations = state.get("needs_human_check_observations") or {}
records, seen = [], set()
for residual in candidates:
_retire_legacy_check_suppression(state, residual)
if _needs_human_residual_error(residual):
continue
identity = _residual_identity(residual)
if identity in seen:
continue
seen.add(identity)
sources = []
for source in residual["sources"]:
kind, source_id = _source_identity(source)
observation = None
record = (threads.get(source_id) if kind == "thread"
else feedback.get(source_id) if kind in ("comment", "review") else None)
if isinstance(record, dict) and record.get("acted_identity") is not None:
observation = record.get("acted_identity")
elif kind == "check":
observation = _migrate_legacy_check_observation(
state, source_id, legacy_check_observations, checks)
else:
observation = _current_source_observation(
state, kind, source_id, threads, feedback, checks)
if observation is None:
sources = []
break
sources.append({"kind": kind, "id": source_id, "observation": observation})
if sources:
records.append(_decision_record(residual, sources, _iso(now)))
return records
def _legacy_answered_decision_records(state, threads, feedback, checks):
"""Preserve only the exact source evidence carried by a legacy answer owner."""
records, seen = [], set()
currency = _load_branch_currency_state(state)
for source_id, item in currency.get("items", {}).items():
answered = item.get("answered_decision")
if not isinstance(answered, dict):
continue
residual, answer = answered.get("residual"), answered.get("answer")
_retire_legacy_check_suppression(state, residual)
if _needs_human_residual_error(residual) or not isinstance(answer, str) or not answer.strip():
continue
owns_residual = any(
_source_identity(source) == ("currency", source_id)
for source in residual["sources"])
observation = _current_source_observation(
state, "currency", source_id, threads, feedback, checks)
if not owns_residual or observation is None:
continue
sources = [{"kind": "currency", "id": source_id, "observation": observation}]
record = {
**_decision_record(residual, sources, answered.get("recorded_at")),
"answer": answer.strip(),
"answered_at": answered.get("recorded_at"),
}
if record["id"] not in seen:
seen.add(record["id"])
records.append(record)
return records
def _normalize_legacy_source_state(state, threads, feedback):
for records in (threads, feedback):
for record in records.values():
if record.get("disposition") == DISPOSITION_NEEDS_HUMAN:
record["disposition"] = DISPOSITION_OPEN
record.pop("needs_human_residual", None)
if record.get("disposition") == DISPOSITION_OPEN:
record.pop("acted_identity", None)
currency = _load_branch_currency_state(state)
for item in currency.get("items", {}).values():
if item.get("disposition") == DISPOSITION_NEEDS_HUMAN:
_reactivate_currency_item(currency, item)
item.pop("requires_decision", None)
item.pop("pending_decision", None)
item.pop("answered_decision", None)
state.pop("needs_human_residuals", None)
state.pop("needs_human_check_observations", None)
def _human_decision_identity(decision):
if not isinstance(decision, dict):
return None
residual, sources = decision.get("residual"), decision.get("sources")
if not isinstance(residual, dict) or not isinstance(sources, list):
return None
return _decision_id(residual, sources)
def _replace_human_decisions(
state, kept, previous=None, threads=None, feedback=None, checks=None):
"""Replace canonical decisions and reactivate exact surviving source evidence."""
previous = previous if isinstance(previous, list) else state.get("human_decisions")
previous = previous if isinstance(previous, list) else []
kept_identities = {
identity for decision in kept
if (identity := _human_decision_identity(decision)) is not None
}
removed = [
decision for decision in previous
if (_human_decision_identity(decision) not in kept_identities)
]
_install_human_decisions(state, kept)
threads = threads if isinstance(threads, dict) else state.get("threads", {})
feedback = feedback if isinstance(feedback, dict) else state.get("feedback", {})
checks = checks if isinstance(checks, dict) else state.get("checks", {})
for decision in removed:
if not isinstance(decision, dict):
continue
for source in decision.get("sources") or []:
if not isinstance(source, dict):
continue
kind, source_id = source.get("kind"), source.get("id")
current = _current_source_observation(
state, kind, source_id, threads, feedback, checks)
if current != source.get("observation") or kind != "currency":
continue
# Review and CI retain ordinary open state; currency has an explicit claim to clear.
currency = _load_branch_currency_state(state)
item = currency.get("items", {}).get(source_id)
if isinstance(item, dict):
_reactivate_currency_item(currency, item)
def _reconcile_human_decisions(state, threads, feedback, checks, now):
"""Let canonical decisions alone suppress exact evidence; activity is never an answer."""
current = state.get("human_decisions")
if not isinstance(current, list):
current = _legacy_decision_records(state, threads, feedback, checks, now)
elif not current and state.get("needs_human_residuals"):
current = _legacy_decision_records(state, threads, feedback, checks, now)
_install_human_decisions(state, current)
answered = state.get("answered_human_decisions")
if not isinstance(answered, list):
answered = _legacy_answered_decision_records(
state, threads, feedback, checks)
elif not answered:
answered = _legacy_answered_decision_records(
state, threads, feedback, checks)
_normalize_legacy_source_state(state, threads, feedback)
kept, seen = [], set()
for decision in current:
if not isinstance(decision, dict) or _needs_human_residual_error(decision.get("residual")):
continue
sources = decision.get("sources")
if not isinstance(sources, list) or len(sources) != len(decision["residual"]["sources"]):
continue
expected = {(source.get("kind"), source.get("id")) for source in sources}
if expected != {_source_identity(source) for source in decision["residual"]["sources"]}:
continue
if not all(
source.get("observation") is not None
and _current_source_observation(
state, source.get("kind"), source.get("id"), threads, feedback, checks)
== source.get("observation")
for source in sources):
continue
normalized = _decision_record(decision["residual"], sources, decision.get("created_at"))
if normalized["id"] in seen:
continue
seen.add(normalized["id"])
kept.append(normalized)
_replace_human_decisions(
state, kept, current, threads=threads, feedback=feedback, checks=checks)
kept_answers = []
for decision in answered:
if not isinstance(decision, dict) or not isinstance(decision.get("answer"), str):
continue
sources = decision.get("sources") or []
matching_actionable = any(
_current_source_observation(
state, source.get("kind"), source.get("id"), threads, feedback, checks)
== source.get("observation")
and _source_is_actionable(
state, source.get("kind"), source.get("id"), threads, feedback, checks)
for source in sources if isinstance(source, dict))
if matching_actionable:
kept_answers.append(decision)
state["answered_human_decisions"] = kept_answers
covered = {
(source["kind"], source["id"])
for decision in kept for source in decision["sources"]
}
public = []
for decision in kept:
public.append(json.loads(json.dumps(decision["residual"])))
return public, covered
def diff(state, cur, now=None, advance_trajectory=True):
"""Pure: given prior state + current snapshot, compute the actionable set
and the persisted observed state. `now` is injectable for tests."""
now = now or _now()
# A transient null/empty head (a gh hiccup) falls back to the last known head,
# so a momentary null does not look like a new head and wipe ci_dispatched.
head = cur["head_sha"] or state.get("head_sha")
head_changed = state.get("head_sha") is not None and head != state["head_sha"]
_apply_stale_merge_computation(state, cur, head, now)
prior_threads = state.get("threads", {})
prior_feedback = state.get("feedback", {})
prior_feedback_activity = {
fid: item.get("edit_id") for fid, item in prior_feedback.items()
}
prior_review_decision = state.get("review_decision")
prior_change_sig = _change_sig(state)
if head_changed:
# SHA-scoped state is meaningless on a new head.
state["ci_dispatched"] = {}
# --- CI: a failing check on the current head is actionable until the agent
# marks it dispatched (recorded in ci_dispatched[head]) or the head moves.
# `checks_terminal` = every check has finished (none IN_PROGRESS/QUEUED).
# Duplicate check keys (same workflow/name) are disambiguated with a #n suffix
# so one never shadows another and silently drops a failing check. ---
new_checks = {}
failing_check_items = {}
has_failing = False
checks_terminal = True
seen_keys = {}
for c in cur["checks"]:
key = c["key"]
if key in seen_keys:
seen_keys[key] += 1
key = f"{key}#{seen_keys[key]}"
else:
seen_keys[key] = 0
new_checks[key] = {
"name": c["name"], "status": c["status"],
"conclusion": c["conclusion"], "head_sha": head,
"details_url": c.get("details_url"),
"created_at": c.get("created_at"),
"started_at": c.get("started_at"),
"completed_at": c.get("completed_at"),
}
if c["status"] != "COMPLETED":
checks_terminal = False
if c["conclusion"] in FAILING:
has_failing = True
failing_check_items[key] = {
"key": key, "name": c["name"], "conclusion": c["conclusion"],
"details_url": c["details_url"],
}
# Review streams keep only ordinary open/dispatched state. Human decisions are a separate
# exact-observation coverage layer reconciled below, so a remote edit invalidates a question
# but is never mistaken for the user's answer.
new_threads, _, _ = _apply_dispositions(
cur["threads"], "thread_id", prior_threads,
identity_fn=lambda t: [t.get("last_comment_id"), t.get("last_comment_at")])
new_feedback, _, _ = _apply_dispositions(
cur.get("feedback") or [], "id", prior_feedback,
identity_fn=lambda c: [c.get("edit_id")], reactivate_dispositions=())
# Reconcile currency before exact decision observations are compared.
branch_currency = _apply_branch_currency(state, cur, head, head_changed, now)
needs_human_residuals, decision_coverage = _reconcile_human_decisions(
state, new_threads, new_feedback, new_checks, now)
dispatched = set(state.get("ci_dispatched", {}).get(head, []))
# Every attention stream is derived from ordinary state minus current decision coverage.
actionable_threads = [
record for source_id, record in new_threads.items()
if record.get("disposition") == DISPOSITION_OPEN
and ("thread", source_id) not in decision_coverage]
actionable_feedback = [
record for source_id, record in new_feedback.items()
if record.get("disposition") == DISPOSITION_OPEN
and (record.get("kind", "comment"), source_id) not in decision_coverage]
actionable_ci = [
item for key, item in failing_check_items.items()
if key not in dispatched and ("check", key) not in decision_coverage]
current_thread_activity = {
item.get("thread_id"): _thread_identity(item)
for item in cur["threads"] if item.get("thread_id")
}
prior_thread_activity = {
tid: _prior_thread_activity_identity(item) for tid, item in prior_threads.items()
if not (item.get("disposition") == DISPOSITION_DISPATCHED
and tid not in current_thread_activity)
}
current_feedback_activity = {
item.get("id"): item.get("edit_id")
for item in (cur.get("feedback") or []) if item.get("id")
}
external_review_moved = (
head_changed
or current_thread_activity != prior_thread_activity
or current_feedback_activity != prior_feedback_activity
or cur.get("review_decision") != prior_review_decision
)
open_needs_human = len(needs_human_residuals)
# Compatibility view for consumers that display source identities. The canonical set above,
# not this flattened list, owns persistence and stop behavior.
needs_human_ids = sorted({
source["id"] for residual in needs_human_residuals for source in residual["sources"]})
state["head_sha"] = head or state.get("head_sha")
state["checks"] = new_checks
state["threads"] = new_threads
state["feedback"] = new_feedback
state["review_decision"] = cur["review_decision"]
state["mergeable"] = cur["mergeable"]
state["merge_state_status"] = cur["merge_state_status"]
observed_awaiting_approval = cur.get("awaiting_approval")
approval_probe_succeeded = observed_awaiting_approval is not None
awaiting_approval = (
observed_awaiting_approval
if approval_probe_succeeded
else state.get("awaiting_approval", 0)
)
state["awaiting_approval"] = awaiting_approval
if awaiting_approval > 0:
gate_is_new_for_head = (
state.get("blocked_external_head_sha") != head
or not state.get("blocked_external_review_last_activity_at")
)
if gate_is_new_for_head:
state["blocked_external_head_sha"] = head
state["blocked_external_first_seen_at"] = _iso(now)
state["blocked_external_review_last_activity_at"] = _iso(now)
elif external_review_moved:
state["blocked_external_review_last_activity_at"] = _iso(now)
elif approval_probe_succeeded:
state["blocked_external_head_sha"] = None
state["blocked_external_first_seen_at"] = None
state["blocked_external_review_last_activity_at"] = None
state["pr_chain"] = cur.get("pr_chain") or _default_pr_chain()
unrequested_base_merge = _detect_unrequested_base_merge(state, cur, head, head_changed, now)
unrequested_base_merge_pending = state.get("unrequested_base_merge_pending") == head
if unrequested_base_merge_pending and isinstance(cur.get("base"), dict) \
and _base_ref_blocker(cur["base"]) is None:
# An unclassified head is a pending merge computation: whatever else the base identity says
# (current or a bounded stale-computation degrade), it is not certain, not ready, not landable.
cur["base"]["identity"] = BASE_REF_MERGEABILITY_PENDING
state["base"] = cur.get("base")
if branch_currency and branch_currency.get("disposition") == DISPOSITION_OPEN:
_prepare_open_currency_item(
state, _load_branch_currency_state(state), branch_currency["key"],
branch_currency, now)
if branch_currency and ("currency", branch_currency.get("key")) in decision_coverage:
branch_currency["attention"] = None
actionable = {"ci": actionable_ci, "threads": actionable_threads, "comments": actionable_feedback}
# Reset on any change to the candidate id-set so a burst arriving over several polls settles
# into one wake; clear when it empties so the next candidate starts a fresh window.
feedback_ids = sorted(str(f.get("id")) for f in actionable_feedback)
if not feedback_ids:
state["feedback_candidate_changed_at"] = None
elif feedback_ids != state.get("feedback_candidate_ids"):
state["feedback_candidate_changed_at"] = _iso(now)
state["feedback_candidate_ids"] = feedback_ids
if advance_trajectory:
state["tick"] = state.get("tick", 0) + 1
# Record CI recurrence transitions on EVERY observation (polls too) so a fail->clear->fail seen
# only between agent ticks is not lost. head_sha for the check-level last_head is the observed
# head; the trajectory-level last_head (for heads_since_progress) is agent-tick-only, inside
# _update_trajectory.
_record_check_history(state, head, new_checks)
if advance_trajectory:
trajectory = _update_trajectory(
state, head, new_checks, new_threads, new_feedback, actionable, decision_coverage)
else:
# A watch poll detects change and advances the settle clock, but must NOT roll the rest of the
# trajectory (tick counter, seen_threads, unresolved_series, heads_since_progress). Advancing
# it would consume new_threads_this_tick — the waking poll marks the just-arrived thread
# "seen", so the agent's real tick reads 0 new arrivals and the review-bot-treadmill
# non-convergence trigger never fires. Only the agent's tick (advance_trajectory=True) rolls it.
trajectory = {}
# --- Settle window: any observable movement resets the quiet clock. ---
changed_this_tick = head_changed or _change_sig(state) != prior_change_sig or state.get("last_change_at") is None
if changed_this_tick:
state["last_change_at"] = _iso(now)
quiet_seconds = _elapsed(state.get("last_change_at"), now)
# Workflow runs awaiting maintainer approval (fork-PR gate) create no check-run, so they are
# invisible to the rollup above — surface them, and never call CI "ok" while the real CI is
# gated. blocked_external = the loop cannot progress (no failing check to fix, but CI can't run)
# and no one in this loop can unblock it — it is up to a maintainer of the base repo.
awaiting_approval = state["awaiting_approval"]
blocked_external_review_quiet_seconds = (
_elapsed(state.get("blocked_external_review_last_activity_at"), now)
if awaiting_approval > 0 else 0
)
# "OK" requires every check terminal, none failing, AND none gated on approval. A still-
# IN_PROGRESS or awaiting-approval check is neither ok nor failing — do not exit green.
all_checks_ok = checks_terminal and not has_failing and bool(cur["checks"]) and awaiting_approval == 0
blocked_external = (awaiting_approval > 0 and not has_failing and checks_terminal
and not actionable_threads and not actionable_feedback)
invocation_elapsed = _active_elapsed(state, now)
return {
"pr_state": cur["pr_state"],
"pr_is_draft": cur.get("pr_is_draft"),
"mergeable": cur["mergeable"],
"merge_state_status": cur["merge_state_status"],
"review_decision": cur["review_decision"],
"head_sha": head,
"head_changed": head_changed,
"base": state.get("base"),
"mergeability_certain": _mergeability_certain(
cur.get("mergeable"), cur.get("merge_state_status"), cur.get("base")),
"base_ref_blocker": _base_ref_blocker(cur.get("base")),
"host_branch_update_capability": cur.get("host_branch_update_capability", "unknown"),
"branch_currency": branch_currency,
"unrequested_base_merge": unrequested_base_merge,
"unrequested_base_merge_pending": unrequested_base_merge_pending,
"branch_currency_blocker": ({
"key": branch_currency.get("key"),
"disposition": branch_currency.get("disposition"),
"recovery_state": branch_currency.get("recovery_state"),
} if branch_currency else None),
"url": cur["url"],
"has_failing_checks": has_failing,
"checks_terminal": checks_terminal,
"checks_present": bool(cur["checks"]),
"all_checks_ok": all_checks_ok,
"checks_awaiting_approval": awaiting_approval,
"blocked_external": blocked_external,
"blocked_external_first_seen_at": state.get("blocked_external_first_seen_at"),
"feedback_candidate_changed_at": state.get("feedback_candidate_changed_at"),
"blocked_external_review_last_activity_at": state.get(
"blocked_external_review_last_activity_at"),
"blocked_external_review_quiet_seconds": blocked_external_review_quiet_seconds,
"blocked_external_review_moved_this_tick": (
awaiting_approval > 0 and external_review_moved),
"pr_chain": state["pr_chain"],
"stack_blocker": _stack_blocker(state["pr_chain"], cur),
"open_needs_human": open_needs_human,
"needs_human_ids": needs_human_ids,
"needs_human_residuals": needs_human_residuals,
"human_decisions": [
{"decision_id": decision.get("id"), "residual": decision.get("residual")}
for decision in state.get("human_decisions", []) if isinstance(decision, dict)
],
"answered_human_decisions": state.get("answered_human_decisions", []),
"actionable": {"threads": actionable_threads, "ci": actionable_ci, "comments": actionable_feedback},
"counts": {"threads": len(actionable_threads), "ci": len(actionable_ci),
"comments": len(actionable_feedback), "needs_human": open_needs_human},
"changed_this_tick": changed_this_tick,
"quiet_seconds": quiet_seconds,
"invocation_id": state.get("invocation_id"),
"invocation_started_at": state.get("started_at"),
"invocation_elapsed_seconds": invocation_elapsed,
"invocation_budget_seconds": state.get("invocation_budget_seconds"),
"invocation_remaining_seconds": max(
0,
(state.get("invocation_budget_seconds") or 0) - invocation_elapsed,
),
"invocation_wall_elapsed_seconds": _elapsed(state.get("started_at"), now),
"invocation_dead_time_seconds": int(state.get("dead_time_seconds") or 0),
"invocation_backstop_seconds": state.get("invocation_backstop_seconds"),
"persisted_state_created_at": state.get("state_created_at"),
"persisted_state_age_seconds": _elapsed(state.get("state_created_at"), now),
# Compatibility aliases for callers migrating from the original single-clock contract.
"session_started_at": state.get("started_at"),
"session_seconds": invocation_elapsed,
"watch_generation": state.get("watch_generation"),
"tick": state["tick"],
"trajectory": trajectory,
}, state
def _settle_base(base):
"""Base fields whose movement should reset the settle clock.
`stale-computation` is a presentation of a still-stable `race` candidate
(same head + live base) plus a report flag. Hashing those presentation
fields would reset last_change_at and add a synthetic settle delay.
Genuine identity values and freshness stay in the signature.
"""
if not isinstance(base, dict):
return base
payload = {k: v for k, v in base.items() if k != "merge_computation_stale"}
if payload.get("identity") == BASE_REF_STALE_COMPUTATION:
payload["identity"] = BASE_REF_RACE
return payload
def _change_sig(state):
"""Everything whose movement should reset the settle clock."""
checks = {k: (v.get("status"), v.get("conclusion")) for k, v in state.get("checks", {}).items()}
threads = {tid: _thread_identity(v) for tid, v in state.get("threads", {}).items()}
# edit_id is part of the signal so a body edit of a silenced item still resets the settle
# clock (review activity happened) even though it no longer reopens the item (#1309).
feedback = {fid: (v.get("disposition"), v.get("edit_id"))
for fid, v in state.get("feedback", {}).items()}
human_decisions = sorted(
item.get("id") for item in state.get("human_decisions", [])
if isinstance(item, dict))
currency = _load_branch_currency_state(state)
current_currency = currency.get("items", {}).get(currency.get("current_key")) or {}
return (checks, threads, feedback, human_decisions,
state.get("review_decision"), state.get("mergeable"),
state.get("merge_state_status"),
json.dumps(state.get("pr_chain") or _default_pr_chain(), sort_keys=True),
json.dumps(_settle_base(state.get("base")), sort_keys=True),
(current_currency.get("host_branch_update_capability")
if current_currency.get("status") == "BEHIND" else None),
currency.get("current_key"), current_currency.get("disposition"),
current_currency.get("attention"), current_currency.get("recovery_state"),
current_currency.get("retry_count"), current_currency.get("mutation_consumed"),
# awaiting-approval clearing is movement: a fork gate lifting must reset the settle clock
# so merge-ready waits for the now-imminent check-runs instead of firing on an empty rollup.
bool(state.get("awaiting_approval")))
def _stack_blocker(chain, cur=None):
"""Return the manager-currency residual that blocks target readiness, if any.
Manager freshness (`needsRebase`) fires on plain trunk drift, which GitHub already prices into
the target's own mergeability against its parent base; so a certain MERGEABLE/CLEAN read defers
to GitHub and only a stale-or-unknown target that GitHub does not clear blocks."""
status = (chain or {}).get("manager_status")
if status == MANAGER_PROBE_ERROR:
return "manager-probe-error"
if (chain or {}).get("relationship_status") == RELATIONSHIP_PROBE_ERROR:
return "relationship-probe-error"
if status == MANAGER_CONFIRMED:
freshness = chain.get("target_needs_rebase")
if freshness is False:
return None
cur = cur or {}
if (cur.get("mergeable") == "MERGEABLE" and cur.get("merge_state_status") == "CLEAN"
and _mergeability_certain(cur.get("mergeable"), cur.get("merge_state_status"),
cur.get("base"))):
return None
if freshness is True:
return "target-needs-rebase"
return "managed-freshness-unknown"
return None
def _parse_iso8601(value):
if isinstance(value, str) and value.endswith("Z"):
value = value[:-1] + "+00:00"
return datetime.fromisoformat(value)
def _elapsed(iso_str, now):
try:
return int((now - _parse_iso8601(iso_str)).total_seconds())
except (ValueError, TypeError):
return 0
def _advance_activity(state, now, threshold=DEAD_TIME_THRESHOLD_SECONDS, accumulate=True):
"""Advance the activity heartbeat, and (for the in-session watch) accumulate suspended time.
Every watch poll and every agent snapshot/mark marks activity. A gap wider than `threshold`
means the whole process was not running — a suspended machine — so, when accumulating, the
excess beyond the threshold is charged to dead time and later excluded from active elapsed.
Agent snapshots/marks bump the heartbeat with `accumulate=False`, so active agent work keeps the
heartbeat fresh and is not refunded as long as the agent touches state (snapshots or marks) within
the threshold; a single silent tick longer than the threshold is the coarse-discriminator limit
the plan defers, bounded by the wall-clock backstop. Dead time only ever grows in the watch loop.
Clock-backward safe: `_elapsed` floors negative deltas at 0 and the accumulator never decreases.
"""
last = state.get("last_activity_at") or state.get("started_at")
if accumulate:
gap = _elapsed(last, now)
if gap > threshold:
state["dead_time_seconds"] = (state.get("dead_time_seconds") or 0) + (gap - threshold)
state["last_activity_at"] = _iso(now)
def _active_elapsed(state, now):
"""Wall-clock elapsed since the invocation anchor, minus accumulated suspended (dead) time."""
wall = _elapsed(state.get("started_at"), now)
return max(0, wall - int(state.get("dead_time_seconds") or 0))
def _session_started_at(value):
"""Parse one invocation-wide, timezone-aware budget anchor for reuse across state dirs."""
try:
parsed = _parse_iso8601(value)
except (ValueError, TypeError):
raise argparse.ArgumentTypeError("must be an ISO-8601 timestamp")
if parsed.tzinfo is None:
raise argparse.ArgumentTypeError("must include a timezone")
return _iso(parsed.astimezone(timezone.utc))
@contextmanager
def locked_state(state_dir, pr, repo, now):
state_path = os.path.join(state_dir, "state.json")
with _state_lock(state_dir):
if os.path.exists(state_path):
with open(state_path) as f:
state = json.load(f)
else:
state = _empty_state(pr, repo, None, now)
# Older state used `started_at` for both durable-state age and the active budget. Preserve
# that original value as the state birth time before a new invocation replaces the clock.
state.setdefault("state_created_at", state.get("started_at") or _iso(now))
# Active-time accounting fields migrate onto legacy state on first observation. Seed the
# heartbeat to *now*, not the old anchor: there is no recorded activity for the pre-migration
# period, so charging that whole span as one supra-threshold gap would refund the entire
# historical invocation (a 9h-old 8h run would read as ~15 min active and never max-runtime).
# Seeding to load time makes the first poll see a ~0 gap and keeps elapsed on real wall-clock.
state.setdefault("last_activity_at", _iso(now))
state.setdefault("dead_time_seconds", 0.0)
state.setdefault("invocation_backstop_seconds", DEFAULT_INVOCATION_BACKSTOP_SECONDS)
box = {"state": state}
yield box
tmp = tempfile.NamedTemporaryFile("w", dir=state_dir, delete=False)
json.dump(box["state"], tmp, indent=2)
tmp.flush()
os.fsync(tmp.fileno())
tmp.close()
_replace_atomic(tmp.name, state_path)
def _fetch_snapshot(args):
"""Fetch current PR state without mutating the persisted babysit state."""
if args.fetch_file:
with open(args.fetch_file) as f:
return json.load(f)
return fetch(args.pr, args.repo)
def _apply_invocation(box, args, now):
"""Start explicitly, or reuse the one fixed clock named by an explicit invocation token."""
state = box["state"]
start_requested = bool(getattr(args, "start_invocation", False)
or getattr(args, "reset_session", False))
continue_requested = bool(getattr(args, "continue_invocation", False))
invocation_id = getattr(args, "invocation_id", None)
requested_started_at = getattr(args, "session_started_at", None)
requested_budget = getattr(args, "invocation_budget_seconds", None)
if start_requested:
invocation_id = uuid.uuid4().hex
requested_started_at = _iso(now)
requested_budget = requested_budget or DEFAULT_INVOCATION_BUDGET_SECONDS
state["invocation_id"] = invocation_id
state["started_at"] = requested_started_at
state["invocation_budget_seconds"] = requested_budget
# A fresh invocation resets the active-time clock: heartbeat at now, no dead time yet.
state["last_activity_at"] = _iso(now)
state["dead_time_seconds"] = 0.0
state.setdefault("invocation_backstop_seconds", DEFAULT_INVOCATION_BACKSTOP_SECONDS)
args.invocation_id = invocation_id
args.session_started_at = requested_started_at
args.invocation_budget_seconds = requested_budget
args._started_new_invocation = True
return
if not invocation_id:
raise SystemExit("pass --start-invocation on the first snapshot or --invocation-id to resume it")
if continue_requested:
if not requested_started_at or not requested_budget:
raise SystemExit("--continue-invocation requires its fixed start and budget")
if state.get("invocation_id") == invocation_id:
try:
same_anchor = (_parse_iso8601(state.get("started_at"))
== _parse_iso8601(requested_started_at))
except (ValueError, TypeError):
same_anchor = False
if not same_anchor or state.get("invocation_budget_seconds") != requested_budget:
raise SystemExit("continuation cannot renew or extend an existing invocation")
# Honor a carry arg on a same-invocation re-continue too, as a monotonic floor: raise the
# dead-time to the carried value without clobbering time this dir already accumulated
# (dead time only ever grows). This makes the carry order-independent — it lands whether
# it arrives on the first adopt or a later re-continue.
continue_dead_time = getattr(args, "continue_dead_time_seconds", None)
if continue_dead_time is not None:
state["dead_time_seconds"] = max(
state.get("dead_time_seconds") or 0, float(continue_dead_time))
return
state["invocation_id"] = invocation_id
state["started_at"] = requested_started_at
state["invocation_budget_seconds"] = requested_budget
# Adopting an invocation into this state dir establishes a fresh active-time clock. Seed the
# heartbeat to now (not the shared anchor, which may be hours old) so the first poll sees a
# ~0 gap and does not refund the span between the anchor and adoption as dead time.
state["last_activity_at"] = _iso(now)
# Dead time is per-state-dir, but a managed-stack continuation shares one active-time budget
# across layers, so it explicitly carries the prior layer's accumulated dead time here. Absent
# that arg (an ordinary adopt of an unrelated dir), reset to 0 rather than inherit stale state.
continue_dead_time = getattr(args, "continue_dead_time_seconds", None)
state["dead_time_seconds"] = float(continue_dead_time) if continue_dead_time is not None else 0.0
state.setdefault("invocation_backstop_seconds", DEFAULT_INVOCATION_BACKSTOP_SECONDS)
return
persisted_id = state.get("invocation_id")
persisted_started_at = state.get("started_at")
persisted_budget = state.get("invocation_budget_seconds")
if persisted_id != invocation_id:
raise SystemExit("invocation token does not match persisted state; start or continue explicitly")
try:
anchors_match = (_parse_iso8601(persisted_started_at)
== _parse_iso8601(requested_started_at))
except (ValueError, TypeError):
anchors_match = False
if not anchors_match:
raise SystemExit("invocation token does not match the persisted budget anchor")
if persisted_budget != requested_budget:
raise SystemExit("invocation token does not match the persisted fixed budget")
def _apply_snapshot(box, args, cur, now, advance_trajectory):
"""Apply one fetched snapshot to a caller-owned locked state box."""
if box["state"].get("pr", {}).get("url") is None and cur.get("url"):
box["state"]["pr"]["url"] = cur["url"]
_apply_invocation(box, args, now)
return diff(box["state"], cur, now, advance_trajectory=advance_trajectory)
def _run_snapshot(args, now, advance_trajectory=True, watch_generation=None):
"""One fetch -> diff -> persist. Returns the actionable/state dict."""
cur = _fetch_snapshot(args)
# Re-capture the clock *after* the (blocking) fetch so activity accounting reflects post-fetch
# time: a suspend that lands during the fetch is then credited as dead time by this poll's
# _advance_activity, instead of leaving the next top-of-loop budget check to fire max-runtime on
# a gap it hasn't yet excluded. It also keeps a new invocation's slow first fetch off its budget.
now = _now()
with locked_state(args.state_dir, args.pr, args.repo, now) as box:
if watch_generation is not None and box["state"].get("watch_generation") != watch_generation:
raise _WatchSuperseded()
if (watch_generation is not None
and box["state"].get("invocation_id") != getattr(args, "invocation_id", None)):
raise _InvocationSuperseded(box["state"].get("invocation_id"))
# Only the self-sustaining in-session watch (watch_generation set) accumulates dead time;
# an agent-driven snapshot bumps the heartbeat without accumulating, so checkpoint/durable
# runs — which have no continuous poll cadence — keep pure wall-clock accounting.
if box["state"].get("started_at"):
_advance_activity(box["state"], now, accumulate=watch_generation is not None)
actionable, box["state"] = _apply_snapshot(box, args, cur, now, advance_trajectory)
return actionable
def cmd_snapshot(args):
result = _run_snapshot(args, _now())
if (getattr(args, "_started_new_invocation", False)
and _elapsed(result.get("invocation_started_at"), _now()) > 60):
raise SystemExit("new invocation unexpectedly inherited more than 60 seconds of elapsed time")
print(json.dumps(result, separators=(",", ":")))
def _wake_reason(a, settle_seconds, feedback_coalesce_seconds=FEEDBACK_COALESCE_SECONDS):
"""Why the in-session agent should wake and run a tick, or None to keep watching.
Ordered by precedence — a terminal/blocked/needs-human state outranks a merge-ready read."""
if a.get("pr_state") in ("MERGED", "CLOSED"):
return "terminal"
c = a.get("counts") or {}
if c.get("threads", 0) or c.get("ci", 0):
return "actionable"
if c.get("comments", 0):
# Coalesced per FEEDBACK_COALESCE_SECONDS above; falling through costs nothing, since the
# set stays actionable and whichever tick runs still sees every candidate. `_elapsed` reads
# a missing or unparseable clock as 0, so unknown must wake rather than hold forever.
changed_at = a.get("feedback_candidate_changed_at")
if changed_at is None or _elapsed(changed_at, _now()) >= feedback_coalesce_seconds:
return "feedback-candidate"
if a.get("stack_blocker"):
return "stack-blocked"
base_ref_blocker = (a.get("base_ref_blocker")
if "base_ref_blocker" in a else _base_ref_blocker(a.get("base")))
if base_ref_blocker:
return "base-ref-blocked"
if a.get("unrequested_base_merge"):
return "unrequested-base-merge"
branch_currency = a.get("branch_currency") or {}
if a.get("has_failing_checks") and a.get("checks_terminal"):
# a dispatched check left terminally red (counts.ci is 0 — nothing new to dispatch) is a
# blocker to hand back, not a reason to idle to max-runtime.
return "blocked-failing"
# Branch maintenance is a third attention stream. It runs only after new review/CI work and
# standing red blockers are clear, but a passive in-progress check or review signal does not
# delay a BEHIND update that will invalidate that work anyway. Claimed observations wake only
# for reconciliation; confirmed and parked residuals remain quiet.
if branch_currency.get("attention") in (CURRENCY_ATTENTION_CLAIM,
CURRENCY_ATTENTION_DECIDE,
CURRENCY_ATTENTION_INSPECT,
CURRENCY_ATTENTION_RECONCILE):
return "branch-currency"
if a.get("blocked_external"):
return "blocked-external"
if a.get("open_needs_human", 0):
return "needs-human" # parked items block ready; surface them
# A widened window is the agent's own judgment about evidence it looked at. When that evidence
# moves, the judgment may no longer hold, so hand the decision back rather than sleeping out a
# window armed against a state that has since changed. Movement is the whole trigger: unchanged
# evidence never fires this, which is what keeps it from re-deciding something the agent already
# settled. The engine reports that something moved and nothing more — whether the movement
# resolved anything is the agent's call at the settle decision.
if settle_seconds > DEFAULT_SETTLE_SECONDS and a.get("changed_this_tick"):
return "review-evidence-moved"
# merge-ready candidate: green, settled, nothing left in the attention set. Whether a review is
# still on its way is the agent's judgment at the settle decision, from live GitHub evidence the
# engine does not model; this wake hands it that decision rather than pre-empting it.
# Interactive merge-ready does NOT require `all_checks_ok`'s "at least one observed check": a repo
# with no configured checks has a CLEAN/MERGEABLE PR that should be callable ready. (That guard
# stays in pipeline success, where a not-yet-created rollup must not read as a pass.)
mergeability_certain = a.get("mergeability_certain")
if mergeability_certain is None:
mergeability_certain = _mergeability_certain(
a.get("mergeable"), a.get("merge_state_status"), a.get("base"))
if (mergeability_certain and a.get("mergeable") == "MERGEABLE"
and a.get("merge_state_status") == "CLEAN"
and a.get("checks_terminal") and not a.get("has_failing_checks")
and a.get("checks_awaiting_approval", 0) == 0
and a.get("branch_currency_blocker") is None
and a.get("quiet_seconds", 0) >= settle_seconds):
return "merge-ready"
return None
def _emit_wake(reason, **fields):
print(json.dumps({"event": "BABYSIT_WAKE", "reason": reason, **fields}), flush=True)
def _persisted_watch_generation(args):
"""Best-effort read of current ownership without creating or mutating watch state."""
try:
with open(os.path.join(args.state_dir, "state.json")) as f:
state = json.load(f)
except (OSError, json.JSONDecodeError):
return None
generation = state.get("watch_generation") if isinstance(state, dict) else None
return generation if isinstance(generation, str) and generation else None
def _blocker_sig(a):
"""Identity of blockers the agent surfaces once but cannot self-clear — parked needs-human
items, a dispatched terminally-red check, and a fork workflow awaiting maintainer approval. The
watch captures this at arm time and does not re-wake on a blocker already in that baseline, so a
parked residual keeps the watch alive (or the blocked-external bounded watch keeps polling for
the gate to clear) instead of busy-waking or terminating on the same condition."""
sig = set(a.get("needs_human_ids") or [])
if a.get("has_failing_checks") and a.get("checks_terminal") and not (a.get("counts") or {}).get("ci"):
sig.add("__terminal_red__")
if a.get("blocked_external"):
sig.add("__blocked_external__")
sig.add(("approval-review-state", a.get("review_decision")))
if a.get("stack_blocker"):
sig.add(f"__stack__:{a['stack_blocker']}")
if a.get("base_ref_blocker"):
sig.add(f"__base_ref__:{a['base_ref_blocker']}")
if a.get("unrequested_base_merge"):
sig.add(("unrequested-base-merge", a["unrequested_base_merge"].get("head")))
branch_currency = a.get("branch_currency") or {}
currency_disposition = branch_currency.get("disposition")
if (currency_disposition == CURRENCY_CONFIRMED
or (currency_disposition == CURRENCY_CLAIMED
and not branch_currency.get("attention"))):
sig.add(("currency", branch_currency.get("key"), currency_disposition,
branch_currency.get("recovery_state")))
# A new head is "context materially changed": a human may have pushed a commit that answers or
# supersedes a parked residual, so the head is part of the baseline — when it moves, the residual
# is no longer "already-surfaced against this state" and the watch wakes to give the agent a tick
# to reopen/reprocess it, instead of parking forever while it still blocks merge-ready.
if sig:
sig.add(("head", a.get("head_sha")))
return frozenset(sig)
def _process_identity(pid):
"""Best-effort PID-reuse guard for replacing a prior watcher. If process identity cannot be
proven, generation invalidation still suppresses its wake and we leave OS cleanup alone."""
if IS_WINDOWS:
return _win_process_identity(pid)
try:
r = subprocess.run(["ps", "-p", str(pid), "-o", "lstart=", "-o", "command="],
capture_output=True, text=True, encoding="utf-8")
except OSError:
return None # no ps on PATH: unproven identity, same as a failed lookup
if r.returncode != 0:
return None
return (r.stdout or "").strip() or None
def _watch_candidate_path(args):
return os.path.join(args.state_dir, "watch-candidate.json")
def _read_watch_candidate(args):
try:
with open(_watch_candidate_path(args)) as f:
candidate = json.load(f)
except (FileNotFoundError, json.JSONDecodeError):
return {}
return candidate if isinstance(candidate, dict) else {}
def _reserve_watch_candidate(args, generation):
"""Make this invocation the newest candidate without displacing the active watcher."""
candidate = {
"generation": generation,
"pid": os.getpid(),
"process_identity": _process_identity(os.getpid()),
}
with _state_lock(args.state_dir, exclusive=True):
previous = _read_watch_candidate(args)
tmp = tempfile.NamedTemporaryFile("w", dir=args.state_dir, delete=False)
json.dump(candidate, tmp)
tmp.flush()
os.fsync(tmp.fileno())
tmp.close()
_replace_atomic(tmp.name, _watch_candidate_path(args))
return previous
def _clear_watch_candidate(args, generation):
with _state_lock(args.state_dir, exclusive=True):
if _read_watch_candidate(args).get("generation") != generation:
return
try:
os.unlink(_watch_candidate_path(args))
except FileNotFoundError:
pass
def _activate_watch(args, generation, now, cur):
"""Atomically activate this watcher and persist its successfully fetched preflight."""
identity = _process_identity(os.getpid())
with locked_state(args.state_dir, args.pr, args.repo, now) as box:
if _read_watch_candidate(args).get("generation") != generation:
return None, None
state = box["state"]
previous = {
"pid": state.get("watch_pid"),
"process_identity": state.get("watch_process_identity"),
}
state["watch_generation"] = generation
state["watch_pid"] = os.getpid()
state["watch_process_identity"] = identity
# The arming poll is a real in-session watch poll: accumulate any suspended span since the
# last activity (e.g. the watch armed right after the machine resumed) before reading elapsed.
if state.get("started_at"):
_advance_activity(state, now, accumulate=True)
actionable, box["state"] = _apply_snapshot(box, args, cur, now, advance_trajectory=False)
try:
os.unlink(_watch_candidate_path(args))
except FileNotFoundError:
pass
return previous, actionable
def _watch_is_current(args, generation):
"""Read the ownership generation under the state lock without rewriting state.json."""
state_path = os.path.join(args.state_dir, "state.json")
with _state_lock(args.state_dir, exclusive=False):
if not os.path.exists(state_path):
return False
with open(state_path) as f:
return json.load(f).get("watch_generation") == generation
def _emit_wake_if_current(args, generation, reason, **fields):
"""Serialize the final ownership check with takeover so an old generation cannot emit after
the new generation has become current. Invocation replacement also supersedes every ordinary
wake from the old budget, even when the replacement kept the same watch generation."""
state_path = os.path.join(args.state_dir, "state.json")
with _state_lock(args.state_dir, exclusive=False):
if not os.path.exists(state_path):
return False
with open(state_path) as f:
state = json.load(f)
if state.get("watch_generation") != generation:
return False
current_invocation_id = state.get("invocation_id")
watcher_invocation_id = getattr(args, "invocation_id", None)
if current_invocation_id != watcher_invocation_id and reason != "invocation-superseded":
_emit_wake(
"invocation-superseded", watch_generation=generation,
superseded_invocation_id=watcher_invocation_id,
current_invocation_id=current_invocation_id,
)
return True
_emit_wake(reason, watch_generation=generation, **fields)
return True
def _terminate_replaced_watch(previous):
"""Promptly stop the replaced process, but never signal a PID whose identity no longer matches.
On Windows os.kill maps to TerminateProcess, so the replaced watcher dies abruptly instead of
unwinding a SIGTERM handler: it cannot clear its own watch candidate, nor finish terminating a
predecessor it was mid-handoff with. Neither leaks. The replacement has already written its own
candidate before calling this, and any watcher left running reads a foreign watch_generation on
its next poll and exits on its own — so the abrupt path costs at most one poll interval."""
pid = previous.get("pid")
identity = previous.get("process_identity")
if not isinstance(pid, int) or pid == os.getpid() or not identity:
return
if _process_identity(pid) != identity:
return
try:
os.kill(pid, signal.SIGTERM)
except OSError:
pass # exited between the identity check and here (Windows reports this as a plain OSError)
def _downstack_targets(args):
"""(pr, fetch_file_or_None) for each --downstack-pr; test fixtures bind a file per PR."""
files = {}
for spec in (getattr(args, "downstack_fetch_file", None) or []):
pr_text, _, path = spec.partition("=")
if pr_text.isdigit() and path:
files[int(pr_text)] = path
return [(pr, files.get(pr)) for pr in (getattr(args, "downstack_pr", None) or [])]
def _downstack_sig(args):
"""Stateless signature of the lower stack layers built from the same identities the engine
uses to decide that a layer has work: the head, each unresolved thread's remote-truth identity
(`_thread_identity` — a reply or an edit both change it), each non-thread comment by id (bodies
are rewritten by status bots and are not new work), failing checks by the canonical `FAILING`
set, the merge state when it demands maintenance, the manager classification tuple, and the PR
state. The watcher captures it at arm time and wakes `downstack-actionable` only when it grows,
so already-parked or dispatched lower-layer items do not re-wake, but any observation that would
make a lower layer non-quiescent does."""
sig = set()
for pr, fetch_file in _downstack_targets(args):
try:
if fetch_file:
with open(fetch_file) as f:
cur = json.load(f)
else:
cur = fetch(pr, args.repo, light=True)
except (OSError, ValueError, SystemExit):
continue # a failed lower-layer probe never wakes or blocks; the next poll retries
head = cur.get("head_sha")
sig.add((pr, "head", head))
sig.add((pr, "state", cur.get("pr_state")))
for t in cur.get("threads") or []:
sig.add((pr, "thread", t.get("thread_id"), _thread_identity(t)))
for c in cur.get("feedback") or []:
sig.add((pr, "comment", c.get("id")))
for c in cur.get("checks") or []:
if (c.get("status") or "").upper() == "COMPLETED" and (c.get("conclusion") or "").upper() in FAILING:
sig.add((pr, "check", head, c.get("key")))
status = (cur.get("merge_state_status") or "").upper()
if status in ("DIRTY", "BEHIND"):
sig.add((pr, "currency", status, head))
chain = cur.get("pr_chain") if isinstance(cur.get("pr_chain"), dict) else {}
sig.add((pr, "manager", chain.get("manager_status"), chain.get("target_needs_rebase")))
return frozenset(sig)
def cmd_watch(args):
"""Deterministic background change-detector (no agent tokens between changes): poll on an
interval, print one wake sentinel line and exit when there is something for the agent to do
(work to inspect or a stop condition), or exit on the stop-signal file / max-runtime.
A residual already present at arm time (a parked needs-human, or a dispatched terminal-red the
agent already handed back) does NOT re-wake the loop — it keeps watching the other streams;
only a *new* blocker (signature grown past the baseline) wakes."""
stop_requested = threading.Event()
interrupt_immediately = True
superseded = False
generation = uuid.uuid4().hex
# An already-requested stop is authoritative before this invocation becomes a candidate. Use
# current ownership for the wake when one exists, but do not fetch, reserve, mutate watch state,
# or disturb that incumbent merely to report the stop condition.
if args.stop_file and os.path.exists(args.stop_file):
_emit_wake("stop-signal", watch_generation=_persisted_watch_generation(args) or generation)
return
def request_stop(_signum, _frame):
nonlocal superseded
superseded = True
stop_requested.set()
if interrupt_immediately:
raise _WatchSuperseded()
# Cooperative supersession. Windows accepts this handler but never delivers SIGTERM across
# processes, so there it is inert and the loop's own watch_generation checks are what retire a
# replaced watcher — the handler is the fast path, not the correctness guarantee.
prior_sigterm = signal.getsignal(signal.SIGTERM)
signal.signal(signal.SIGTERM, request_stop)
try:
previous_candidate = _reserve_watch_candidate(args, generation)
_terminate_replaced_watch(previous_candidate)
# Preflight before takeover: an invalid fetch/auth/config must not displace a healthy watcher.
cur = _fetch_snapshot(args)
if stop_requested.is_set():
return
# Finish the handoff even if an even newer watcher arrives: otherwise that watcher could
# stop us between activation and terminating our predecessor, orphaning the oldest process.
interrupt_immediately = False
previous, actionable = _activate_watch(args, generation, _now(), cur)
if previous is None:
superseded = True
return
_terminate_replaced_watch(previous)
interrupt_immediately = True
if stop_requested.is_set():
return
armed = _blocker_sig(actionable) # blockers already surfaced when this generation armed
downstack_targets = _downstack_targets(args)
downstack_armed = _downstack_sig(args) if downstack_targets else frozenset()
while True:
if stop_requested.is_set():
return
if not _watch_is_current(args, generation):
superseded = True
return
if args.stop_file and os.path.exists(args.stop_file):
_emit_wake_if_current(args, generation, "stop-signal")
return
now = _now()
# Active elapsed = raw wall-clock minus accumulated suspended time (dead_time was
# updated by the previous poll's _run_snapshot, including any gap from a resumed
# suspension). The 8h budget spends active time; the 3-day backstop is raw wall-clock.
wall_elapsed = _elapsed(actionable.get("invocation_started_at"), now)
dead_time = int(actionable.get("invocation_dead_time_seconds") or 0)
invocation_elapsed = max(0, wall_elapsed - dead_time)
budget = actionable.get("invocation_budget_seconds") or 0
backstop = actionable.get("invocation_backstop_seconds") or 0
actionable["invocation_elapsed_seconds"] = invocation_elapsed
actionable["invocation_remaining_seconds"] = max(0, budget - invocation_elapsed)
actionable["invocation_wall_elapsed_seconds"] = wall_elapsed
actionable["persisted_state_age_seconds"] = _elapsed(
actionable.get("persisted_state_created_at"), now)
reason = _wake_reason(
actionable, args.settle_seconds,
getattr(args, "feedback_coalesce_seconds", FEEDBACK_COALESCE_SECONDS))
if reason in ("terminal", "merge-ready"):
_emit_wake_if_current(args, generation, reason, url=actionable.get("url"),
pr_state=actionable.get("pr_state"), counts=actionable.get("counts"))
return
active_cap_hit = bool(budget) and invocation_elapsed >= budget
backstop_hit = bool(backstop) and wall_elapsed >= backstop
if active_cap_hit or backstop_hit:
_emit_wake_if_current(
args, generation, "max-runtime", url=actionable.get("url"),
invocation_id=actionable.get("invocation_id"),
invocation_started_at=actionable.get("invocation_started_at"),
invocation_elapsed_seconds=actionable.get("invocation_elapsed_seconds"),
invocation_budget_seconds=budget,
invocation_wall_elapsed_seconds=wall_elapsed,
invocation_backstop_seconds=backstop or None,
max_runtime_ceiling=("backstop" if backstop_hit and not active_cap_hit
else "active-budget"),
persisted_state_age_seconds=actionable.get("persisted_state_age_seconds"),
)
return
drain_seconds = getattr(args, "blocked_external_drain_seconds", None)
drain_quiet = actionable.get("blocked_external_review_quiet_seconds", 0)
if (drain_seconds is not None and actionable.get("blocked_external")
and drain_quiet >= drain_seconds):
_emit_wake_if_current(
args, generation, "blocked-external-drained", url=actionable.get("url"),
pr_state=actionable.get("pr_state"), counts=actionable.get("counts"),
blocked_external_review_quiet_seconds=drain_quiet,
blocked_external_drain_seconds=drain_seconds,
)
return
approval_review_moved = (
reason == "blocked-external"
and actionable.get("blocked_external_review_moved_this_tick"))
if (reason in ("needs-human", "blocked-failing", "blocked-external",
"stack-blocked", "base-ref-blocked", "unrequested-base-merge")
and not approval_review_moved and _blocker_sig(actionable) <= armed):
reason = None # already-surfaced residual — keep watching, do not re-wake or terminate
if reason:
_emit_wake_if_current(args, generation, reason, url=actionable.get("url"),
pr_state=actionable.get("pr_state"), counts=actionable.get("counts"))
return
if downstack_targets:
grown = _downstack_sig(args) - downstack_armed
if grown:
_emit_wake_if_current(
args, generation, "downstack-actionable", url=actionable.get("url"),
pr_state=actionable.get("pr_state"), counts=actionable.get("counts"),
downstack_prs=sorted({item[0] for item in grown}))
return
remaining = max(0, budget - invocation_elapsed) if budget else args.interval
wait_seconds = min(args.interval, remaining)
if stop_requested.wait(wait_seconds):
return
if not _watch_is_current(args, generation):
superseded = True
return
actionable = _run_snapshot(args, _now(), advance_trajectory=False,
watch_generation=generation)
except _WatchSuperseded:
superseded = True
return
except _InvocationSuperseded as exc:
superseded = True
_emit_wake_if_current(
args, generation, "invocation-superseded",
superseded_invocation_id=getattr(args, "invocation_id", None),
current_invocation_id=exc.current_invocation_id,
)
return
finally:
interrupt_immediately = False
_clear_watch_candidate(args, generation)
# A newer owner can still signal our recorded PID after observing us stale but before this
# process has exited. Keep that late takeover signal harmless; ordinary wake/timeout/stop
# returns restore the embedding caller's handler as before.
signal.signal(signal.SIGTERM, signal.SIG_IGN if superseded else prior_sigterm)
def _mark_thread_records(item_ids, args, state):
"""Best-effort current records for every thread covered by one atomic mark."""
wanted = set(item_ids)
if not wanted:
return {}
try:
if getattr(args, "fetch_file", None):
with open(args.fetch_file) as f:
threads = json.load(f).get("threads", [])
else:
owner, name, host = _resolve_repo_ref(args.repo, (state.get("pr") or {}).get("url"))
threads = fetch_threads(args.pr, owner, name, host)
except (SystemExit, Exception): # SystemExit (from _run_checked) is not an Exception subclass
return {}
return {
thread.get("thread_id"): thread
for thread in threads if thread.get("thread_id") in wanted
}
def _mark_thread_baselines(item_ids, args, state):
"""Best-effort current identities for every thread covered by one atomic mark."""
return {
thread_id: [thread.get("last_comment_id"), thread.get("last_comment_at")]
for thread_id, thread in _mark_thread_records(item_ids, args, state).items()
}
def _mark_thread_baseline(item_id, args, state):
"""Compatibility wrapper for an ordinary single-thread dispatched mark."""
return _mark_thread_baselines([item_id], args, state).get(item_id)
def _mark_branch_currency(args, state, now):
currency = _load_branch_currency_state(state)
key = args.currency_key
if not key or currency.get("current_key") != key:
raise SystemExit("currency mark requires the exact current observation key")
item = currency.get("items", {}).get(key)
if not isinstance(item, dict):
raise SystemExit("currency mark requires a current observed item")
prior = item.get("disposition", DISPOSITION_OPEN)
outcome = args.currency_outcome
inspected_fingerprint = args.currency_inspected_fingerprint
if outcome:
if prior != CURRENCY_CLAIMED:
raise SystemExit("currency outcomes require a claimed observation")
item["reconciled_invocation_id"] = state.get("invocation_id")
item["reconciled_at"] = _iso(now)
if outcome == CURRENCY_OUTCOME_MUTATION_OBSERVED:
item["mutation_consumed"] = True
item["mutation_observed_at"] = _iso(now)
item["recovery_state"] = CURRENCY_OUTCOME_MUTATION_OBSERVED
elif outcome == CURRENCY_OUTCOME_AMBIGUOUS:
item["recovery_state"] = CURRENCY_OUTCOME_AMBIGUOUS
elif outcome == CURRENCY_OUTCOME_PROVEN_NO_MUTATION:
if item.get("mutation_consumed"):
raise SystemExit("cannot record no mutation after mutation start was observed")
retries = int(item.get("retry_count", 0))
if retries < 1:
item["retry_count"] = retries + 1
item["disposition"] = DISPOSITION_OPEN
item["recovery_state"] = "retry-authorized"
item["retry_not_before"] = _iso(
now + timedelta(seconds=CURRENCY_RETRY_BACKOFF_SECONDS))
item.pop("claimed_invocation_id", None)
item.pop("reconciled_invocation_id", None)
else:
item["disposition"] = DISPOSITION_OPEN
item["recovery_state"] = "retry-exhausted"
state["last_action"] = f"{outcome} currency {key}"
return item.get("recovery_state") in (
CURRENCY_OUTCOME_AMBIGUOUS, "retry-exhausted")
if inspected_fingerprint:
if prior != DISPOSITION_OPEN:
raise SystemExit("currency inspection requires an open observation")
parks = currency.get("semantic_parks") or {}
if not parks:
raise SystemExit("currency inspection requires carried semantic conflict evidence")
item["inspected_semantic_conflict_fingerprint"] = inspected_fingerprint
item["inspected_at"] = _iso(now)
if inspected_fingerprint in parks:
item["semantic_conflict_fingerprint"] = inspected_fingerprint
item["recovery_state"] = "semantic-unchanged"
item["inspection_result"] = "unchanged"
else:
item["inspection_required"] = False
item["inspection_result"] = "changed"
item["recovery_state"] = "semantic-changed"
item.pop("semantic_conflict_fingerprint", None)
# The preview describes the current conflict set for this head. Once it differs, the
# carried fingerprints are no longer standing residuals and must not tax every later
# base generation with inspection of stale evidence.
currency["semantic_parks"] = {}
item["parked_semantic_fingerprints"] = []
state["last_action"] = f"inspected currency {key}"
return inspected_fingerprint in parks
disposition = args.currency_disposition
allowed = {
DISPOSITION_OPEN: {DISPOSITION_OPEN, CURRENCY_CLAIMED},
CURRENCY_CLAIMED: {CURRENCY_CLAIMED, CURRENCY_CONFIRMED},
CURRENCY_CONFIRMED: {CURRENCY_CONFIRMED, DISPOSITION_OPEN},
}
if disposition not in allowed.get(prior, set()):
raise SystemExit(f"invalid currency transition: {prior} -> {disposition}")
if (prior == DISPOSITION_OPEN and disposition == CURRENCY_CLAIMED
and item.get("inspection_required")):
raise SystemExit("currency claim requires semantic-fingerprint inspection first")
if (prior == DISPOSITION_OPEN and disposition == CURRENCY_CLAIMED
and _needs_human_source_is_covered(state, key, "currency")):
raise SystemExit("currency claim requires the current human decision to be answered")
if (prior == DISPOSITION_OPEN and disposition == CURRENCY_CLAIMED
and _currency_mutation_requires_answer(state, key, item)):
raise SystemExit("currency claim requires an exact human answer")
if prior == DISPOSITION_OPEN and disposition == CURRENCY_CLAIMED:
retry_not_before = item.get("retry_not_before")
if retry_not_before:
try:
if _parse_iso8601(retry_not_before) > now:
raise SystemExit("currency retry backoff has not elapsed")
except (ValueError, TypeError):
raise SystemExit("currency retry backoff is invalid")
budget = state.get("invocation_budget_seconds")
backstop = state.get("invocation_backstop_seconds")
if ((budget is not None and _active_elapsed(state, now) >= budget)
or (backstop is not None and _elapsed(state.get("started_at"), now) >= backstop)):
raise SystemExit("currency claim cannot start after max-runtime")
item["disposition"] = disposition
item["transitioned_at"] = _iso(now)
if disposition == CURRENCY_CLAIMED:
# A repeated claim is deliberately idempotent. It cannot transfer or renew an existing
# claim; a later invocation must reconcile it and record an explicit outcome.
if prior == DISPOSITION_OPEN:
item["mutation_requires_answer"] = False
item["claimed_invocation_id"] = state.get("invocation_id")
item["claimed_at"] = _iso(now)
item["attempt_number"] = int(item.get("retry_count", 0)) + 1
item["mutation_consumed"] = False
item["recovery_state"] = "claimed"
item.pop("reconciled_invocation_id", None)
elif disposition == CURRENCY_CONFIRMED:
item["confirmed_invocation_id"] = state.get("invocation_id")
elif disposition == DISPOSITION_OPEN:
_reactivate_currency_item(currency, item)
else:
item["parked_semantic_fingerprints"] = sorted(currency.get("semantic_parks") or {})
state["last_action"] = f"{disposition} currency {key}"
return False
def _load_needs_human_residual(path, item_id=None, expected_kind=None):
if not path:
raise SystemExit("needs-human marks require --residual-file")
try:
with open(path, encoding="utf-8") as f:
residual = json.load(f)
except (OSError, json.JSONDecodeError) as exc:
raise SystemExit(f"cannot read --residual-file: {exc}")
error = _needs_human_residual_error(residual)
if error:
raise SystemExit(f"invalid --residual-file: {error}")
if item_id is not None and not any(
source["id"] == item_id and source["kind"] == expected_kind
for source in residual["sources"]):
raise SystemExit("--residual-file sources must include the marked item ID and kind")
return residual
def _load_human_answer(path):
if not path:
raise SystemExit("--answer-decision requires --answer-file")
try:
with open(path, encoding="utf-8") as f:
answer = f.read().strip()
except OSError as exc:
raise SystemExit(f"cannot read --answer-file: {exc}")
if not answer:
raise SystemExit("--answer-file must contain the human's answer")
return answer
def _park_human_decision(state, residual, args, now):
"""Freeze one complete decision over ordinary source observations, then publish it."""
threads = state.get("threads", {})
feedback = state.get("feedback", {})
checks = state.get("checks", {})
decisions = state.setdefault("human_decisions", [])
existing = next((
decision for decision in decisions
if isinstance(decision, dict)
and _residual_identity(decision.get("residual")) == _residual_identity(residual)
), None)
if existing is not None:
_install_human_decisions(state, decisions)
return existing
covered = {
(source.get("kind"), source.get("id"))
for decision in decisions if isinstance(decision, dict)
for source in decision.get("sources") or [] if isinstance(source, dict)
}
overlap = covered.intersection(_source_identity(source) for source in residual["sources"])
if overlap:
raise SystemExit("needs-human source already belongs to another current decision")
for source in residual["sources"]:
kind, source_id = _source_identity(source)
if not _source_can_enter_human_decision(
state, kind, source_id, threads, feedback, checks):
raise SystemExit(f"needs-human source is not current decision evidence: {kind} {source_id}")
if kind == "currency" and args.semantic_conflict_fingerprint:
currency = _load_branch_currency_state(state)
item = currency["items"][source_id]
fingerprint = args.semantic_conflict_fingerprint
item["semantic_conflict_fingerprint"] = fingerprint
currency.setdefault("semantic_parks", {})[fingerprint] = {
"head_sha": item.get("head_sha"),
"status": item.get("status"),
"route": item.get("route"),
"observation_key": source_id,
}
thread_ids = [
source["id"] for source in residual["sources"] if source["kind"] == "thread"]
thread_records = _mark_thread_records(thread_ids, args, state)
thread_baselines = {
thread_id: [thread.get("last_comment_id"), thread.get("last_comment_at")]
for thread_id, thread in thread_records.items()
}
missing_thread_baselines = sorted(set(thread_ids) - set(thread_baselines))
if missing_thread_baselines:
raise SystemExit(
"cannot freeze the post-reply thread observation: "
+ ", ".join(missing_thread_baselines))
authoritative_thread_urls = {
thread_id: thread.get("url") for thread_id, thread in thread_records.items()
}
if (any(not isinstance(url, str) or not url.strip()
for url in authoritative_thread_urls.values())
or set(residual["thread_urls"]) != set(authoritative_thread_urls.values())):
raise SystemExit("thread_urls must match the authoritative URL of every thread source")
sources = []
for source in residual["sources"]:
kind, source_id = _source_identity(source)
if kind == "thread":
observation = thread_baselines[source_id]
elif kind in ("comment", "review"):
observation = _current_source_observation(
state, kind, source_id, threads, feedback, checks)
if args.comment == source_id and args.acted_edit_id:
observation = [args.acted_edit_id]
else:
observation = _current_source_observation(
state, kind, source_id, threads, feedback, checks)
if observation is None:
raise SystemExit(f"cannot freeze needs-human observation: {kind} {source_id}")
sources.append({"kind": kind, "id": source_id, "observation": observation})
if any(
_answered_human_source_matches(
state, source["kind"], source["id"], source["observation"])
for source in sources):
raise SystemExit("source observations already have a recorded human answer")
_release_ordinary_suppression_for_decision(state, sources, threads, feedback)
decision = _decision_record(residual, sources, _iso(now))
if any(item.get("id") == decision["id"] for item in decisions if isinstance(item, dict)):
return decision
decisions.append(decision)
_install_human_decisions(state, decisions)
return decision
def _invalidate_human_source(state, item_id, kind):
kept = [
decision for decision in state.get("human_decisions") or []
if not any(source.get("id") == item_id and source.get("kind") == kind
for source in decision.get("sources") or [] if isinstance(source, dict))
]
_replace_human_decisions(state, kept)
def _needs_human_source_is_covered(state, item_id, kind):
return any(
source.get("id") == item_id and source.get("kind") == kind
for decision in state.get("human_decisions") or []
if isinstance(decision, dict)
for source in decision.get("sources") or []
if isinstance(source, dict)
)
def _answer_human_decision(state, decision_id, answer_file, now):
answer = _load_human_answer(answer_file)
current = next((
decision for decision in state.get("human_decisions") or []
if isinstance(decision, dict) and decision.get("id") == decision_id
), None)
if current is None:
existing = next((
decision for decision in state.get("answered_human_decisions") or []
if isinstance(decision, dict) and decision.get("id") == decision_id
), None)
if existing and existing.get("answer") == answer:
return existing
raise SystemExit("--answer-decision must name an exact current decision ID")
kept = [
decision for decision in state.get("human_decisions") or []
if not isinstance(decision, dict) or decision.get("id") != decision_id
]
_replace_human_decisions(state, kept)
answered = {**json.loads(json.dumps(current)), "answer": answer, "answered_at": _iso(now)}
state.setdefault("answered_human_decisions", []).append(answered)
state["last_action"] = f"answered human decision {decision_id}"
return answered
def cmd_mark(args):
now = _now()
marked = None
with locked_state(args.state_dir, args.pr, args.repo, now) as box:
_apply_invocation(box, args, now)
if (args.check or args.thread or args.comment) and args.residual_file \
and args.disposition != DISPOSITION_NEEDS_HUMAN:
raise SystemExit("--residual-file requires --disposition needs-human for source marks")
if (args.disposition == DISPOSITION_NEEDS_HUMAN and not args.residual_file
and not (args.check or args.thread or args.comment)):
raise SystemExit("needs-human marks require --residual-file")
state = box["state"]
# An agent-driven mark is activity: bump the heartbeat (no accumulation) so a long tick
# that only marks — never snapshots — still never has its active time refunded as dead time.
if state.get("started_at"):
_advance_activity(state, now, accumulate=False)
if args.answer_decision or args.answer_file:
if not args.answer_decision or not args.answer_file:
raise SystemExit("answer marks require --answer-decision and --answer-file")
if any((
args.check,
args.thread,
args.comment,
args.residual_file,
args.acted_edit_id,
args.currency_key,
args.currency_disposition,
args.semantic_conflict_fingerprint,
args.currency_outcome,
args.currency_inspected_fingerprint,
args.disposition != DISPOSITION_DISPATCHED,
)):
raise SystemExit("answer marks cannot be combined with source actions")
marked = _answer_human_decision(
state, args.answer_decision, args.answer_file, now)
elif (args.currency_key or args.currency_disposition or args.currency_outcome
or args.currency_inspected_fingerprint):
actions = sum(bool(value) for value in (
args.currency_disposition, args.currency_outcome,
args.currency_inspected_fingerprint))
if not args.currency_key or actions != 1:
raise SystemExit("currency marks require --currency-key and exactly one currency action")
decision_required = _mark_branch_currency(args, state, now)
if decision_required:
residual = _load_needs_human_residual(
args.residual_file, args.currency_key, "currency")
marked = _park_human_decision(state, residual, args, now)
elif args.residual_file:
raise SystemExit(
"--residual-file requires a currency outcome that needs a human decision")
if args.currency_disposition == DISPOSITION_OPEN:
_invalidate_human_source(state, args.currency_key, "currency")
marked = marked or args.currency_key
elif args.residual_file and not (args.check or args.thread or args.comment):
if args.disposition != DISPOSITION_NEEDS_HUMAN:
raise SystemExit("a residual-only mark requires --disposition needs-human")
residual = _load_needs_human_residual(args.residual_file)
marked = _park_human_decision(state, residual, args, now)
state["last_action"] = "needs-human residual"
elif args.check:
head = state.get("head_sha")
if not head:
raise SystemExit("mark --check requires a prior snapshot (state has no head_sha)")
state.setdefault("ci_dispatched", {}).setdefault(head, [])
if args.disposition == DISPOSITION_OPEN:
_invalidate_human_source(state, args.check, "check")
state["ci_dispatched"][head] = [
key for key in state["ci_dispatched"][head] if key != args.check]
elif args.disposition == DISPOSITION_NEEDS_HUMAN:
residual = _load_needs_human_residual(
args.residual_file, args.check, "check")
marked = _park_human_decision(state, residual, args, now)
else:
if _needs_human_source_is_covered(state, args.check, "check"):
raise SystemExit(
"cannot dispatch a source covered by a current human decision")
if args.check not in state["ci_dispatched"][head]:
state["ci_dispatched"][head].append(args.check)
state["last_action"] = f"{args.disposition} check {args.check}"
marked = marked or args.check
elif args.thread or args.comment:
item_id = args.thread or args.comment
collection, id_field, label = (
("threads", "thread_id", "thread") if args.thread else ("feedback", "id", "comment"))
entry = state.setdefault(collection, {}).setdefault(item_id, {id_field: item_id})
kind = "thread" if args.thread else entry.get("kind", "comment")
if args.disposition == DISPOSITION_OPEN:
entry["disposition"] = DISPOSITION_OPEN
_invalidate_human_source(state, item_id, kind)
entry.pop("acted_identity", None) # reopened -> next dispatch/park re-baselines
elif args.disposition == DISPOSITION_NEEDS_HUMAN:
residual = _load_needs_human_residual(
args.residual_file, item_id, kind)
marked = _park_human_decision(state, residual, args, now)
else:
if _needs_human_source_is_covered(state, item_id, kind):
raise SystemExit(
"cannot dispatch a source covered by a current human decision")
entry["disposition"] = DISPOSITION_DISPATCHED
if args.thread:
# our reply moved the thread's last comment, so re-read it now as the baseline.
ident = _mark_thread_baseline(item_id, args, state)
if ident is not None:
entry["acted_identity"] = ident
elif args.comment and args.acted_edit_id:
entry["acted_identity"] = [args.acted_edit_id]
state["last_action"] = f"{args.disposition} {label} {item_id}"
marked = marked or item_id
if getattr(args, "invariant_key", None) and args.disposition == DISPOSITION_DISPATCHED:
_record_invariant_round(state, args.invariant_key, state.get("head_sha"))
if isinstance(marked, dict):
print(json.dumps({"marked": marked.get("id"), "decision": marked}, separators=(",", ":")))
else:
print(json.dumps({"marked": marked}))
WATCH_BOOTSTRAP_FLAGS = ("--state-dir", "--invocation-id", "--session-started-at",
"--invocation-budget-seconds")
WATCH_BOOTSTRAP_HINT = """\
watch cannot start a babysit run: it only arms the change detector for an
invocation that was already bootstrapped by a snapshot --start-invocation run
(the ce-babysit-pr Step 2 bootstrap).
To recover, invoke the ce-babysit-pr skill through your harness's callable
skill mechanism with this PR: its instructions own the bootstrap, arming,
marks, and stop protocol. Do not drive pr-snapshot directly outside that skill.
Never mint the bootstrap values yourself; they come from the skill's
bootstrap snapshot."""
class _WatchHintingParser(argparse.ArgumentParser):
# A bare `watch --pr N` is the observed illegal start path (an agent arming the
# detector without the Step 2 bootstrap). Keep the fail-closed exit 2, but make
# the refusal carry its own recovery path instead of a raw usage dump.
def error(self, message):
if (self.prog.endswith(" watch") and "required" in message
and any(flag in message for flag in WATCH_BOOTSTRAP_FLAGS)):
self.print_usage(sys.stderr)
self.exit(2, f"{self.prog}: error: {message}\n\n{WATCH_BOOTSTRAP_HINT}\n")
super().error(message)
def main():
p = argparse.ArgumentParser(prog="pr-snapshot")
sub = p.add_subparsers(dest="cmd", required=True, parser_class=_WatchHintingParser)
s = sub.add_parser("snapshot")
s.add_argument("--pr", type=int, required=True)
s.add_argument("--repo", default=None)
s.add_argument("--state-dir", required=True)
s.add_argument("--fetch-file", default=None)
sg = s.add_mutually_exclusive_group()
sg.add_argument("--start-invocation", action="store_true",
help="mint one new fixed budget (first snapshot only)")
sg.add_argument("--reset-session", action="store_true",
help=argparse.SUPPRESS) # deprecated alias for --start-invocation
sg.add_argument("--continue-invocation", action="store_true",
help="carry the current fixed budget into a managed-stack layer state dir")
s.add_argument("--invocation-id", default=None,
help="token emitted by the first snapshot; required on every resume")
s.add_argument("--session-started-at", type=_session_started_at,
help="fixed anchor emitted by the first snapshot; required for a new state dir")
s.add_argument("--invocation-budget-seconds", type=float,
help=f"fixed total budget; first snapshot default {DEFAULT_INVOCATION_BUDGET_SECONDS}")
s.add_argument("--continue-dead-time-seconds", type=float, default=None,
help="carry the prior layer's accumulated dead time into a managed-stack layer "
"(with --continue-invocation), so the shared active-time budget stays correct")
s.set_defaults(func=cmd_snapshot)
m = sub.add_parser("mark")
m.add_argument("--pr", type=int, default=0)
m.add_argument("--repo", default=None)
m.add_argument("--state-dir", required=True)
m.add_argument("--invocation-id", required=True,
help="token emitted by the invocation's first snapshot")
m.add_argument("--session-started-at", type=_session_started_at, required=True,
help="fixed anchor emitted by the invocation's first snapshot")
m.add_argument("--invocation-budget-seconds", type=float, required=True,
help="fixed budget emitted by the invocation's first snapshot")
m.add_argument("--thread", default=None)
# `open` cancels a covering decision and returns the source to ordinary actionable state.
m.add_argument("--disposition", choices=[DISPOSITION_NEEDS_HUMAN, DISPOSITION_DISPATCHED, DISPOSITION_OPEN], default=DISPOSITION_DISPATCHED)
m.add_argument("--check", default=None)
m.add_argument("--comment", default=None)
m.add_argument("--fetch-file", default=None) # reuse the tick's fetch for the at-mark baseline
m.add_argument("--acted-edit-id", default=None) # snapshot-time identity for a covered comment
m.add_argument("--residual-file", default=None)
m.add_argument("--currency-key", default=None,
help="exact branch-currency observation key emitted by snapshot")
m.add_argument("--currency-disposition",
choices=[DISPOSITION_OPEN, CURRENCY_CLAIMED, CURRENCY_CONFIRMED],
default=None)
m.add_argument("--semantic-conflict-fingerprint", default=None,
help="semantic conflict identity preserved with a currency decision")
m.add_argument("--currency-outcome",
choices=[CURRENCY_OUTCOME_MUTATION_OBSERVED,
CURRENCY_OUTCOME_PROVEN_NO_MUTATION,
CURRENCY_OUTCOME_AMBIGUOUS], default=None)
m.add_argument("--currency-inspected-fingerprint", default=None)
m.add_argument("--answer-decision", default=None,
help="exact current decision ID whose human answer is being recorded")
m.add_argument("--answer-file", default=None,
help="file containing the human answer to preserve until covered work moves")
m.add_argument("--invariant-key", default=None,
help="resolver-supplied review invariant; counted per unique head, never inferred")
m.set_defaults(func=cmd_mark)
w = sub.add_parser("watch")
w.add_argument("--pr", type=int, required=True)
w.add_argument("--repo", default=None)
w.add_argument("--state-dir", required=True)
w.add_argument("--interval", type=float, default=150.0, help="poll cadence seconds")
w.add_argument("--downstack-pr", type=int, action="append", default=[],
help="lower managed-stack layer to keep probing; a new thread/comment/failing check/head there wakes downstack-actionable")
w.add_argument("--downstack-fetch-file", action="append", default=[], help="N=path fixture for --downstack-pr N (tests)")
w.add_argument("--settle-seconds", type=float, default=DEFAULT_SETTLE_SECONDS,
help="quiet window before a merge-ready wake")
w.add_argument("--feedback-coalesce-seconds", type=float, default=FEEDBACK_COALESCE_SECONDS,
help="quiet window before a comment-only wake, so a burst costs one dispatch")
w.add_argument("--blocked-external-drain-seconds", type=float, default=None,
help="head-scoped review-quiet bound before a gated-CI handback")
w.add_argument("--stop-file", default=None, help="path whose existence stops the watch")
w.add_argument("--fetch-file", default=None)
w.add_argument("--invocation-id", required=True,
help="token emitted by the invocation's first snapshot")
w.add_argument("--session-started-at", type=_session_started_at, required=True,
help="fixed anchor emitted by the first snapshot")
w.add_argument("--invocation-budget-seconds", type=float, required=True,
help="fixed budget emitted by the invocation's first snapshot")
w.set_defaults(func=cmd_watch)
args = p.parse_args()
if (args.cmd in ("snapshot", "watch")
and args.invocation_budget_seconds is not None
and args.invocation_budget_seconds <= 0):
p.error("--invocation-budget-seconds must be positive")
if (args.cmd == "watch" and args.blocked_external_drain_seconds is not None
and args.blocked_external_drain_seconds <= 0):
p.error("--blocked-external-drain-seconds must be positive")
if args.cmd == "snapshot":
starting = args.start_invocation or args.reset_session
continuing = args.continue_invocation
if starting and (args.invocation_id or args.session_started_at):
p.error("the first snapshot cannot combine --start-invocation with resume fields")
if continuing and not (args.invocation_id and args.session_started_at
and args.invocation_budget_seconds):
p.error("--continue-invocation requires --invocation-id, --session-started-at, and --invocation-budget-seconds")
if not starting and not continuing and not (args.invocation_id and args.session_started_at
and args.invocation_budget_seconds):
p.error("snapshot requires --start-invocation or --invocation-id")
args.func(args)
if __name__ == "__main__":
main()
SHA-256: c53fb1f6b778460785267d1a0686089561a3229ac7d56bb7ad346a5d8fa5d474