diff --git a/src/git_agent/quilt_emit.py b/src/git_agent/quilt_emit.py new file mode 100644 index 0000000..7e7e270 --- /dev/null +++ b/src/git_agent/quilt_emit.py @@ -0,0 +1,245 @@ +""" +Vessel-Quilt emitter — git-agent's lifecycle events as quilt 5-opcode WAL records. + +The fleet's quilt kernel speaks five opcodes (BIND / LINK / EFFECT / VIEW / TICK, +plus FORGET). git-agent's vessel already produces structured lifecycle events +(session start, worklog entries, task outcomes, career promotions, fences, +skills, snapshots, heartbeats, session end) — this module is the translation +and validation layer that makes git-agent the first quilt-native fleet agent. + +Every ingested event is: + 1. validated against the event schema (enforced here in stdlib code; the + canonical JSON Schema lives beside the producer at + ``git_agent/schemas/event.schema.json`` for external consumers), + 2. mapped to exactly one quilt opcode line, + 3. appended to a fnv1a hash-chained JSONL WAL (default ``~/.git-agent/quilt.jsonl``), + 4. folded into a reducer state that ``replay()`` can reconstruct from the WAL alone. + +Guarantees: invalid events raise ``QuiltValidationError`` (never silently +dropped, never crash the ingest loop); duplicate ``event_id`` deliveries are +idempotent (at-least-once ingest → exactly-once WAL). +""" + +from __future__ import annotations + +import datetime +import json +import re +from pathlib import Path +from typing import Any, Dict, List, Optional + +SCHEMA_PATH = Path(__file__).parent / "schemas" / "event.schema.json" + +DEFAULT_WAL = Path.home() / ".git-agent" / "quilt.jsonl" + +EVENT_TYPES = ( + "session_start", "worklog", "task_completion", "promotion", + "fence", "skill", "snapshot", "heartbeat", "session_end", +) + +# per-type required fields beyond the universal (event_id, type, timestamp) +REQUIRED: Dict[str, tuple] = { + "session_start": ("name", "designation", "version"), + "worklog": ("action", "target", "summary", "outcome"), + "task_completion": ("success",), + "promotion": ("from_stage", "to_stage"), + "fence": ("fence_name",), + "skill": ("skill",), + "snapshot": ("stage",), + "heartbeat": (), + "session_end": (), +} + +OUTCOMES = ("success", "failure", "partial") +_TS_RE = re.compile(r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}") + + +class QuiltValidationError(ValueError): + """Raised when an event fails schema validation. Precise about the why.""" + + +def fnv1a(s: str) -> str: + """64-bit FNV-1a, hex-encoded (16 chars). Same integrity family the fleet's + rate limiter and quilt WAL use — fast, deterministic, non-cryptographic.""" + h = 0xCBF29CE484222325 + for b in s.encode("utf-8"): + h ^= b + h = (h * 0x100000001B3) & 0xFFFFFFFFFFFFFFFF + return f"{h:016x}" + + +def _canonical(obj: Dict[str, Any]) -> str: + return json.dumps(obj, sort_keys=True, separators=(",", ":")) + + +def _validate(event: Dict[str, Any]) -> None: + if not isinstance(event, dict): + raise QuiltValidationError("event must be a JSON object") + for field in ("event_id", "type", "timestamp"): + if field not in event: + raise QuiltValidationError(f"missing required field: {field!r}") + if not isinstance(event["event_id"], str) or not event["event_id"]: + raise QuiltValidationError("event_id must be a non-empty string") + if event["type"] not in EVENT_TYPES: + raise QuiltValidationError( + f"unknown event type: {event['type']!r} (expected one of {', '.join(EVENT_TYPES)})") + if not (isinstance(event["timestamp"], str) and _TS_RE.match(event["timestamp"])): + raise QuiltValidationError( + f"timestamp must be ISO-8601 (YYYY-MM-DDTHH:MM:SS...), got {event['timestamp']!r}") + for field in REQUIRED[event["type"]]: + if field not in event: + raise QuiltValidationError( + f"{event['type']} event missing required field: {field!r}") + if event["type"] == "worklog": + if event["outcome"] not in OUTCOMES: + raise QuiltValidationError( + f"worklog outcome must be one of {OUTCOMES}, got {event['outcome']!r}") + + +def _map(event: Dict[str, Any]) -> Dict[str, Any]: + """One vessel event → one quilt opcode line (minus chain fields).""" + t = event["type"] + if t == "session_start": + return {"op": "BIND", "cell": "vessel/identity", + "args": {"name": event["name"], "designation": event["designation"], + "version": event["version"], "domains": event.get("domains", [])}} + if t == "worklog": + return {"op": "LINK", "cell": "vessel/worklog", + "args": {"action": event["action"], "target": event["target"], + "summary": event["summary"], "outcome": event["outcome"]}} + if t == "task_completion": + return {"op": "EFFECT", "cell": "vessel/career", + "args": {"kind": "task", "success": bool(event["success"])}} + if t == "promotion": + return {"op": "EFFECT", "cell": "vessel/career", + "args": {"kind": "promotion", "from_stage": event["from_stage"], + "to_stage": event["to_stage"]}} + if t == "fence": + return {"op": "EFFECT", "cell": "vessel/career", + "args": {"kind": "fence", "fence_name": event["fence_name"]}} + if t == "skill": + return {"op": "BIND", "cell": "vessel/skills", "args": {"skill": event["skill"]}} + if t == "snapshot": + return {"op": "VIEW", "cell": "vessel/state", + "args": {"stage": event["stage"], + "total_tasks_completed": event.get("total_tasks_completed"), + "total_tasks_failed": event.get("total_tasks_failed"), + "worklog_len": event.get("worklog_len")}} + if t == "heartbeat": + return {"op": "TICK", "cell": "vessel/heartbeat", "args": {}} + if t == "session_end": + return {"op": "FORGET", "cell": "vessel/session", + "args": {"archived": bool(event.get("archived", True))}} + raise QuiltValidationError(f"unmapped event type: {t}") # pragma: no cover + + +class QuiltEmitter: + """Ingest vessel lifecycle events → append hash-chained quilt WAL lines.""" + + def __init__(self, wal_path: Optional[Path] = None): + self.wal_path = Path(wal_path) if wal_path else DEFAULT_WAL + self._seen: set = set() + self._state: Dict[str, Any] = self._empty_state() + if self.wal_path.exists(): + for line in self._read_lines(): + self._seen.add(line["event_id"]) + self._fold(line, self._state) + + @staticmethod + def _empty_state() -> Dict[str, Any]: + return {"identity": None, "tasks_completed": 0, "tasks_failed": 0, + "promotions": [], "fences": [], "skills": [], "worklog": 0, + "last_tick": None, "archived": False, "lines": 0} + + def _read_lines(self) -> List[Dict[str, Any]]: + if not self.wal_path.exists(): + return [] + out = [] + for raw in self.wal_path.read_text().splitlines(): + if raw.strip(): + out.append(json.loads(raw)) + return out + + def wal(self) -> List[Dict[str, Any]]: + return self._read_lines() + + def ingest(self, event: Dict[str, Any]) -> Dict[str, Any]: + _validate(event) + if event["event_id"] in self._seen: + for line in self._read_lines(): # idempotent no-op: return the existing line + if line["event_id"] == event["event_id"]: + return line + raise QuiltValidationError( # pragma: no cover + f"event_id {event['event_id']!r} seen but absent from WAL") + line = _map(event) + line["event_id"] = event["event_id"] + line["timestamp"] = event["timestamp"] + prev = self._read_lines() + line["seq"] = len(prev) + line["prev_hash"] = prev[-1]["hash"] if prev else "0" * 16 + line["hash"] = fnv1a(_canonical({k: v for k, v in line.items() if k != "hash"})) + self.wal_path.parent.mkdir(parents=True, exist_ok=True) + with open(self.wal_path, "a", encoding="utf-8") as f: + f.write(json.dumps(line, sort_keys=True) + "\n") + self._seen.add(event["event_id"]) + self._fold(line, self._state) + return line + + @staticmethod + def _fold(line: Dict[str, Any], st: Dict[str, Any]) -> None: + st["lines"] += 1 + op, args = line["op"], line.get("args", {}) + if op == "BIND" and line["cell"] == "vessel/identity": + st["identity"] = args + elif op == "BIND" and line["cell"] == "vessel/skills": + if args.get("skill") not in st["skills"]: + st["skills"].append(args.get("skill")) + elif op == "LINK": + st["worklog"] += 1 + elif op == "EFFECT": + k = args.get("kind") + if k == "task": + st["tasks_completed" if args.get("success") else "tasks_failed"] += 1 + elif k == "promotion": + st["promotions"].append((args.get("from_stage"), args.get("to_stage"))) + elif k == "fence": + if args.get("fence_name") not in st["fences"]: + st["fences"].append(args.get("fence_name")) + elif op == "TICK": + st["last_tick"] = line.get("timestamp") + elif op == "FORGET": + st["archived"] = bool(args.get("archived")) + + def state(self) -> Dict[str, Any]: + return self._state + + def replay(self) -> Dict[str, Any]: + st = self._empty_state() + for line in self._read_lines(): + self._fold(line, st) + return st + + def verify(self) -> Dict[str, Any]: + """Replay the WAL from disk and diff against expectations. + + Checks: (a) per-line hash recomputation, (b) chain linkage, (c) seq + continuity, (d) replay equivalence with the live reducer state. + """ + divergences: List[str] = [] + lines = self._read_lines() + for i, line in enumerate(lines): + want = fnv1a(_canonical({k: v for k, v in line.items() if k != "hash"})) + if line.get("hash") != want: + divergences.append(f"line {i}: hash mismatch (tampered)") + if i == 0: + if line.get("prev_hash") != "0" * 16: + divergences.append("line 0: genesis prev_hash must be 16 zeros") + else: + if line.get("prev_hash") != lines[i - 1].get("hash"): + divergences.append(f"line {i}: chain link broken (truncation or reorder)") + if line.get("seq") != i: + divergences.append(f"line {i}: seq gap (expected {i}, got {line.get('seq')})") + replayed = self.replay() + if replayed != self._state: + divergences.append("replay state diverges from live reducer state") + return {"ok": not divergences, "divergences": divergences} diff --git a/src/git_agent/schemas/event.schema.json b/src/git_agent/schemas/event.schema.json new file mode 100644 index 0000000..8d1b0cb --- /dev/null +++ b/src/git_agent/schemas/event.schema.json @@ -0,0 +1,101 @@ +{ + "$schema": "http://json-schema.org/draft-07/schema#", + "title": "git-agent vessel lifecycle event", + "description": "Schema for events ingested by git_agent.quilt_emit. Lives with the producer (git-agent), per fleet doctrine: schemas belong with the producer, not the linter.", + "type": "object", + "required": ["event_id", "type", "timestamp"], + "additionalProperties": true, + "definitions": { + "timestamp": { + "type": "string", + "format": "date-time", + "pattern": "^\\d{4}-\\d{2}-\\d{2}T\\d{2}:\\d{2}:\\d{2}" + } + }, + "allOf": [ + { + "if": { "properties": { "type": { "const": "session_start" } } }, + "then": { + "required": ["event_id", "type", "timestamp", "name", "designation", "version"], + "properties": { + "name": { "type": "string", "minLength": 1 }, + "designation": { "type": "string", "minLength": 1 }, + "version": { "type": "string", "minLength": 1 }, + "domains": { "type": "array", "items": { "type": "string" } } + } + } + }, + { + "if": { "properties": { "type": { "const": "worklog" } } }, + "then": { + "required": ["event_id", "type", "timestamp", "action", "target", "summary", "outcome"], + "properties": { + "action": { "type": "string", "minLength": 1 }, + "target": { "type": "string", "minLength": 1 }, + "summary": { "type": "string", "minLength": 1 }, + "outcome": { "type": "string", "enum": ["success", "failure", "partial"] } + } + } + }, + { + "if": { "properties": { "type": { "const": "task_completion" } } }, + "then": { + "required": ["event_id", "type", "timestamp", "success"], + "properties": { "success": { "type": "boolean" } } + } + }, + { + "if": { "properties": { "type": { "const": "promotion" } } }, + "then": { + "required": ["event_id", "type", "timestamp", "from_stage", "to_stage"], + "properties": { + "from_stage": { "type": "string", "minLength": 1 }, + "to_stage": { "type": "string", "minLength": 1 } + } + } + }, + { + "if": { "properties": { "type": { "const": "fence" } } }, + "then": { + "required": ["event_id", "type", "timestamp", "fence_name"], + "properties": { "fence_name": { "type": "string", "minLength": 1 } } + } + }, + { + "if": { "properties": { "type": { "const": "skill" } } }, + "then": { + "required": ["event_id", "type", "timestamp", "skill"], + "properties": { "skill": { "type": "string", "minLength": 1 } } + } + }, + { + "if": { "properties": { "type": { "const": "snapshot" } } }, + "then": { + "required": ["event_id", "type", "timestamp", "stage"], + "properties": { "stage": { "type": "string", "minLength": 1 } } + } + }, + { + "if": { "properties": { "type": { "const": "session_end" } } }, + "then": { + "required": ["event_id", "type", "timestamp"], + "properties": { "archived": { "type": "boolean" } } + } + }, + { + "if": { "properties": { "type": { "const": "heartbeat" } } }, + "then": { "required": ["event_id", "type", "timestamp"] } + } + ], + "properties": { + "event_id": { "type": "string", "minLength": 1 }, + "type": { + "type": "string", + "enum": [ + "session_start", "worklog", "task_completion", "promotion", + "fence", "skill", "snapshot", "heartbeat", "session_end" + ] + }, + "timestamp": { "$ref": "#/definitions/timestamp" } + } +} diff --git a/tests/_capture_fixtures.py b/tests/_capture_fixtures.py new file mode 100644 index 0000000..df04741 --- /dev/null +++ b/tests/_capture_fixtures.py @@ -0,0 +1,88 @@ +"""One-shot fixture capture: emit the REAL lifecycle events a git-agent vessel +produces, save them as test fixtures. Run: python3 _capture_fixtures.py