Files
ThothII/harness/tht/decisions.py

177 lines
5.9 KiB
Python

import fcntl
import os
import tempfile
from datetime import UTC, datetime
from pathlib import Path
from typing import Literal
from pydantic import BaseModel
DECISIONS_FILE = "review_decisions.jsonl"
DECISIONS_LOCK_FILE = ".review_decisions.lock"
# 22 tipi di the reference implementation (verified leggendo session/decisions.py) + 1 nuovo (D15):
# `decision_retracted` per il rollback a granularità step (ritira una decisione
# senza cancellarne la riga dal log di audit; effective_decisions la onora).
DecisionType = Literal[
"concept_clarified",
"question_rewritten",
"table_promoted",
"table_excluded",
"column_promoted",
"column_excluded",
"column_corrected",
"join_modified",
"evidence_accepted",
"evidence_rejected",
"ambiguity_open",
"memory_rejected",
"cte_approved",
"cte_corrected",
"cte_rejected",
"sql_revised",
"sql_approved",
"sql_rejected",
"phase_approved",
"phase_auto_approved",
"phase_reopened",
"phase_skipped",
"datamart_requested",
"datamart_declined",
# F8: promozione memorie riusabili al gate reviewer_memory_promote. subject =
# subject della decisione originale, detail = "seq:<decision_seq>" (usato da
# declined_promotion_seqs per non riproporre i candidati rifiutati).
"memory_promoted",
"memory_promotion_declined",
# D15: marker di ritrazione. subject = "phase:N", retracts = decision_seq ritirata.
# Resta nel log di audit (append-only); effective_decisions() la esclude dalla vista.
"decision_retracted",
# D14a: valore citato nella domanda ancorato a una o piu' colonne. subject =
# "phase:4", detail = il valore (es. "ablazione"), rationale = la/e colonna/e scelta/e
# dal reviewer (aggregate_lsh_multi le espone tutte senza collassare al miglior match).
"value_grounded",
# D14b: formula di concetto approvata/rifiutata dal reviewer. subject = "phase:4",
# detail = il concetto (es. "fascia pediatrica"), rationale = la/e colonna/e o il motivo.
# La ricerca nelle Formula Evidence pubblicate restituisce i candidati; queste decisioni
# registrano la scelta della proposta nella sessione.
"concept_formula_approved",
"concept_formula_rejected",
]
class DecisionRecord(BaseModel):
seq: int
ts: datetime
type: DecisionType
subject: str
detail: str = ""
rationale: str = ""
# D15: se type == "decision_retracted", indica quale seq viene ritirata.
retracts: int | None = None
# D15: la fase corrente al momento della scrittura (high-water-mark). Permette a
# effective_decisions() di marcare stale le decisioni il cui subject NON e' "phase:N"
# (es. cte_approved usa il nome del CTE) dopo un rollback. None per record storici
# (pre-fix) o costruiti a mano: in quel caso si ricade sulla logica basata sul subject.
phase: int | None = None
class DecisionInput(BaseModel):
type: DecisionType
subject: str
detail: str = ""
rationale: str = ""
retracts: int | None = None
def list_decisions(session_dir: Path) -> list[DecisionRecord]:
path = session_dir / DECISIONS_FILE
if not path.exists():
return []
return [
DecisionRecord.model_validate_json(line)
for line in path.read_text().splitlines()
if line.strip()
]
def append_decision(
session_dir: Path,
*,
type: str,
subject: str,
detail: str = "",
rationale: str = "",
retracts: int | None = None,
) -> DecisionRecord:
return append_decisions(
session_dir,
[{
"type": type,
"subject": subject,
"detail": detail,
"rationale": rationale,
"retracts": retracts,
}],
)[0]
def append_decisions(
session_dir: Path,
decisions: list[DecisionInput | dict],
) -> list[DecisionRecord]:
"""Validate and append a decision set through one atomic file replacement."""
inputs = [DecisionInput.model_validate(decision) for decision in decisions]
if not inputs:
return []
session_dir.mkdir(parents=True, exist_ok=True)
lock_path = session_dir / DECISIONS_LOCK_FILE
with lock_path.open("a+") as lock:
fcntl.flock(lock, fcntl.LOCK_EX)
try:
return _append_decisions_locked(session_dir, inputs)
finally:
fcntl.flock(lock, fcntl.LOCK_UN)
def _append_decisions_locked(
session_dir: Path,
inputs: list[DecisionInput],
) -> list[DecisionRecord]:
"""Append while the caller holds the session's cross-process ledger lock."""
# Fase corrente PRIMA dell'append (high-water-mark D15). Import lazy: phase.py
# importa decisions.py (ciclo). Per i marker di fase (subject "phase:N") il valore
# e' ridondante col subject; per le decisioni sostanziali con subject "a nome"
# (cte_approved, ...) e' l'unico modo per filtrarle dopo un reopen.
from tht.phase import current_phase
existing = list_decisions(session_dir)
phase = current_phase(session_dir)
records = [
DecisionRecord(
seq=len(existing) + index,
ts=datetime.now(UTC),
phase=phase,
**decision.model_dump(),
)
for index, decision in enumerate(inputs, start=1)
]
path = session_dir / DECISIONS_FILE
path.parent.mkdir(parents=True, exist_ok=True)
previous = path.read_text() if path.exists() else ""
separator = "" if not previous or previous.endswith("\n") else "\n"
content = previous + separator + "".join(record.model_dump_json() + "\n" for record in records)
fd, temporary_name = tempfile.mkstemp(prefix=f".{DECISIONS_FILE}.", dir=path.parent)
temporary = Path(temporary_name)
try:
with os.fdopen(fd, "w") as handle:
handle.write(content)
handle.flush()
os.fsync(handle.fileno())
os.replace(temporary, path)
finally:
temporary.unlink(missing_ok=True)
return records