import json import os import re import shutil import uuid from datetime import UTC, datetime from pathlib import Path import yake from tht.config import DatabaseConfig from tht.session.models import SessionManifest MANIFEST = "session_manifest.yaml" MAX_SLUG_CHARS = 40 MAX_SUMMARY_CHARS = 120 MAX_NAME_CHARS = 60 MAX_NAME_WORDS = 5 # Generic query verbs/words dropped from a derived session name (lowercased compare); # they carry no content and would just crowd out the salient keywords. _STOPNAME = { "crea", "creare", "elenca", "elencare", "elencami", "mostra", "mostrami", "mostrare", "dammi", "dai", "trova", "trovami", "cerca", "cercami", "voglio", "vorrei", "fammi", "quanti", "quante", "quanto", "quanta", "conta", "numero", "lista", "elenco", "vedi", "visualizza", "ottieni", "restituisci", "calcola", "esponi", "riporta", "fornisci", } class SessionError(Exception): pass def current_author() -> str: """Autore della sessione (D6). Nel profilo `none` dell'MVP l'operatore e' sulla propria macchina: default `dev@local`. Con auth reale il backend passa l'id via THT_AUTHOR. Una sola fonte per author/updated_by.""" return os.environ.get("THT_AUTHOR", "dev@local").strip() or "dev@local" def _summarize(question: str) -> str: """Domanda sintetica per il manifest/la sidebar (D6 §5.2): prima riga, troncata.""" first = question.strip().splitlines()[0].strip() if question.strip() else "" return first[:MAX_SUMMARY_CHARS] def _extract_name(question: str) -> str: """Nome-sintesi (3-5 parole-chiave) della domanda per la sidebar, senza LLM (YAKE). Le parole-chiave restano nella lingua della domanda (è contenuto, non chrome UI). Ripiega su `_summarize` se YAKE non è disponibile o non estrae nulla di utile.""" q = (question or "").strip() if not q: return "" try: extractor = yake.KeywordExtractor(lan="it", n=1, top=8, dedupLim=0.9) ranked = [k for k, _ in extractor.extract_keywords(q)] except Exception: # noqa: BLE001 - keyword extraction has a deterministic fallback return _summarize(question) seen: set[str] = set() picked: list[str] = [] for k in ranked: kl = k.lower() if kl in _STOPNAME or kl in seen: continue seen.add(kl) picked.append(k) if len(picked) >= MAX_NAME_WORDS: break if not picked: return _summarize(question) ql = q.lower() picked.sort(key=lambda k: ql.find(k.lower())) # natural reading order name = " ".join(picked) return (name[0].upper() + name[1:])[:MAX_NAME_CHARS] def render_question_md(question: str, assumptions: list[str] | None = None) -> str: """Rende question.md in modo deterministico: domanda + assunzioni opzionali. Unica fonte di formattazione per question.md (riusata da create_session e set_question), così il gate non deve costruire markdown a mano. """ body = f"# Domanda\n\n{question.strip()}\n" items = [a.strip() for a in (assumptions or []) if a.strip()] if items: body += "\n## Assunzioni\n\n" + "".join(f"- {a}\n" for a in items) return body def _new_id(question: str, sessions_root: Path, stamp: str) -> str: """Create an opaque UUIDv4 session identifier. ``question``, ``sessions_root`` and ``stamp`` remain accepted temporarily so existing workflow callers do not need to change as the repository boundary is introduced. """ return str(uuid.uuid4()) def create_session( question: str, db: DatabaseConfig, sessions_root: Path, *, author: str | None = None, summary: str | None = None, provider: str | None = None, model: str | None = None, thinking: str | None = None, workspace_id: str | None = None, workspace_revision: str | None = None, name: str | None = None, ) -> SessionManifest: now = datetime.now(UTC) # stamp con ora/min/sec: identifica univocamente sessioni dello stesso giorno # sulla stessa domanda. Il contatore -n resta come rete per collisioni nello # stesso secondo. session_id = _new_id(question, sessions_root, now.strftime("%Y-%m-%d-%H%M%S")) who = author or current_author() # schema_version del workflow usato (D6/§5.3): consente di interpretare il ledger # secondo la versione anche se il workflow evolve. try: from tht.workflow import load_workflow schema_version = load_workflow().schema_version except Exception: # noqa: BLE001 - legacy sessions may predate workflow metadata schema_version = None manifest = SessionManifest( id=session_id, created_at=now, question=question, database=db.database, schema=db.db_schema, author=who, summary=summary or _summarize(question), updated_at=now, updated_by=who, schema_version=schema_version, provider=provider, model=model, thinking=thinking, workspace_id=workspace_id, workspace_revision=workspace_revision, name=name, ) session_dir = sessions_root / session_id manifest.to_yaml(session_dir / MANIFEST) (session_dir / "question.md").write_text(render_question_md(question)) return manifest def new_session_manifest( question: str, db: DatabaseConfig, *, provider=None, model=None, thinking=None, workspace_id=None, workspace_revision=None, name=None, interaction_language=None, ) -> SessionManifest: """Create an unsaved UUIDv4 manifest for a repository-owned session.""" now = datetime.now(UTC) return SessionManifest( id=str(uuid.uuid4()), created_at=now, question=question, database=db.database, schema=db.db_schema, author=current_author(), summary=_summarize(question), updated_at=now, updated_by=current_author(), provider=provider, model=model, thinking=thinking, workspace_id=workspace_id, workspace_revision=workspace_revision, name=name, interaction_language=interaction_language, ) def pin_interaction_language(manifest: SessionManifest, workspace_language: str) -> bool: """Apply legacy compatibility once, while the repository holds its writer lock.""" from tht.session.language import question_interaction_language if manifest.interaction_language is not None: return False if manifest.status == "finalized" or manifest.archived: raise SessionError("Session is read-only (finalized or archived)") manifest.interaction_language = question_interaction_language(manifest.question, workspace_language) manifest.updated_at = datetime.now(UTC) manifest.updated_by = current_author() return True _ASSUMPTIONS_HEADING = re.compile( r"^#{1,6}\s+(?:assunzioni|assumptions)\s*$", re.IGNORECASE | re.MULTILINE ) _PREVIEW_HEADING = re.compile( r"^##\s+preview(?:\s*/\s*aggregato)?\s*$", re.IGNORECASE | re.MULTILINE ) _NEXT_H2 = re.compile(r"^##\s+", re.MULTILINE) _MEMORY_TYPES = frozenset( {"concept_clarified", "memory_promoted", "memory_promotion_declined", "memory_rejected"} ) _TABLE_DECISION_TYPES = frozenset({"table_approved", "table_promoted", "table_excluded"}) _SUMMARY_HIDDEN_TYPES = _MEMORY_TYPES | frozenset( { "phase_approved", "phase_auto_approved", "table_approved", "table_promoted", "column_promoted", "cte_approved", "sql_approved", "datamart_requested", "datamart_declined", } ) def _split_question_sections(content: str) -> tuple[str, str]: """Split legacy combined question Markdown without rewriting the artifact.""" match = _ASSUMPTIONS_HEADING.search(content) if match is None: return content.strip(), "" return content[:match.start()].strip(), content[match.start():].strip() def _split_validation_preview(content: str) -> tuple[str, str]: """Return the persisted preview section and the remaining validation report.""" match = _PREVIEW_HEADING.search(content) if match is None: return "", content.strip() body_start = match.end() next_heading = _NEXT_H2.search(content, body_start) section_end = next_heading.start() if next_heading else len(content) preview = content[body_start:section_end].strip() remaining = (content[:match.start()] + content[section_end:]).strip() return preview, remaining def _referenced_decision(marker, by_seq: dict[int, object]): match = re.fullmatch(r"seq:(\d+)", marker.detail.strip()) return by_seq.get(int(match.group(1))) if match else None def _normalized_identifier(value: str) -> str: return value.strip().strip("`\"'").lower() def _schema_linking_tables(content: str | None, decisions: list) -> set[str]: names = { _normalized_identifier(decision.subject) for decision in decisions if decision.type in _TABLE_DECISION_TYPES and decision.subject.strip() } if not content: return names try: schema_linking = json.loads(content) except (TypeError, json.JSONDecodeError): return names if not isinstance(schema_linking, dict): return names for key in ("candidates", "excluded"): values = schema_linking.get(key, []) if not isinstance(values, list): continue for value in values: if not isinstance(value, dict) or value.get("kind") != "table": continue name = value.get("name") if isinstance(name, str) and name.strip(): names.add(_normalized_identifier(name)) return names def _is_standalone_table_subject(subject: str, schema_tables: set[str]) -> bool: normalized = _normalized_identifier(subject) leaf = normalized.rsplit(".", 1)[-1] schema_leaves = {name.rsplit(".", 1)[-1] for name in schema_tables} return ( normalized in schema_tables or leaf in schema_leaves or re.fullmatch(r"(?:fact|dim)(?:_[a-z0-9_]+)?", leaf) is not None ) def _memory_projection(decisions: list, schema_tables: set[str]) -> list[dict[str, str]]: by_seq = {decision.seq: decision for decision in decisions} approved = [decision for decision in decisions if decision.type == "memory_promoted"] declined = [ decision for decision in decisions if decision.type in {"memory_promotion_declined", "memory_rejected"} ] items: list[dict[str, str]] = [] for status, markers in (("approved", approved), ("declined", declined)): for marker in markers: origin = _referenced_decision(marker, by_seq) if origin is not None and origin.type != "concept_clarified": continue subject = origin.subject if origin is not None else marker.subject if _is_standalone_table_subject(subject, schema_tables): continue detail = origin.detail if origin is not None else marker.detail if re.fullmatch(r"seq:\d+", detail.strip()): detail = "" rationale = ( origin.rationale if origin is not None and origin.rationale else marker.rationale ) if rationale.strip() == detail.strip(): rationale = "" items.append( { "status": status, "subject": subject, "detail": detail, "rationale": rationale, } ) return items def _document_bundle(manifest: SessionManifest, artifacts: dict[str, str], decisions: list) -> list[dict]: docs: list[dict] = [{ "phase": "—", "key": "question", "title": "Original question", "format": "text", "content": manifest.question, }] sql = artifacts.get("sql_final") if sql is not None: docs.append({ "phase": "F7", "key": "sql", "title": "Final SQL", "format": "sql", "content": sql, }) validation = artifacts.get("validation_report") preview, validation_without_preview = _split_validation_preview(validation or "") if preview: docs.append({ "phase": "finalize", "key": "preview", "title": "Data preview", "format": "markdown", "content": preview, }) revised, assumptions = _split_question_sections(artifacts.get("question", "")) if revised: docs.append({ "phase": "F3", "key": "revised_question", "title": "Revised question", "format": "markdown", "content": revised, }) if assumptions: docs.append({ "phase": "F3", "key": "assumptions", "title": "Assumptions", "format": "markdown", "content": assumptions, }) schema_linking = artifacts.get("schema_linking") memories = _memory_projection( decisions, _schema_linking_tables(schema_linking, decisions), ) if memories: docs.append({ "phase": "F8", "key": "memories", "title": "Memories", "format": "memories", "content": json.dumps(memories, ensure_ascii=False), }) if schema_linking is not None: docs.append({ "phase": "F4", "key": "schema_linking", "title": "Schema linking", "format": "schema-linking", "content": schema_linking, }) if validation_without_preview: docs.append({ "phase": "finalize", "key": "validation_report", "title": "Validation report", "format": "markdown", "content": validation_without_preview, }) visible_decisions = [d for d in decisions if d.type not in _SUMMARY_HIDDEN_TYPES] if visible_decisions: docs.append({ "phase": "—", "key": "decisions", "title": "Decisions", "format": "decisions", "content": "\n".join(d.model_dump_json() for d in visible_decisions) + "\n", }) return docs def build_snapshot_documents(snapshot) -> list[dict]: from tht.phase import effective_decisions return _document_bundle( snapshot.manifest, snapshot.artifacts, effective_decisions(snapshot), ) def touch_manifest( session_id: str, sessions_root: Path, *, updated_by: str | None = None ) -> SessionManifest: """Aggiorna updated_at/updated_by del manifest (D6 §5.2). Da chiamare ad ogni mutazione della sessione (set-question, finalize, close).""" manifest = load_session(session_id, sessions_root) manifest.updated_at = datetime.now(UTC) manifest.updated_by = updated_by or current_author() manifest.to_yaml(sessions_root / session_id / MANIFEST) return manifest def set_question( session_id: str, question: str, assumptions: list[str], sessions_root: Path, ) -> Path: """Riscrive question.md (Fase 3) in modo deterministico, senza edit tool. Valida l'esistenza della sessione (SessionError se assente) e ritorna il path scritto. Non tocca il manifest né le decisioni. """ load_session(session_id, sessions_root) path = sessions_root / session_id / "question.md" path.write_text(render_question_md(question, assumptions)) touch_manifest(session_id, sessions_root) return path def set_schema_linking( session_id: str, data: dict, sessions_root: Path, ) -> Path: """Valida e scrive schema_linking.json (Fase 4) in modo deterministico. Valida `data` contro il modello SchemaLinking (ValidationError se invalido) PRIMA di scrivere, così un artefatto malformato non tocca mai il disco. Ritorna il path. """ from tht.session.models import SchemaLinking load_session(session_id, sessions_root) model = SchemaLinking.model_validate(data) path = sessions_root / session_id / "schema_linking.json" path.write_text( json.dumps(model.model_dump(by_alias=True), indent=2, ensure_ascii=False) ) touch_manifest(session_id, sessions_root) return path def sync_schema_linking(session_id: str, sessions_root: Path) -> Path: """Project the effective F4 ledger decisions into schema_linking.json. candidates/excluded are rebuilt from table_promoted/table_excluded + column_promoted/column_excluded (last decision per subject wins). question, joins, concept_formulas and open_questions are preserved from the existing file when present. The reviewer's curation is thus authoritative and deterministic (no model transcription).""" from tht.phase import effective_decisions session_dir = sessions_root / session_id manifest = load_session(session_id, sessions_root) existing: dict = {} sl_path = session_dir / "schema_linking.json" if sl_path.exists(): existing = json.loads(sl_path.read_text()) # last decision per subject wins (handles a re-run of the gate). latest: dict[str, str] = {} for d in effective_decisions(session_dir): if d.type in ("table_promoted", "table_excluded", "column_promoted", "column_excluded"): latest[d.subject] = d.type candidates: list[dict] = [] excluded: list[dict] = [] for subject, dtype in latest.items(): is_column = "." in subject kind = "column" if is_column else "table" if dtype in ("table_promoted", "column_promoted"): candidates.append({"kind": kind, "name": subject, "decision": "promoted"}) else: excluded.append({"kind": kind, "name": subject}) data = { "question": existing.get("question") or manifest.question, "candidates": candidates, "joins": existing.get("joins", []), "excluded": excluded, "open_questions": existing.get("open_questions", []), "concept_formulas": existing.get("concept_formulas", []), } return set_schema_linking(session_id, data, sessions_root) def sync_schema_linking_snapshot(snapshot) -> str: """Repository projection of effective F4 decisions into schema_linking JSON.""" from tht.phase import effective_decisions existing = json.loads(snapshot.artifacts.get("schema_linking", "{}")) latest: dict[str, str] = {} for decision in effective_decisions(snapshot): if decision.type in {"table_promoted", "table_excluded", "column_promoted", "column_excluded"}: latest[decision.subject] = decision.type candidates, excluded = [], [] for subject, kind in latest.items(): entity = "column" if "." in subject else "table" if kind.endswith("promoted"): candidates.append({"kind": entity, "name": subject, "decision": "promoted"}) else: excluded.append({"kind": entity, "name": subject}) from tht.session.models import SchemaLinking model = SchemaLinking.model_validate({ "question": existing.get("question") or snapshot.manifest.question, "candidates": candidates, "excluded": excluded, "joins": existing.get("joins", []), "open_questions": existing.get("open_questions", []), "concept_formulas": existing.get("concept_formulas", []), }) return json.dumps(model.model_dump(by_alias=True), indent=2, ensure_ascii=False) def load_session(session_id: str, sessions_root: Path) -> SessionManifest: path = sessions_root / session_id / MANIFEST if not path.exists(): raise SessionError(f"Sessione non trovata: {session_id} (atteso {path})") return SessionManifest.from_yaml(path) def close_session(session_id: str, sessions_root: Path) -> SessionManifest: manifest = load_session(session_id, sessions_root) manifest.status = "closed" manifest.updated_at = datetime.now(UTC) manifest.updated_by = current_author() manifest.to_yaml(sessions_root / session_id / MANIFEST) return manifest def fail_session(session_id: str, sessions_root: Path) -> SessionManifest: """Record a fatal managed-runtime failure without losing phase artifacts.""" manifest = load_session(session_id, sessions_root) manifest.status = "failed" manifest.updated_at = datetime.now(UTC) manifest.updated_by = current_author() manifest.to_yaml(sessions_root / session_id / MANIFEST) return manifest def reopen_session(session_id: str, sessions_root: Path) -> SessionManifest: """Mark a manually resumed session active again.""" manifest = load_session(session_id, sessions_root) manifest.status = "open" manifest.updated_at = datetime.now(UTC) manifest.updated_by = current_author() manifest.to_yaml(sessions_root / session_id / MANIFEST) return manifest def _save_touched(manifest: SessionManifest, sessions_root: Path) -> SessionManifest: """Persist `manifest` updating updated_at/updated_by (single save path for mutations).""" manifest.updated_at = datetime.now(UTC) manifest.updated_by = current_author() manifest.to_yaml(sessions_root / manifest.id / MANIFEST) return manifest def set_name(session_id: str, name: str | None, sessions_root: Path) -> SessionManifest: """Set the descriptive name (empty/blank clears it back to None).""" manifest = load_session(session_id, sessions_root) manifest.name = (name or "").strip() or None return _save_touched(manifest, sessions_root) def set_group(session_id: str, group: str | None, sessions_root: Path) -> SessionManifest: """Set the group (empty/blank clears it back to None).""" manifest = load_session(session_id, sessions_root) manifest.group = (group or "").strip() or None return _save_touched(manifest, sessions_root) def set_archived(session_id: str, archived: bool, sessions_root: Path) -> SessionManifest: """Flip the archived flag. Archiving does NOT change resumability (a finalized session stays read-only); unarchive only moves it back to the active list.""" manifest = load_session(session_id, sessions_root) manifest.archived = archived return _save_touched(manifest, sessions_root) def delete_session(session_id: str, sessions_root: Path) -> None: """Hard-delete the session directory. SessionError if it does not exist.""" load_session(session_id, sessions_root) # raises SessionError if absent shutil.rmtree(sessions_root / session_id) def persist_verified_finalization( repository, session_id: str, *, validation_report: str, evidence: str, ) -> SessionManifest: """Publish DWH-verified final artifacts and status at one repository boundary. The caller must complete static validation, EXPLAIN and preview before this function is entered. It intentionally does not touch solved-question indexing: that derivative is best-effort and happens after the durable commit. """ snapshot = repository.get(session_id) manifest = snapshot.manifest.model_copy(deep=True) manifest.status = "finalized" manifest.updated_at = datetime.now(UTC) manifest.updated_by = current_author() return repository.finalize( manifest, {"validation_report": validation_report, "evidence": evidence}, ).manifest def build_documents(manifest: SessionManifest, session_dir: Path) -> list[dict]: """Build the canonical read-only bundle for legacy filesystem sessions.""" from tht.phase import effective_decisions artifact_files = { "question": "question.md", "schema_linking": "schema_linking.json", "sql_final": "sql_final.sql", "validation_report": "validation_report.md", } artifacts = { key: path.read_text() for key, filename in artifact_files.items() if (path := session_dir / filename).exists() } return _document_bundle(manifest, artifacts, effective_decisions(session_dir))