Add sssf skill, installable via the skills CLI
Port the sssf skill from ~/.agents/skills/sssf into this repo so it can be distributed and installed with the skills CLI (skills add INDigitalStudio/skills --skill sssf). - Copy the skill (SKILL.md, cookbooks, references, scripts, templates, and the visualizer app source) into sssf/. - Gitignore build/runtime artifacts: the visualizer's node_modules/ and dist/, Python bytecode, and the machine-specific repos.json. - Make the skill location-independent: install.py now stamps the skill's real path into the stamped justfile's skill_dir (replacing the hardcoded ~/.agents/skills/sssf), so 'just obs' finds the visualizer wherever the CLI installed the skill. - Update cookbooks to use <skill>/scripts/... instead of the hardcoded path, and document the skills CLI install command. - Update the repo README with install instructions.
This commit is contained in:
parent
a42608f602
commit
2cc766aabe
98 changed files with 12508 additions and 0 deletions
0
sssf/templates/adws/adw_modules/__init__.py
Normal file
0
sssf/templates/adws/adw_modules/__init__.py
Normal file
15
sssf/templates/adws/adw_modules/agent_cc.py
Normal file
15
sssf/templates/adws/adw_modules/agent_cc.py
Normal file
|
|
@ -0,0 +1,15 @@
|
|||
"""Claude Code interface — STUB in v1. The factory is Pi-only for now.
|
||||
|
||||
The config schema accepts `coding_agent: claude_code` so nothing breaks at the
|
||||
schema level, but selecting it raises until v2 implements this interface
|
||||
(`claude -p --output-format stream-json --resume <session_id>`).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
|
||||
def run(*args, **kwargs):
|
||||
raise NotImplementedError(
|
||||
"coding_agent 'claude_code' is not implemented in v1 — SSSF v1 runs the "
|
||||
"Pi coding agent only. Set coding_agent: pi (or omit it) in sssf.config.yaml."
|
||||
)
|
||||
210
sssf/templates/adws/adw_modules/agent_omp.py
Normal file
210
sssf/templates/adws/adw_modules/agent_omp.py
Normal file
|
|
@ -0,0 +1,210 @@
|
|||
"""OMP coding agent interface — an alternative to Pi for SSSF.
|
||||
|
||||
Runs `omp -p --mode json` and tails its JSONL stdout line by line, forwarding
|
||||
each event to a callback WHILE the agent works. OMP's event stream is the same
|
||||
shape as Pi's (session / message_start / message_end / turn_end /
|
||||
tool_execution_start / tool_execution_end / agent_end), so the shared
|
||||
ToolCallTracker folds tool calls identically.
|
||||
|
||||
The one real difference is session identity. Pi accepts `--session-id` and
|
||||
creates-or-continues that exact id. OMP mints its own session id and resumes by
|
||||
`--resume <id-prefix>`. To keep SSSF's retry/correction loop (which re-enters
|
||||
the SAME context window) working, this adapter persists the real OMP session id
|
||||
in a `session_id.txt` next to the session dir: the first send creates a session
|
||||
and records its id; every later send with the same session_dir resumes it.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import subprocess
|
||||
from functools import lru_cache
|
||||
from pathlib import Path
|
||||
from typing import Callable, Optional
|
||||
|
||||
from .agent_pi import (ToolCallTracker, _context_tokens, _label, _text_of)
|
||||
from .data_types import PiRequest, PiResult
|
||||
from .utils import now_iso, operator_env
|
||||
|
||||
OMP_PATH = os.environ.get("OMP_PATH", "omp")
|
||||
|
||||
RESULT_SNIPPET_CHARS = 20_000 # tool output rides along whole; clip only guards pathological cases
|
||||
ARG_VALUE_CHARS = 20_000 # args too — the UI scrolls, it must not be handed cut-off data
|
||||
LABEL_CHARS = 80 # "bash: <command>" shown as the event name
|
||||
|
||||
# The arg that identifies a call at a glance, in the order tools tend to use.
|
||||
PRIMARY_ARGS = ("command", "path", "file_path", "pattern", "query", "url")
|
||||
|
||||
# The config's `tools` list is written in pi's vocabulary. OMP's tool set
|
||||
# overlaps but is not identical, so translate the names that differ and drop
|
||||
# the ones OMP has no equivalent for (the `writes` boundary, enforced in code,
|
||||
# is what actually keeps an agent read-only — the tool list is a hint).
|
||||
TOOL_TRANSLATION = {
|
||||
"find": "glob", # pi's find -> omp's glob
|
||||
"ls": None, # omp reads directories via `read`
|
||||
"subagent_create": "task", # omp's subagent tool
|
||||
"subagent_continue": "task",
|
||||
"subagent_list": "task",
|
||||
"subagent_remove": "task",
|
||||
}
|
||||
|
||||
|
||||
def _translate_tools(tools: list[str]) -> list[str]:
|
||||
"""Map pi tool names to omp's vocabulary, dropping unmappable ones."""
|
||||
out = []
|
||||
for tool in tools:
|
||||
mapped = TOOL_TRANSLATION.get(tool, tool)
|
||||
if mapped and mapped not in out:
|
||||
out.append(mapped)
|
||||
return out
|
||||
|
||||
|
||||
@lru_cache(maxsize=1)
|
||||
def _omp_catalog() -> list[dict]:
|
||||
"""Read omp's merged model catalog (`omp models --json`)."""
|
||||
try:
|
||||
result = subprocess.run(
|
||||
[OMP_PATH, "models", "--json"], capture_output=True, text=True,
|
||||
timeout=30, env=operator_env(), check=False,
|
||||
)
|
||||
except (OSError, subprocess.TimeoutExpired):
|
||||
return []
|
||||
if result.returncode != 0:
|
||||
return []
|
||||
try:
|
||||
return json.loads(result.stdout).get("models", [])
|
||||
except json.JSONDecodeError:
|
||||
return []
|
||||
|
||||
|
||||
def resolve_model(pattern: str) -> tuple[str, str]:
|
||||
"""Resolve a model pattern to an explicit ``(provider, model_id)`` pair.
|
||||
|
||||
OMP's catalog exposes a `selector` of the form ``provider/model-id``. A
|
||||
pattern with a slash must match a selector exactly; a bare pattern matches
|
||||
by model id (substring, then exact), failing on ambiguity.
|
||||
"""
|
||||
catalog = _omp_catalog()
|
||||
selectors = {m.get("selector"): (m.get("provider", ""), m.get("id", ""))
|
||||
for m in catalog if m.get("selector")}
|
||||
if "/" in pattern:
|
||||
if pattern in selectors:
|
||||
return selectors[pattern]
|
||||
raise ValueError(f"model pattern {pattern!r} not found in omp models — "
|
||||
"run `omp models` to see available selectors")
|
||||
matches = [(m.get("provider", ""), m.get("id", "")) for m in catalog
|
||||
if pattern == m.get("id") or pattern in m.get("id", "")]
|
||||
exact = [match for match in matches if match[1] == pattern]
|
||||
if len(exact) == 1:
|
||||
return exact[0]
|
||||
if len(matches) == 1:
|
||||
return matches[0]
|
||||
if not matches:
|
||||
raise ValueError(f"model pattern {pattern!r} not found in omp models — "
|
||||
"run `omp models` to see available models")
|
||||
raise ValueError(f"model pattern {pattern!r} is ambiguous: {matches}")
|
||||
|
||||
|
||||
def context_window(provider: str, model_id: str) -> int:
|
||||
"""The model's context ceiling from omp's merged model catalog."""
|
||||
for model in _omp_catalog():
|
||||
if model.get("provider") == provider and model.get("id") == model_id:
|
||||
return int(model.get("contextWindow") or 0)
|
||||
return 0
|
||||
|
||||
|
||||
def _session_id_file(session_dir: Path) -> Path:
|
||||
return session_dir / "session_id.txt"
|
||||
|
||||
|
||||
def run(request: PiRequest, on_event: Optional[Callable[[dict], None]] = None,
|
||||
on_spawn: Optional[Callable[[int], None]] = None,
|
||||
on_exit: Optional[Callable[[int], None]] = None) -> PiResult:
|
||||
"""Run one non-interactive omp turn.
|
||||
|
||||
`on_spawn(pid)` and `on_exit(pid)` bracket the child process so the caller
|
||||
can record it as killable — a hung coding agent is otherwise a pid you have
|
||||
to hunt for in `ps` while the run sits there.
|
||||
"""
|
||||
provider, model_id = resolve_model(request.model)
|
||||
selector = f"{provider}/{model_id}"
|
||||
|
||||
session_dir = Path(request.session_dir)
|
||||
session_dir.mkdir(parents=True, exist_ok=True)
|
||||
sid_file = _session_id_file(session_dir)
|
||||
|
||||
cmd = [
|
||||
OMP_PATH, "-p", "--mode", "json",
|
||||
"--model", selector,
|
||||
"--thinking", request.thinking,
|
||||
"--system-prompt", request.system_prompt,
|
||||
"--session-dir", str(session_dir),
|
||||
"--auto-approve", # agents run tools autonomously; no approval prompts
|
||||
]
|
||||
# Create-or-continue: the first send mints a session and records its id;
|
||||
# every later send with the same session_dir resumes that context window.
|
||||
if sid_file.exists():
|
||||
cmd += ["--resume", sid_file.read_text().strip()]
|
||||
if request.tools:
|
||||
cmd += ["--tools", ",".join(_translate_tools(request.tools))]
|
||||
for extension in request.extensions:
|
||||
cmd += ["-e", extension]
|
||||
cmd.append(request.prompt)
|
||||
|
||||
raw_path = Path(request.raw_output_path)
|
||||
raw_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
result = PiResult(session_id="", context_window=context_window(provider, model_id))
|
||||
# stdin is DEVNULL, deliberately. The prompt travels in argv, so the child
|
||||
# never needs stdin — but inheriting the parent's means omp sees a non-TTY
|
||||
# and can sit forever waiting for piped input that will never arrive or EOF.
|
||||
process = subprocess.Popen(cmd, stdin=subprocess.DEVNULL,
|
||||
stdout=subprocess.PIPE, stderr=subprocess.PIPE,
|
||||
text=True, bufsize=1, cwd=request.cwd,
|
||||
env=operator_env())
|
||||
if on_spawn:
|
||||
on_spawn(process.pid)
|
||||
with raw_path.open("a") as raw:
|
||||
assert process.stdout is not None
|
||||
for line in process.stdout:
|
||||
raw.write(line)
|
||||
raw.flush() # events land on disk as they happen
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
event = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
if event.get("type") == "session":
|
||||
sid = event.get("id")
|
||||
if sid:
|
||||
result.session_id = sid
|
||||
sid_file.write_text(sid) # persist for later resumes
|
||||
if event.get("type") == "message_end":
|
||||
message = event.get("message", {})
|
||||
if message.get("role") == "assistant":
|
||||
text = _text_of(message)
|
||||
if text:
|
||||
result.text = text # last assistant message wins
|
||||
usage = message.get("usage", {}) or {}
|
||||
turn = _context_tokens(usage)
|
||||
result.tokens += turn
|
||||
result.usage.add_turn(usage, turn)
|
||||
# Occupancy is read off the last VALID assistant turn, the
|
||||
# way omp does it — an aborted or errored turn reports usage
|
||||
# you can't trust, so it must not overwrite a good reading.
|
||||
if turn and message.get("stopReason") not in ("aborted", "error"):
|
||||
result.context_tokens = turn
|
||||
result.cost += (usage.get("cost", {}) or {}).get("total", 0.0) or 0.0
|
||||
if on_event:
|
||||
on_event(event)
|
||||
|
||||
stderr = process.stderr.read() if process.stderr else ""
|
||||
result.returncode = process.wait()
|
||||
if on_exit:
|
||||
on_exit(process.pid)
|
||||
if result.returncode != 0 and not result.text:
|
||||
raise RuntimeError(f"omp exited {result.returncode}: {stderr.strip()[-800:]}")
|
||||
return result
|
||||
286
sssf/templates/adws/adw_modules/agent_pi.py
Normal file
286
sssf/templates/adws/adw_modules/agent_pi.py
Normal file
|
|
@ -0,0 +1,286 @@
|
|||
"""Pi coding agent interface — v1's only coding agent.
|
||||
|
||||
Runs `pi -p --mode json` and tails its JSONL stdout line by line, forwarding
|
||||
each event to a callback WHILE the agent works (the streaming crack, solved
|
||||
by construction). `--session-id` creates-or-continues, so running and
|
||||
continuing an agent are the same call: same session id = same context window.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import subprocess
|
||||
import time
|
||||
from functools import lru_cache
|
||||
from pathlib import Path
|
||||
from typing import Callable, Optional
|
||||
|
||||
from .data_types import PiRequest, PiResult
|
||||
from .utils import now_iso, operator_env
|
||||
|
||||
PI_PATH = os.environ.get("PI_PATH", "pi")
|
||||
MODELS_JSON = os.environ.get("PI_MODELS_PATH",
|
||||
str(Path.home() / ".pi" / "agent" / "models.json"))
|
||||
|
||||
RESULT_SNIPPET_CHARS = 20_000 # tool output rides along whole; clip only guards pathological cases
|
||||
ARG_VALUE_CHARS = 20_000 # args too — the UI scrolls, it must not be handed cut-off data
|
||||
LABEL_CHARS = 80 # "bash: <command>" shown as the event name
|
||||
|
||||
# The arg that identifies a call at a glance, in the order tools tend to use.
|
||||
PRIMARY_ARGS = ("command", "path", "file_path", "pattern", "query", "url")
|
||||
|
||||
|
||||
def _count(value: str) -> int:
|
||||
"""Parse pi's compact model-list counts (`272K`, `1.0M`)."""
|
||||
suffixes = {"K": 1_000, "M": 1_000_000}
|
||||
suffix = value[-1:].upper()
|
||||
if suffix in suffixes:
|
||||
return int(float(value[:-1]) * suffixes[suffix])
|
||||
return int(value)
|
||||
|
||||
|
||||
@lru_cache(maxsize=1)
|
||||
def _pi_catalog() -> list[tuple[str, str, int]]:
|
||||
"""Read pi's merged catalog, including built-in providers and custom models."""
|
||||
try:
|
||||
result = subprocess.run(
|
||||
[PI_PATH, "--list-models"], capture_output=True, text=True,
|
||||
timeout=30, env=operator_env(), check=False,
|
||||
)
|
||||
except (OSError, subprocess.TimeoutExpired):
|
||||
return []
|
||||
if result.returncode != 0:
|
||||
return []
|
||||
rows = []
|
||||
for line in result.stdout.splitlines()[1:]:
|
||||
columns = line.split()
|
||||
if len(columns) < 3:
|
||||
continue
|
||||
try:
|
||||
rows.append((columns[0], columns[1], _count(columns[2])))
|
||||
except ValueError:
|
||||
continue
|
||||
return rows
|
||||
|
||||
|
||||
def resolve_model(pattern: str) -> tuple[str, str]:
|
||||
"""Resolve a model pattern to an explicit ``(provider, model_id)`` pair.
|
||||
|
||||
Pi's catalog merges built-in models with ``~/.pi/agent/models.json``. Using
|
||||
that same merged view lets SSSF target direct providers such as
|
||||
``openai/gpt-5.6-terra`` without re-registering built-in models locally.
|
||||
"""
|
||||
catalog = [(provider, model_id) for provider, model_id, _ in _pi_catalog()]
|
||||
if "/" in pattern:
|
||||
provider, model_id = pattern.split("/", 1)
|
||||
if (provider, model_id) in catalog:
|
||||
return provider, model_id
|
||||
matches = [(provider, model_id) for provider, model_id in catalog
|
||||
if pattern == model_id or pattern in model_id]
|
||||
exact = [match for match in matches
|
||||
if match[1] == pattern or match[1].endswith("/" + pattern)]
|
||||
if len(exact) == 1:
|
||||
return exact[0]
|
||||
if len(matches) == 1:
|
||||
return matches[0]
|
||||
if not matches:
|
||||
raise ValueError(f"model pattern {pattern!r} not found in pi --list-models — "
|
||||
"authenticate/register it or fix the config")
|
||||
raise ValueError(f"model pattern {pattern!r} is ambiguous: {matches}")
|
||||
|
||||
|
||||
def _context_tokens(usage: dict) -> int:
|
||||
"""Tokens occupying the window after a turn.
|
||||
|
||||
Mirrors pi's own `calculateContextTokens` (coding-agent
|
||||
`core/compaction/compaction.ts`), which is what pi compacts against and
|
||||
shows in its footer: prefer the provider's `totalTokens`, else sum the
|
||||
parts. Cache reads count — cached prompt is still prompt.
|
||||
"""
|
||||
total = usage.get("totalTokens") or 0
|
||||
if total:
|
||||
return int(total)
|
||||
return int(sum(usage.get(part) or 0
|
||||
for part in ("input", "output", "cacheRead", "cacheWrite")))
|
||||
|
||||
|
||||
def context_window(provider: str, model_id: str) -> int:
|
||||
"""The model's context ceiling from pi's merged model catalog."""
|
||||
registry = json.loads(Path(MODELS_JSON).read_text())
|
||||
for model in registry.get("providers", {}).get(provider, {}).get("models", []):
|
||||
if model.get("id") == model_id:
|
||||
return int(model.get("contextWindow") or 0)
|
||||
for listed_provider, listed_model, window in _pi_catalog():
|
||||
if listed_provider == provider and listed_model == model_id:
|
||||
return window
|
||||
return 0
|
||||
|
||||
|
||||
def _text_of(container: dict) -> str:
|
||||
"""Join the text blocks of anything pi shapes as {content: [...]} — a
|
||||
message or a tool result."""
|
||||
return "".join(part.get("text", "") for part in container.get("content", []) or []
|
||||
if isinstance(part, dict) and part.get("type") == "text")
|
||||
|
||||
|
||||
def _clip(text: str, limit: int) -> str:
|
||||
return text if len(text) <= limit else text[:limit].rstrip() + "…"
|
||||
|
||||
|
||||
def _label(tool: str, args: dict) -> str:
|
||||
"""One-line human name for a tool call: `bash: ls -la src`."""
|
||||
value = next((args[key] for key in PRIMARY_ARGS
|
||||
if isinstance(args.get(key), str) and args[key].strip()), "")
|
||||
if not value:
|
||||
value = next((v for v in args.values() if isinstance(v, str) and v.strip()), "")
|
||||
value = " ".join(str(value).split())
|
||||
return f"{tool}: {_clip(value, LABEL_CHARS)}" if value else tool
|
||||
|
||||
|
||||
class ToolCallTracker:
|
||||
"""Folds pi's tool stream into ONE normalized record per completed call.
|
||||
|
||||
pi announces a call as a `toolCall` content block, then emits
|
||||
tool_execution_start / _update / _end for it. Only the end carries the
|
||||
result, so that is where a record is emitted — one trace event per real
|
||||
tool call, the moment it returns, instead of three shapeless ones.
|
||||
|
||||
The record carries the call's real span (`started_at`/`ended_at`), which the
|
||||
tracer writes to columns so the UI can lay tool calls on a time axis without
|
||||
parsing every payload.
|
||||
"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self._open: dict[str, dict] = {}
|
||||
|
||||
def observe(self, event: dict) -> Optional[dict]:
|
||||
"""Returns the record for a finished tool call, else None."""
|
||||
etype = event.get("type", "")
|
||||
if etype == "message_end":
|
||||
for block in event.get("message", {}).get("content", []) or []:
|
||||
if isinstance(block, dict) and block.get("type") == "toolCall":
|
||||
self._announce(block.get("id"), block.get("name"),
|
||||
block.get("arguments"))
|
||||
return None
|
||||
if etype == "tool_execution_start":
|
||||
self._announce(event.get("toolCallId"), event.get("toolName"),
|
||||
event.get("args"))
|
||||
return None
|
||||
if etype != "tool_execution_end":
|
||||
return None
|
||||
|
||||
call_id = str(event.get("toolCallId") or "")
|
||||
opened = self._open.pop(call_id, {})
|
||||
tool = str(event.get("toolName") or opened.get("tool") or "tool")
|
||||
args = event.get("args") or opened.get("args") or {}
|
||||
record = {
|
||||
"tool": tool,
|
||||
"tool_call_id": call_id,
|
||||
"args": {key: _clip(value, ARG_VALUE_CHARS) if isinstance(value, str) else value
|
||||
for key, value in args.items()},
|
||||
"ok": not event.get("isError", False),
|
||||
"label": _label(tool, args),
|
||||
}
|
||||
result_text = _text_of(event.get("result") or {})
|
||||
if result_text:
|
||||
record["result_snippet"] = _clip(result_text, RESULT_SNIPPET_CHARS)
|
||||
record["ended_at"] = now_iso()
|
||||
if opened.get("clock"):
|
||||
record["duration_ms"] = int((time.monotonic() - opened["clock"]) * 1000)
|
||||
if opened.get("started_at"):
|
||||
record["started_at"] = opened["started_at"]
|
||||
return record
|
||||
|
||||
def _announce(self, call_id, tool, args) -> None:
|
||||
"""First sighting starts the clock; a later sighting only fills gaps."""
|
||||
if not call_id:
|
||||
return
|
||||
known = self._open.get(str(call_id), {})
|
||||
self._open[str(call_id)] = {
|
||||
"tool": tool or known.get("tool", ""),
|
||||
"args": args or known.get("args", {}),
|
||||
"started_at": known.get("started_at") or now_iso(), # wall clock, for the row
|
||||
"clock": known.get("clock") or time.monotonic(), # monotonic, for duration
|
||||
}
|
||||
|
||||
|
||||
def run(request: PiRequest, on_event: Optional[Callable[[dict], None]] = None,
|
||||
on_spawn: Optional[Callable[[int], None]] = None,
|
||||
on_exit: Optional[Callable[[int], None]] = None) -> PiResult:
|
||||
"""Run one non-interactive pi turn.
|
||||
|
||||
`on_spawn(pid)` and `on_exit(pid)` bracket the child process so the caller
|
||||
can record it as killable — a hung coding agent is otherwise a pid you have
|
||||
to hunt for in `ps` while the run sits there.
|
||||
"""
|
||||
provider, model_id = resolve_model(request.model)
|
||||
cmd = [
|
||||
PI_PATH, "-p", "--mode", "json",
|
||||
"--provider", provider, "--model", model_id,
|
||||
"--thinking", request.thinking,
|
||||
"--session-id", request.session_id,
|
||||
"--session-dir", request.session_dir,
|
||||
"--system-prompt", request.system_prompt,
|
||||
]
|
||||
if request.tools:
|
||||
cmd += ["--tools", ",".join(request.tools)]
|
||||
for extension in request.extensions:
|
||||
cmd += ["-e", extension]
|
||||
cmd.append(request.prompt)
|
||||
|
||||
raw_path = Path(request.raw_output_path)
|
||||
raw_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
result = PiResult(session_id=request.session_id,
|
||||
context_window=context_window(provider, model_id))
|
||||
# stdin is DEVNULL, deliberately. The prompt travels in argv, so the child
|
||||
# never needs stdin — but inheriting the parent's means pi sees a non-TTY
|
||||
# and can sit forever waiting for piped input that will never arrive or
|
||||
# EOF. That failure is silent and total: no request goes out, no bytes come
|
||||
# back, and the ADW blocks on a read loop with nothing to read. Observed as
|
||||
# a run that sat idle at 0% CPU with an empty raw_output.jsonl.
|
||||
process = subprocess.Popen(cmd, stdin=subprocess.DEVNULL,
|
||||
stdout=subprocess.PIPE, stderr=subprocess.PIPE,
|
||||
text=True, bufsize=1, cwd=request.cwd,
|
||||
env=operator_env())
|
||||
if on_spawn:
|
||||
on_spawn(process.pid)
|
||||
with raw_path.open("a") as raw:
|
||||
assert process.stdout is not None
|
||||
for line in process.stdout:
|
||||
raw.write(line)
|
||||
raw.flush() # events land on disk as they happen
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
event = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
if event.get("type") == "message_end":
|
||||
message = event.get("message", {})
|
||||
if message.get("role") == "assistant":
|
||||
text = _text_of(message)
|
||||
if text:
|
||||
result.text = text # last assistant message wins
|
||||
usage = message.get("usage", {}) or {}
|
||||
turn = _context_tokens(usage)
|
||||
result.tokens += turn
|
||||
result.usage.add_turn(usage, turn)
|
||||
# Occupancy is read off the last VALID assistant turn, the
|
||||
# way pi does it — an aborted or errored turn reports usage
|
||||
# you can't trust, so it must not overwrite a good reading.
|
||||
if turn and message.get("stopReason") not in ("aborted", "error"):
|
||||
result.context_tokens = turn
|
||||
result.cost += (usage.get("cost", {}) or {}).get("total", 0.0) or 0.0
|
||||
if on_event:
|
||||
on_event(event)
|
||||
|
||||
stderr = process.stderr.read() if process.stderr else ""
|
||||
result.returncode = process.wait()
|
||||
if on_exit:
|
||||
on_exit(process.pid)
|
||||
if result.returncode != 0 and not result.text:
|
||||
raise RuntimeError(f"pi exited {result.returncode}: {stderr.strip()[-800:]}")
|
||||
return result
|
||||
320
sssf/templates/adws/adw_modules/agents.py
Normal file
320
sssf/templates/adws/adw_modules/agents.py
Normal file
|
|
@ -0,0 +1,320 @@
|
|||
"""Config loading/validation and agent execution.
|
||||
|
||||
Every ADW validates its agents before running (fail fast, nothing spawns
|
||||
against a half-valid config). Every agent call parses against a concrete
|
||||
output type; parse failures and gate violations re-prompt the SAME session
|
||||
with a correction — context intact, bounded retries. Agent proposes, code
|
||||
disposes.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
|
||||
import yaml
|
||||
|
||||
from . import agent_omp, agent_pi, permissions, prompts
|
||||
from .data_types import (AgentCall, AgentConfig, EnvelopeBase, EventRecord,
|
||||
GateCheck, GateReport, Phase, PiRequest, SSSFConfig,
|
||||
UsageBreakdown)
|
||||
from .utils import new_id
|
||||
|
||||
JSON_FIX_ATTEMPTS = 2 # continue-with-correction attempts for malformed JSON
|
||||
|
||||
|
||||
class GateFailure(RuntimeError):
|
||||
pass
|
||||
|
||||
|
||||
# ── config ───────────────────────────────────────────────────────────────────
|
||||
|
||||
def load_config(path: str = "adws/adw_sssf_config/sssf.config.yaml") -> SSSFConfig:
|
||||
raw = yaml.safe_load(Path(path).read_text()) or {}
|
||||
defaults = raw.get("defaults", {}) or {}
|
||||
for agent in raw.get("agents", []) or []:
|
||||
for key in ("coding_agent", "model", "thinking", "color", "tools", "writes"):
|
||||
if key in defaults:
|
||||
agent.setdefault(key, defaults[key])
|
||||
agent.setdefault("harness_engineering", defaults.get("harness_engineering", []))
|
||||
return SSSFConfig(**raw)
|
||||
|
||||
|
||||
def resolve(cfg: SSSFConfig, name: str) -> AgentConfig:
|
||||
for agent in cfg.agents:
|
||||
if agent.name == name:
|
||||
return agent
|
||||
raise SystemExit(f"agent {name!r} is not defined in the config — "
|
||||
f"available: {[a.name for a in cfg.agents]}")
|
||||
|
||||
|
||||
def validate(cfg: SSSFConfig, required: list[str]) -> None:
|
||||
"""Fail fast: every required name must resolve to a usable agent."""
|
||||
problems = []
|
||||
for name in required:
|
||||
try:
|
||||
agent = resolve(cfg, name)
|
||||
except SystemExit as e:
|
||||
problems.append(str(e))
|
||||
continue
|
||||
if agent.coding_agent not in ("pi", "omp"):
|
||||
problems.append(f"agent {name!r}: coding_agent {agent.coding_agent!r} "
|
||||
f"is not implemented (pi and omp are)")
|
||||
for label, ref in (("system", agent.prompt_engineering.system),
|
||||
("user", agent.prompt_engineering.user)):
|
||||
if not Path(ref).is_file():
|
||||
problems.append(f"agent {name!r}: {label} prompt not found: {ref}")
|
||||
try:
|
||||
_resolve_model(agent)
|
||||
except ValueError as e:
|
||||
problems.append(f"agent {name!r}: {e}")
|
||||
if problems:
|
||||
raise SystemExit("config validation failed:\n- " + "\n- ".join(problems))
|
||||
|
||||
|
||||
# ── execution ────────────────────────────────────────────────────────────────
|
||||
|
||||
def execute(run, phase: Phase, call: AgentCall) -> EnvelopeBase:
|
||||
"""One agent call: render prompts -> pi run -> typed parse -> gates -> envelope."""
|
||||
agent = resolve(run.cfg, phase.params.owner)
|
||||
agent_dir = run.session_dir / agent.name
|
||||
agent_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
variables = {
|
||||
"prompt": call.prompt,
|
||||
"previous_envelope": call.previous.model_dump_json(indent=2) if call.previous else "(none)",
|
||||
"context_handoff_dir": str(run.context_handoff_dir),
|
||||
}
|
||||
system_text = prompts.render(agent.prompt_engineering.system, variables)
|
||||
user_text = prompts.render(agent.prompt_engineering.user, variables)
|
||||
prompts.save(agent_dir / "prompts", "system.md", system_text)
|
||||
prompts.save(agent_dir / "prompts", "user.md", user_text)
|
||||
|
||||
session_id = _agent_session_id(run, agent)
|
||||
run.tracer.event(EventRecord(adw_id=run.adw_id, phase_id=phase.phase_id,
|
||||
type="agent_start", name=agent.name,
|
||||
payload={"model": agent.model, "thinking": agent.thinking,
|
||||
"color": agent.color,
|
||||
"session_id": session_id,
|
||||
"coding_agent": agent.coding_agent,
|
||||
"purpose": agent.purpose,
|
||||
"tools": agent.tools, # None = all tools
|
||||
"harness_engineering": agent.harness_engineering}))
|
||||
run.console.agent_started(agent.name, agent.model, session_id)
|
||||
|
||||
# Parse retries and gate corrections re-enter the SAME pi session, so the
|
||||
# last send is the one whose context occupancy is current — while spend is
|
||||
# the opposite: every send costs, so usage accumulates across all of them.
|
||||
latest: agent_pi.PiResult | None = None
|
||||
spent = UsageBreakdown()
|
||||
|
||||
def send(prompt_text: str) -> agent_pi.PiResult:
|
||||
nonlocal latest
|
||||
request = PiRequest(
|
||||
prompt=prompt_text,
|
||||
system_prompt=system_text,
|
||||
model=agent.model,
|
||||
thinking=agent.thinking,
|
||||
session_id=session_id,
|
||||
# absolute: these are read by the coding-agent subprocess, which runs in repo_root
|
||||
session_dir=_agent_session_dir(agent_dir, agent.coding_agent),
|
||||
raw_output_path=str((agent_dir / "raw_output.jsonl").resolve()),
|
||||
tools=agent.tools,
|
||||
extensions=agent.harness_engineering,
|
||||
cwd=str(run.repo_root),
|
||||
)
|
||||
result = _agent_runner(agent)(
|
||||
request,
|
||||
on_event=_event_forwarder(run, phase, agent.name),
|
||||
on_spawn=lambda pid: run.tracer.process_start(
|
||||
run.adw_id, "agent", agent.name, pid,
|
||||
f"{agent.coding_agent} {agent.name} {agent.model}"),
|
||||
on_exit=lambda pid: run.tracer.process_end(run.adw_id, pid))
|
||||
run.add_usage(result.tokens, result.cost)
|
||||
spent.merge(result.usage)
|
||||
latest = result
|
||||
return result
|
||||
|
||||
# What the tree looked like before this agent got its hands on it. Every
|
||||
# send in this phase — first prompt, JSON retries, gate corrections — is
|
||||
# measured against this one baseline.
|
||||
tree_before = permissions.snapshot(run)
|
||||
|
||||
result = send(user_text)
|
||||
envelope, attempt = _parse_with_retries(run, phase, call, result, send)
|
||||
|
||||
# claim gates — violations flow back into the SAME session as corrections
|
||||
for gate_attempt in range(1, max(1, phase.params.retries + 1) + 1):
|
||||
violations = []
|
||||
for gate in call.gates:
|
||||
report = _as_report(gate(envelope, run))
|
||||
found = report.violations
|
||||
run.tracer.gate_row(phase, gate.__name__, report, gate_attempt)
|
||||
run.tracer.event(EventRecord(
|
||||
adw_id=run.adw_id, phase_id=phase.phase_id,
|
||||
type="gate_fail" if found else "gate_pass", name=gate.__name__,
|
||||
payload={"attempt": gate_attempt, "violations": found,
|
||||
"checks": [c.model_dump() for c in report.checks]}))
|
||||
run.console.gate_result(gate.__name__, report)
|
||||
violations.extend(found)
|
||||
if not violations:
|
||||
break
|
||||
if gate_attempt > phase.params.retries:
|
||||
raise GateFailure(f"{agent.name} failed gates after {gate_attempt} attempt(s):\n- "
|
||||
+ "\n- ".join(violations))
|
||||
phase.attempt = gate_attempt
|
||||
run.console.retry(agent.name, gate_attempt, phase.params.retries,
|
||||
f"{len(violations)} gate violation(s)")
|
||||
correction = ("Your previous response failed validation:\n- "
|
||||
+ "\n- ".join(violations)
|
||||
+ "\n\nFix these problems, then re-emit ONLY your Report JSON.")
|
||||
result = send(correction)
|
||||
envelope, attempt = _parse_with_retries(run, phase, call, result, send)
|
||||
|
||||
# Permission is checked after every send is done, and before the envelope is
|
||||
# accepted: an agent does not get to report success on a phase in which it
|
||||
# wrote somewhere it was not allowed to.
|
||||
try:
|
||||
touched = permissions.enforce(run, phase, agent, tree_before)
|
||||
except permissions.PermissionBreach as breach:
|
||||
run.tracer.event(EventRecord(adw_id=run.adw_id, phase_id=phase.phase_id,
|
||||
type="error", name="permission_breach",
|
||||
payload={"agent": agent.name, "error": str(breach),
|
||||
"writes": agent.writes,
|
||||
"protected_files": run.cfg.defaults.protected_files}))
|
||||
raise
|
||||
if touched:
|
||||
run.tracer.event(EventRecord(adw_id=run.adw_id, phase_id=phase.phase_id,
|
||||
type="log", name="paths_touched",
|
||||
payload={"agent": agent.name, "paths": touched}))
|
||||
|
||||
_persist_envelope(run, phase, agent.name, call, envelope, attempt, valid=True)
|
||||
run.console.envelope_summary(envelope)
|
||||
context = latest or result
|
||||
run.tracer.agent_session_row(run.adw_id, agent, session_id,
|
||||
context_tokens=context.context_tokens,
|
||||
context_window=context.context_window)
|
||||
run.save_agent_map(agent.name, {"session_id": session_id, "model": agent.model,
|
||||
"coding_agent": agent.coding_agent})
|
||||
run.tracer.event(EventRecord(adw_id=run.adw_id, phase_id=phase.phase_id,
|
||||
type="handoff", name=agent.name,
|
||||
payload={"artifacts": envelope.artifacts,
|
||||
"summary": envelope.summary}))
|
||||
run.tracer.event(EventRecord(adw_id=run.adw_id, phase_id=phase.phase_id,
|
||||
type="agent_end", name=agent.name,
|
||||
# Phase totals, not the last send's: a retried
|
||||
# phase paid for every attempt.
|
||||
tokens=spent.total_tokens,
|
||||
payload={"cost": spent.total_cost,
|
||||
"usage": spent.model_dump(),
|
||||
"context_tokens": context.context_tokens,
|
||||
"context_window": context.context_window}))
|
||||
run.console.agent_finished(agent.name, spent.total_tokens, spent.total_cost)
|
||||
if envelope.status != "success":
|
||||
raise RuntimeError(f"{agent.name} reported status={envelope.status!r}: {envelope.summary}")
|
||||
return envelope
|
||||
|
||||
|
||||
# ── internals ────────────────────────────────────────────────────────────────
|
||||
|
||||
def _as_report(result) -> GateReport:
|
||||
"""Accept a GateReport, or a legacy gate that returned a violations list."""
|
||||
if isinstance(result, GateReport):
|
||||
return result
|
||||
return GateReport(checks=[GateCheck(item=str(v), ok=False) for v in (result or [])])
|
||||
|
||||
|
||||
def _resolve_model(agent: AgentConfig) -> tuple[str, str]:
|
||||
"""Resolve an agent's model pattern against its coding agent's catalog."""
|
||||
if agent.coding_agent == "omp":
|
||||
return agent_omp.resolve_model(agent.model)
|
||||
return agent_pi.resolve_model(agent.model)
|
||||
|
||||
|
||||
def _agent_runner(agent: AgentConfig):
|
||||
"""The run() callable for an agent's coding agent."""
|
||||
if agent.coding_agent == "omp":
|
||||
return agent_omp.run
|
||||
return agent_pi.run
|
||||
|
||||
|
||||
def _agent_session_dir(agent_dir, coding_agent: str) -> str:
|
||||
"""Absolute session dir for the coding agent's subprocess."""
|
||||
sub = "omp_sessions" if coding_agent == "omp" else "pi_sessions"
|
||||
return str((agent_dir / sub).resolve())
|
||||
|
||||
|
||||
def _agent_session_id(run, agent: AgentConfig) -> str:
|
||||
entry = run.agent_map.get(agent.name)
|
||||
if entry and entry.get("model") == agent.model:
|
||||
return entry["session_id"] # rejoin the existing context window
|
||||
return f"sssf-{run.adw_id}-{agent.name}-{new_id(4)}"
|
||||
|
||||
|
||||
def _event_forwarder(run, phase: Phase, agent_name: str):
|
||||
"""One tool_call event per real tool call, with its exact args and result."""
|
||||
tracker = agent_pi.ToolCallTracker()
|
||||
|
||||
def forward(event: dict) -> None:
|
||||
record = tracker.observe(event)
|
||||
if record is None:
|
||||
return
|
||||
# The call's span rides the columns; duration_ms stays in the payload as
|
||||
# pi's own authoritative number.
|
||||
run.tracer.event(EventRecord(adw_id=run.adw_id, phase_id=phase.phase_id,
|
||||
type="tool_call", name=record.pop("label"),
|
||||
started_at=record.pop("started_at", None),
|
||||
ended_at=record.pop("ended_at", None),
|
||||
payload={**record, "agent": agent_name}))
|
||||
return forward
|
||||
|
||||
|
||||
def _extract_json(text: str) -> dict:
|
||||
candidate = text
|
||||
if "```" in text:
|
||||
for block in text.split("```")[1::2]:
|
||||
block = block.removeprefix("json").strip()
|
||||
if block.startswith("{"):
|
||||
candidate = block
|
||||
break
|
||||
start, end = candidate.find("{"), candidate.rfind("}")
|
||||
if start == -1 or end <= start:
|
||||
raise ValueError("no JSON object found in the response")
|
||||
return json.loads(candidate[start:end + 1])
|
||||
|
||||
|
||||
def _parse_with_retries(run, phase: Phase, call: AgentCall, result, send):
|
||||
"""Parse the final response against the declared output type; on failure,
|
||||
continue the SAME session with a correction (bounded)."""
|
||||
for attempt in range(1, JSON_FIX_ATTEMPTS + 2):
|
||||
try:
|
||||
payload = _extract_json(result.text)
|
||||
return call.output_type.model_validate(payload), attempt
|
||||
except Exception as error:
|
||||
_persist_envelope(run, phase, phase.params.owner, call, None, attempt,
|
||||
valid=False, raw=result.text)
|
||||
if attempt > JSON_FIX_ATTEMPTS:
|
||||
raise RuntimeError(
|
||||
f"{phase.params.owner} never produced valid "
|
||||
f"{call.output_type.__name__} JSON: {error}") from error
|
||||
run.console.retry(phase.params.owner, attempt, JSON_FIX_ATTEMPTS,
|
||||
f"invalid {call.output_type.__name__} JSON: {error}")
|
||||
fields = ", ".join(call.output_type.model_fields.keys())
|
||||
result = send(
|
||||
f"Your response was not valid JSON for the required structure "
|
||||
f"({error}). Respond again with ONLY a JSON object with these "
|
||||
f"fields: {fields}. No prose, no code fences.")
|
||||
|
||||
|
||||
def _persist_envelope(run, phase: Phase, agent_name: str, call: AgentCall,
|
||||
envelope: Optional[EnvelopeBase], attempt: int,
|
||||
valid: bool, raw: str = "") -> None:
|
||||
payload_json = envelope.model_dump_json(indent=2) if envelope else json.dumps({"raw": raw[-2000:]})
|
||||
run.tracer.envelope_row(phase, agent_name, call.output_type.__name__,
|
||||
payload_json, valid, attempt)
|
||||
if envelope:
|
||||
record = {"agent_name": agent_name, "purpose": resolve(run.cfg, agent_name).purpose,
|
||||
"output_type": call.output_type.__name__, "attempt": attempt,
|
||||
**envelope.model_dump()}
|
||||
(run.session_dir / agent_name / "envelope.json").write_text(json.dumps(record, indent=2))
|
||||
103
sssf/templates/adws/adw_modules/changes.py
Normal file
103
sssf/templates/adws/adw_modules/changes.py
Normal file
|
|
@ -0,0 +1,103 @@
|
|||
"""Deterministic change capture: what was built, straight from git.
|
||||
|
||||
"What changed since main" is not a judgement call — it is two git commands and
|
||||
a subtraction. So it is code, and an agent is only handed the result. The
|
||||
capture writes the full diff into `context_handoff/` and returns a ChangeSet;
|
||||
`as_envelope` adapts that into the one door every agent handoff uses.
|
||||
|
||||
The base is resolved, not assumed. Off the base branch the diff covers the
|
||||
whole branch plus the working tree; on it, the uncommitted tree; and on a clean
|
||||
tree, the last commit — because "document the work that was just done" still
|
||||
has an answer right after a chain committed. Whichever it picked rides along in
|
||||
`BaseRef.reason`, so the trace never leaves you guessing what a diff was
|
||||
measured against.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from . import git_helper
|
||||
from .data_types import BaseRef, ChangeCapture, ChangeSet, ChangesOutput
|
||||
|
||||
DIFF_FILENAME = "changes.diff"
|
||||
|
||||
|
||||
def resolve_base(ref: str) -> BaseRef:
|
||||
"""Pick the commit the work is measured from, and record why that one."""
|
||||
if not git_helper.is_repo():
|
||||
raise RuntimeError(
|
||||
"not a git repository — change capture needs one. Run `git init` in "
|
||||
"the repo root before running an ADW that documents a change.")
|
||||
if not git_helper.ref_exists(ref):
|
||||
raise RuntimeError(
|
||||
f"base ref {ref!r} does not exist in this repository — pass --base "
|
||||
f"with a ref that does (e.g. --base master, --base HEAD~1).")
|
||||
|
||||
# Built first, then given its reason: BaseRef.label knows how to print a
|
||||
# pinned sha, and the reason is the line a human reads in the trace.
|
||||
base = BaseRef(ref=ref, commit=git_helper.merge_base(ref, "HEAD"))
|
||||
if git_helper.short_sha(base.commit) != git_helper.short_sha("HEAD"):
|
||||
base.reason = (f"HEAD is ahead of {base.label} — diffing every commit since, "
|
||||
f"plus the working tree")
|
||||
elif git_helper.is_dirty():
|
||||
base.reason = f"HEAD is on {base.label} — diffing the uncommitted working tree"
|
||||
elif git_helper.ref_exists("HEAD~1"):
|
||||
base.commit = git_helper.rev("HEAD~1")
|
||||
base.reason = (f"HEAD is on {base.label} with a clean tree — falling back to "
|
||||
f"the last commit")
|
||||
else:
|
||||
base.reason = f"HEAD is on {base.label} with a clean tree and no parent commit"
|
||||
return base
|
||||
|
||||
|
||||
def capture(run, params: ChangeCapture) -> ChangeSet:
|
||||
"""Diff the working tree against the resolved base and persist the evidence."""
|
||||
base = resolve_base(params.base)
|
||||
files = git_helper.diff_files(base.commit)
|
||||
untracked = git_helper.untracked_files() if params.include_untracked else []
|
||||
insertions, deletions = git_helper.diff_counts(base.commit)
|
||||
stat = git_helper.diff_stat(base.commit)
|
||||
|
||||
text = git_helper.diff_text(base.commit)
|
||||
lines = text.splitlines()
|
||||
truncated = len(lines) > params.max_diff_lines
|
||||
if truncated:
|
||||
text = "\n".join(lines[:params.max_diff_lines])
|
||||
text += (f"\n\n[truncated at {params.max_diff_lines} lines of "
|
||||
f"{len(lines)} — run `git diff {base.commit}` for the rest]")
|
||||
|
||||
# Untracked files are absent from `git diff` by construction, so they are
|
||||
# named here rather than silently missing from the record. The reader has
|
||||
# `read` and can open any of them.
|
||||
untracked_block = ("\n".join(f" {f}" for f in untracked) if untracked
|
||||
else " (none)")
|
||||
diff_path = run.context_handoff_dir / DIFF_FILENAME
|
||||
diff_path.write_text(
|
||||
f"# changes since {base.label} @ {git_helper.short_sha(base.commit)}\n"
|
||||
f"# {base.reason}\n"
|
||||
f"# +{insertions} -{deletions} across {len(files)} tracked file(s)\n\n"
|
||||
f"## stat\n{stat or ' (no tracked changes)'}\n\n"
|
||||
f"## untracked files\n{untracked_block}\n\n"
|
||||
f"## diff\n{text}\n")
|
||||
|
||||
return ChangeSet(base=base, files=files, untracked=untracked,
|
||||
insertions=insertions, deletions=deletions, stat=stat,
|
||||
diff_path=str(diff_path), truncated=truncated)
|
||||
|
||||
|
||||
def as_envelope(changes: ChangeSet, notes: str = "") -> ChangesOutput:
|
||||
"""Wrap a captured change so an agent can be handed it directly."""
|
||||
total = len(changes.files) + len(changes.untracked)
|
||||
return ChangesOutput(
|
||||
status="success",
|
||||
summary=(f"{total} file(s) changed since {changes.base.label} "
|
||||
f"(+{changes.insertions} -{changes.deletions})"),
|
||||
artifacts=[changes.diff_path],
|
||||
notes_for_next_agent=notes,
|
||||
base=f"{changes.base.label} @ {git_helper.short_sha(changes.base.commit)} "
|
||||
f"— {changes.base.reason}",
|
||||
changed_files=changes.files + changes.untracked,
|
||||
insertions=changes.insertions,
|
||||
deletions=changes.deletions,
|
||||
stat=changes.stat,
|
||||
diff_path=changes.diff_path,
|
||||
)
|
||||
131
sssf/templates/adws/adw_modules/console.py
Normal file
131
sssf/templates/adws/adw_modules/console.py
Normal file
|
|
@ -0,0 +1,131 @@
|
|||
"""Console reporter: one narrative, two destinations.
|
||||
|
||||
Every line an ADW prints ALSO lands in the db as a `log` event, so the swim-lane
|
||||
UI reads the same story the terminal does. Both go through `_emit` — print and
|
||||
trace cannot drift. Plain sequential lines only: no spinners, no live displays,
|
||||
so a CI log reads exactly like a terminal.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from rich.console import Console as RichConsole
|
||||
from rich.markup import escape
|
||||
from rich.panel import Panel
|
||||
from rich.text import Text
|
||||
|
||||
from .data_types import EnvelopeBase, EventRecord, Phase
|
||||
|
||||
KIND_COLOR = {"engineer": "cyan", "agent": "magenta", "code": "yellow"}
|
||||
MAX_LINE = 160 # dynamic text (summaries, violations, errors) is clipped
|
||||
|
||||
|
||||
def _clip(text: str, limit: int = MAX_LINE) -> str:
|
||||
text = " ".join(str(text).split())
|
||||
return text if len(text) <= limit else text[: limit - 1] + "…"
|
||||
|
||||
|
||||
class Console:
|
||||
"""Bound to one run's tracer. Reachable as `run.console` everywhere."""
|
||||
|
||||
def __init__(self, tracer, adw_id: str):
|
||||
self.tracer = tracer
|
||||
self.adw_id = adw_id
|
||||
self.phase_id = "" # current lane — log events attach to it
|
||||
self.phase_name = ""
|
||||
self.results: list[str] = [] # phase statuses, for the summary
|
||||
self._finished = False # the summary panel prints once
|
||||
self._out = RichConsole(highlight=False, soft_wrap=True)
|
||||
|
||||
# ── the one helper: print AND trace, always together ────────────────────
|
||||
def _emit(self, markup: str, level: str = "info", renderable=None) -> None:
|
||||
text = Text.from_markup(markup)
|
||||
self._out.print(renderable if renderable is not None else text)
|
||||
self.tracer.event(EventRecord(
|
||||
adw_id=self.adw_id, phase_id=self.phase_id, type="log",
|
||||
name=self.phase_name or "console",
|
||||
payload={"message": text.plain, "level": level}))
|
||||
|
||||
# ── session ─────────────────────────────────────────────────────────────
|
||||
def session_started(self, adw_id: str, engineer: str) -> None:
|
||||
self._emit(f"[bold cyan]adw_id:[/bold cyan] [bold]{escape(adw_id)}[/bold]"
|
||||
f" [dim]engineer[/dim] {escape(engineer)}")
|
||||
|
||||
def session_finished(self, ok: bool, tokens: int, cost: float, db_path: str) -> None:
|
||||
if self._finished:
|
||||
return
|
||||
self._finished = True
|
||||
passed = sum(1 for r in self.results if r == "success")
|
||||
status = "[green]✓ success[/green]" if ok else "[red]✗ fail[/red]"
|
||||
rows = [f" [dim]status[/dim] {status}",
|
||||
f" [dim]phases[/dim] {passed}/{len(self.results)} passed",
|
||||
f" [dim]tokens[/dim] {tokens:,}",
|
||||
f" [dim]cost[/dim] ${cost:.4f}",
|
||||
f" [dim]adw_id[/dim] {escape(self.adw_id)}",
|
||||
f" [dim]db[/dim] {escape(str(db_path))}",
|
||||
f" [dim]next[/dim] [bold]just phases {escape(self.adw_id)}[/bold]"]
|
||||
panel = Panel(Text.from_markup("\n".join(rows)),
|
||||
title="[bold]ADW complete[/bold]",
|
||||
border_style="green" if ok else "red", expand=False)
|
||||
plain = (f"session {self.adw_id} {'success' if ok else 'fail'} · "
|
||||
f"{passed}/{len(self.results)} phases · {tokens:,} tokens · ${cost:.4f}")
|
||||
self._emit(escape(plain), level="info" if ok else "error", renderable=panel)
|
||||
|
||||
# ── phases ──────────────────────────────────────────────────────────────
|
||||
def phase_started(self, phase: Phase) -> None:
|
||||
self.phase_id, self.phase_name = phase.phase_id, phase.params.name
|
||||
p = phase.params
|
||||
color = KIND_COLOR.get(p.kind, "white")
|
||||
line = (f"[bold {color}]▶ {phase.seq:02d} {escape(p.name)}[/bold {color}]"
|
||||
f" [{color}]{p.kind}[/{color}] [dim]· {escape(p.owner)}[/dim]")
|
||||
if p.description:
|
||||
line += f" [dim]{escape(_clip(p.description))}[/dim]"
|
||||
self._emit(line)
|
||||
|
||||
def phase_ended(self, phase: Phase, seconds: float) -> None:
|
||||
ok = phase.status == "success"
|
||||
self.results.append(phase.status)
|
||||
line = (f" {'[green]✓[/green]' if ok else '[red]✗[/red]'} "
|
||||
f"{escape(phase.params.name)} [dim]{seconds:.1f}s[/dim]")
|
||||
if not ok and phase.error:
|
||||
line += f" [red]{escape(_clip(phase.error))}[/red]"
|
||||
self._emit(line, level="info" if ok else "error")
|
||||
self.phase_id, self.phase_name = "", ""
|
||||
|
||||
def note(self, message: str) -> None:
|
||||
"""Free-form detail inside the current phase — what `ph.log()` recorded."""
|
||||
self._emit(f" [dim]· {escape(_clip(message))}[/dim]")
|
||||
|
||||
# ── agents ──────────────────────────────────────────────────────────────
|
||||
def agent_started(self, name: str, model: str, session_id: str) -> None:
|
||||
self._emit(f" [magenta]▸[/magenta] {escape(name)} [dim]{escape(model)}[/dim]"
|
||||
f" [dim]session {escape(session_id)}[/dim]")
|
||||
|
||||
def agent_finished(self, name: str, tokens: int, cost: float) -> None:
|
||||
self._emit(f" [dim]└ {escape(name)} used {tokens:,} tokens · ${cost:.4f}[/dim]")
|
||||
|
||||
def retry(self, name: str, attempt: int, limit: int, reason: str) -> None:
|
||||
self._emit(f" [yellow]⟳[/yellow] {escape(name)} retry {attempt}/{limit} "
|
||||
f"[dim]— same session · {escape(_clip(reason))}[/dim]", level="warn")
|
||||
|
||||
# ── verification ────────────────────────────────────────────────────────
|
||||
def gate_result(self, name: str, report) -> None:
|
||||
"""A gate reports WHAT it checked, not just whether it passed."""
|
||||
ok = report.passed
|
||||
mark = "[green]✓[/green]" if ok else "[red]✗[/red]"
|
||||
summary = (f"{len(report.checks)} checked" if ok
|
||||
else f"[red]{len(report.violations)} of {len(report.checks)} failed[/red]")
|
||||
self._emit(f" {mark} gate [dim]{escape(name)}[/dim] [dim]{summary}[/dim]",
|
||||
level="info" if ok else "error")
|
||||
for check in report.checks:
|
||||
style = "dim" if check.ok else "dim red"
|
||||
detail = f" — {_clip(check.note)}" if check.note else ""
|
||||
self._emit(f" [{style}]{'·' if check.ok else '✗'} {escape(_clip(check.item))}"
|
||||
f"{escape(detail)}[/{style}]", level="info" if check.ok else "error")
|
||||
|
||||
def envelope_summary(self, envelope: EnvelopeBase) -> None:
|
||||
ok = envelope.status == "success"
|
||||
line = (f" {'[green]✓[/green]' if ok else '[red]✗[/red]'} "
|
||||
f"{type(envelope).__name__} [dim]{escape(_clip(envelope.summary))}[/dim]")
|
||||
self._emit(line, level="info" if ok else "error")
|
||||
if envelope.artifacts:
|
||||
self._emit(f" [dim]artifacts: {escape(_clip(', '.join(envelope.artifacts)))}[/dim]")
|
||||
447
sssf/templates/adws/adw_modules/data_types.py
Normal file
447
sssf/templates/adws/adw_modules/data_types.py
Normal file
|
|
@ -0,0 +1,447 @@
|
|||
"""Concrete data types for the SSSF ADW system.
|
||||
|
||||
RULE (four-param rule): any function that takes more than 4 parameters takes
|
||||
ONE of these objects instead. AgentCall and PhaseParams are the pattern.
|
||||
|
||||
Every agent call declares a concrete output type — an EnvelopeBase subclass —
|
||||
that its final JSON response is parsed against. No untyped handoffs.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any, Callable, Literal, Optional, Type
|
||||
|
||||
from pydantic import BaseModel, Field, ValidationInfo, field_validator
|
||||
|
||||
PhaseKind = Literal["engineer", "agent", "code"]
|
||||
PhaseStatus = Literal["queued", "running", "success", "fail"]
|
||||
|
||||
|
||||
# ── Phases ────────────────────────────────────────────────────────────────────
|
||||
|
||||
class PhaseParams(BaseModel):
|
||||
"""Everything run.phase() needs. Passed as one object, never loose params."""
|
||||
|
||||
name: str # short id, unique within the run: "plan", "build"
|
||||
kind: PhaseKind # which lane the block renders in
|
||||
owner: str # engineer's name, "git", or an agent name from config
|
||||
description: str # REQUIRED: what this phase does and why — see below
|
||||
retries: int = 0 # agent phases: gate-failure retries via continue
|
||||
|
||||
@field_validator("description")
|
||||
@classmethod
|
||||
def _description_must_be_earned(cls, value: str, info: ValidationInfo) -> str:
|
||||
"""A phase name identifies; a description explains. Both are required.
|
||||
|
||||
The description is the only sentence the trace, the console, and the
|
||||
phase block in the UI ever show about intent — everything else is ids,
|
||||
statuses, and timings. `commit_plan: "Commit the plan"` tells a reader
|
||||
nothing they could not already see, so an echo is rejected the same way
|
||||
a blank one is. This is a construction-time error on purpose: it fires
|
||||
before the phase opens, not after a run is already in the trace.
|
||||
"""
|
||||
text = " ".join(value.split())
|
||||
name = str(info.data.get("name", "?"))
|
||||
if not text:
|
||||
raise ValueError(
|
||||
f"phase {name!r}: description is required — one sentence on what this "
|
||||
f"phase does and why. It is what the trace and the UI show.")
|
||||
if text.rstrip(".").casefold() == name.replace("_", " ").casefold():
|
||||
raise ValueError(
|
||||
f"phase {name!r}: description {text!r} only restates the phase name — "
|
||||
f"say what it does and why instead.")
|
||||
return text
|
||||
|
||||
|
||||
class Phase(BaseModel):
|
||||
"""The persisted phase record — PhaseParams plus lifecycle."""
|
||||
|
||||
phase_id: str
|
||||
adw_id: str
|
||||
seq: int
|
||||
params: PhaseParams
|
||||
status: PhaseStatus = "fail" # success must be earned
|
||||
attempt: int = 0
|
||||
error: Optional[str] = None
|
||||
started_at: Optional[str] = None
|
||||
ended_at: Optional[str] = None
|
||||
|
||||
|
||||
# ── Envelopes (agent output types) ───────────────────────────────────────────
|
||||
|
||||
class EnvelopeBase(BaseModel):
|
||||
"""Base of every agent's final JSON response. Output types extend this."""
|
||||
|
||||
status: Literal["success", "fail"]
|
||||
summary: str = ""
|
||||
artifacts: list[str] = Field(default_factory=list)
|
||||
notes_for_next_agent: str = ""
|
||||
|
||||
|
||||
class GenericOutput(EnvelopeBase):
|
||||
pass
|
||||
|
||||
|
||||
class PlanOutput(EnvelopeBase):
|
||||
# Subject for committing the PLAN — the spec file the planner wrote, not the
|
||||
# implementation it describes. Each agent's commit_message covers its own
|
||||
# work product, so a chain that commits per step never reuses one agent's
|
||||
# words for another agent's diff.
|
||||
commit_message: str = ""
|
||||
|
||||
|
||||
class BuildOutput(EnvelopeBase):
|
||||
changed_files: list[str] = Field(default_factory=list)
|
||||
commit_message: str = "" # consumed by the git commit phase
|
||||
|
||||
|
||||
class ScoutFinding(BaseModel):
|
||||
file: str
|
||||
note: str = ""
|
||||
|
||||
|
||||
class ScoutOutput(EnvelopeBase):
|
||||
findings: list[ScoutFinding] = Field(default_factory=list)
|
||||
|
||||
|
||||
class ReviewFinding(BaseModel):
|
||||
"""One thing the request (or plan) asked for, and whether it is there."""
|
||||
|
||||
requirement: str # the ask, in the requester's words
|
||||
met: bool
|
||||
evidence: str = "" # where it lives, or what is missing
|
||||
|
||||
|
||||
class ReviewOutput(EnvelopeBase):
|
||||
"""Confirmation that what was built is what was asked for — not a test run."""
|
||||
|
||||
approved: bool = False
|
||||
findings: list[ReviewFinding] = Field(default_factory=list)
|
||||
blocking: list[str] = Field(default_factory=list) # what must change before approval
|
||||
|
||||
|
||||
class DocumentOutput(EnvelopeBase):
|
||||
"""Where the write-up of a completed change landed."""
|
||||
|
||||
document_path: str = "" # the doc in the repo, e.g. app_docs/<adw_id>_<slug>.md
|
||||
documented_files: list[str] = Field(default_factory=list)
|
||||
commit_message: str = ""
|
||||
|
||||
|
||||
# ── Deterministic quality blocks ─────────────────────────────────────────────
|
||||
|
||||
QualityArea = Literal["frontend", "backend"]
|
||||
QualityOperation = Literal["lint", "typecheck", "build"]
|
||||
|
||||
|
||||
class QualityCheckSpec(BaseModel):
|
||||
"""One deterministic quality command."""
|
||||
|
||||
name: str
|
||||
area: QualityArea
|
||||
operation: QualityOperation
|
||||
argv: list[str]
|
||||
timeout_seconds: int = 120
|
||||
|
||||
|
||||
class QualityCheckResult(BaseModel):
|
||||
"""Captured evidence from one quality command."""
|
||||
|
||||
name: str
|
||||
area: QualityArea
|
||||
operation: QualityOperation
|
||||
command: str
|
||||
returncode: int
|
||||
passed: bool
|
||||
duration_seconds: float
|
||||
output_artifact: str
|
||||
# The tail of stdout+stderr, verbatim and unparsed. A failure has to travel
|
||||
# back to the builder as an envelope, and the builder cannot open a log file
|
||||
# it was never handed — so the evidence rides along. Deliberately raw: every
|
||||
# runner formats failures differently and a generic parser would be
|
||||
# confidently wrong. The full log is always at output_artifact.
|
||||
output_tail: str = ""
|
||||
|
||||
|
||||
class QualityResult(BaseModel):
|
||||
"""Aggregate result from a quality block: every check it ran, and the verdict."""
|
||||
|
||||
passed: bool
|
||||
checks: list[QualityCheckResult] = Field(default_factory=list)
|
||||
failures: list[str] = Field(default_factory=list)
|
||||
artifacts: list[str] = Field(default_factory=list)
|
||||
|
||||
|
||||
# ── Change capture (git diff, deterministic) ─────────────────────────────────
|
||||
|
||||
class ChangeCapture(BaseModel):
|
||||
"""Everything documentation.capture() needs. One object, never loose params."""
|
||||
|
||||
base: str = "main" # the ref the work is measured against
|
||||
max_diff_lines: int = 2000 # the diff artifact is truncated past this
|
||||
include_untracked: bool = True # a brand-new file is part of the change
|
||||
|
||||
|
||||
class BaseRef(BaseModel):
|
||||
"""The commit a change is measured from, and why that one.
|
||||
|
||||
`reason` is the line the trace shows. A diff is only as trustworthy as the
|
||||
thing it was taken against, so the ADW records that choice instead of
|
||||
leaving the reader to infer it.
|
||||
"""
|
||||
|
||||
ref: str # what was asked for: "main", or a pinned sha
|
||||
commit: str # the commit actually diffed against
|
||||
reason: str = ""
|
||||
|
||||
@property
|
||||
def label(self) -> str:
|
||||
"""Display form — a named ref as itself, a pinned raw sha shortened."""
|
||||
if len(self.ref) == 40 and all(c in "0123456789abcdef" for c in self.ref):
|
||||
return self.ref[:7]
|
||||
return self.ref
|
||||
|
||||
|
||||
class ChangeSet(BaseModel):
|
||||
"""What changed since the base commit — pure git facts, no judgement."""
|
||||
|
||||
base: BaseRef
|
||||
files: list[str] = Field(default_factory=list)
|
||||
untracked: list[str] = Field(default_factory=list)
|
||||
insertions: int = 0
|
||||
deletions: int = 0
|
||||
stat: str = "" # `git diff --stat` output, verbatim
|
||||
diff_path: str = "" # the full diff, written into context_handoff/
|
||||
truncated: bool = False
|
||||
|
||||
@property
|
||||
def empty(self) -> bool:
|
||||
return not (self.files or self.untracked)
|
||||
|
||||
|
||||
class ChangesOutput(EnvelopeBase):
|
||||
"""A ChangeSet shaped as an envelope so an agent can be handed it directly.
|
||||
|
||||
Same adapter idea as VerifyOutput: code computes the diff, the documenter
|
||||
consumes it through the one door every agent handoff uses.
|
||||
"""
|
||||
|
||||
base: str = "" # "<ref> @ <commit> — <reason>"
|
||||
changed_files: list[str] = Field(default_factory=list)
|
||||
insertions: int = 0
|
||||
deletions: int = 0
|
||||
stat: str = ""
|
||||
diff_path: str = "" # read this for the full diff
|
||||
|
||||
|
||||
class VerifyOutput(EnvelopeBase):
|
||||
"""A deterministic result, shaped as an envelope so an agent can consume it.
|
||||
|
||||
Agents hand each other typed envelopes; code blocks return QualityResult.
|
||||
This is the adapter, so a failing lint or test run flows back into the
|
||||
builder through exactly the same door a tester agent's report used to —
|
||||
the ADW script is the only thing that knows the difference.
|
||||
"""
|
||||
|
||||
passed: bool = False
|
||||
failures: list[str] = Field(default_factory=list)
|
||||
|
||||
|
||||
# ── Agent calls ──────────────────────────────────────────────────────────────
|
||||
|
||||
class GateCheck(BaseModel):
|
||||
"""One thing a gate looked at, and what it found.
|
||||
|
||||
`note` is the evidence — "exists, 2.1KB", "exit 0", "not in the diff". On a
|
||||
failed check it doubles as the reason, so it is what the agent is told.
|
||||
"""
|
||||
|
||||
item: str # what was checked: a path, a command, a test
|
||||
ok: bool
|
||||
note: str = ""
|
||||
|
||||
|
||||
class GateReport(BaseModel):
|
||||
"""What every gate returns: the checks it ran. Violations are derived.
|
||||
|
||||
Authoring stays a one-liner per item — `report.check(...)` appends and
|
||||
returns self, so a gate is a loop and a return.
|
||||
"""
|
||||
|
||||
checks: list[GateCheck] = Field(default_factory=list)
|
||||
|
||||
def check(self, item: str, ok: bool, note: str = "") -> "GateReport":
|
||||
self.checks.append(GateCheck(item=item, ok=ok, note=note))
|
||||
return self
|
||||
|
||||
@property
|
||||
def violations(self) -> list[str]:
|
||||
return [f"{c.item}: {c.note or 'failed'}" for c in self.checks if not c.ok]
|
||||
|
||||
@property
|
||||
def passed(self) -> bool:
|
||||
return not self.violations
|
||||
|
||||
|
||||
class AgentCall(BaseModel):
|
||||
"""One agent invocation: prompt in, typed envelope out, gates verified."""
|
||||
|
||||
model_config = {"arbitrary_types_allowed": True}
|
||||
|
||||
output_type: Type[EnvelopeBase]
|
||||
prompt: str
|
||||
previous: Optional[EnvelopeBase] = None
|
||||
gates: list[Callable] = Field(default_factory=list) # gate(envelope, run) -> list[str]
|
||||
|
||||
|
||||
# ── Config ───────────────────────────────────────────────────────────────────
|
||||
|
||||
class PromptEngineering(BaseModel):
|
||||
system: str # path to system.md
|
||||
user: str # path to user.md
|
||||
|
||||
|
||||
class AgentConfig(BaseModel):
|
||||
name: str
|
||||
coding_agent: Literal["pi", "omp", "claude_code"] = "pi"
|
||||
model: str = "google/gemini-3.6-flash"
|
||||
thinking: str = "medium" # off | minimal | low | medium | high | xhigh | max
|
||||
color: str = "" # hex swatch for this agent's lane in the UI
|
||||
purpose: str = ""
|
||||
prompt_engineering: PromptEngineering
|
||||
harness_engineering: list[str] = Field(default_factory=list)
|
||||
tools: Optional[list[str]] = None # allowlist; None = all tools usable
|
||||
# What this agent may MODIFY in the repo, enforced in code after every call
|
||||
# (see adw_modules/permissions.py). `tools` cannot express this: `bash` runs
|
||||
# anything and `write` reaches any path, so an agent's capability list is a
|
||||
# statement of intent that nothing checks.
|
||||
# None -> unrestricted, except the roster-wide `protected_files` paths
|
||||
# [] -> read-only: may modify nothing tracked
|
||||
# [...] -> only these. A trailing "/" means a directory prefix; a "*"
|
||||
# makes it a glob; anything else is an exact path.
|
||||
writes: Optional[list[str]] = None
|
||||
|
||||
|
||||
class ConfigDefaults(BaseModel):
|
||||
coding_agent: Literal["pi", "omp", "claude_code"] = "pi"
|
||||
model: str = "google/gemini-3.6-flash"
|
||||
thinking: str = "medium"
|
||||
color: str = ""
|
||||
harness_engineering: list[str] = Field(default_factory=list)
|
||||
tools: Optional[list[str]] = None # roster-wide allowlist; None = all tools usable
|
||||
# Off-limits to every agent that has not named them in its own `writes`.
|
||||
# The factory's own code is the default: an agent must not be able to edit
|
||||
# the machinery that decides whether its work passed.
|
||||
protected_files: list[str] = Field(default_factory=lambda: [
|
||||
"adws/adw_modules/", "adws/adw_sssf_config/", "adws/adw_*.py",
|
||||
])
|
||||
data_dir: str = "adws/adw_data"
|
||||
|
||||
|
||||
class ObservabilityConfig(BaseModel):
|
||||
db: str = "adws/adw_data/sssf.db"
|
||||
poll_ms: int = 500
|
||||
|
||||
|
||||
class SSSFConfig(BaseModel):
|
||||
defaults: ConfigDefaults = Field(default_factory=ConfigDefaults)
|
||||
observability: ObservabilityConfig = Field(default_factory=ObservabilityConfig)
|
||||
agents: list[AgentConfig] = Field(default_factory=list)
|
||||
|
||||
|
||||
# ── Tracing ──────────────────────────────────────────────────────────────────
|
||||
|
||||
class EventRecord(BaseModel):
|
||||
"""One traced event, always logged against adw_id + phase."""
|
||||
|
||||
adw_id: str
|
||||
phase_id: str = ""
|
||||
type: str # phase_start | agent_start | tool_call | handoff | gate_pass | gate_fail | log | agent_end | phase_end | error
|
||||
name: str = ""
|
||||
payload: dict[str, Any] = Field(default_factory=dict)
|
||||
parent_id: str = ""
|
||||
tokens: Optional[int] = None
|
||||
# Spans: set both when an event covers real elapsed time (a tool call), so
|
||||
# the UI lays it out on a time axis without parsing payload JSON. Left unset,
|
||||
# the tracer stamps started_at with the moment the event was recorded.
|
||||
started_at: Optional[str] = None
|
||||
ended_at: Optional[str] = None
|
||||
|
||||
|
||||
# ── Pi coding agent interface ────────────────────────────────────────────────
|
||||
|
||||
class PiRequest(BaseModel):
|
||||
"""Everything one non-interactive pi run needs."""
|
||||
|
||||
prompt: str
|
||||
system_prompt: str
|
||||
model: str # registry pattern, resolved to provider + id
|
||||
thinking: str = "medium"
|
||||
session_id: str # pi --session-id: creates or continues
|
||||
session_dir: str
|
||||
raw_output_path: str # JSONL stream lands here
|
||||
tools: Optional[list[str]] = None
|
||||
extensions: list[str] = Field(default_factory=list)
|
||||
cwd: str = "." # set from run.repo_root — the codebase root agents work in
|
||||
|
||||
|
||||
class UsageBreakdown(BaseModel):
|
||||
"""Tokens and the dollars they cost, per component, summed over a call.
|
||||
|
||||
Mirrors pi's `usage` shape one-for-one so the numbers reconcile with what
|
||||
pi itself reports: `input` EXCLUDES cache reads, which bill at their own
|
||||
(cheaper) rate — add them to learn the size of the prompt that was sent.
|
||||
"""
|
||||
input_tokens: int = 0
|
||||
output_tokens: int = 0
|
||||
cache_read_tokens: int = 0
|
||||
cache_write_tokens: int = 0
|
||||
# Thinking tokens. NOT a fifth component: measured across every session on
|
||||
# disk, reasoning is always <= output and the four components above always
|
||||
# sum to totalTokens, so reasoning is the thinking SHARE of output, billed
|
||||
# at the output rate. Report it nested under output, never added to it.
|
||||
reasoning_tokens: int = 0
|
||||
total_tokens: int = 0
|
||||
input_cost: float = 0.0
|
||||
output_cost: float = 0.0
|
||||
cache_read_cost: float = 0.0
|
||||
cache_write_cost: float = 0.0
|
||||
total_cost: float = 0.0
|
||||
|
||||
def add_turn(self, usage: dict, total_tokens: int) -> None:
|
||||
"""Fold in one pi `message_end` usage object.
|
||||
|
||||
`total_tokens` is passed in rather than re-derived: the caller already
|
||||
computes it pi's way (totalTokens, else the sum of the parts).
|
||||
"""
|
||||
cost = usage.get("cost") or {}
|
||||
self.input_tokens += usage.get("input") or 0
|
||||
self.output_tokens += usage.get("output") or 0
|
||||
self.cache_read_tokens += usage.get("cacheRead") or 0
|
||||
self.cache_write_tokens += usage.get("cacheWrite") or 0
|
||||
self.reasoning_tokens += usage.get("reasoning") or 0
|
||||
self.total_tokens += total_tokens
|
||||
self.input_cost += cost.get("input") or 0.0
|
||||
self.output_cost += cost.get("output") or 0.0
|
||||
self.cache_read_cost += cost.get("cacheRead") or 0.0
|
||||
self.cache_write_cost += cost.get("cacheWrite") or 0.0
|
||||
self.total_cost += cost.get("total") or 0.0
|
||||
|
||||
def merge(self, other: "UsageBreakdown") -> None:
|
||||
"""Add another call's usage — a phase that retries spends more than once."""
|
||||
for field in self.model_fields:
|
||||
setattr(self, field, getattr(self, field) + getattr(other, field))
|
||||
|
||||
|
||||
class PiResult(BaseModel):
|
||||
text: str = ""
|
||||
returncode: int = 0
|
||||
session_id: str = ""
|
||||
tokens: int = 0
|
||||
cost: float = 0.0
|
||||
usage: UsageBreakdown = Field(default_factory=UsageBreakdown)
|
||||
# Context occupancy after the LAST turn — not a sum. `tokens` bills every
|
||||
# turn; this is how full the window is right now, which is what the
|
||||
# visualizer's context bar measures against `context_window`.
|
||||
context_tokens: int = 0
|
||||
context_window: int = 0 # 0 when the registry declares no ceiling
|
||||
108
sssf/templates/adws/adw_modules/gates.py
Normal file
108
sssf/templates/adws/adw_modules/gates.py
Normal file
|
|
@ -0,0 +1,108 @@
|
|||
"""Validation gates: verify the envelope's CLAIMS, never guesses.
|
||||
|
||||
A gate is `gate(envelope, run) -> GateReport` — one check per item it looked at.
|
||||
Violations are derived from the failed checks and sent back to the SAME agent
|
||||
session as a correction. Every check is recorded either way, so a green gate
|
||||
says WHAT it verified instead of only that it passed.
|
||||
|
||||
Gates check what is mechanically checkable; plan quality is a reviewer's job.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
|
||||
from .data_types import EnvelopeBase, GateReport
|
||||
|
||||
TAIL_CHARS = 1000 # command output kept as evidence on a failure
|
||||
|
||||
|
||||
def _size(path: Path) -> str:
|
||||
n = path.stat().st_size
|
||||
return f"{n}B" if n < 1024 else f"{n / 1024:.1f}KB"
|
||||
|
||||
|
||||
def artifacts_exist(envelope: EnvelopeBase, run) -> GateReport:
|
||||
report = GateReport()
|
||||
for a in envelope.artifacts:
|
||||
p = Path(a)
|
||||
report.check(a, p.exists(),
|
||||
f"exists, {_size(p)}" if p.exists() else "declared artifact does not exist")
|
||||
return report
|
||||
|
||||
|
||||
def files_non_empty(envelope: EnvelopeBase, run) -> GateReport:
|
||||
report = GateReport()
|
||||
for a in envelope.artifacts:
|
||||
p = Path(a)
|
||||
if not (p.exists() and p.is_file()):
|
||||
continue # existence is artifacts_exist's job
|
||||
empty = p.stat().st_size == 0
|
||||
report.check(a, not empty, "declared artifact is empty" if empty else _size(p))
|
||||
return report
|
||||
|
||||
|
||||
def json_parses(envelope: EnvelopeBase, run) -> GateReport:
|
||||
report = GateReport()
|
||||
for a in envelope.artifacts:
|
||||
p = Path(a)
|
||||
if p.suffix != ".json" or not p.exists():
|
||||
continue
|
||||
try:
|
||||
parsed = json.loads(p.read_text())
|
||||
report.check(a, True, f"parses, {type(parsed).__name__}")
|
||||
except json.JSONDecodeError as e:
|
||||
report.check(a, False, f"declared JSON artifact does not parse: {e}")
|
||||
return report
|
||||
|
||||
|
||||
def diff_matches_claims(envelope: EnvelopeBase, run) -> GateReport:
|
||||
"""Every file claimed changed must exist on disk."""
|
||||
report = GateReport()
|
||||
for f in getattr(envelope, "changed_files", []):
|
||||
p = Path(f)
|
||||
report.check(f, p.exists(),
|
||||
f"exists, {_size(p)}" if p.exists() else "claimed changed file does not exist")
|
||||
return report
|
||||
|
||||
|
||||
def verdict_consistent(envelope: EnvelopeBase, run) -> GateReport:
|
||||
"""A review's verdict must agree with the findings it just wrote down.
|
||||
|
||||
Nothing here judges the code — that is the reviewer's job. This checks the
|
||||
envelope against itself: an approval that ships blocking items, or a
|
||||
rejection that names no problem, is a claim the harness can refute without
|
||||
reading a line of the diff.
|
||||
"""
|
||||
report = GateReport()
|
||||
approved = bool(getattr(envelope, "approved", False))
|
||||
blocking = list(getattr(envelope, "blocking", []))
|
||||
unmet = [f.requirement for f in getattr(envelope, "findings", []) if not f.met]
|
||||
|
||||
report.check("approved vs blocking", not (approved and blocking),
|
||||
"no blocking items" if not blocking
|
||||
else f"{len(blocking)} blocking item(s) while approved=true"
|
||||
if approved else f"{len(blocking)} blocking item(s), not approved")
|
||||
report.check("approved vs findings", not (approved and unmet),
|
||||
"every requirement met" if not unmet
|
||||
else f"{len(unmet)} unmet requirement(s) while approved=true"
|
||||
if approved else f"{len(unmet)} unmet requirement(s), not approved")
|
||||
report.check("rejection names a problem", approved or bool(blocking or unmet),
|
||||
"verdict is supported" if approved or blocking or unmet
|
||||
else "approved=false but no blocking item or unmet requirement was given")
|
||||
return report
|
||||
|
||||
|
||||
def tests_pass(command: str):
|
||||
"""Gate factory: the given shell command must exit 0."""
|
||||
def gate(envelope: EnvelopeBase, run) -> GateReport:
|
||||
result = subprocess.run(command, shell=True, capture_output=True, text=True)
|
||||
ok = result.returncode == 0
|
||||
note = f"exit {result.returncode}"
|
||||
if not ok:
|
||||
note += "\n" + (result.stdout + result.stderr)[-TAIL_CHARS:]
|
||||
return GateReport().check(command, ok, note)
|
||||
gate.__name__ = f"tests_pass({command})"
|
||||
return gate
|
||||
120
sssf/templates/adws/adw_modules/git_helper.py
Normal file
120
sssf/templates/adws/adw_modules/git_helper.py
Normal file
|
|
@ -0,0 +1,120 @@
|
|||
"""Low-level git operations for code phases. All low-level logic lives in adw_modules."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
|
||||
|
||||
def _git(*args: str) -> str:
|
||||
result = subprocess.run(["git", *args], capture_output=True, text=True)
|
||||
if result.returncode != 0:
|
||||
raise RuntimeError(f"git {' '.join(args)} failed: {result.stderr.strip()}")
|
||||
return result.stdout.strip()
|
||||
|
||||
|
||||
def current_branch() -> str:
|
||||
return _git("rev-parse", "--abbrev-ref", "HEAD")
|
||||
|
||||
|
||||
def create_branch(name: str) -> str:
|
||||
_git("checkout", "-b", name)
|
||||
return name
|
||||
|
||||
|
||||
def is_repo() -> bool:
|
||||
result = subprocess.run(["git", "rev-parse", "--git-dir"],
|
||||
capture_output=True, text=True)
|
||||
return result.returncode == 0
|
||||
|
||||
|
||||
def repo_root() -> Path:
|
||||
"""Absolute root of the codebase — where agents are spawned to work.
|
||||
|
||||
The git toplevel when there is one, else the process cwd (ADWs run fine in a
|
||||
non-git dir; only a commit phase requires a repo). Always absolute, so it is
|
||||
safe to hand to a subprocess regardless of where the ADW was launched from.
|
||||
"""
|
||||
if is_repo():
|
||||
return Path(_git("rev-parse", "--show-toplevel")).resolve()
|
||||
return Path.cwd().resolve()
|
||||
|
||||
|
||||
def commit_all(message: str) -> str:
|
||||
"""Stage the working tree and commit it. Returns the new short sha."""
|
||||
if not is_repo():
|
||||
raise RuntimeError(
|
||||
"not a git repository — a commit phase needs one. Run `git init` in the "
|
||||
"repo root (and make a first commit) before running an ADW that commits.")
|
||||
_git("add", "-A")
|
||||
if not _git("status", "--porcelain"):
|
||||
raise RuntimeError("nothing to commit — the preceding phases changed no files")
|
||||
_git("commit", "-m", message)
|
||||
return _git("rev-parse", "--short", "HEAD")
|
||||
|
||||
|
||||
def changed_files() -> list[str]:
|
||||
out = _git("status", "--porcelain")
|
||||
return [line[3:] for line in out.splitlines() if line]
|
||||
|
||||
|
||||
# ── diff plumbing (composed into a ChangeSet by documentation.py) ────────────
|
||||
|
||||
def ref_exists(ref: str) -> bool:
|
||||
"""True when `ref` resolves to a commit. Never raises — this is a question."""
|
||||
result = subprocess.run(["git", "rev-parse", "--verify", "--quiet", f"{ref}^{{commit}}"],
|
||||
capture_output=True, text=True)
|
||||
return result.returncode == 0
|
||||
|
||||
|
||||
def rev(ref: str = "HEAD") -> str:
|
||||
return _git("rev-parse", ref)
|
||||
|
||||
|
||||
def short_sha(ref: str = "HEAD") -> str:
|
||||
return _git("rev-parse", "--short", ref)
|
||||
|
||||
|
||||
def merge_base(ref: str, other: str = "HEAD") -> str:
|
||||
"""The commit where `ref` and `other` diverged — the honest base of a branch.
|
||||
|
||||
On the base branch itself this returns HEAD, which makes the diff exactly
|
||||
"what is not committed yet". Off it, the diff is the whole branch plus the
|
||||
working tree. One command covers both cases, so no ADW has to branch on it.
|
||||
"""
|
||||
return _git("merge-base", ref, other)
|
||||
|
||||
|
||||
def is_dirty() -> bool:
|
||||
return bool(_git("status", "--porcelain"))
|
||||
|
||||
|
||||
def untracked_files() -> list[str]:
|
||||
out = _git("ls-files", "--others", "--exclude-standard")
|
||||
return [line for line in out.splitlines() if line]
|
||||
|
||||
|
||||
def diff_files(base: str) -> list[str]:
|
||||
"""Tracked files that differ between `base` and the working tree."""
|
||||
out = _git("diff", "--name-only", base)
|
||||
return [line for line in out.splitlines() if line]
|
||||
|
||||
|
||||
def diff_stat(base: str) -> str:
|
||||
return _git("diff", "--stat", base)
|
||||
|
||||
|
||||
def diff_counts(base: str) -> tuple[int, int]:
|
||||
"""(insertions, deletions) across the diff. Binary files count as neither."""
|
||||
insertions = deletions = 0
|
||||
for line in _git("diff", "--numstat", base).splitlines():
|
||||
added, removed, *_ = line.split("\t")
|
||||
if added.isdigit():
|
||||
insertions += int(added)
|
||||
if removed.isdigit():
|
||||
deletions += int(removed)
|
||||
return insertions, deletions
|
||||
|
||||
|
||||
def diff_text(base: str) -> str:
|
||||
return _git("diff", base)
|
||||
185
sssf/templates/adws/adw_modules/permissions.py
Normal file
185
sssf/templates/adws/adw_modules/permissions.py
Normal file
|
|
@ -0,0 +1,185 @@
|
|||
"""What an agent may CHANGE, enforced in code after the fact.
|
||||
|
||||
`tools:` is a capability list, not a sandbox, and two holes make it
|
||||
unenforceable on its own:
|
||||
|
||||
* `bash` runs anything. A builder handed bash to run a test suite can also
|
||||
run `git checkout adws/` — which is not hypothetical: one did, discarding
|
||||
uncommitted changes to the very quality check it was about to be judged by.
|
||||
* `write` reaches any path, not just the one report file an agent was given
|
||||
it for. A reviewer configured with "no edit, so it cannot quietly fix"
|
||||
could still rewrite the code it was reviewing.
|
||||
|
||||
So permission is verified the way every other claim in this system is —
|
||||
after the fact, against the repo itself. `snapshot()` fingerprints the working
|
||||
tree's change-set before an agent runs; `enforce()` compares it afterwards and
|
||||
fails the phase if the agent touched anything outside its allowlist.
|
||||
|
||||
Comparing change-sets, rather than watching for writes, is what catches the
|
||||
`git checkout` case: a path that was modified before the agent ran and is clean
|
||||
afterwards has been reverted, and a reversion is a modification. Appearing,
|
||||
disappearing, and changing all count.
|
||||
|
||||
A breach is NOT a gate violation. Gates are for work an agent can be asked to
|
||||
redo; a breach cannot be corrected by re-prompting, because the write already
|
||||
happened. It aborts the phase and names every offending path.
|
||||
|
||||
Two keys drive it, both in sssf.config.yaml:
|
||||
defaults.protected_files paths no agent may touch unless it names them itself
|
||||
agents[].writes None = unrestricted · [] = read-only · [...] = only these
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
|
||||
from .data_types import AgentConfig, SSSFConfig
|
||||
|
||||
|
||||
class PermissionBreach(RuntimeError):
|
||||
"""An agent modified a path it was not permitted to modify."""
|
||||
|
||||
|
||||
def _git(args: list[str], cwd) -> str:
|
||||
result = subprocess.run(["git", *args], cwd=cwd, capture_output=True, text=True)
|
||||
return result.stdout if result.returncode == 0 else ""
|
||||
|
||||
|
||||
def snapshot(run) -> dict[str, str]:
|
||||
"""Fingerprint every path the working tree currently differs on.
|
||||
|
||||
Tracked files carry their numstat counts, so an edit to an already-dirty
|
||||
file still registers as a change. Untracked files are listed by name.
|
||||
Gitignored paths never appear, which is why the session runtime under
|
||||
`data_dir` — where handoff files legitimately land — needs no special case.
|
||||
"""
|
||||
fingerprints: dict[str, str] = {}
|
||||
for line in _git(["diff", "HEAD", "--numstat"], run.repo_root).splitlines():
|
||||
fields = line.split("\t")
|
||||
if len(fields) >= 3:
|
||||
path = fields[-1].strip()
|
||||
fingerprints[path] = f"{fields[0]},{fields[1]}"
|
||||
for path in _git(["ls-files", "--others", "--exclude-standard"],
|
||||
run.repo_root).splitlines():
|
||||
if path.strip():
|
||||
fingerprints[path.strip()] = "untracked"
|
||||
return fingerprints
|
||||
|
||||
|
||||
def changed_paths(before: dict[str, str], after: dict[str, str]) -> list[str]:
|
||||
"""Every path whose state differs — appeared, vanished, or was rewritten."""
|
||||
return sorted({p for p in set(before) | set(after)
|
||||
if before.get(p) != after.get(p)})
|
||||
|
||||
|
||||
def _glob(pattern: str) -> re.Pattern:
|
||||
"""Translate a pattern, with `*` stopping at a path separator.
|
||||
|
||||
fnmatch would let `*` cross `/`, which quietly widens every pattern:
|
||||
`adws/adw_*.py` would match `adws/adw_data/sessions/x/y.py` as well as the
|
||||
ADW scripts it means. `**` is the way to say "cross directories".
|
||||
"""
|
||||
out, i = [], 0
|
||||
while i < len(pattern):
|
||||
char = pattern[i]
|
||||
if pattern.startswith("**", i):
|
||||
out.append(".*")
|
||||
i += 2
|
||||
elif char == "*":
|
||||
out.append("[^/]*")
|
||||
i += 1
|
||||
elif char == "?":
|
||||
out.append("[^/]")
|
||||
i += 1
|
||||
else:
|
||||
out.append(re.escape(char))
|
||||
i += 1
|
||||
return re.compile("".join(out))
|
||||
|
||||
|
||||
def _matches(path: str, pattern: str) -> bool:
|
||||
if pattern.endswith("/"): # directory prefix
|
||||
return path.startswith(pattern)
|
||||
if "*" in pattern or "?" in pattern:
|
||||
return _glob(pattern).fullmatch(path) is not None
|
||||
return path == pattern
|
||||
|
||||
|
||||
def always_writable(cfg: SSSFConfig) -> list[str]:
|
||||
"""The session runtime, which EVERY agent must be able to write.
|
||||
|
||||
`context_handoff/` is the one place agents hand work to each other, and an
|
||||
agent's own prompts, raw_output.jsonl, and envelope.json land beside it.
|
||||
Scout writes its findings there, the reviewer its review, the planner its
|
||||
plan — a read-only agent is read-only with respect to the REPO, never with
|
||||
respect to its own report.
|
||||
|
||||
This is granted from `data_dir` rather than left to .gitignore. The runtime
|
||||
is normally ignored, so it never even appears in a snapshot — but an agent's
|
||||
ability to record its work must not hang on a gitignore entry that someone
|
||||
can delete or that a changed `data_dir` can outgrow.
|
||||
"""
|
||||
return [cfg.defaults.data_dir.rstrip("/") + "/"]
|
||||
|
||||
|
||||
def permitted(path: str, agent: AgentConfig, cfg: SSSFConfig) -> bool:
|
||||
"""Session runtime first, then the agent's own list, then what is protected."""
|
||||
if any(_matches(path, p) for p in always_writable(cfg)):
|
||||
return True
|
||||
if any(_matches(path, p) for p in (agent.writes or [])):
|
||||
return True # naming a path is what unlocks a protected one
|
||||
if any(_matches(path, p) for p in cfg.defaults.protected_files):
|
||||
return False
|
||||
return agent.writes is None # None = unrestricted, [] = no repo writes
|
||||
|
||||
|
||||
def _roll_back(run, path: str, before: dict[str, str], after: dict[str, str]) -> str:
|
||||
"""Undo one unauthorized change. Returns a word describing what happened.
|
||||
|
||||
Only changes the agent INTRODUCED are undone. A path that was already dirty
|
||||
when the agent started is left exactly as it is: the operator had
|
||||
uncommitted work there, and discarding it to tidy up would be the same harm
|
||||
this module exists to prevent, committed by the cleanup instead of the agent.
|
||||
"""
|
||||
if path in before:
|
||||
# Already dirty beforehand. If it is gone from the diff now, the agent
|
||||
# reverted an engineer's uncommitted work and the content is not ours
|
||||
# to reconstruct — say so loudly rather than pretend it was handled.
|
||||
return "REVERTED-BY-AGENT (uncommitted work lost, cannot restore)" \
|
||||
if path not in after else "left as-is (was already modified)"
|
||||
if after.get(path) == "untracked":
|
||||
try:
|
||||
(Path(run.repo_root) / path).unlink()
|
||||
return "deleted"
|
||||
except OSError as error:
|
||||
return f"could not delete ({error})"
|
||||
result = subprocess.run(["git", "checkout", "--", path],
|
||||
cwd=run.repo_root, capture_output=True, text=True)
|
||||
return "rolled back" if result.returncode == 0 else "could not roll back"
|
||||
|
||||
|
||||
def enforce(run, phase, agent: AgentConfig, before: dict[str, str]) -> list[str]:
|
||||
"""Compare the tree against `before`; undo and raise if the agent overstepped.
|
||||
|
||||
Returns the paths it legitimately changed, so the trace records what an
|
||||
agent actually touched rather than only what it claimed in its envelope.
|
||||
|
||||
Detection alone would leave the repo holding the unauthorized change while
|
||||
reporting a failure, so anything the agent introduced outside its allowlist
|
||||
is rolled back before the phase dies. What it cannot undo, it names.
|
||||
"""
|
||||
after = snapshot(run)
|
||||
touched = changed_paths(before, after)
|
||||
breaches = [p for p in touched if not permitted(p, agent, run.cfg)]
|
||||
if not breaches:
|
||||
return touched
|
||||
|
||||
outcomes = {p: _roll_back(run, p, before, after) for p in breaches}
|
||||
scope = ("read-only" if agent.writes == []
|
||||
else f"limited to {agent.writes}" if agent.writes
|
||||
else f"barred from {run.cfg.defaults.protected_files}")
|
||||
detail = "\n".join(f" - {p} — {outcome}" for p, outcome in outcomes.items())
|
||||
raise PermissionBreach(
|
||||
f"{agent.name} is {scope} but modified {len(breaches)} path(s):\n{detail}")
|
||||
21
sssf/templates/adws/adw_modules/prompts.py
Normal file
21
sssf/templates/adws/adw_modules/prompts.py
Normal file
|
|
@ -0,0 +1,21 @@
|
|||
"""Prompt rendering: load system/user refs from config, replace {{placeholders}}."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from pathlib import Path
|
||||
|
||||
|
||||
def render(template_path: str | Path, variables: dict[str, str]) -> str:
|
||||
text = Path(template_path).read_text()
|
||||
for key, value in variables.items():
|
||||
text = text.replace("{{" + key + "}}", value)
|
||||
return text
|
||||
|
||||
|
||||
def save(directory: str | Path, name: str, content: str) -> Path:
|
||||
"""Save the exact prompt sent, before execution — the audit copy."""
|
||||
directory = Path(directory)
|
||||
directory.mkdir(parents=True, exist_ok=True)
|
||||
path = directory / name
|
||||
path.write_text(content)
|
||||
return path
|
||||
240
sssf/templates/adws/adw_modules/quality.py
Normal file
240
sssf/templates/adws/adw_modules/quality.py
Normal file
|
|
@ -0,0 +1,240 @@
|
|||
"""Deterministic lint, typecheck, build, and test blocks.
|
||||
|
||||
A known command is not a judgement call. Anything whose invocation you can write
|
||||
down belongs here as code — it runs in milliseconds, costs nothing, and returns
|
||||
the same answer every time. Agents are for the parts that need reading and
|
||||
deciding.
|
||||
|
||||
╔══════════════════════════════════════════════════════════════════════════════╗
|
||||
║ REPLACE THE PLACEHOLDER COMMANDS BELOW. ║
|
||||
║ ║
|
||||
║ Every block ships as an `echo` that exits 0 and announces it is fake. They ║
|
||||
║ are placeholders on purpose: a stamped repo has no way to guess your test ║
|
||||
║ runner, and a wrong-but-plausible command that silently passes is worse ║
|
||||
║ than one that says so out loud. ║
|
||||
║ ║
|
||||
║ For each block you want: swap `_placeholder(...)` for the real argv, e.g. ║
|
||||
║ argv=["bun", "test", "apps/web/server.test.ts"] ║
|
||||
║ argv=["uv", "run", "pytest", "-q"] ║
|
||||
║ argv=["npm", "run", "lint"] ║
|
||||
║ Delete the blocks you don't need, and drop them from run_quality()'s list. ║
|
||||
║ ║
|
||||
║ Two rules when you write the real command: ║
|
||||
║ 1. argv LIST, never a shell string — no quoting bugs, no shell injection. ║
|
||||
║ 2. Call binaries by BARE NAME. These blocks inherit the operator's ║
|
||||
║ environment (see utils.operator_env), so `bun`, `uv`, `pytest` resolve ║
|
||||
║ exactly as they do in their terminal. Never hard-code an absolute path ║
|
||||
║ like /Users/you/.bun/bin/bun — that bakes your machine into the trace. ║
|
||||
╚══════════════════════════════════════════════════════════════════════════════╝
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import shlex
|
||||
import subprocess
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Callable
|
||||
|
||||
from .data_types import (EventRecord, QualityCheckResult, QualityCheckSpec, QualityResult,
|
||||
VerifyOutput)
|
||||
from .utils import now_iso, operator_env
|
||||
|
||||
# How much of a failing command's output rides back inside the envelope. Enough
|
||||
# for a builder to act on without opening the artifact; bounded so a runaway
|
||||
# stack trace can't swamp the next agent's context.
|
||||
TAIL_CHARS = 4_000
|
||||
|
||||
|
||||
def _placeholder(name: str) -> list[str]:
|
||||
"""A command that does nothing and admits it. Replace every call to this."""
|
||||
return ["echo", f"PLACEHOLDER {name}: edit adws/adw_modules/quality.py and "
|
||||
f"replace this echo with the real {name} command"]
|
||||
|
||||
|
||||
def _check_dir(run, name: str) -> Path:
|
||||
seq = run.phases[-1].seq if run.phases else 0
|
||||
path = run.context_handoff_dir / "quality" / f"{seq:02d}_{name}"
|
||||
path.mkdir(parents=True, exist_ok=True)
|
||||
return path
|
||||
|
||||
|
||||
def _run(spec: QualityCheckSpec, run) -> QualityCheckResult:
|
||||
phase = run.phases[-1]
|
||||
output_dir = _check_dir(run, spec.name)
|
||||
output_artifact = output_dir / "command.log"
|
||||
command = shlex.join(spec.argv)
|
||||
env = operator_env() # the engineer's own shell environment
|
||||
|
||||
run.console.note(f"quality {spec.name}: {command}")
|
||||
started_at = now_iso()
|
||||
clock = time.monotonic()
|
||||
stdout = ""
|
||||
stderr = ""
|
||||
try:
|
||||
completed = subprocess.run(
|
||||
spec.argv,
|
||||
cwd=run.repo_root,
|
||||
env=env,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=spec.timeout_seconds,
|
||||
)
|
||||
returncode = completed.returncode
|
||||
stdout = completed.stdout
|
||||
stderr = completed.stderr
|
||||
except subprocess.TimeoutExpired as error:
|
||||
returncode = 124
|
||||
stdout = error.stdout or ""
|
||||
stderr = (error.stderr or "") + f"\nTimed out after {spec.timeout_seconds}s."
|
||||
except OSError as error:
|
||||
# A missing binary lands here as exit 127 with the real message — no
|
||||
# pre-flight probe needed, and none wanted.
|
||||
returncode = 127
|
||||
stderr = str(error)
|
||||
|
||||
duration = time.monotonic() - clock
|
||||
output_artifact.write_text(
|
||||
f"$ {command}\nexit: {returncode}\nduration_seconds: {duration:.3f}\n"
|
||||
f"\n--- stdout ---\n{stdout}\n--- stderr ---\n{stderr}\n"
|
||||
)
|
||||
passed = returncode == 0
|
||||
run.tracer.event(EventRecord(
|
||||
adw_id=run.adw_id,
|
||||
phase_id=phase.phase_id,
|
||||
type="tool_call",
|
||||
name=f"quality:{spec.name}",
|
||||
payload={
|
||||
"area": spec.area,
|
||||
"operation": spec.operation,
|
||||
"command": command,
|
||||
"returncode": returncode,
|
||||
"passed": passed,
|
||||
"output_artifact": str(output_artifact),
|
||||
},
|
||||
started_at=started_at,
|
||||
ended_at=now_iso(),
|
||||
))
|
||||
run.console.note(
|
||||
f"quality {spec.name}: {'passed' if passed else 'failed'} "
|
||||
f"(exit {returncode}, {duration:.1f}s)"
|
||||
)
|
||||
return QualityCheckResult(
|
||||
name=spec.name,
|
||||
area=spec.area,
|
||||
operation=spec.operation,
|
||||
command=command,
|
||||
returncode=returncode,
|
||||
passed=passed,
|
||||
duration_seconds=duration,
|
||||
output_artifact=str(output_artifact),
|
||||
output_tail=(stdout + stderr)[-TAIL_CHARS:],
|
||||
)
|
||||
|
||||
|
||||
# ── Blocks ────────────────────────────────────────────────────────────────────
|
||||
# Replace every argv below. See the banner at the top of this file.
|
||||
|
||||
def test(run) -> QualityCheckResult:
|
||||
"""Run the project's test suite. The highest-value block to wire up first."""
|
||||
return _run(QualityCheckSpec(
|
||||
name="test",
|
||||
area="backend",
|
||||
operation="build",
|
||||
argv=_placeholder("test"), # e.g. ["bun", "test"] or ["uv", "run", "pytest", "-q"]
|
||||
timeout_seconds=600,
|
||||
), run)
|
||||
|
||||
|
||||
def lint(run) -> QualityCheckResult:
|
||||
return _run(QualityCheckSpec(
|
||||
name="lint",
|
||||
area="backend",
|
||||
operation="lint",
|
||||
argv=_placeholder("lint"), # e.g. ["bun", "x", "oxlint@1.36.0", "src"]
|
||||
), run)
|
||||
|
||||
|
||||
def typecheck(run) -> QualityCheckResult:
|
||||
return _run(QualityCheckSpec(
|
||||
name="typecheck",
|
||||
area="backend",
|
||||
operation="typecheck",
|
||||
argv=_placeholder("typecheck"), # e.g. ["bun", "x", "tsc", "--noEmit"]
|
||||
), run)
|
||||
|
||||
|
||||
def build(run) -> QualityCheckResult:
|
||||
output_dir = _check_dir(run, "build") / "bundle"
|
||||
return _run(QualityCheckSpec(
|
||||
name="build",
|
||||
area="backend",
|
||||
operation="build",
|
||||
argv=_placeholder("build"), # e.g. ["bun", "build", "src/index.ts", "--outdir", str(output_dir)]
|
||||
), run)
|
||||
|
||||
|
||||
def run_tests(run) -> QualityResult:
|
||||
"""The test suite alone, as a QualityResult — the deterministic test phase.
|
||||
|
||||
This is what replaces a `tester` agent once the command is written down. An
|
||||
agent rediscovering the runner on every run costs a fortune to learn what a
|
||||
subprocess already knows; the repair loop is unchanged, because a failure
|
||||
still reaches the builder through `as_envelope` below.
|
||||
"""
|
||||
check = test(run)
|
||||
failures = ([] if check.passed else
|
||||
[f"{check.name}: `{check.command}` exited {check.returncode}\n"
|
||||
f"{check.output_tail}".rstrip()])
|
||||
return QualityResult(passed=check.passed, checks=[check], failures=failures,
|
||||
artifacts=[check.output_artifact])
|
||||
|
||||
|
||||
def as_envelope(result: QualityResult, what: str) -> VerifyOutput:
|
||||
"""Wrap a deterministic result so an agent can be handed it directly.
|
||||
|
||||
Agents hand each other typed envelopes; code blocks return QualityResult.
|
||||
This is the adapter, so a failing lint or test run flows back into the
|
||||
builder through exactly the same door an agent's report would — the ADW
|
||||
script is the only thing that knows the difference.
|
||||
"""
|
||||
return VerifyOutput(
|
||||
status="success" if result.passed else "fail",
|
||||
summary=(f"{what}: all {len(result.checks)} check(s) passed" if result.passed
|
||||
else f"{what}: {len(result.failures)} of {len(result.checks)} check(s) failed"),
|
||||
artifacts=result.artifacts,
|
||||
notes_for_next_agent=("" if result.passed else
|
||||
"Fix every failure below. The output is verbatim from the "
|
||||
"command — trust it over any summary."),
|
||||
passed=result.passed,
|
||||
failures=result.failures,
|
||||
)
|
||||
|
||||
|
||||
def run_quality(run) -> QualityResult:
|
||||
"""Run every block and collect ALL failures — one pass tells you everything.
|
||||
|
||||
Ordering contract for the caller: a failing block does NOT fail the phase.
|
||||
The runner did its job; the CODE is what failed. Hand this result to the
|
||||
builder and let the bounded repair loop decide the run's fate.
|
||||
"""
|
||||
blocks: list[Callable] = [
|
||||
test,
|
||||
lint,
|
||||
typecheck,
|
||||
build,
|
||||
]
|
||||
checks = [block(run) for block in blocks]
|
||||
# A failure is the command, its exit code, and what it actually printed —
|
||||
# everything a builder needs to repair without opening a log or being told
|
||||
# what the error "means" by a parser that guessed.
|
||||
failures = [
|
||||
f"{check.name}: `{check.command}` exited {check.returncode}\n{check.output_tail}".rstrip()
|
||||
for check in checks if not check.passed
|
||||
]
|
||||
return QualityResult(
|
||||
passed=not failures,
|
||||
checks=checks,
|
||||
failures=failures,
|
||||
artifacts=[check.output_artifact for check in checks],
|
||||
)
|
||||
142
sssf/templates/adws/adw_modules/runner.py
Normal file
142
sssf/templates/adws/adw_modules/runner.py
Normal file
|
|
@ -0,0 +1,142 @@
|
|||
"""The Run object: config + adw_id + agent_map + tracer + console, bound once.
|
||||
|
||||
`run.phase(PhaseParams(...))` is the ONE phase primitive — a context manager
|
||||
for all three kinds (engineer, agent, code). Success must be earned: every
|
||||
phase defaults to fail; only a clean exit flips it (agent phases additionally
|
||||
require a parsed envelope + green gates, enforced inside ph.call).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import time
|
||||
from contextlib import contextmanager
|
||||
from pathlib import Path
|
||||
|
||||
from . import agents, git_helper
|
||||
from .console import Console
|
||||
from .data_types import AgentCall, EnvelopeBase, EventRecord, Phase, PhaseParams
|
||||
from .utils import ensure_dir, now_iso
|
||||
|
||||
|
||||
class PhaseHandle:
|
||||
def __init__(self, run: "Run", phase: Phase):
|
||||
self.run = run
|
||||
self.phase = phase
|
||||
|
||||
def log(self, **payload) -> None:
|
||||
self.run.tracer.event(EventRecord(adw_id=self.run.adw_id,
|
||||
phase_id=self.phase.phase_id,
|
||||
type="log", name=self.phase.params.name,
|
||||
payload=payload))
|
||||
self.run.console.note(", ".join(f"{k}: {v}" for k, v in payload.items()))
|
||||
if self.phase.params.kind == "engineer" and "input" in payload:
|
||||
self.run.tracer.session_request(self.run.adw_id, str(payload["input"]))
|
||||
|
||||
def call(self, call: AgentCall) -> EnvelopeBase:
|
||||
if self.phase.params.kind != "agent":
|
||||
raise RuntimeError("ph.call() is only valid inside an agent phase")
|
||||
return agents.execute(self.run, self.phase, call)
|
||||
|
||||
|
||||
class Run:
|
||||
def __init__(self, cfg, adw_id: str, tracer, engineer: str):
|
||||
self.cfg = cfg
|
||||
self.adw_id = adw_id
|
||||
self.tracer = tracer
|
||||
self.console = Console(tracer, adw_id)
|
||||
self.engineer = engineer
|
||||
self.phases: list[Phase] = []
|
||||
self.tokens = 0
|
||||
self.cost = 0.0
|
||||
self._seq = tracer.max_phase_seq(adw_id) # a joined run continues the sequence
|
||||
self.repo_root = git_helper.repo_root() # where every agent is spawned to work
|
||||
self.session_dir = ensure_dir(Path(cfg.defaults.data_dir) / "sessions" / adw_id)
|
||||
self.context_handoff_dir = ensure_dir(self.session_dir / "context_handoff")
|
||||
self._agent_map_path = self.session_dir / "agent_map.json"
|
||||
self.agent_map: dict = (json.loads(self._agent_map_path.read_text())
|
||||
if self._agent_map_path.exists() else {})
|
||||
|
||||
# ── agent map (adw_id -> per-agent coding-agent session ids) ────────────
|
||||
def save_agent_map(self, agent: str, entry: dict) -> None:
|
||||
self.agent_map[agent] = entry
|
||||
self._agent_map_path.write_text(json.dumps(self.agent_map, indent=2))
|
||||
|
||||
# ── usage (run totals mirror what the tracer accumulates in sqlite) ─────
|
||||
def add_usage(self, tokens: int, cost: float) -> None:
|
||||
self.tokens += tokens
|
||||
self.cost += cost
|
||||
self.tracer.session_add_usage(self.adw_id, tokens, cost)
|
||||
|
||||
# ── the phase primitive ─────────────────────────────────────────────────
|
||||
@contextmanager
|
||||
def phase(self, params: PhaseParams):
|
||||
self._seq += 1
|
||||
phase = Phase(phase_id=f"{self.adw_id}_{self._seq:02d}_{params.name}",
|
||||
adw_id=self.adw_id, seq=self._seq, params=params,
|
||||
status="running", started_at=now_iso())
|
||||
self.phases.append(phase)
|
||||
self.tracer.phase_upsert(phase)
|
||||
self.tracer.event(EventRecord(adw_id=self.adw_id, phase_id=phase.phase_id,
|
||||
type="phase_start", name=params.name,
|
||||
payload={"kind": params.kind, "owner": params.owner,
|
||||
"description": params.description}))
|
||||
self.console.phase_started(phase)
|
||||
clock = time.monotonic()
|
||||
try:
|
||||
yield PhaseHandle(self, phase)
|
||||
except BaseException as error:
|
||||
phase.status = "fail" # success must be earned
|
||||
phase.error = str(error)[:1000]
|
||||
phase.ended_at = now_iso()
|
||||
self.tracer.event(EventRecord(adw_id=self.adw_id, phase_id=phase.phase_id,
|
||||
type="error", name=params.name,
|
||||
payload={"error": phase.error}))
|
||||
self.tracer.event(EventRecord(adw_id=self.adw_id, phase_id=phase.phase_id,
|
||||
type="phase_end", name=params.name,
|
||||
payload={"status": "fail"}))
|
||||
self.tracer.phase_upsert(phase)
|
||||
self.tracer.session_finish(self.adw_id, ok=False)
|
||||
self.console.phase_ended(phase, time.monotonic() - clock)
|
||||
self.console.session_finished(False, self.tokens, self.cost,
|
||||
self.cfg.observability.db)
|
||||
raise
|
||||
else:
|
||||
phase.status = "success"
|
||||
phase.ended_at = now_iso()
|
||||
self.tracer.event(EventRecord(adw_id=self.adw_id, phase_id=phase.phase_id,
|
||||
type="phase_end", name=params.name,
|
||||
payload={"status": "success"}))
|
||||
self.tracer.phase_upsert(phase)
|
||||
self.console.phase_ended(phase, time.monotonic() - clock)
|
||||
|
||||
# ── run outcome ─────────────────────────────────────────────────────────
|
||||
def finish(self, accepted: bool = True, reason: str = "") -> int:
|
||||
"""Finalize the run and return its exit code. Call this exactly once.
|
||||
|
||||
Two criteria, not one. Every phase must have passed, AND the ADW's own
|
||||
acceptance test must hold. They are different questions on purpose: a
|
||||
test phase that ran the suite did its job even when the suite came back
|
||||
red, so the PHASE succeeds while the RUN must not.
|
||||
|
||||
This replaces a `succeeded` property that answered only the first
|
||||
question — and, being a property with side effects, wrote the session
|
||||
status and printed the banner before the caller's `and test.passed` was
|
||||
ever evaluated. A run whose suite never passed was recorded green in the
|
||||
db, on the terminal, and in the UI while exiting 1. Anyone reading the
|
||||
trace saw success; only a CI job checking `$?` saw the truth. One call
|
||||
now settles the db, the banner, and the exit code together, so the three
|
||||
cannot disagree.
|
||||
"""
|
||||
phases_ok = bool(self.phases) and all(p.status == "success" for p in self.phases)
|
||||
ok = phases_ok and accepted
|
||||
if phases_ok and not accepted:
|
||||
note = reason or "the run's acceptance criterion was not met"
|
||||
self.tracer.event(EventRecord(
|
||||
adw_id=self.adw_id,
|
||||
phase_id=self.phases[-1].phase_id if self.phases else "",
|
||||
type="error", name="not_accepted", payload={"reason": note}))
|
||||
self.console.note(f"not accepted: {note}")
|
||||
self.tracer.session_finish(self.adw_id, ok=ok)
|
||||
self.console.session_finished(ok, self.tokens, self.cost, self.cfg.observability.db)
|
||||
return 0 if ok else 1
|
||||
50
sssf/templates/adws/adw_modules/session.py
Normal file
50
sssf/templates/adws/adw_modules/session.py
Normal file
|
|
@ -0,0 +1,50 @@
|
|||
"""Session lifecycle: pin-or-create an adw_id, build the Run object.
|
||||
|
||||
`ensure(cfg, adw_id)` joins the session if it exists or creates it under
|
||||
exactly that id (pinned ids for repeatable runs); omitted, a fresh id is
|
||||
minted and printed so the next ADW can pick it up.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import signal
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
from .data_types import SSSFConfig
|
||||
from .runner import Run
|
||||
from .tracer import Tracer
|
||||
from .utils import engineer_name, new_id
|
||||
|
||||
|
||||
def _finalize_when_killed(run: Run) -> None:
|
||||
"""A killed run still closes its own trace.
|
||||
|
||||
Python's default SIGTERM handling exits without unwinding, so `just kill`
|
||||
(or any `kill <pid>`) would leave the session reading `running` forever and
|
||||
its process rows open — the trace would claim work is in flight that is
|
||||
already dead. Turning the signal into SystemExit both finalizes here and
|
||||
lets the phase context manager record the phase as failed on the way out.
|
||||
"""
|
||||
def handler(signum, _frame):
|
||||
run.tracer.session_finish(run.adw_id, ok=False) # also closes process rows
|
||||
raise SystemExit(128 + signum)
|
||||
|
||||
for sig in (signal.SIGTERM, signal.SIGINT):
|
||||
signal.signal(sig, handler)
|
||||
|
||||
|
||||
def ensure(cfg: SSSFConfig, adw_id: str | None = None) -> Run:
|
||||
adw_id = adw_id or new_id(8)
|
||||
tracer = Tracer(cfg.observability.db,
|
||||
f"{cfg.defaults.data_dir}/sessions/{adw_id}/events.jsonl")
|
||||
run = Run(cfg=cfg, adw_id=adw_id, tracer=tracer, engineer=engineer_name())
|
||||
tracer.session_start(adw_id, run.engineer, adw_name=Path(sys.argv[0]).stem)
|
||||
# This process is the run. Record it before any phase opens, so a run that
|
||||
# hangs in its first agent call is still killable by adw_id.
|
||||
tracer.process_start(adw_id, "adw", "", os.getpid(),
|
||||
" ".join([Path(sys.argv[0]).name, *sys.argv[1:]]))
|
||||
_finalize_when_killed(run)
|
||||
run.console.session_started(adw_id, run.engineer)
|
||||
return run
|
||||
271
sssf/templates/adws/adw_modules/tracer.py
Normal file
271
sssf/templates/adws/adw_modules/tracer.py
Normal file
|
|
@ -0,0 +1,271 @@
|
|||
"""Tracer: every event lands in JSONL and SQLite AS IT HAPPENS.
|
||||
|
||||
Files are the raw record; sssf.db is the queryable mirror the UI polls.
|
||||
No push transport — the flow is always: agents -> sqlite -> web ui.
|
||||
WAL mode so the UI can read while ADW processes write.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import sqlite3
|
||||
from pathlib import Path
|
||||
|
||||
from .data_types import AgentConfig, EventRecord, GateReport, Phase
|
||||
from .utils import ensure_dir, new_id, now_iso
|
||||
|
||||
SCHEMA = """
|
||||
CREATE TABLE IF NOT EXISTS sessions (
|
||||
adw_id TEXT PRIMARY KEY,
|
||||
adw_name TEXT, -- ADW script(s) run, e.g. "adw_plan + adw_build_test"
|
||||
request TEXT,
|
||||
status TEXT,
|
||||
engineer TEXT,
|
||||
started_at TEXT, ended_at TEXT,
|
||||
total_tokens INTEGER DEFAULT 0, total_cost REAL DEFAULT 0,
|
||||
archived INTEGER DEFAULT 0 -- review triage, set by the UI; never by a run
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS phases (
|
||||
phase_id TEXT PRIMARY KEY,
|
||||
adw_id TEXT REFERENCES sessions,
|
||||
seq INTEGER,
|
||||
name TEXT, kind TEXT, owner TEXT, description TEXT,
|
||||
status TEXT DEFAULT 'fail',
|
||||
attempt INTEGER DEFAULT 0, retries INTEGER DEFAULT 0,
|
||||
error TEXT,
|
||||
started_at TEXT, ended_at TEXT
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS events (
|
||||
event_id TEXT PRIMARY KEY,
|
||||
adw_id TEXT REFERENCES sessions,
|
||||
phase_id TEXT REFERENCES phases,
|
||||
parent_id TEXT,
|
||||
type TEXT,
|
||||
name TEXT,
|
||||
payload_json TEXT,
|
||||
tokens INTEGER,
|
||||
started_at TEXT, ended_at TEXT
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS envelopes (
|
||||
envelope_id TEXT PRIMARY KEY,
|
||||
adw_id TEXT REFERENCES sessions,
|
||||
phase_id TEXT REFERENCES phases,
|
||||
agent TEXT,
|
||||
output_type TEXT,
|
||||
payload_json TEXT,
|
||||
valid INTEGER,
|
||||
attempt INTEGER,
|
||||
created_at TEXT
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS gate_results (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
adw_id TEXT REFERENCES sessions,
|
||||
phase_id TEXT REFERENCES phases,
|
||||
attempt INTEGER,
|
||||
gate TEXT,
|
||||
passed INTEGER,
|
||||
violations_json TEXT,
|
||||
checks_json TEXT, -- [{item, ok, note}] — WHAT the gate verified
|
||||
created_at TEXT
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS processes (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
adw_id TEXT REFERENCES sessions,
|
||||
kind TEXT, -- 'adw' (the workflow process) | 'agent' (a coding-agent child)
|
||||
name TEXT, -- '' for the adw, the agent name for a child
|
||||
pid INTEGER,
|
||||
command TEXT, -- what the pid was, so a recycled pid is not killed by mistake
|
||||
started_at TEXT, ended_at TEXT -- ended_at NULL = believed alive
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS agent_sessions (
|
||||
adw_id TEXT REFERENCES sessions,
|
||||
agent TEXT,
|
||||
coding_agent TEXT, model TEXT, color TEXT,
|
||||
session_id TEXT,
|
||||
context_tokens INTEGER, -- window occupancy after the agent's last turn
|
||||
context_window INTEGER, -- the model's ceiling; 0/NULL = unknown
|
||||
created_at TEXT, last_used_at TEXT,
|
||||
PRIMARY KEY (adw_id, agent)
|
||||
);
|
||||
"""
|
||||
|
||||
# Columns added after a schema shipped. CREATE TABLE IF NOT EXISTS never
|
||||
# revisits an existing table, so additive changes need an explicit ALTER.
|
||||
MIGRATIONS = [("agent_sessions", "color", "TEXT"),
|
||||
("gate_results", "checks_json", "TEXT"),
|
||||
("sessions", "adw_name", "TEXT"),
|
||||
("agent_sessions", "context_tokens", "INTEGER"),
|
||||
("agent_sessions", "context_window", "INTEGER"),
|
||||
("sessions", "archived", "INTEGER DEFAULT 0")]
|
||||
|
||||
|
||||
class Tracer:
|
||||
def __init__(self, db_path: str | Path, events_jsonl: str | Path):
|
||||
ensure_dir(Path(db_path).parent)
|
||||
self.db_path = str(db_path)
|
||||
self.events_jsonl = Path(events_jsonl)
|
||||
ensure_dir(self.events_jsonl.parent)
|
||||
self.conn = sqlite3.connect(self.db_path, isolation_level=None)
|
||||
self.conn.execute("PRAGMA journal_mode=WAL;")
|
||||
self.conn.execute("PRAGMA synchronous=NORMAL;")
|
||||
self.conn.execute("PRAGMA busy_timeout=5000;")
|
||||
self.conn.executescript(SCHEMA)
|
||||
self._migrate()
|
||||
|
||||
def _migrate(self) -> None:
|
||||
"""Additive column migrations, so a db from an older SSSF still opens."""
|
||||
for table, column, decl in MIGRATIONS:
|
||||
columns = {row[1] for row in self.conn.execute(f"PRAGMA table_info({table})")}
|
||||
if column not in columns:
|
||||
self.conn.execute(f"ALTER TABLE {table} ADD COLUMN {column} {decl}")
|
||||
|
||||
# ── events ──────────────────────────────────────────────────────────────
|
||||
def event(self, record: EventRecord) -> str:
|
||||
event_id = f"evt_{new_id(12)}"
|
||||
ts = now_iso()
|
||||
line = {"event_id": event_id, "ts": ts, **record.model_dump()}
|
||||
with self.events_jsonl.open("a") as f:
|
||||
f.write(json.dumps(line) + "\n")
|
||||
self.conn.execute(
|
||||
"INSERT INTO events (event_id, adw_id, phase_id, parent_id, type, name,"
|
||||
" payload_json, tokens, started_at, ended_at) VALUES (?,?,?,?,?,?,?,?,?,?)",
|
||||
(event_id, record.adw_id, record.phase_id, record.parent_id, record.type,
|
||||
record.name, json.dumps(record.payload), record.tokens,
|
||||
record.started_at or ts, record.ended_at),
|
||||
)
|
||||
return event_id
|
||||
|
||||
# ── sessions ────────────────────────────────────────────────────────────
|
||||
def session_start(self, adw_id: str, engineer: str, adw_name: str | None = None) -> None:
|
||||
self.conn.execute(
|
||||
"INSERT INTO sessions (adw_id, status, engineer, started_at) VALUES (?,?,?,?) "
|
||||
"ON CONFLICT(adw_id) DO UPDATE SET status='running'",
|
||||
(adw_id, "running", engineer, now_iso()),
|
||||
)
|
||||
if not adw_name:
|
||||
return
|
||||
# A joined session chains ADWs — record each distinct one, in run order.
|
||||
row = self.conn.execute("SELECT adw_name FROM sessions WHERE adw_id=?",
|
||||
(adw_id,)).fetchone()
|
||||
names = row[0].split(" + ") if row and row[0] else []
|
||||
if adw_name not in names:
|
||||
names.append(adw_name)
|
||||
self.conn.execute("UPDATE sessions SET adw_name=? WHERE adw_id=?",
|
||||
(" + ".join(names), adw_id))
|
||||
|
||||
def session_request(self, adw_id: str, request: str) -> None:
|
||||
self.conn.execute("UPDATE sessions SET request=? WHERE adw_id=?",
|
||||
(request[:500], adw_id))
|
||||
|
||||
def session_finish(self, adw_id: str, ok: bool) -> None:
|
||||
self.conn.execute(
|
||||
"UPDATE sessions SET status=?, ended_at=? WHERE adw_id=?",
|
||||
("success" if ok else "fail", now_iso(), adw_id),
|
||||
)
|
||||
self.processes_end_all(adw_id) # nothing of this run is alive any more
|
||||
|
||||
def session_add_usage(self, adw_id: str, tokens: int, cost: float) -> None:
|
||||
self.conn.execute(
|
||||
"UPDATE sessions SET total_tokens=total_tokens+?, total_cost=total_cost+? WHERE adw_id=?",
|
||||
(tokens, cost, adw_id),
|
||||
)
|
||||
|
||||
# ── processes (adw_id → pid, so a hung run can be found and killed) ─────
|
||||
def process_start(self, adw_id: str, kind: str, name: str, pid: int,
|
||||
command: str) -> None:
|
||||
"""Record a live process for this run.
|
||||
|
||||
A coding agent that hangs produces no events at all, which is exactly
|
||||
when you need its pid — and `ps` cannot tell you which adw_id it
|
||||
belongs to. Writing it here makes the trace the answer to "what is this
|
||||
run running, and how do I stop it".
|
||||
"""
|
||||
self.conn.execute(
|
||||
"INSERT INTO processes (adw_id, kind, name, pid, command, started_at)"
|
||||
" VALUES (?,?,?,?,?,?)",
|
||||
(adw_id, kind, name, pid, command[:500], now_iso()),
|
||||
)
|
||||
|
||||
def process_end(self, adw_id: str, pid: int) -> None:
|
||||
"""Mark the newest live row for this pid as finished."""
|
||||
self.conn.execute(
|
||||
"UPDATE processes SET ended_at=? WHERE id = ("
|
||||
" SELECT id FROM processes WHERE adw_id=? AND pid=? AND ended_at IS NULL"
|
||||
" ORDER BY id DESC LIMIT 1)",
|
||||
(now_iso(), adw_id, pid),
|
||||
)
|
||||
|
||||
def processes_end_all(self, adw_id: str) -> None:
|
||||
"""Close out every live row for a run — called when the session ends."""
|
||||
self.conn.execute(
|
||||
"UPDATE processes SET ended_at=? WHERE adw_id=? AND ended_at IS NULL",
|
||||
(now_iso(), adw_id),
|
||||
)
|
||||
|
||||
# ── phases ──────────────────────────────────────────────────────────────
|
||||
def max_phase_seq(self, adw_id: str) -> int:
|
||||
"""Highest seq already recorded for this session; 0 when it is new.
|
||||
|
||||
A joined run continues the sequence instead of restarting at 1 — which
|
||||
would collide with the first run's phases on both `seq` (breaking
|
||||
ordering) and `phase_id` (silently overwriting a row through the
|
||||
phase_upsert conflict clause).
|
||||
"""
|
||||
row = self.conn.execute("SELECT MAX(seq) FROM phases WHERE adw_id = ?",
|
||||
(adw_id,)).fetchone()
|
||||
return row[0] if row and row[0] is not None else 0
|
||||
|
||||
def phase_upsert(self, phase: Phase) -> None:
|
||||
p = phase.params
|
||||
self.conn.execute(
|
||||
"INSERT INTO phases (phase_id, adw_id, seq, name, kind, owner, description,"
|
||||
" status, attempt, retries, error, started_at, ended_at)"
|
||||
" VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)"
|
||||
" ON CONFLICT(phase_id) DO UPDATE SET status=excluded.status,"
|
||||
" attempt=excluded.attempt, error=excluded.error, ended_at=excluded.ended_at",
|
||||
(phase.phase_id, phase.adw_id, phase.seq, p.name, p.kind, p.owner,
|
||||
p.description, phase.status, phase.attempt, p.retries, phase.error,
|
||||
phase.started_at, phase.ended_at),
|
||||
)
|
||||
|
||||
# ── envelopes / gates / agent sessions ──────────────────────────────────
|
||||
def envelope_row(self, phase: Phase, agent: str, output_type: str,
|
||||
payload_json: str, valid: bool, attempt: int) -> None:
|
||||
self.conn.execute(
|
||||
"INSERT INTO envelopes (envelope_id, adw_id, phase_id, agent, output_type,"
|
||||
" payload_json, valid, attempt, created_at) VALUES (?,?,?,?,?,?,?,?,?)",
|
||||
(f"env_{new_id(12)}", phase.adw_id, phase.phase_id, agent, output_type,
|
||||
payload_json, int(valid), attempt, now_iso()),
|
||||
)
|
||||
|
||||
def gate_row(self, phase: Phase, gate: str, report: GateReport, attempt: int) -> None:
|
||||
"""The report carries both the verdict and the evidence behind it."""
|
||||
self.conn.execute(
|
||||
"INSERT INTO gate_results (adw_id, phase_id, attempt, gate, passed,"
|
||||
" violations_json, checks_json, created_at) VALUES (?,?,?,?,?,?,?,?)",
|
||||
(phase.adw_id, phase.phase_id, attempt, gate, int(report.passed),
|
||||
json.dumps(report.violations),
|
||||
json.dumps([c.model_dump() for c in report.checks]), now_iso()),
|
||||
)
|
||||
|
||||
def agent_session_row(self, adw_id: str, agent: AgentConfig, session_id: str,
|
||||
context_tokens: int = 0, context_window: int = 0) -> None:
|
||||
"""The agent's config row is the source of truth for its label and color.
|
||||
|
||||
Context is carried here rather than derived from events because the lane
|
||||
wants one number per agent — the latest — and a session that runs the
|
||||
same agent twice overwrites it, exactly like model and session_id.
|
||||
"""
|
||||
ts = now_iso()
|
||||
self.conn.execute(
|
||||
"INSERT INTO agent_sessions (adw_id, agent, coding_agent, model, color,"
|
||||
" session_id, context_tokens, context_window, created_at, last_used_at)"
|
||||
" VALUES (?,?,?,?,?,?,?,?,?,?)"
|
||||
" ON CONFLICT(adw_id, agent) DO UPDATE SET model=excluded.model,"
|
||||
" color=excluded.color, session_id=excluded.session_id,"
|
||||
" context_tokens=excluded.context_tokens,"
|
||||
" context_window=excluded.context_window,"
|
||||
" last_used_at=excluded.last_used_at",
|
||||
(adw_id, agent.name, agent.coding_agent, agent.model, agent.color,
|
||||
session_id, context_tokens, context_window, ts, ts),
|
||||
)
|
||||
77
sssf/templates/adws/adw_modules/utils.py
Normal file
77
sssf/templates/adws/adw_modules/utils.py
Normal file
|
|
@ -0,0 +1,77 @@
|
|||
"""Small shared helpers. Anything bigger belongs in its own module."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import secrets
|
||||
import subprocess
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
from dotenv import load_dotenv
|
||||
|
||||
load_dotenv()
|
||||
|
||||
|
||||
def operator_env() -> dict[str, str]:
|
||||
"""The engineer's own environment, as their shell would hand it over.
|
||||
|
||||
Agents and quality blocks are meant to see exactly what the operator sees:
|
||||
their PATH, their toolchains, their globally installed packages. Copying
|
||||
os.environ gets almost all the way there — but ADWs launch under `uv run`,
|
||||
which prepends its ephemeral venv's bin to PATH and sets VIRTUAL_ENV. That
|
||||
venv holds the ADW's OWN dependencies (pydantic, pyyaml), not the
|
||||
operator's, so anything a subprocess resolves through it — `python3`,
|
||||
`pip`, every globally pip-installed CLI — silently becomes the wrong one.
|
||||
|
||||
Stripping the venv restores parity: `python3` in an agent's bash is the
|
||||
same `python3` the engineer gets in their terminal. The ADW's own imports
|
||||
are unaffected; this env is only ever handed to child processes.
|
||||
"""
|
||||
env = os.environ.copy()
|
||||
venv = env.pop("VIRTUAL_ENV", "")
|
||||
if not venv:
|
||||
return env
|
||||
venv_bin = str(Path(venv) / "bin")
|
||||
parts = [p for p in env.get("PATH", "").split(os.pathsep) if p and p != venv_bin]
|
||||
env["PATH"] = os.pathsep.join(parts)
|
||||
return env
|
||||
|
||||
|
||||
def new_id(length: int = 8) -> str:
|
||||
return secrets.token_hex(length // 2)
|
||||
|
||||
|
||||
def now_iso() -> str:
|
||||
return datetime.now(timezone.utc).isoformat(timespec="milliseconds")
|
||||
|
||||
|
||||
def ensure_dir(path: str | Path) -> Path:
|
||||
p = Path(path)
|
||||
p.mkdir(parents=True, exist_ok=True)
|
||||
return p
|
||||
|
||||
|
||||
def resolve_prompt(arg: str) -> str:
|
||||
"""CLI prompt arg: a file path resolves to its contents, else inline text."""
|
||||
try:
|
||||
p = Path(arg)
|
||||
if p.is_file():
|
||||
return p.read_text()
|
||||
except OSError:
|
||||
pass
|
||||
return arg
|
||||
|
||||
|
||||
def engineer_name() -> str:
|
||||
name = os.environ.get("ENGINEER_NAME", "").strip()
|
||||
if name:
|
||||
return name
|
||||
try:
|
||||
out = subprocess.run(["git", "config", "user.name"],
|
||||
capture_output=True, text=True, timeout=5)
|
||||
if out.returncode == 0 and out.stdout.strip():
|
||||
return out.stdout.strip()
|
||||
except OSError:
|
||||
pass
|
||||
return os.environ.get("USER", "engineer")
|
||||
Loading…
Add table
Add a link
Reference in a new issue