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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,12 @@ Versions follow [SemVer](https://semver.org).

### Added

- Phase 4 loop plumbing: `compute` (Slurm submit/status/cancel behind an
injectable runner; afterany wake jobs; query-failure ≠ job-gone),
`runstate` (atomic run records, six endings, expiring wake leases with
handoff), `tick` (pause sentinel, heartbeat, the five-layer fail-safe
sweep), and the self-resubmitting chain script.

- `autoresearch.brief` (typed, bounded, replayable session briefs — the
context-engineering artifact) and `autoresearch.harness` (backend seam;
Claude Code adapter with scrubbed session env, timeouts, and key-redacted
Expand Down
4 changes: 3 additions & 1 deletion docs/design/scaling.md
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,9 @@ concurrent seats).
immediately (promotes to the requested lane).
2. **Consolidation**: accumulation-triggered on K benchmark-consequential
merges. Maintenance merges are valuable but don't count toward K — see
Part 2b.
Part 2b. **Round 2 (2026-08-06): K defaults to 5–10 and the
consequential threshold to a 10% relative metric delta (ε) — both
contract-configurable, never hard-coded.**
3. **Credentials**: API keys for now (seats cost more up front); start a seat
when multi-agent testing is actually ready. The `setup-token` spike waits
until then.
Expand Down
24 changes: 15 additions & 9 deletions docs/roadmap.md
Original file line number Diff line number Diff line change
Expand Up @@ -61,15 +61,21 @@ standing instruction they supersede — a resumed agent honors stale constraints

## Phase 4 — Torch tick + compute

- [ ] Wake delivery per the architecture's fail-safe layers: afterany
dependency job (mechanism verified live 2026-08-06, incl. on failed
experiments) + run lease + tick sweep + deadline floor + `stuck` state

- [ ] The chain: two queued successors, `--dependency=singleton`, absolute
`--begin` cadence grid, sbatch retry/backoff
- [ ] Deploy shim: submit-successors → pull main → `uv sync --locked` → exec;
pull uses the bot PAT (bot has a Read collaborator role here; add this
repo to the token's selection)
- [x] Wake delivery per the architecture's fail-safe layers: afterany
dependency job (`compute.submit_after`; mechanism verified live
2026-08-06, incl. on failed experiments) + expiring run lease with
handoff + tick sweep + deadline floor (pending-cancel / gone-on-good-
query / defer-on-query-failure) + `stuck` state. Session dispatch
behind the `WakeDispatcher` seam (real dispatcher lands with phase 5)

- [x] The chain: `scripts/tick_chain.sbatch` — successor top-up to depth 2,
`--dependency=singleton`, absolute `--begin` cadence grid, sbatch
retry/backoff; successors submitted FIRST so nothing below can break
the chain
- [x] Deploy shim (same script): pull main with the bot PAT → `uv sync
--locked` → exec tick; every step best-effort (bad merges crash ticks,
never the chain). Bot PAT on Torch + repo in the token's selection
still owed by Mengye before live deploy
- [ ] Pause sentinel + lease (compare-and-swap on the state branch); heartbeat at
tick start; stale-lease reaping
- [ ] `compute`: sbatch submit / squeue poll, jobs tagged + `--nice`, GPU-hour caps
Expand Down
100 changes: 100 additions & 0 deletions scripts/tick_chain.sbatch
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
#!/bin/bash
# The self-resubmitting tick chain (docs/design/architecture.md, Scheduling).
#
# Order matters and is the fail-safety:
# 1. top up successors FIRST — a crash anywhere below never breaks the chain
# 2. pull main + sync (deploy-at-tick-start; failures logged, tick continues
# on the previous code — a bad merge crashes ticks, never the chain)
# 3. exec the tick
#
# Deploy/config knobs come from the environment (set once via
# `sbatch --export=...` when starting the chain; successors inherit):
# AUTORESEARCH_HOME checkout to run from (required)
# AUTORESEARCH_ROOT state root on the shared FS (required)
# AUTORESEARCH_ACCOUNT / AUTORESEARCH_PARTITION slurm placement (required)
# AUTORESEARCH_CADENCE_MIN tick cadence, default 30
# AUTORESEARCH_PAT_FILE bot PAT for the deploy pull (optional; no file
# means run whatever code is already deployed)
#
#SBATCH --job-name=autoresearch-tick
#SBATCH --time=15
#SBATCH --mem=4G
#SBATCH --cpus-per-task=2
#SBATCH --output=/dev/null

# NO set -u before the top-up: a config mistake must fail loudly AFTER the
# chain has secured its successors, not break the chain silently. Required
# vars are checked by hand instead.
CADENCE_MIN="${AUTORESEARCH_CADENCE_MIN:-30}"
JOB_NAME="autoresearch-tick"
missing=""
for var in AUTORESEARCH_HOME AUTORESEARCH_ROOT AUTORESEARCH_ACCOUNT AUTORESEARCH_PARTITION; do
eval "val=\${$var:-}"
[ -z "$val" ] && missing="$missing $var"
done

# --- 1. keep two successors queued (singleton serializes same-name jobs) ---
# Successors land on slots AFTER the ones already-queued jobs occupy —
# otherwise every top-up doubles the tick rate on the next slot.
if squeue_out=$(squeue -u "$USER" --name="$JOB_NAME" -h -t PENDING 2>/dev/null); then
pending=$(printf '%s' "$squeue_out" | grep -c '[0-9]')
need=$((2 - pending))
else
# A failed squeue must not read as "zero pending" — that would stack
# extra successors onto occupied slots every outage. Skip the top-up;
# the (very likely still-queued) successors keep the chain alive and
# the next tick retries.
pending=0
need=0
fi
epoch_now=$(date +%s)
cadence_s=$((CADENCE_MIN * 60))
next_slot=$(( (epoch_now / cadence_s + 1) * cadence_s ))
for i in $(seq 1 "$need"); do
begin_epoch=$((next_slot + (pending + i - 1) * cadence_s))
begin=$(date -d "@$begin_epoch" +%Y-%m-%dT%H:%M:%S 2>/dev/null \
|| date -r "$begin_epoch" +%Y-%m-%dT%H:%M:%S)
for attempt in 1 2 3; do
if sbatch --dependency=singleton --begin="$begin" \
--account="${AUTORESEARCH_ACCOUNT:-}" --partition="${AUTORESEARCH_PARTITION:-}" \
"${AUTORESEARCH_HOME:-}/scripts/tick_chain.sbatch"; then
break
fi
echo "sbatch retry $attempt failed; backing off"
sleep $((attempt * 20))
done
done

# Only NOW may configuration problems kill this tick — successors are queued.
set -u
if [ -n "$missing" ]; then
echo "tick misconfigured; missing:$missing" >&2
exit 1
fi
LOG_DIR="$AUTORESEARCH_ROOT/logs"
mkdir -p "$LOG_DIR" || true
if [ -w "$LOG_DIR" ]; then
exec >>"$LOG_DIR/tick-$(date +%Y%m%d).log" 2>&1
fi
echo "=== tick $(date -Is) on $(hostname -s) job=${SLURM_JOB_ID:-none}"

# --- 2. deploy: pull main with the bot PAT, sync deps (best-effort) ---
# The PAT never appears in argv (argv is world-readable via /proc on shared
# nodes): git asks for it through GIT_ASKPASS instead.
if [ -n "${AUTORESEARCH_PAT_FILE:-}" ] && [ -r "$AUTORESEARCH_PAT_FILE" ]; then
ASKPASS=$(mktemp 2>/dev/null || echo "") && [ -n "$ASKPASS" ] && chmod 700 "$ASKPASS"
printf '#!/bin/sh\ncat "%s"\n' "$AUTORESEARCH_PAT_FILE" > "$ASKPASS"
if GIT_ASKPASS="$ASKPASS" GIT_TERMINAL_PROMPT=0 git -C "$AUTORESEARCH_HOME" fetch --quiet \
"https://x-access-token@github.com/agentic-learning-ai-lab/autoresearch.git" main; then
git -C "$AUTORESEARCH_HOME" reset --hard --quiet FETCH_HEAD || echo "deploy: reset failed"
else
echo "deploy: fetch failed; running previous code"
fi
rm -f "$ASKPASS"
fi
(cd "$AUTORESEARCH_HOME" && uv sync --locked --quiet) || echo "deploy: uv sync failed"

# --- 3. the tick itself (scrubbed of the PAT path; sessions are scrubbed
# further down in the harness) ---
cd "$AUTORESEARCH_HOME"
exec env -u AUTORESEARCH_PAT_FILE uv run python -m autoresearch.tick --root "$AUTORESEARCH_ROOT"
196 changes: 196 additions & 0 deletions src/autoresearch/compute.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,196 @@
"""Slurm behind one small interface: submit, status, cancel.

Everything the loop knows about the cluster goes through here, so a CI
runner or cloud backend is a new implementation of the same three verbs.
Commands run through an injectable runner (tests use a fake; nothing here
requires a cluster).

The status query preserves a distinction the fail-safe design depends on
(docs/design/architecture.md, "Wake delivery and fail-safety"): a FAILED
query ("Slurm unknown") is not the same as a successful query that finds
nothing ("job gone") — misreading an outage as a vanished job would
terminate healthy runs.
"""

from __future__ import annotations

import logging
import shlex
import subprocess
from collections.abc import Callable, Sequence
from dataclasses import dataclass, field

log = logging.getLogger(__name__)

# Terminal Slurm states (prefix-matched: sacct reports e.g. "CANCELLED by 123").
TERMINAL_STATES = (
"COMPLETED",
"FAILED",
"CANCELLED",
"TIMEOUT",
"OUT_OF_MEMORY",
"NODE_FAIL",
"PREEMPTED",
"BOOT_FAIL",
"DEADLINE",
)
# A successful query that returns no record: the job left Slurm's memory.
GONE = "GONE"


class SlurmError(RuntimeError):
"""A Slurm command failed (submit/cancel), with its stderr."""


class SlurmQueryError(RuntimeError):
"""A status query failed — the answer is UNKNOWN, not 'job gone'.

Callers must treat this as "defer and retry", never as a terminal state.
"""


@dataclass(frozen=True)
class CommandResult:
returncode: int
stdout: str
stderr: str


Runner = Callable[[Sequence[str], int], CommandResult]


def _subprocess_runner(argv: Sequence[str], timeout_s: int) -> CommandResult:
completed = subprocess.run(list(argv), capture_output=True, text=True, timeout=timeout_s)
return CommandResult(completed.returncode, completed.stdout, completed.stderr)


@dataclass(frozen=True)
class JobSpec:
"""One sbatch submission. `command` is run via --wrap; a script path can
be passed as `script` instead (mutually exclusive).

--wrap executes under a shell on the compute node: `command` must be
built from trusted parts, with anything variable passed through
`quote_command`. Never interpolate agent- or contract-supplied text."""

job_name: str
account: str
partition: str
time_minutes: int
command: str = ""
script: str = ""
script_args: tuple[str, ...] = ()
cpus: int = 1
mem: str = "2G"
gpus: int = 0
qos: str = ""
output: str = "/dev/null"
# Slurm scheduling controls
dependency: str = "" # e.g. "afterany:12345" or "singleton"
begin: str = "" # e.g. "now+30" or an absolute "YYYY-MM-DDTHH:MM:SS"
extra: tuple[str, ...] = ()

def to_argv(self) -> list[str]:
if bool(self.command) == bool(self.script):
raise ValueError("exactly one of command/script must be set")
argv = [
"sbatch",
"--parsable",
f"--job-name={self.job_name}",
f"--account={self.account}",
f"--partition={self.partition}",
f"--time={self.time_minutes}",
f"--cpus-per-task={self.cpus}",
f"--mem={self.mem}",
f"--output={self.output}",
]
if self.gpus:
argv.append(f"--gpus={self.gpus}")
if self.qos:
argv.append(f"--qos={self.qos}")
if self.dependency:
argv.append(f"--dependency={self.dependency}")
if self.begin:
argv.append(f"--begin={self.begin}")
argv.extend(self.extra)
if self.command:
argv.append(f"--wrap={self.command}")
else:
argv.append(self.script)
argv.extend(self.script_args)
return argv


@dataclass
class SlurmCompute:
"""The three verbs, plus afterany for wake jobs."""

runner: Runner = field(default=_subprocess_runner)
command_timeout_s: int = 60

def submit(self, spec: JobSpec) -> str:
"""Submit; returns the job id. Raises SlurmError on failure."""
result = self.runner(spec.to_argv(), self.command_timeout_s)
if result.returncode != 0:
raise SlurmError(f"sbatch failed ({result.returncode}): {result.stderr.strip()}")
job_id = result.stdout.strip().split(";")[0]
if not job_id.isdigit():
raise SlurmError(f"sbatch returned no job id: {result.stdout.strip()!r}")
log.info("submitted %s as job %s", spec.job_name, job_id)
return job_id

def submit_after(self, spec: JobSpec, after_job_id: str) -> str:
"""Submit `spec` to run when `after_job_id` terminates — however it
terminates (afterany: the wake-on-failure semantics the fail-safe
design requires; verified live on Torch 2026-08-06). Refuses a spec
that already carries a dependency rather than silently replacing it."""
if not after_job_id.isdigit():
raise ValueError(f"not a job id: {after_job_id!r}")
if spec.dependency:
raise ValueError(f"spec already has dependency {spec.dependency!r}")
import dataclasses

return self.submit(dataclasses.replace(spec, dependency=f"afterany:{after_job_id}"))

def status(self, job_id: str) -> str:
"""The job's Slurm state, or GONE when a *successful* query finds no
record. Raises SlurmQueryError when the query itself fails."""
if not job_id.isdigit():
raise ValueError(f"not a job id: {job_id!r}")
try:
result = self.runner(
["sacct", "-j", job_id, "--parsable2", "--noheader", "-X", "-o", "State"],
self.command_timeout_s,
)
except (OSError, subprocess.TimeoutExpired) as exc:
raise SlurmQueryError(f"sacct did not run: {exc}") from exc
if result.returncode != 0:
raise SlurmQueryError(f"sacct failed ({result.returncode}): {result.stderr.strip()}")
state = result.stdout.strip().splitlines()[0].strip() if result.stdout.strip() else ""
return state if state else GONE

def cancel(self, job_id: str) -> None:
"""Cancel; idempotent (cancelling a finished job is not an error)."""
if not job_id.isdigit():
raise ValueError(f"not a job id: {job_id!r}")
result = self.runner(["scancel", job_id], self.command_timeout_s)
if result.returncode != 0:
log.warning("scancel %s: %s", job_id, result.stderr.strip())


def is_terminal(state: str) -> bool:
"""Whether a state string from `status` means the job is over.

GONE is deliberately NOT terminal here: it means "no record", and the
deadline-floor logic decides what that implies — not this predicate.
"""
return any(state.startswith(t) for t in TERMINAL_STATES)


def is_pending(state: str) -> bool:
return state.startswith("PENDING")


def quote_command(parts: Sequence[str]) -> str:
"""Shell-quote a command for JobSpec.command (--wrap takes a string)."""
return " ".join(shlex.quote(p) for p in parts)
Loading
Loading