- 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)
638 lines
27 KiB
Python
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()
|