trading-system-docs/a2/rq_mission.py
Red Queen 968094498a feat(a5): Controlled Heartbeat & Resume v1
- a5/rq_heartbeat.py: HeartbeatStore + Heartbeat.run_tick (thin scheduler/resume layer)
- a5/rq_heartbeat_cli.py: status/tick/enable/disable/resume/approve/deny/priority
- a5/test_a5.py: 31 deterministic tests (tick-lock incl. ownership, eligibility, kill-switch, approval, circuit, git-conflict, bounded single A4)
- a5/DESIGN.md, a5/README.md, a5/scripts/a5_heartbeat_tick.sh
- a2/rq_mission.py: add read-only mission_list() for A5 enumeration (additive)
- Fresh checker: PASS after tick-lock ownership repair (Regressionschutz lock_owned)
2026-08-24 23:04:16 +00:00

638 lines
27 KiB
Python

#!/usr/bin/env python3
"""
Red Queen — A2: Persistentes Mission-State & Work-Package-State Management.
Schichtung (siehe A2 §8 NATIVE vs MARKDOWN):
* MISSION-Ebene -> dieses Modul + `missions.db` (SQLite, stdlib sqlite3).
Missions-IDs, Titel/Goal/Scope/Constraints/Acceptance,
Mission-State (A1 State Machine), Transition-Evidence.
Der Mission-State und die WPs werden hier deterministisch
validiert. Die `missions.db` ist der PRIMAERE, durable
State fuer MISSION & WORK PACKAGE.
* WORK-PACKAGE -> Wird in `missions.db` persistiert UND (Best Effort) als
natives Kanban-Task gespiegelt, damit Dependencies und
Block nativen Hermes-Mechanismen nutzen. Die Kanban-Tasks
sind ein DURABLE MIRROR; die Authoritat (Fach-Transitions)
liegt in `missions.db`. Der Mirror ist optional aktivierbar
(kanban_db Pfad), so dass Tests ihn komplett isolieren oder
deaktivieren koennen.
Idempotenz & Restart-Stabilitaet:
* Missions-ID `RQ-MISSION-YYYYMMDD-NNN` (NNN laufend pro Tag, Max+1 restart-stabil).
* WP-ID `<mission_id>-WP-NNN` (NNN laufend pro Mission).
* create mit gleicher ID / Idempotency-Key -> kein Duplikat.
* DONE->DONE erzeugt KEIN zweites Transition-Event.
* Doppelte Dependency -> kein Duplikat-Eintrag.
* Jede persistente Mutation laeuft in einer SQLite-Transaktion.
Fail-Closed:
* Bei State-Inkonsistenz / unbekanntem State -> `RqError` (Code STATE_ERROR),
KEINE Mutation, Evidence wird gemeldet.
* Ungueltige Transition -> `RqError` (Code INVALID_TRANSITION), keine Mutation.
KEINE Orchestrierung: kein Loop/Retry/Heartbeat/Cron in diesem Modul.
"""
from __future__ import annotations
import datetime
import json
import os
import sqlite3
import subprocess
from pathlib import Path
from typing import Any, Dict, List, Optional
from rq_state_machine import (
is_mission_state,
is_wp_state,
mission_transition_allowed,
wp_transition_allowed,
)
# --------------------------------------------------------------------------- #
# Fehler-Codes (maschinenlesbar)
# --------------------------------------------------------------------------- #
class RqError(Exception):
"""Geworfener, maschinenlesbarer Fehler mit stabilem `code`."""
def __init__(self, code: str, message: str, detail: Optional[Dict[str, Any]] = None):
super().__init__(message)
self.code = code
self.message = message
self.detail = detail or {}
def to_dict(self) -> Dict[str, Any]:
return {"code": self.code, "message": self.message, "detail": self.detail}
def _err(code: str, msg: str, **detail) -> RqError:
return RqError(code, msg, detail)
# --------------------------------------------------------------------------- #
# DB-Hilfsfunktionen
# --------------------------------------------------------------------------- #
def _connect(path: str) -> sqlite3.Connection:
conn = sqlite3.connect(path)
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("PRAGMA foreign_keys=ON")
return conn
def _utcnow() -> str:
return datetime.datetime.now(datetime.timezone.utc).isoformat(timespec="seconds")
# --------------------------------------------------------------------------- #
# Kanban-Mirror (Best Effort, isoliert per HERMES_KANBAN_DB)
# --------------------------------------------------------------------------- #
class KanbanMirror:
"""Reflektiert WPs als native Kanban-Tasks (durable mirror).
Nur aktiv, wenn `db` gesetzt und nicht 'off'. Aufrufe verwenden die
`hermes kanban` CLI mit `HERMES_KANBAN_DB` auf einen dedizierten DB-Pfad.
Fehler im Mirror sind NICHT fatal fuer die Missions-Transaktion: die
`missions.db` ist immer authoritativ.
"""
def __init__(self, db: Optional[str]):
self.enabled = bool(db) and db.strip().lower() != "off"
self.db = Path(db).expanduser() if self.enabled else None
self._env = dict(os.environ)
if self.enabled:
self._env["HERMES_KANBAN_DB"] = str(self.db)
def _run(self, args: List[str]) -> str:
cmd = ["hermes", "kanban", *args]
try:
proc = subprocess.run(cmd, capture_output=True, text=True, timeout=60, env=self._env)
return (proc.stdout or "").strip()
except Exception:
return ""
def _show_task(self, task_id: str) -> Optional[Dict[str, Any]]:
out = self._run(["show", task_id, "--json"])
try:
return json.loads(out).get("task", {})
except Exception:
return None
def create(self, title: str, mission_id: str, wp_id: str) -> Optional[str]:
out = self._run(
["create", title, "--idempotency-key", f"rq:{mission_id}:{wp_id}", "--json"]
)
try:
return json.loads(out).get("id")
except Exception:
return None
def link(self, parent_task: str, child_task: str) -> None:
if parent_task and child_task:
self._run(["link", parent_task, child_task])
def block(self, task_id: str, reason: str) -> None:
if task_id:
self._run(["block", task_id, "--kind", "needs_input", reason])
def unblock(self, task_id: str) -> None:
if task_id:
self._run(["unblock", task_id, "--reason", "unblocked by state machine"])
def complete(self, task_id: str, result: str) -> None:
if task_id:
self._run(["complete", task_id, "--result", result])
# --------------------------------------------------------------------------- #
# Missions-Store
# --------------------------------------------------------------------------- #
class MissionStore:
"""Persistenter, deterministischer MISSION & WP Store auf einer SQLite-DB."""
def __init__(self, db_path: str, mirror_db: Optional[str] = None):
self.db_path = db_path
self._ensure_schema()
self.mirror = KanbanMirror(mirror_db)
# -- Schema -------------------------------------------------------------
def _ensure_schema(self) -> None:
conn = _connect(self.db_path)
try:
conn.executescript(
"""
CREATE TABLE IF NOT EXISTS missions (
id TEXT PRIMARY KEY,
title TEXT NOT NULL,
goal TEXT,
scope TEXT,
constraints TEXT,
acceptance_criteria TEXT,
state TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS work_packages (
id TEXT PRIMARY KEY,
mission_id TEXT NOT NULL REFERENCES missions(id),
title TEXT NOT NULL,
state TEXT NOT NULL,
kanban_task_id TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS wp_dependencies (
mission_id TEXT NOT NULL,
wp_id TEXT NOT NULL,
depends_on TEXT NOT NULL,
PRIMARY KEY (mission_id, wp_id, depends_on)
);
CREATE TABLE IF NOT EXISTS mission_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
mission_id TEXT NOT NULL,
transition TEXT NOT NULL,
actor TEXT NOT NULL,
evidence TEXT,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS wp_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
wp_id TEXT NOT NULL,
mission_id TEXT NOT NULL,
transition TEXT NOT NULL,
actor TEXT NOT NULL,
evidence TEXT,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_wp_mission ON work_packages(mission_id);
CREATE INDEX IF NOT EXISTS idx_me_mission ON mission_events(mission_id);
CREATE INDEX IF NOT EXISTS idx_we_wp ON wp_events(wp_id);
"""
)
conn.commit()
finally:
conn.close()
# -- row helpers ----------------------------------------------------------
def _mission_row(self, conn, mission_id: str):
return conn.execute("SELECT * FROM missions WHERE id = ?", (mission_id,)).fetchone()
def _wp_row(self, conn, wp_id: str):
return conn.execute("SELECT * FROM work_packages WHERE id = ?", (wp_id,)).fetchone()
def _mission_exists(self, conn, mission_id: str) -> bool:
return conn.execute("SELECT 1 FROM missions WHERE id = ?", (mission_id,)).fetchone() is not None
def _wp_exists(self, conn, wp_id: str) -> bool:
return conn.execute("SELECT 1 FROM work_packages WHERE id = ?", (wp_id,)).fetchone() is not None
# -- ID-Erzeugung ----------------------------------------------------------
def _new_mission_id(self, conn) -> str:
"""Restart-stabile Mission-ID: RQ-MISSION-YYYYMMDD-NNN (Max+1 pro Tag).
Der laufende Zaehler wird aus den bereits vorhandenen IDs des heutigen
Tages berechnet (Max+1). Das Parsen des Suffixes erfolgt in Python,
da SQLite hier kein `reverse()` bietet.
"""
today = datetime.date.today().strftime("%Y%m%d")
prefix = f"RQ-MISSION-{today}-"
rows = conn.execute("SELECT id FROM missions WHERE id LIKE ?", (prefix + "%",)).fetchall()
n = 0
for r in rows:
suffix = r["id"][len(prefix):]
if suffix.isdigit():
n = max(n, int(suffix))
return f"{prefix}{n + 1:03d}"
def _new_wp_id(self, conn, mission_id: str) -> str:
"""Restart-stabile WP-ID: <mission_id>-WP-NNN (Max+1 pro Mission)."""
row = conn.execute(
"SELECT COALESCE(MAX(CAST(substr(id, instr(id,'-WP-') + 4) AS INTEGER)), 0) AS m "
"FROM work_packages WHERE mission_id = ?",
(mission_id,),
).fetchone()
n = int(row["m"]) + 1 if row and row["m"] is not None else 1
return f"{mission_id}-WP-{n:03d}"
# -- Mission: Create ------------------------------------------------------
def mission_create(
self,
title: str,
goal: Optional[str] = None,
scope: Optional[str] = None,
constraints: Optional[str] = None,
acceptance_criteria: Optional[str] = None,
mission_id: Optional[str] = None,
actor: str = "red-queen",
) -> Dict[str, Any]:
if not title or not title.strip():
raise _err("INVALID_ARGS", "mission title required")
conn = _connect(self.db_path)
try:
if mission_id:
if self._mission_exists(conn, mission_id):
return self._mission_read(conn, mission_id) # idempotent
mid = mission_id
else:
mid = self._new_mission_id(conn)
now = _utcnow()
conn.execute(
"INSERT INTO missions "
"(id,title,goal,scope,constraints,acceptance_criteria,state,created_at,updated_at) "
"VALUES (?,?,?,?,?,?,?,?,?)",
(mid, title.strip(), goal, scope, constraints, acceptance_criteria, "CREATED", now, now),
)
conn.execute(
"INSERT INTO mission_events (mission_id,transition,actor,evidence,created_at) "
"VALUES (?,?,?,?,?)",
(mid, "CREATED", actor, "mission created", now),
)
conn.commit()
return self._mission_read(conn, mid)
finally:
conn.close()
def _mission_read(self, conn, mission_id: str) -> Dict[str, Any]:
m = self._mission_row(conn, mission_id)
if m is None:
raise _err("UNKNOWN_MISSION", f"unknown mission id {mission_id!r}", mission_id=mission_id)
wps = [
{"id": r["id"], "title": r["title"], "state": r["state"], "kanban_task_id": r["kanban_task_id"]}
for r in conn.execute(
"SELECT id,title,state,kanban_task_id FROM work_packages WHERE mission_id=? ORDER BY id",
(mission_id,),
)
]
deps = [
{"wp_id": r["wp_id"], "depends_on": r["depends_on"]}
for r in conn.execute(
"SELECT wp_id, depends_on FROM wp_dependencies WHERE mission_id=? ORDER BY wp_id, depends_on",
(mission_id,),
)
]
return {
"id": m["id"],
"title": m["title"],
"goal": m["goal"],
"scope": m["scope"],
"constraints": m["constraints"],
"acceptance_criteria": m["acceptance_criteria"],
"state": m["state"],
"created_at": m["created_at"],
"updated_at": m["updated_at"],
"work_packages": wps,
"dependencies": deps,
}
def mission_read(self, mission_id: str) -> Dict[str, Any]:
conn = _connect(self.db_path)
try:
return self._mission_read(conn, mission_id)
finally:
conn.close()
def mission_list(self) -> List[Dict[str, Any]]:
"""Read-only: alle Missionen (id, title, state, created_at, updated_at),
deterministisch sortiert (created_at ASC, id ASC). Reine Enumeration —
keine Mutation, keine Transition. Dient A5 der Heartbeat-Eligibility."""
conn = _connect(self.db_path)
try:
rows = conn.execute(
"SELECT id,title,state,created_at,updated_at "
"FROM missions ORDER BY created_at ASC, id ASC"
).fetchall()
return [dict(r) for r in rows]
finally:
conn.close()
def mission_state(self, mission_id: str) -> Dict[str, Any]:
conn = _connect(self.db_path)
try:
m = self._mission_row(conn, mission_id)
if m is None:
raise _err("UNKNOWN_MISSION", f"unknown mission {mission_id!r}", mission_id=mission_id)
return {"id": m["id"], "state": m["state"]}
finally:
conn.close()
# -- Mission: transition --------------------------------------------------
def mission_transition(
self,
mission_id: str,
target: str,
actor: str = "red-queen",
evidence: Optional[str] = None,
) -> Dict[str, Any]:
if not is_mission_state(target):
raise _err("UNKNOWN_STATE", f"unknown mission state {target!r}", target=target)
conn = _connect(self.db_path)
try:
m = self._mission_row(conn, mission_id)
if m is None:
raise _err("UNKNOWN_MISSION", f"unknown mission {mission_id!r}", mission_id=mission_id)
current = m["state"]
if not is_mission_state(current):
raise _err(
"STATE_ERROR",
f"corrupt mission state {current!r}",
mission_id=mission_id, current=current,
)
if current == target:
return {
"mission_id": mission_id,
"state": target,
"transition": f"{current}->{target}",
"idempotent": True,
}
if not mission_transition_allowed(current, target):
raise _err(
"INVALID_TRANSITION",
f"forbidden mission transition {current} -> {target}",
mission_id=mission_id, current=current, target=target,
)
now = _utcnow()
conn.execute("UPDATE missions SET state=?, updated_at=? WHERE id=?", (target, now, mission_id))
conn.execute(
"INSERT INTO mission_events (mission_id,transition,actor,evidence,created_at) "
"VALUES (?,?,?,?,?)",
(mission_id, f"{current}->{target}", actor, evidence, now),
)
conn.commit()
return {
"mission_id": mission_id,
"state": target,
"transition": f"{current}->{target}",
"idempotent": False,
}
finally:
conn.close()
def mission_pause(self, mission_id, actor="red-queen", evidence=None):
return self.mission_transition(mission_id, "PAUSED", actor, evidence)
def mission_block(self, mission_id, actor="red-queen", evidence=None):
return self.mission_transition(mission_id, "BLOCKED", actor, evidence)
def mission_resume(self, mission_id, actor="red-queen", evidence=None):
return self.mission_transition(mission_id, "RUNNING", actor, evidence)
def mission_complete(self, mission_id, actor="red-queen", evidence=None):
"""REVIEW -> COMPLETED. Rejected, solange ein Pflicht-WP nicht DONE ist."""
conn = _connect(self.db_path)
try:
m = self._mission_row(conn, mission_id)
if m is None:
raise _err("UNKNOWN_MISSION", f"unknown mission {mission_id!r}", mission_id=mission_id)
if m["state"] == "COMPLETED":
# idempotent: erneutes complete ist ein No-op
return {
"mission_id": mission_id,
"state": "COMPLETED",
"transition": "COMPLETED->COMPLETED",
"idempotent": True,
}
if m["state"] != "REVIEW":
raise _err(
"INVALID_TRANSITION",
f"mission_complete requires REVIEW, got {m['state']}",
mission_id=mission_id, current=m["state"],
)
open_wps = conn.execute(
"SELECT id,state FROM work_packages WHERE mission_id=? AND state <> 'DONE' ORDER BY id",
(mission_id,),
).fetchall()
if open_wps:
raise _err(
"MISSION_INCOMPLETE_WPS",
"mission has open (non-DONE) work packages",
mission_id=mission_id,
open_wps=[{"id": r["id"], "state": r["state"]} for r in open_wps],
)
return self.mission_transition(mission_id, "COMPLETED", actor, evidence)
finally:
conn.close()
# -- WP: create ------------------------------------------------------------
def wp_create(self, mission_id: str, title: str, wp_id: Optional[str] = None, actor: str = "red-queen"):
if not title or not title.strip():
raise _err("INVALID_ARGS", "wp title required")
conn = _connect(self.db_path)
try:
if not self._mission_exists(conn, mission_id):
raise _err("UNKNOWN_MISSION", f"unknown mission {mission_id!r}", mission_id=mission_id)
if wp_id:
if self._wp_exists(conn, wp_id):
return self._wp_read(conn, wp_id) # idempotent
wid = wp_id
else:
wid = self._new_wp_id(conn, mission_id)
now = _utcnow()
conn.execute(
"INSERT INTO work_packages (id,mission_id,title,state,created_at,updated_at) "
"VALUES (?,?,?,?,?,?)",
(wid, mission_id, title.strip(), "TODO", now, now),
)
conn.execute(
"INSERT INTO wp_events (wp_id,mission_id,transition,actor,evidence,created_at) "
"VALUES (?,?,?,?,?,?)",
(wid, mission_id, "TODO", actor, "wp created", now),
)
conn.commit()
# Kanban-Mirror (best effort)
if self.mirror.enabled:
ktask = self.mirror.create(title.strip(), mission_id, wid)
if ktask:
conn.execute("UPDATE work_packages SET kanban_task_id=? WHERE id=?", (ktask, wid))
conn.commit()
return self._wp_read(conn, wid)
finally:
conn.close()
def _wp_read(self, conn, wp_id: str) -> Dict[str, Any]:
w = self._wp_row(conn, wp_id)
if w is None:
raise _err("UNKNOWN_WP", f"unknown wp {wp_id!r}", wp_id=wp_id)
return {
"id": w["id"],
"mission_id": w["mission_id"],
"title": w["title"],
"state": w["state"],
"kanban_task_id": w["kanban_task_id"],
"created_at": w["created_at"],
"updated_at": w["updated_at"],
}
def wp_read(self, wp_id: str) -> Dict[str, Any]:
conn = _connect(self.db_path)
try:
return self._wp_read(conn, wp_id)
finally:
conn.close()
def wp_state(self, wp_id: str) -> Dict[str, Any]:
conn = _connect(self.db_path)
try:
w = self._wp_row(conn, wp_id)
if w is None:
raise _err("UNKNOWN_WP", f"unknown wp {wp_id!r}", wp_id=wp_id)
return {"id": w["id"], "state": w["state"]}
finally:
conn.close()
# -- WP: transition ---------------------------------------------------------
def wp_transition(self, wp_id, target, actor="red-queen", evidence=None):
if not is_wp_state(target):
raise _err("UNKNOWN_STATE", f"unknown wp state {target!r}", target=target)
conn = _connect(self.db_path)
try:
w = self._wp_row(conn, wp_id)
if w is None:
raise _err("UNKNOWN_WP", f"unknown wp {wp_id!r}", wp_id=wp_id)
current = w["state"]
if not is_wp_state(current):
raise _err(
"STATE_ERROR",
f"corrupt wp state {current!r}",
wp_id=wp_id, current=current,
)
if current == target:
return {"id": wp_id, "state": target, "transition": f"{current}->{target}", "idempotent": True}
if not wp_transition_allowed(current, target):
raise _err(
"INVALID_TRANSITION",
f"forbidden wp transition {current} -> {target}",
wp_id=wp_id, current=current, target=target,
)
if target == "DONE" and current != "CHECKING":
raise _err(
"INVALID_TRANSITION",
f"wp can only reach DONE via CHECKING (current {current})",
wp_id=wp_id, current=current, target="DONE",
)
if target == "DONE":
open_deps = self._open_deps(conn, w["mission_id"], wp_id)
if open_deps:
raise _err(
"WP_OPEN_DEPENDENCY",
"wp has open dependency",
wp_id=wp_id, dependencies=open_deps,
)
now = _utcnow()
conn.execute("UPDATE work_packages SET state=?, updated_at=? WHERE id=?", (target, now, wp_id))
conn.execute(
"INSERT INTO wp_events (wp_id,mission_id,transition,actor,evidence,created_at) "
"VALUES (?,?,?,?,?,?)",
(wp_id, w["mission_id"], f"{current}->{target}", actor, evidence, now),
)
conn.commit()
self._mirror_sync(w["mission_id"], wp_id, current, target, evidence)
return {"id": wp_id, "state": target, "transition": f"{current}->{target}", "idempotent": False}
finally:
conn.close()
def wp_block(self, wp_id, reason="blocked", actor="red-queen"):
return self.wp_transition(wp_id, "BLOCKED", actor, reason)
def wp_complete(self, wp_id, result=None, actor="red-queen"):
return self.wp_transition(wp_id, "DONE", actor, result)
# -- WP: dependency ----------------------------------------------------------
def wp_dependency(self, mission_id: str, wp_id: str, depends_on: str) -> Dict[str, Any]:
conn = _connect(self.db_path)
try:
if not self._wp_exists(conn, wp_id):
raise _err("UNKNOWN_WP", f"unknown wp {wp_id!r}", wp_id=wp_id)
if not self._wp_exists(conn, depends_on):
raise _err("UNKNOWN_WP", f"unknown dependency wp {depends_on!r}", wp_id=depends_on)
w = self._wp_row(conn, wp_id)
d = self._wp_row(conn, depends_on)
if w["mission_id"] != mission_id or d["mission_id"] != mission_id:
raise _err(
"INVALID_DEPENDENCY",
"wp and dependency must belong to same mission",
mission_id=mission_id, wp_id=wp_id, depends_on=depends_on,
)
if wp_id == depends_on:
raise _err("INVALID_DEPENDENCY", "wp cannot depend on itself", wp_id=wp_id)
conn.execute(
"INSERT OR IGNORE INTO wp_dependencies (mission_id,wp_id,depends_on) VALUES (?,?,?)",
(mission_id, wp_id, depends_on),
)
conn.commit()
if self.mirror.enabled:
wt = self._wp_row(conn, wp_id)
dt = self._wp_row(conn, depends_on)
if wt["kanban_task_id"] and dt["kanban_task_id"]:
self.mirror.link(dt["kanban_task_id"], wt["kanban_task_id"])
return {"mission_id": mission_id, "wp_id": wp_id, "depends_on": depends_on}
finally:
conn.close()
# -- interne Helfer -----------------------------------------------------------
def _open_deps(self, conn, mission_id: str, wp_id: str) -> List[Dict[str, Any]]:
rows = conn.execute(
"SELECT wd.depends_on, w.state FROM wp_dependencies wd "
"JOIN work_packages w ON w.id = wd.depends_on "
"WHERE wd.mission_id=? AND wd.wp_id=? AND w.state <> 'DONE'",
(mission_id, wp_id),
).fetchall()
return [{"id": r["depends_on"], "state": r["state"]} for r in rows]
def _mirror_sync(self, mission_id: str, wp_id: str, current: str, target: str, evidence: Optional[str]):
if not self.mirror.enabled:
return
conn = _connect(self.db_path)
try:
w = self._wp_row(conn, wp_id)
if w is None or not w["kanban_task_id"]:
return
k = w["kanban_task_id"]
if target == "BLOCKED":
self.mirror.block(k, evidence or "blocked by state machine")
elif current == "BLOCKED" and target != "BLOCKED":
self.mirror.unblock(k)
elif target == "DONE":
self.mirror.complete(k, evidence or "done")
finally:
conn.close()