## CHAPTER requirements-and-contracts The task runner, called `taskrun`, is a local coding-agent workhorse: it takes a declarative tasks file, builds a dependency graph, executes each task's command on a bounded thread pool, records everything in an append-only journal, and survives crashes by combining atomic state snapshots with journal replay. This chapter fixes the vocabulary used by all later chapters so the scheduler in chapter 9, the recovery logic in chapter 14, and the HTTP facade in chapter 17 all speak the same language. The functional requirements are: (1) parse a `.tasks` DSL file into `TaskSpec` objects; (2) execute tasks respecting `needs` ordering with parallelism; (3) bound every subprocess with timeouts, output caps, and OS rlimits; (4) retry transient failures under an explicit policy; (5) support cooperative cancellation from the CLI or HTTP interface; (6) persist an event journal with sequence numbers and checksums; (7) recover interrupted runs on restart; (8) expose `validate`, `run`, `resume`, `status`, `cancel`, and `journal` operations through both a CLI and a local HTTP server. Non-goals are deliberate: no distributed execution, no remote queue, no authentication beyond localhost binding, and no plugin loading beyond the tool protocol in chapter 16. The contract types are `TaskSpec` (static definition parsed from the DSL), `TaskState` (runtime status, attempt count, exit code, and human-readable detail), `CommandSpec` (either an `argv` vector executed directly or a named `tool` resolved through the chapter 16 protocol), and `RetrySpec` (attempt budget and backoff shape). Every deliberate failure raises a subclass of `TaskrunError`, so the CLI can map errors to exit codes without parsing message strings. **Invariants.** (I1) A `TaskState` moves only forward through the transition table: `PENDING → RUNNING → {SUCCEEDED, FAILED, CANCELLED}`, with `PENDING → SKIPPED` when an upstream task fails or is skipped, and `RUNNING → RETRY_WAIT` as a transient state used by the retry layer. (I2) All subprocess execution is argv-vector based; `shell=True` is banned project-wide. (I3) The journal is append-only with contiguous per-run sequence numbers starting at 1. (I4) Once cancellation is requested, no task may transition from `PENDING` to `RUNNING`. (I5) Configuration objects are frozen dataclasses; only `TaskState` mutates, and only under the scheduler lock. **Error handling.** `ConfigError` covers unusable workspaces and unknown tools; `ValidationError` aggregates every static problem in a tasks file rather than failing on the first; `CycleError` carries the offending path; `ToolError` records the exit code and whether the failure is retryable; `CancelledError` unwinds cooperative workers; `StorageError` stops the runner when snapshots cannot be trusted. **Worked example.** The file below drives every later chapter: `lint` runs first, `test` waits on `lint` and is gated by `env.CI`, and `package` needs `test`. Given this file, `taskrun run demo.tasks` produces a run id such as `188f2a41c0d39b1c4e7f` from `make_run_id`, executes `lint`, then `test` and `package`, and leaves a journal whose last line is `RUN_FINISHED`. ```python # taskrun/core.py """Shared vocabulary for every taskrun subsystem.""" from __future__ import annotations import enum import random from dataclasses import dataclass, field from typing import Any class TaskStatus(enum.Enum): PENDING = "pending" READY = "ready" RUNNING = "running" RETRY_WAIT = "retry_wait" SUCCEEDED = "succeeded" FAILED = "failed" CANCELLED = "cancelled" SKIPPED = "skipped" TERMINAL = frozenset({TaskStatus.SUCCEEDED, TaskStatus.FAILED, TaskStatus.CANCELLED, TaskStatus.SKIPPED}) class TaskrunError(Exception): """Base class for every deliberate taskrun failure.""" class ConfigError(TaskrunError): """Workspace, registry, or CLI inputs are unusable.""" class ValidationError(TaskrunError): """A tasks file violates the static contract.""" def __init__(self, errors: list[str]) -> None: super().__init__("; ".join(errors)) self.errors = list(errors) class CycleError(TaskrunError): """The dependency graph contains a cycle.""" def __init__(self, cycle: list[str]) -> None: super().__init__("dependency cycle: " + " -> ".join(cycle + cycle[:1])) self.cycle = list(cycle) class StorageError(TaskrunError): """Snapshot persistence failed in a way that must stop the runner.""" class ToolError(TaskrunError): def __init__(self, tool: str, message: str, exit_code: int | None = None, retryable: bool = False) -> None: super().__init__(f"tool {tool!r}: {message}") self.tool, self.exit_code, self.retryable = tool, exit_code, retryable class CancelledError(TaskrunError): """Raised inside work that observed a cancellation token.""" @dataclass(frozen=True) class RetrySpec: attempts: int = 0 delay_s: float = 0.0 backoff: float = 1.0 max_delay_s: float = 30.0 @dataclass(frozen=True) class CommandSpec: kind: str # "tool" or "argv" argv: tuple[str, ...] = () tool: str = "" params: dict[str, Any] = field(default_factory=dict) @dataclass class TaskSpec: name: str needs: tuple[str, ...] = () command: CommandSpec | None = None timeout_s: float = 60.0 retry: RetrySpec = field(default_factory=RetrySpec) when: Any = None # expression AST from chapter 4 description: str = "" @dataclass class TaskState: status: TaskStatus = TaskStatus.PENDING attempts: int = 0 exit_code: int | None = None detail: str = "" def make_run_id(now: float, rng: random.Random) -> str: """20 hex chars: 12 from the millisecond clock, 8 random.""" return f"{int(now * 1000) & 0xFFFFFFFFFFFF:012x}{rng.randrange(1 << 32):08x}" ``` ## CHAPTER workspace-layout Everything the runner persists lives under one directory so backups, permissions, and cleanup have a single target. The default root is `./.taskrun`; the `--home` flag overrides it, and the `TASKRUN_HOME` environment variable provides a middle layer. Resolution order is therefore CLI flag, then environment, then cwd-relative default, and the resolved path is always canonicalized with `Path.resolve()` so later containment checks in chapter 20 compare absolute, symlink-free paths. The layout has six directories with strict responsibilities. `journal/` holds one append-only JSONL file per run, named `.jsonl`, plus `.1` rotation successors. `state/` holds `state-.json` snapshots written atomically by chapter 13, with `.bak` predecessors. `logs/` holds the rotating structured log from chapter 19. `artifacts/` is the only directory tasks may write into voluntarily, addressed through the sandbox. `tools/` holds `tools.json`, the allowlist mapping tool names to argv vectors used by chapter 16. Keeping journal, state, and logs separate matters operationally: the journal is the recovery source of truth, snapshots are an optimization over it, and logs are human diagnostics that may be deleted freely. All directories are created with mode `0o700` because task runs frequently embed tokens in the environment, and a world-readable journal would leak command output. A write probe (`write`, `unlink` of a dotfile) runs at startup so a read-only mount or a squashed NFS mount fails immediately with `ConfigError` instead of corrupting a run twenty minutes later. **Invariants.** (W1) No component writes outside the workspace root; temporary files are created inside the destination directory so `os.replace` stays on one filesystem. (W2) Journal, snapshot, and log files are never renamed across directories. (W3) Directory creation is idempotent; `ensure_workspace` is safe to call from every entrypoint. (W4) The tools allowlist file is optional; when absent, only `argv` commands are executable. **Error handling.** Every filesystem failure during setup is converted to `ConfigError` with the errno message embedded, because callers (CLI chapter 18, HTTP chapter 17) need one exception type to map to exit code 3. `mkdir(parents=True, exist_ok=True)` tolerates concurrent creation by two runner processes, which the load tests in chapter 23 actually exercise. `chmod` failures on filesystems that ignore modes are tolerated silently, since the write probe is the real safety check. **Worked example.** Running the module below against a fresh temporary directory prints exactly the six directories, then the `tools.json` file we placed there. After a real run of `demo.tasks`, the same listing would additionally show `journal/.jsonl` and `state/state-.json`, which is precisely what the operations script in chapter 24 scans for. ```python # taskrun/workspace.py from __future__ import annotations import os from dataclasses import dataclass from pathlib import Path @dataclass(frozen=True) class Workspace: root: Path journal_dir: Path state_dir: Path logs_dir: Path artifacts_dir: Path tools_dir: Path def resolve_home(explicit: str | None = None, env: dict[str, str] | None = None, cwd: Path | None = None) -> Path: """CLI flag beats TASKRUN_HOME beats ./.taskrun next to the cwd.""" env = os.environ if env is None else env cwd = Path.cwd() if cwd is None else cwd if explicit: return Path(explicit).expanduser().resolve() if env.get("TASKRUN_HOME"): return Path(env["TASKRUN_HOME"]).expanduser().resolve() return (cwd / ".taskrun").resolve() def ensure_workspace(home: Path) -> Workspace: home = Path(home).resolve() ws = Workspace(root=home, journal_dir=home / "journal", state_dir=home / "state", logs_dir=home / "logs", artifacts_dir=home / "artifacts", tools_dir=home / "tools") try: for path in (ws.root, ws.journal_dir, ws.state_dir, ws.logs_dir, ws.artifacts_dir, ws.tools_dir): path.mkdir(parents=True, exist_ok=True) try: os.chmod(path, 0o700) except OSError: pass # modeless filesystems are fine; the probe below decides except OSError as exc: raise ConfigError(f"cannot create workspace under {home}: {exc}") from exc probe = ws.root / ".write-probe" try: probe.write_text("ok", encoding="utf-8") probe.unlink() except OSError as exc: raise ConfigError(f"workspace {home} is not writable: {exc}") from exc return ws if __name__ == "__main__": import tempfile with tempfile.TemporaryDirectory() as td: ws = ensure_workspace(Path(td) / "agent-home") (ws.tools_dir / "tools.json").write_text("{}\n", encoding="utf-8") for path in sorted(ws.root.rglob("*")): kind = "dir " if path.is_dir() else "file" print(kind, path.relative_to(ws.root)) ``` ## CHAPTER lexer The DSL is lexed by a hand-written scanner rather than `re.finditer` or a third-party tokenizer for three reasons: precise line and column on every token so parser errors can point at a column, first-class duration literals like `30s` and `1m` that would otherwise need post-processing, and deterministic maximal-munch behavior for two-character operators such as `==` and `&&`. The lexer is a single pass with one character of lookahead and produces a flat token list ending in an `EOF` sentinel, which frees the parser from bounds checking. Token kinds are: `IDENT` for `[A-Za-z_][A-Za-z0-9_-]*`, `STRING` for double-quoted text with escapes `\n`, `\t`, `\"`, `\\`, `NUMBER` for integer and decimal digits, `DURATION` for a number suffixed with `ms`, `s`, `m`, or `h`, `OP` for the six comparison and logical operators, `PUNCT` for structural characters `{ } = , [ ] ( ) : < > ! +`, `NEWLINE`, and `EOF`. Keywords such as `task`, `true`, and `env` are deliberately *not* reserved at lex time; the parser interprets `IDENT` values in context, which keeps the lexer reusable for the expression grammar in chapter 5. Comments start with `#` and run to end of line; they are skipped without emitting a token, but the newline that terminates them is still emitted so line-oriented field parsing works. Whitespace between tokens is skipped. An unterminated string, a bad escape, a bare number glued to letters (`30x`), or any unexpected character raises `LexError` carrying line, column, and a human-readable reason. **Invariants.** (L1) Concatenating token source text plus skipped whitespace and comments reproduces the input exactly, which makes column arithmetic trustworthy. (L2) The final token is always `EOF`; the parser can never read past it because `next()` refuses to advance at `EOF`. (L3) A `DURATION` token is always followed by a non-identifier character or end of input; `30sleep` is a lex error, not a number plus identifier. **Error handling.** `LexError` subclasses `ValueError` and formats as `line:col: message`, matching `ParseError` so chapter 6 can merge both into one error list. The lexer fails fast on the first error; the parser, not the lexer, implements multi-error recovery by resynchronizing at the next newline or `task` keyword. **Worked example.** For the three-line sample in the demo, the stream begins with `IDENT 'task'`, `IDENT 'lint'`, `PUNCT ':'`, then on line two `IDENT 'timeout'`, `PUNCT '='`, `DURATION '30s'`; the comment contributes nothing, and line three yields `IDENT 'when'`, `PUNCT '='`, `IDENT 'env'`, `PUNCT '.'`, `IDENT 'CI'`, `OP '!='`, `STRING 'false'`. The trailing token is `EOF` at line 4, column 1. ```python # taskrun/lexer.py from __future__ import annotations import re from dataclasses import dataclass class LexError(ValueError): def __init__(self, msg: str, line: int, col: int) -> None: super().__init__(f"{line}:{col}: {msg}") self.line, self.col = line, col @dataclass class Token: kind: str # IDENT STRING NUMBER DURATION OP PUNCT NEWLINE EOF value: str line: int col: int IDENT_RE = re.compile(r"[A-Za-z_][A-Za-z0-9_\-]*") NUMBER_RE = re.compile(r"\d+(?:\.\d+)?") DURATION_RE = re.compile(r"\d+(?:\.\d+)?(?:ms|s|m|h)") TWO_CHAR_OPS = ("==", "!=", "<=", ">=", "&&", "||") ONE_CHAR = set("{}[]=(),:<>!") class Lexer: def __init__(self, text: str) -> None: self.text, self.i, self.line, self.col = text, 0, 1, 1 def _peek(self, off: int = 0) -> str: j = self.i + off return self.text[j] if j < len(self.text) else "" def _advance(self, n: int = 1) -> None: for _ in range(n): if self.i < len(self.text): if self.text[self.i] == "\n": self.line, self.col = self.line + 1, 1 else: self.col += 1 self.i += 1 def tokens(self) -> list[Token]: out: list[Token] = [] while self.i < len(self.text): ch = self._peek() if ch in " \t\r": self._advance() elif ch == "#": while self._peek() not in ("", "\n"): self._advance() elif ch == "\n": out.append(Token("NEWLINE", "\\n", self.line, self.col)) self._advance() elif ch == '"': out.append(self._string()) elif ch.isdigit(): out.append(self._number_or_duration()) elif ch.isalpha() or ch == "_": m = IDENT_RE.match(self.text, self.i) out.append(Token("IDENT", m.group(0), self.line, self.col)) self._advance(len(m.group(0))) elif self._peek(0) + self._peek(1) in TWO_CHAR_OPS: out.append(Token("OP", self._peek(0) + self._peek(1), self.line, self.col)) self._advance(2) elif ch in ONE_CHAR: out.append(Token("PUNCT", ch, self.line, self.col)) self._advance() else: raise LexError(f"unexpected character {ch!r}", self.line, self.col) out.append(Token("EOF", "", self.line, self.col)) return out def _string(self) -> Token: line, col = self.line, self.col self._advance() buf: list[str] = [] while True: ch = self._peek() if ch in ("", "\n"): raise LexError("unterminated string", line, col) if ch == '"': self._advance() return Token("STRING", "".join(buf), line, col) if ch == "\\": nxt = self._peek(1) mapping = {"n": "\n", "t": "\t", '"': '"', "\\": "\\"} if nxt not in mapping: raise LexError(f"bad escape \\{nxt}", self.line, self.col) buf.append(mapping[nxt]) self._advance(2) else: buf.append(ch) self._advance() def _number_or_duration(self) -> Token: line, col = self.line, self.col m = DURATION_RE.match(self.text, self.i) if m and not self.text[m.end():m.end() + 1].isalnum(): self._advance(len(m.group(0))) return Token("DURATION", m.group(0), line, col) m = NUMBER_RE.match(self.text, self.i) nxt = self.text[m.end():m.end() + 1] if nxt and (nxt.isalpha() or nxt == "_"): raise LexError(f"bare number followed by {nxt!r}; units are ms s m h", line, col) self._advance(len(m.group(0))) return Token("NUMBER", m.group(0), line, col) if __name__ == "__main__": sample = 'task lint:\n timeout = 30s # seconds\n' for tok in Lexer(sample).tokens(): print(f"{tok.line}:{tok.col:>2} {tok.kind:<9} {tok.value!r}") ``` ## CHAPTER parser The parser is a recursive-descent translator from the token stream to `TaskSpec` objects plus an expression AST. The grammar is line-oriented: a file is a sequence of `task :` headers, each followed by `key = expression` fields until the next header or EOF. Field keys are restricted to `needs`, `command`, `timeout`, `retry`, `when`, and `description`; unknown keys are rejected with a `difflib` "did you mean" hint, which catches the classic `timeoutt` typo. Expression precedence, lowest to highest, is `||`, `&&`, unary `!`, a single non-associative comparison (`== != < <= > >=`), then primaries: string, number, duration, `true`, `false`, `env.`, function calls `name(args, key=value)`, bracketed lists, and parenthesized groups. A bare identifier parses as `Const(name)` so `needs = [lint, test]` works; the validator in chapter 6 rejects bare identifiers where they make no sense. Function calls accept positional and keyword arguments, which is exactly what the `tool(...)` and `retry` builders need. Structural errors, such as a missing colon after the task name, are fatal to that task but not to the file: `parse_file` records the `ParseError`, then `_sync_to_next_task` skips tokens to the next `task NAME` pair, so one broken task does not hide five others. Field-level errors are capped at 20 to keep diagnostics bounded. Two DSL constructs are translated eagerly: `retry` builders (`none()`, `fixed(attempts, delay, backoff)`, `exponential(...)` with `max_delay`) become `RetrySpec`, and `command` builders (`tool(name, **params)` and `argv([...])`) become `CommandSpec`. Durations are converted to float seconds at parse time, so `timeout = 90s` and `timeout = 1.5m` both arrive as `90.0`. **Invariants.** (P1) The parser never evaluates `env` references or spawns anything; `when` remains an AST until the scheduler evaluates it in chapter 9. (P2) Every `TaskSpec` reaching chapter 6 has syntactically valid needs (a `ListLit` of strings) and a converted or absent `command` and `retry`. (P3) Parse errors always include line and column taken from the offending token. **Error handling.** `ParseError` wraps the failing token, and field errors are re-raised prefixed with the task name so a message reads `task 'test': 4:11: unknown dependency syntax`. The `errors` list on `TaskFile` is merged with semantic errors in chapter 6, giving the user the complete problem list in one run. **Worked example.** Parsing the demo file yields two tasks: `lint` with `needs=()`, `timeout_s=30.0`, `command=CommandSpec(kind='tool', tool='lint', params={'target': 'src/'})`, and `retry=RetrySpec(attempts=2, delay_s=2.0)`; and `test` with `needs=('lint',)`, `when=Binary('!=', EnvRef('CI'), Const('false'))`. The printed errors list is empty, and the AST printout shows the un-evaluated expression tree. ```python # taskrun/parser.py from __future__ import annotations import difflib from dataclasses import dataclass from taskrun.core import CommandSpec, RetrySpec, TaskSpec from taskrun.lexer import Lexer, Token FIELD_KEYS = {"needs", "command", "timeout", "retry", "when", "description"} COMPARISONS = {"==", "!=", "<", "<=", ">", ">="} MAX_ERRORS = 20 class ParseError(Exception): def __init__(self, msg: str, tok: Token) -> None: super().__init__(f"{tok.line}:{tok.col}: {msg}") @dataclass(frozen=True) class Const: value: object @dataclass(frozen=True) class EnvRef: name: str @dataclass(frozen=True) class ListLit: items: tuple @dataclass(frozen=True) class Unary: op: str operand: object @dataclass(frozen=True) class Binary: op: str left: object right: object @dataclass(frozen=True) class Call: name: str args: tuple kwargs: dict @dataclass class TaskFile: tasks: list errors: list[str] class Parser: def __init__(self, tokens: list[Token]) -> None: self.toks, self.pos = tokens, 0 def peek(self, off: int = 0) -> Token: return self.toks[min(self.pos + off, len(self.toks) - 1)] def next(self) -> Token: tok = self.toks[self.pos] if tok.kind != "EOF": self.pos += 1 return tok def expect(self, kind: str, value: str | None = None) -> Token: tok = self.peek() if tok.kind != kind or (value is not None and tok.value != value): raise ParseError(f"expected {value or kind!r}, found {tok.value or tok.kind!r}", tok) return self.next() def parse_file(self) -> TaskFile: tasks, errors = [], [] while self.peek().kind != "EOF": while self.peek().kind == "NEWLINE": self.next() if self.peek().kind == "EOF": break try: tasks.append(self.parse_task()) except ParseError as exc: errors.append(str(exc)) if len(errors) >= MAX_ERRORS: break self._sync_to_next_task() return TaskFile(tasks=tasks, errors=errors) def _sync_to_next_task(self) -> None: while self.peek().kind != "EOF": if (self.peek().kind == "IDENT" and self.peek().value == "task" and self.peek(1).kind == "IDENT"): return self.next() def parse_task(self) -> "TaskSpec": self.expect("IDENT", "task") name = self.expect("IDENT") self.expect("PUNCT", ":") spec = TaskSpec(name=name.value) while not (self.peek().kind == "IDENT" and self.peek().value == "task" and self.peek(1).kind == "IDENT") and self.peek().kind != "EOF": if self.peek().kind == "NEWLINE": self.next() continue self.parse_field(spec) return spec def parse_field(self, spec: "TaskSpec") -> None: key = self.expect("IDENT") self.expect("PUNCT", "=") if key.value not in FIELD_KEYS: close = difflib.get_close_matches(key.value, FIELD_KEYS, 1) hint = f" (did you mean {close[0]!r}?)" if close else "" raise ParseError(f"unknown field {key.value!r}{hint}", key) expr = self.parse_or() self._end