176 lines
5.8 KiB
Python
176 lines
5.8 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.
|
|
# retrieve_formula restituisce i candidati; queste decisioni registrano la scelta.
|
|
"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
|