From c1cddaa667a2b7b1d826fd65e0a3ab4fd2222f66 Mon Sep 17 00:00:00 2001 From: User Date: Thu, 16 Jul 2026 17:55:34 +0200 Subject: [PATCH] refactor(harness): route workflow persistence through repositories --- .superpowers/sdd/task-3-report.md | 101 ++++----- .../session-repository-writes.test.js | 41 ++++ harness/.pi/extensions/tht-gate.js | 32 +++ harness/.pi/skills/tht-sessione/SKILL.md | 7 +- .../tests/test_adapter_command_regressions.py | 11 +- harness/tests/test_decision_join_set_cli.py | 21 +- harness/tests/test_decision_retract_cli.py | 19 +- harness/tests/test_promoted_columns_for.py | 6 +- .../tests/test_repository_memory_sql_paths.py | 16 ++ harness/tests/test_require_phase_or_exit.py | 9 +- harness/tests/test_session_repository.py | 54 ++++- .../tests/test_session_repository_workflow.py | 60 +++++ harness/tests/test_sql_preview_json.py | 8 +- harness/tht/cli/cte_cmd.py | 83 ++++--- harness/tht/cli/decision_cmd.py | 38 ++-- harness/tht/cli/memory_cmd.py | 38 ++-- harness/tht/cli/phase_cmd.py | 51 +++-- harness/tht/cli/search_cmd.py | 7 +- harness/tht/cli/session_cmd.py | 206 ++++++++++-------- harness/tht/cli/sql_cmd.py | 58 +++-- harness/tht/ctetest.py | 10 +- harness/tht/memory.py | 20 ++ harness/tht/phase.py | 85 ++++---- harness/tht/session/filesystem_repository.py | 55 ++++- harness/tht/session/postgres_repository.py | 50 +++++ harness/tht/session/repository.py | 42 +++- harness/tht/session/store.py | 80 +++++++ harness/tht/solved.py | 17 ++ harness/tht/taskdoc.py | 20 +- harness/tht/teardown.py | 28 +++ 30 files changed, 925 insertions(+), 348 deletions(-) create mode 100644 harness/.pi/extensions/gate/__tests__/session-repository-writes.test.js create mode 100644 harness/tests/test_repository_memory_sql_paths.py create mode 100644 harness/tests/test_session_repository_workflow.py diff --git a/.superpowers/sdd/task-3-report.md b/.superpowers/sdd/task-3-report.md index 0228878d..5fb7ae6d 100644 --- a/.superpowers/sdd/task-3-report.md +++ b/.superpowers/sdd/task-3-report.md @@ -1,65 +1,56 @@ -# Task 3 report — reconnect SSE on same-session Resume +# Task 3 report — workflow repository migration -## Status +## RED -Complete. A successful Resume of the currently active session now replaces its existing -`EventSource` connection. Resuming a different session continues to reconnect through the -session ID change only, without a generation-driven second connection. +- `harness/tests/test_session_repository_workflow.py` initially failed at collection: + `persist_verified_finalization` did not exist. +- The new gate test initially failed because `write_cte_sql` and `write_final_sql` + were not registered. Its first run also exposed the worktree-local missing + Node dependency (`typebox`); `npm ci` installed the lockfile dependency. +- After the principal/legacy policy was clarified, the resolver tests initially + failed because `resolve_principal` did not exist. -## Implementation +## GREEN evidence -- `useSessionStream` accepts an optional `generation` argument (default `0`) and includes it - in the stream effect dependencies. A generation change therefore runs the existing cleanup, - closes the old source, and opens the same URL again. -- `AppShell` captures whether the requested Resume ID is already active before its existing - optimistic state updates. It increments the stream generation only after `resumeSession(id)` - succeeds and only for that same-ID case. -- The existing optimistic session switch, phase refresh, and failed-Resume rollback remain - unchanged. A failed POST cannot increment the generation. +- Focused Python regression set: `66 passed`: + `test_session_repository_workflow`, `test_session_repository`, session mutation/list/ + documents/schema-linking, CTE plan/next, decision phase gate, and phase requirement tests. +- Gate suite: `127 passed`, including + `session-repository-writes.test.js`. +- Changed-source Ruff checks pass. `git diff --check` passes. -## TDD evidence +## Implemented boundary -- RED command: - `cd frontend && npx vitest run src/stream/useSessionStream.test.tsx src/shell/AppShell.session-mgmt.test.tsx` -- RED result: 2 expected failures and 15 passes. The hook test observed - `first.closed === false`; the AppShell test observed one `FakeEventSource` instead of two - after the second same-ID Resume. -- GREEN focused result: the same command passed 2/2 files and 17/17 tests after the minimal - production wiring. +- Added `resolve_principal`: PostgreSQL session storage requires trusted + `THT_PRINCIPAL_ISSUER` and `THT_PRINCIPAL_SUBJECT`, optional display name, and + strict admin parsing (`1`/`true`). It fails closed and never substitutes a local + identity. Filesystem storage uses `local_principal()`. +- Filesystem repository creates UUIDv4 sessions only and permits safe historical + timestamp IDs (`YYYY-MM-DD-HHMMSS`) for read/mutate compatibility. PostgreSQL + remains UUIDv4 only. +- Phase helpers fold `SessionSnapshot` ledger/artifacts; decision, phase, CTE, + session mutation/list/document paths, retrieval-pack persistence, SQL promotion + lookup, and task-doc/CTE test helpers gained repository/snapshot paths. +- Finalization now publishes report, evidence, and finalized manifest through + `repository.finalize`: one PostgreSQL transaction; filesystem writes artifacts + before the finalized manifest commit marker. Solved-question indexing stays + best-effort after this durable write. +- Added `tht cte save --session --name --file -` and + `tht sql set-final --session --file -`; Pi tools and SKILL.md now use them. -## Full verification +## Outstanding in-scope migration work -- Baseline before edits: `cd frontend && npx vitest run` — 42/42 files and 251/251 tests passed. -- Focused tests: 2/2 files and 17/17 tests passed. -- Full frontend suite: `cd frontend && npx vitest run` — 42/42 files and 253/253 tests passed. -- Typecheck: `cd frontend && npx tsc -b` — exit 0. -- Production build: `cd frontend && npm run build` — exit 0; Vite transformed 4,835 modules - and completed the production bundle. -- `git diff --check` — passed. +Do not treat this task as complete yet. Remaining direct session path consumers are: -The suite and build retained the pre-existing MSW unhandled-request, React ref/`act`, Node type -stripping, and Vite chunk-size warnings. This task introduced no new warning category. +- `harness/tht/cli/memory_cmd.py`: lines 60, 93, 165, 400, 458. +- `harness/tht/cli/sql_cmd.py`: `_session_sql_file` at line 254 remains a legacy + Path-returning bridge for preview/save/export. +- `harness/tht/cli/session_cmd.py:session_dir` remains only as a compatibility + bridge for the out-of-scope datamart command and the still-unmigrated memory/ + SQL consumers; workflow mutations in session_cmd do not call it. -## Files - -- `frontend/src/stream/useSessionStream.ts` -- `frontend/src/stream/useSessionStream.test.tsx` -- `frontend/src/shell/AppShell.tsx` -- `frontend/src/shell/AppShell.session-mgmt.test.tsx` -- `.superpowers/sdd/task-3-report.md` - -## Self-review - -- Confirmed the old EventSource is closed before the replacement is retained by React's effect - lifecycle, and the replacement uses the identical session URL. -- Confirmed same-ID detection happens before the optimistic `setActiveSessionId(id)` call. -- Confirmed the generation increments only after a successful Resume POST; the catch/rollback - branch is unchanged. -- Confirmed a different ID leaves the generation unchanged, so the existing session-ID effect - change creates exactly one replacement connection. -- Confirmed the diff is frontend-only apart from this report and contains no backend, Docker, - configuration, or session changes. - -## Concerns - -None. +The full Python suite has not been conclusively re-run to completion after the +latest changes. An earlier root-directory invocation failed only because a +pre-existing test expects `workflow.yaml` relative to `harness/`. Full gate tests +are green. Full-repo Ruff currently fails on pre-existing test-file lint findings; +changed-source Ruff passes. diff --git a/harness/.pi/extensions/gate/__tests__/session-repository-writes.test.js b/harness/.pi/extensions/gate/__tests__/session-repository-writes.test.js new file mode 100644 index 00000000..58a6cff2 --- /dev/null +++ b/harness/.pi/extensions/gate/__tests__/session-repository-writes.test.js @@ -0,0 +1,41 @@ +const test = require("node:test"); +const assert = require("node:assert"); +const cp = require("node:child_process"); +const { createRequire } = require("node:module"); +const path = require("node:path"); + +const GATE = path.join(__dirname, "..", "..", "tht-gate.js"); +if (typeof globalThis.require === "undefined") globalThis.require = createRequire(GATE); + +test("repository write tools persist CTE and final SQL through deterministic tht commands", async () => { + const calls = []; + const original = cp.execFileSync; + cp.execFileSync = (file, args, options) => { + calls.push({ args, input: options?.input }); + return ""; + }; + try { + const gate = require(GATE); + const { createFakePi } = require("./fake_pi_runtime.js"); + const { pi, ctx, tools } = createFakePi(); + ctx.cwd = "/nonexistent-thothii-test-cwd"; + gate.default(pi); + + const cte = tools.get("write_cte_sql"); + const final = tools.get("write_final_sql"); + assert.ok(cte, "write_cte_sql must be registered"); + assert.ok(final, "write_final_sql must be registered"); + + await cte.def.execute("cte", { session: "s1", name: "base", sql: "WITH base AS (SELECT 1)" }, null, null, ctx); + await final.def.execute("sql", { session: "s1", sql: "SELECT * FROM base" }, null, null, ctx); + + assert.deepEqual(calls.map((call) => call.args), [ + ["cte", "save", "--session", "s1", "--name", "base", "--file", "-"], + ["sql", "set-final", "--session", "s1", "--file", "-"], + ]); + assert.deepEqual(calls.map((call) => call.input), ["WITH base AS (SELECT 1)", "SELECT * FROM base"]); + assert.equal(JSON.stringify(calls).includes("sessions/s1"), false); + } finally { + cp.execFileSync = original; + } +}); diff --git a/harness/.pi/extensions/tht-gate.js b/harness/.pi/extensions/tht-gate.js index eac2c836..357af984 100644 --- a/harness/.pi/extensions/tht-gate.js +++ b/harness/.pi/extensions/tht-gate.js @@ -1401,6 +1401,38 @@ export default function (pi) { }, }); + pi.registerTool({ + name: "write_cte_sql", + label: "Scrittura CTE SQL", + description: "Persiste un blocco CTE tramite tht cte save; non scrivere mai file di sessione direttamente.", + parameters: Type.Object({ session: Type.String(), name: Type.String(), sql: Type.String() }), + async execute(_id, params, _signal, _onUpdate, ctx) { + const err = relayIfThtFails( + ctx, + ["cte", "save", "--session", params.session, "--name", params.name, "--file", "-"], + "Correggi il blocco CTE e riprova.", + params.sql, + ); + return err || textResult(`CTE ${params.name} salvato per la sessione ${params.session}.`); + }, + }); + + pi.registerTool({ + name: "write_final_sql", + label: "Scrittura SQL finale", + description: "Persiste il SQL finale tramite tht sql set-final; non scrivere mai file di sessione direttamente.", + parameters: Type.Object({ session: Type.String(), sql: Type.String() }), + async execute(_id, params, _signal, _onUpdate, ctx) { + const err = relayIfThtFails( + ctx, + ["sql", "set-final", "--session", params.session, "--file", "-"], + "Correggi il SQL finale e riprova.", + params.sql, + ); + return err || textResult(`SQL finale salvato per la sessione ${params.session}.`); + }, + }); + // --- slash command: /torna [session_id] [N] (rollback to a previous phase) -- pi.registerCommand("torna", { description: diff --git a/harness/.pi/skills/tht-sessione/SKILL.md b/harness/.pi/skills/tht-sessione/SKILL.md index 9a59551e..d32079ad 100644 --- a/harness/.pi/skills/tht-sessione/SKILL.md +++ b/harness/.pi/skills/tht-sessione/SKILL.md @@ -353,7 +353,9 @@ Prerequisite: Phase 5 closed. "output_columns":["cod_paz","data_ricovero"]} ]} ``` -3. For each CTE (in plan order): write `sessions//ctes/.sql` (ONLY the +3. For each CTE (in plan order): call `write_cte_sql` with the session id, CTE name, + and SQL block (the tool invokes `tht cte save --session --name --file -`). + Persist ONLY the `WITH ... AS (...)` block, NO trailing SELECT), test with `tht cte test --session ` (with an **ok** outcome), then present it with `reviewer_confirm kind:"cte_result"`. Pass ONLY the thin v2 data @@ -385,7 +387,8 @@ Prerequisite: Phase 6 closed. columns (`dt.year`, `dt.month`, …), NEVER arithmetic on the key. 3. `tht sql validate` + `tht sql preview` (max 10 rows). On errors / suspicious results, apply the `sql-generation.md` checklist and correct with the reviewer. -4. `tht sql save` writes `sessions//sql_final.sql` (ONLY clean SQL, no comments). +4. Call `write_final_sql` with the session id and clean SQL; it invokes + `tht sql set-final --session --file -` (ONLY clean SQL, no comments). Approve with `reviewer_confirm kind:"sql"` (records `sql_approved`), then advance to Phase 8 with `reviewer_confirm kind:"phase"` — `kind:"sql"` alone does NOT advance F7. diff --git a/harness/tests/test_adapter_command_regressions.py b/harness/tests/test_adapter_command_regressions.py index 28ee362c..ebadf06a 100644 --- a/harness/tests/test_adapter_command_regressions.py +++ b/harness/tests/test_adapter_command_regressions.py @@ -67,17 +67,17 @@ def test_memory_command_writes_through_factory_vector_store(monkeypatch): # configure the workstation-only REST writer key. cfg = SimpleNamespace(profile="server", embeddings=object(), vector_write_rest=None) manifest = SimpleNamespace(id="s1") + snapshot = SimpleNamespace(manifest=manifest, decisions=[], artifacts={}) record = MemoryRecord(id="m1", ts=datetime(2026, 1, 1), session_id="s1", decision_seq=7, type="table_promoted", subject="t", question_context="q") monkeypatch.setattr(memory_cmd, "_load_config_or_exit", lambda path: cfg) - monkeypatch.setattr(memory_cmd, "load_session_or_exit", lambda cfg, session: manifest) - monkeypatch.setattr(memory_cmd, "session_dir", lambda *args: None) + monkeypatch.setattr(memory_cmd, "load_snapshot_or_exit", lambda cfg, session: snapshot) monkeypatch.setattr(memory_cmd, "registry_path", lambda cfg: None) monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: store) monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda cfg: SimpleNamespace(embed_documents=lambda texts: [[0.1]])) - monkeypatch.setattr("tht.memory.promote", lambda *args, **kwargs: None) + monkeypatch.setattr("tht.memory.promote_snapshot", lambda *args, **kwargs: None) monkeypatch.setattr("tht.memory.load_registry", lambda path: [record]) memory_cmd.save_one_cmd(session="s1", decision=7, json_out=True) from tht.ports.vector import VectorWriteRecord @@ -96,14 +96,13 @@ def test_solved_index_writes_through_writer_only_factory_store(monkeypatch): calls = [] monkeypatch.setattr(memory_cmd, "has_vector_write_rest", lambda cfg: True) - monkeypatch.setattr(memory_cmd, "load_session_or_exit", lambda cfg, session: manifest) - monkeypatch.setattr(memory_cmd, "session_dir", lambda *args: None) + monkeypatch.setattr(memory_cmd, "load_snapshot_or_exit", lambda cfg, session: SimpleNamespace(manifest=manifest, decisions=[], artifacts={})) monkeypatch.setattr( "tht.adapters.factory.build_vector_store", lambda cfg, require_write: calls.append(require_write) or writer_only_store, ) monkeypatch.setattr("tht.cli.sql_cmd.promoted_tables_for", lambda *args: []) - monkeypatch.setattr("tht.solved.build_solved_record", lambda *args: solved_record) + monkeypatch.setattr("tht.solved.build_solved_snapshot", lambda *args: solved_record) monkeypatch.setattr( "tht.solved.save_solved_question", lambda record, *, store, embedder: int( diff --git a/harness/tests/test_decision_join_set_cli.py b/harness/tests/test_decision_join_set_cli.py index 3865197a..f14807d7 100644 --- a/harness/tests/test_decision_join_set_cli.py +++ b/harness/tests/test_decision_join_set_cli.py @@ -28,15 +28,28 @@ def _walk_to_phase(session, target): def _configure_command(monkeypatch, sessions): import tht.cli.decision_cmd as mod + import tht.cli.session_cmd as session_mod + from tht.session.models import SessionManifest, SessionSnapshot + + class _Repository: + def get(self, session_id): + return SessionSnapshot( + manifest=SessionManifest(id=session_id, created_at="2026-01-01T00:00:00Z", question="q", database="d", schema="s"), + decisions=list_decisions(sessions / session_id), + ) + + def append_decisions(self, session_id, decisions): + return append_decisions(sessions / session_id, list(decisions)) class _Cfg: - class paths: - pass + pass - _Cfg.paths.sessions = sessions + repository = _Repository() monkeypatch.setattr(mod, "_load_config_or_exit", lambda _c: _Cfg()) monkeypatch.setattr(mod, "load_session_or_exit", lambda _cfg, _s: None) - monkeypatch.setattr(mod, "session_dir", lambda _cfg, sid: sessions / sid) + monkeypatch.setattr(mod, "load_snapshot_or_exit", lambda _cfg, sid: repository.get(sid)) + monkeypatch.setattr(mod, "session_repository", lambda _cfg: repository) + monkeypatch.setattr(session_mod, "load_snapshot_or_exit", lambda _cfg, sid: repository.get(sid)) def test_add_join_set_rejects_the_whole_batch_when_one_item_is_invalid(tmp_path, monkeypatch): diff --git a/harness/tests/test_decision_retract_cli.py b/harness/tests/test_decision_retract_cli.py index 4991d2aa..c77ee42b 100644 --- a/harness/tests/test_decision_retract_cli.py +++ b/harness/tests/test_decision_retract_cli.py @@ -25,6 +25,8 @@ def test_retract_drops_last_substantive_decision(tmp_path, monkeypatch): # stub config + session loading (the command only needs a session dir) import tht.cli.decision_cmd as mod + import tht.cli.session_cmd as session_mod + from tht.session.models import SessionManifest, SessionSnapshot class _Cfg: class paths: @@ -32,7 +34,22 @@ def test_retract_drops_last_substantive_decision(tmp_path, monkeypatch): monkeypatch.setattr(mod, "_load_config_or_exit", lambda _c: _Cfg()) monkeypatch.setattr(mod, "load_session_or_exit", lambda _cfg, _s: None) - monkeypatch.setattr(mod, "session_dir", lambda _cfg, sid: tmp_path / sid) + class _Repository: + def get(self, session_id): + return SessionSnapshot( + manifest=SessionManifest(id=session_id, created_at="2026-01-01T00:00:00Z", question="q", database="d", schema="s"), + decisions=list_decisions(tmp_path / session_id), + ) + + def append_decisions(self, session_id, decisions): + from tht.decisions import append_decisions + + return append_decisions(tmp_path / session_id, list(decisions)) + + repository = _Repository() + monkeypatch.setattr(mod, "load_snapshot_or_exit", lambda _cfg, sid: repository.get(sid)) + monkeypatch.setattr(mod, "session_repository", lambda _cfg: repository) + monkeypatch.setattr(session_mod, "load_snapshot_or_exit", lambda _cfg, sid: repository.get(sid)) retract_cmd(session="2026-01-01-000000-x", config=Path("x")) diff --git a/harness/tests/test_promoted_columns_for.py b/harness/tests/test_promoted_columns_for.py index 3aa40f54..57ff20f7 100644 --- a/harness/tests/test_promoted_columns_for.py +++ b/harness/tests/test_promoted_columns_for.py @@ -3,9 +3,13 @@ from tht.cli.sql_cmd import promoted_columns_for def test_promoted_columns_for(tmp_path): - sid = "sess1" + sid = "2026-07-16-120000" sdir = tmp_path / sid sdir.mkdir(parents=True) + (sdir / "session_manifest.yaml").write_text( + f"id: {sid}\nquestion: q\ndatabase: d\nschema: s\n" + "created_at: 2026-01-01T00:00:00+00:00\nstatus: open\n" + ) (sdir / "schema_linking.json").write_text(json.dumps({ "question": "q", "candidates": [ diff --git a/harness/tests/test_repository_memory_sql_paths.py b/harness/tests/test_repository_memory_sql_paths.py new file mode 100644 index 00000000..31a6edfb --- /dev/null +++ b/harness/tests/test_repository_memory_sql_paths.py @@ -0,0 +1,16 @@ +"""Repository-only workflow commands must not recover server state via session_dir.""" + +import inspect + +from tht.cli import memory_cmd, sql_cmd + + +def test_memory_workflow_commands_do_not_import_or_call_session_dir(): + source = inspect.getsource(memory_cmd) + assert "session_dir" not in source + + +def test_sql_session_commands_read_snapshot_artifacts_not_session_paths(): + source = inspect.getsource(sql_cmd) + assert "_session_sql_file" not in source + assert "session_dir" not in source diff --git a/harness/tests/test_require_phase_or_exit.py b/harness/tests/test_require_phase_or_exit.py index 75f5c8b2..57ff6221 100644 --- a/harness/tests/test_require_phase_or_exit.py +++ b/harness/tests/test_require_phase_or_exit.py @@ -16,8 +16,13 @@ def _make_session(tmp_path, current_phase_num: int) -> str: approving phases 1..N-1 puts the session at phase N.""" import json - s = tmp_path / "sess" + session_id = "2026-07-16-120000" + s = tmp_path / session_id s.mkdir() + (s / "session_manifest.yaml").write_text( + f"id: {session_id}\nquestion: q\ndatabase: d\nschema: s\n" + "created_at: 2026-01-01T00:00:00+00:00\nstatus: open\n" + ) decisions = [ {"seq": n, "type": "phase_approved", "subject": f"phase:{n}", "ts": "2025-01-01T00:00:00"} for n in range(1, current_phase_num) @@ -25,7 +30,7 @@ def _make_session(tmp_path, current_phase_num: int) -> str: (s / "review_decisions.jsonl").write_text( "\n".join(json.dumps(d) for d in decisions) + ("\n" if decisions else "") ) - return "sess" + return session_id class _StubConfig: diff --git a/harness/tests/test_session_repository.py b/harness/tests/test_session_repository.py index 1d2f654f..9afc93d6 100644 --- a/harness/tests/test_session_repository.py +++ b/harness/tests/test_session_repository.py @@ -1,10 +1,13 @@ import uuid +import pytest + from tht.config import DatabaseConfig, load_config from tht.decisions import DecisionInput from tht.session.filesystem_repository import FilesystemSessionRepository from tht.session.models import PrincipalContext, SessionManifest -from tht.session.repository import build_session_repository +from tht.session.repository import build_session_repository, resolve_principal +from tht.session.store import SessionError from tht.session.store import create_session @@ -100,3 +103,52 @@ def test_filesystem_repository_keeps_preferences_per_principal(tmp_path): assert alice.get_preferences() == {"model": "glm"} assert bob.get_preferences() == {} + + +@pytest.mark.parametrize("session_id", ["2026-01-01-000000-test", "2026-01-01-000000-x", "s1"]) +def test_filesystem_repository_reads_safe_legacy_session_ids(tmp_path, session_id): + repository = FilesystemSessionRepository( + tmp_path / "home", "demo", PrincipalContext(issuer="local", subject="alice") + ) + directory = repository.root / session_id + directory.mkdir(parents=True) + _manifest(session_id).to_yaml(directory / "session_manifest.yaml") + + assert repository.get(session_id).manifest.id == session_id + + +@pytest.mark.parametrize("session_id", ["..", "a/b", "/absolute", "space id"]) +def test_filesystem_repository_rejects_unsafe_legacy_session_ids(tmp_path, session_id): + repository = FilesystemSessionRepository( + tmp_path / "home", "demo", PrincipalContext(issuer="local", subject="alice") + ) + + with pytest.raises(SessionError): + repository.get(session_id) + + +def test_postgres_session_storage_requires_trusted_principal_environment(tmp_path, monkeypatch): + config = _config(tmp_path).model_copy(update={"session_storage": { + "type": "postgres_direct", + "connection": _db().model_dump(by_alias=True), + }}) + monkeypatch.delenv("THT_PRINCIPAL_ISSUER", raising=False) + monkeypatch.delenv("THT_PRINCIPAL_SUBJECT", raising=False) + + with pytest.raises(SessionError, match="THT_PRINCIPAL_ISSUER"): + resolve_principal(config) + + +def test_postgres_session_storage_uses_only_explicit_trusted_principal(tmp_path, monkeypatch): + config = _config(tmp_path).model_copy(update={"session_storage": { + "type": "postgres_direct", + "connection": _db().model_dump(by_alias=True), + }}) + monkeypatch.setenv("THT_PRINCIPAL_ISSUER", "portal") + monkeypatch.setenv("THT_PRINCIPAL_SUBJECT", "alice") + monkeypatch.setenv("THT_PRINCIPAL_DISPLAY_NAME", "Alice") + monkeypatch.setenv("THT_PRINCIPAL_IS_ADMIN", "TRUE") + + assert resolve_principal(config) == PrincipalContext( + issuer="portal", subject="alice", display_name="Alice", is_admin=True + ) diff --git a/harness/tests/test_session_repository_workflow.py b/harness/tests/test_session_repository_workflow.py new file mode 100644 index 00000000..8608ce25 --- /dev/null +++ b/harness/tests/test_session_repository_workflow.py @@ -0,0 +1,60 @@ +import uuid + +from tht.decisions import DecisionInput +from tht.phase import current_phase, cte_plan, next_cte +from tht.session.filesystem_repository import FilesystemSessionRepository +from tht.session.models import PrincipalContext, SessionManifest +from tht.session.store import persist_verified_finalization + + +def _manifest(session_id: str) -> SessionManifest: + return SessionManifest( + id=session_id, + created_at="2026-07-16T10:00:00Z", + question="Which patients had an ablation?", + database="testdb", + schema="public", + ) + + +def test_phase_helpers_fold_the_repository_snapshot_not_a_session_path(tmp_path): + repository = FilesystemSessionRepository( + tmp_path / "home", "demo", PrincipalContext(issuer="local", subject="alice") + ) + session_id = str(uuid.uuid4()) + repository.create(_manifest(session_id)) + repository.write_artifact(session_id, "cte_plan", '["base_patients"]') + repository.append_decisions( + session_id, + [ + DecisionInput(type="phase_approved", subject="phase:1"), + DecisionInput(type="phase_auto_approved", subject="phase:2"), + DecisionInput(type="cte_approved", subject="base_patients"), + ], + ) + + snapshot = repository.get(session_id) + + assert current_phase(snapshot) == 3 + assert cte_plan(snapshot) == ["base_patients"] + assert next_cte(snapshot) is None + + +def test_verified_finalization_commits_report_evidence_and_status_through_repository(tmp_path): + repository = FilesystemSessionRepository( + tmp_path / "home", "demo", PrincipalContext(issuer="local", subject="alice") + ) + session_id = str(uuid.uuid4()) + repository.create(_manifest(session_id)) + + persist_verified_finalization( + repository, + session_id, + validation_report="# Validation\n\nverified against DWH\n", + evidence='[{"source":"review"}]\n', + ) + + snapshot = repository.get(session_id) + assert snapshot.manifest.status == "finalized" + assert snapshot.artifacts["validation_report"] == "# Validation\n\nverified against DWH\n" + assert snapshot.artifacts["evidence"] == '[{"source":"review"}]\n' diff --git a/harness/tests/test_sql_preview_json.py b/harness/tests/test_sql_preview_json.py index a19b34c7..a1d36997 100644 --- a/harness/tests/test_sql_preview_json.py +++ b/harness/tests/test_sql_preview_json.py @@ -108,16 +108,12 @@ def test_do_run_offset_zero_path_unchanged(monkeypatch): def test_preview_session_no_file_resolves_sql_final(monkeypatch, tmp_path, capsys): - """--session without a positional FILE resolves sql_final.sql via _session_sql_file.""" + """--session without a positional FILE resolves repository SQL text.""" from types import SimpleNamespace from tht.cli import sql_cmd - sql_file = tmp_path / "sql_final.sql" - sql_file.write_text("SELECT session_resolved") - - # Patch _session_sql_file to return our tmp file (no real workspace/DB needed). - monkeypatch.setattr(sql_cmd, "_session_sql_file", lambda cfg, sid: sql_file) + monkeypatch.setattr(sql_cmd, "_session_sql", lambda cfg, sid: "SELECT session_resolved") captured_sql = {} diff --git a/harness/tht/cli/cte_cmd.py b/harness/tht/cli/cte_cmd.py index 94cd8ca2..d7b39a58 100644 --- a/harness/tht/cli/cte_cmd.py +++ b/harness/tht/cli/cte_cmd.py @@ -6,7 +6,7 @@ import typer from tht.cli.config_cmd import CONFIG_OPT from tht.cli.schema_cmd import _load_config_or_exit -from tht.cli.session_cmd import load_session_or_exit, session_dir +from tht.cli.session_cmd import load_session_or_exit, load_snapshot_or_exit, session_repository from tht.cli.sql_cmd import promoted_tables_for, require_action _RULE6_HINT = ( @@ -17,6 +17,27 @@ _RULE6_HINT = ( cte_app = typer.Typer(help="Test controllato dei CTE proposti (Agent View Generation)") +@cte_app.command("save") +def save_cmd( + session: str = typer.Option(..., "--session"), + name: str = typer.Option(..., "--name"), + file: str = typer.Option(..., "--file", help="File SQL, oppure '-' per stdin."), + config: Path = CONFIG_OPT, +) -> None: + """Persist one CTE SQL block through the configured session repository.""" + import sys + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session) + raw = sys.stdin.read() if file == "-" else Path(file).read_text() + try: + session_repository(cfg).write_artifact(session, f"cte_sql:{name}", raw) + except ValueError as exc: + typer.secho(f"ERRORE: {exc}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from None + typer.secho(f"OK: CTE {name} salvato.", fg=typer.colors.GREEN) + + @cte_app.command("test") def test_cmd( name: str = typer.Argument(..., help="Nome del CTE (file sessions//ctes/.sql)."), @@ -33,7 +54,7 @@ def test_cmd( CteError, CteTestRecord, _jsonable, - append_cte_test, + append_cte_test_snapshot, build_test_sql, has_trailing_select, ) @@ -43,13 +64,13 @@ def test_cmd( cfg = _load_config_or_exit(config) require_action(cfg, "cte_test") load_session_or_exit(cfg, session) - sdir = session_dir(cfg, session) + snapshot = load_snapshot_or_exit(cfg, session) from tht.cli.phase_cmd import require_phase_or_exit from tht.phase import next_cte require_phase_or_exit(cfg, session, 6) - nxt = next_cte(sdir) + nxt = next_cte(snapshot) if nxt is not None and name != nxt: typer.secho( f"ERRORE: ordine CTE. Ora tocca a '{nxt}' (il primo CTE del piano non " @@ -60,11 +81,10 @@ def test_cmd( ) raise typer.Exit(code=5) - cte_file = sdir / "ctes" / f"{name}.sql" - if not cte_file.exists(): - typer.secho(f"ERRORE: file CTE non trovato: {cte_file}", fg=typer.colors.RED, err=True) + cte_sql = snapshot.artifacts.get(f"cte_sql:{name}") + if cte_sql is None: + typer.secho(f"ERRORE: file CTE non trovato: {name}", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) - cte_sql = cte_file.read_text() sql_hash = hashlib.sha256(cte_sql.encode()).hexdigest() def _record_error(message: str) -> None: @@ -72,7 +92,7 @@ def test_cmd( name=name, ts=datetime.now(UTC), sql_hash=sql_hash, status="error", error=message, ) - append_cte_test(sdir, record) + append_cte_test_snapshot(session_repository(cfg), snapshot, record) if json_out: typer.echo(record.model_dump_json()) else: @@ -118,7 +138,7 @@ def test_cmd( execution_ms=result.execution_ms, warnings=warnings, preview_rows=[[_jsonable(cell) for cell in row] for row in result.rows], ) - append_cte_test(sdir, record) + append_cte_test_snapshot(session_repository(cfg), snapshot, record) if json_out: typer.echo(record.model_dump_json()) @@ -133,7 +153,7 @@ def test_cmd( Console().print(table) for w in warnings: typer.secho(f" warning: {w}", fg=typer.colors.YELLOW) - typer.secho(f"OK: esito registrato in {sdir / 'cte_tests.json'}", fg=typer.colors.GREEN) + typer.secho("OK: esito registrato in cte_tests.json", fg=typer.colors.GREEN) CTE_PLAN_DOC_FILE = "cte_plan_doc.json" @@ -156,13 +176,10 @@ def plan_cmd( import json import sys - from tht.phase import CTE_PLAN_FILE - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session) from tht.cli.phase_cmd import require_phase_or_exit require_phase_or_exit(cfg, session, 6) - sdir = session_dir(cfg, session) doc_data = None if doc is not None: @@ -180,11 +197,11 @@ def plan_cmd( ) raise typer.Exit(code=1) - sdir.mkdir(parents=True, exist_ok=True) - (sdir / CTE_PLAN_FILE).write_text(json.dumps(name, ensure_ascii=False)) + repository = session_repository(cfg) + repository.write_artifact(session, "cte_plan", json.dumps(name, ensure_ascii=False)) if doc_data is not None: - (sdir / CTE_PLAN_DOC_FILE).write_text(json.dumps(doc_data, ensure_ascii=False)) - typer.secho(f"OK: piano CTE salvato ({len(name)} CTE) in {sdir / CTE_PLAN_FILE}.", + repository.write_artifact(session, "cte_plan_doc", json.dumps(doc_data, ensure_ascii=False)) + typer.secho(f"OK: piano CTE salvato ({len(name)} CTE).", fg=typer.colors.GREEN) @@ -201,7 +218,7 @@ def next_cmd( cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session) - nxt = next_cte(session_dir(cfg, session)) + nxt = next_cte(load_snapshot_or_exit(cfg, session)) if nxt: typer.echo(nxt) @@ -222,9 +239,9 @@ def info_cmd( cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session) - sdir = session_dir(cfg, session) + snapshot = load_snapshot_or_exit(cfg, session) - plan = cte_plan(sdir) + plan = cte_plan(snapshot) if not plan: typer.secho(f"ERRORE: {CTE_PLAN_FILE} assente o vuoto per la sessione '{session}'.", fg=typer.colors.RED, err=True) @@ -234,24 +251,24 @@ def info_cmd( fg=typer.colors.RED, err=True) raise typer.Exit(code=1) - cte_file = sdir / "ctes" / f"{name}.sql" - if not cte_file.exists(): - typer.secho(f"ERRORE: file CTE non trovato: {cte_file}", fg=typer.colors.RED, err=True) + cte_sql = snapshot.artifacts.get(f"cte_sql:{name}") + if cte_sql is None: + typer.secho(f"ERRORE: file CTE non trovato: {name}", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) doc = None - doc_path = sdir / CTE_PLAN_DOC_FILE - if doc_path.exists(): - full_doc = _json.loads(doc_path.read_text()) + raw_doc = snapshot.artifacts.get("cte_plan_doc") + if raw_doc: + full_doc = _json.loads(raw_doc) for c in full_doc.get("ctes", []): if c.get("name") == name: doc = {k: v for k, v in c.items() if k != "name"} break - from tht.ctetest import CteError, load_cte_tests + from tht.ctetest import CteError, load_cte_tests_text try: - records = [r for r in load_cte_tests(sdir) if r.name == name] + records = [r for r in load_cte_tests_text(snapshot.artifacts.get("cte_tests", "")) if r.name == name] except CteError as e: typer.secho(f"ERRORE: impossibile leggere cte_tests.json: {e}", fg=typer.colors.RED, err=True) @@ -263,8 +280,8 @@ def info_cmd( "index": plan.index(name) + 1, "total": len(plan), "plan": plan, - "sql": cte_file.read_text(), - "approved": name in approved_ctes(sdir), + "sql": cte_sql, + "approved": name in approved_ctes(snapshot), "doc": doc, "last_test": last_test, } @@ -282,12 +299,12 @@ def list_cmd( config: Path = CONFIG_OPT, ) -> None: """Ultimo esito registrato per ogni CTE della sessione.""" - from tht.ctetest import CteError, load_cte_tests + from tht.ctetest import CteError, load_cte_tests_text cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session) try: - records = load_cte_tests(session_dir(cfg, session)) + records = load_cte_tests_text(load_snapshot_or_exit(cfg, session).artifacts.get("cte_tests", "")) except CteError as e: typer.secho(f"ERRORE: impossibile leggere cte_tests.json: {e}", fg=typer.colors.RED, err=True) diff --git a/harness/tht/cli/decision_cmd.py b/harness/tht/cli/decision_cmd.py index 0f13d9bf..0c695a87 100644 --- a/harness/tht/cli/decision_cmd.py +++ b/harness/tht/cli/decision_cmd.py @@ -8,7 +8,7 @@ from pydantic import ValidationError from tht.cli.config_cmd import CONFIG_OPT from tht.cli.schema_cmd import _load_config_or_exit -from tht.cli.session_cmd import load_session_or_exit, session_dir +from tht.cli.session_cmd import load_session_or_exit, load_snapshot_or_exit, session_repository decision_app = typer.Typer(help="Decisioni del reviewer (per sessione, append-only)") @@ -20,7 +20,7 @@ def add_join_set_cmd( config: Path = CONFIG_OPT, ) -> None: """Registra un insieme completo di join con un'unica sostituzione atomica del ledger.""" - from tht.decisions import DecisionInput, append_decisions + from tht.decisions import DecisionInput cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session) @@ -40,7 +40,7 @@ def add_join_set_cmd( from tht.workflow import load_workflow require_phase_or_exit(cfg, session, load_workflow().decision_min_phase("join_modified")) - records = append_decisions(session_dir(cfg, session), decisions) + records = session_repository(cfg).append_decisions(session, decisions) typer.secho( f"OK: registrato set atomico di {len(records)} join.", fg=typer.colors.GREEN, @@ -61,7 +61,7 @@ def add_cmd( config: Path = CONFIG_OPT, ) -> None: """Registra una decisione del reviewer nella sessione.""" - from tht.decisions import DecisionType, append_decision + from tht.decisions import DecisionType cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session) @@ -82,7 +82,7 @@ def add_cmd( from tht.workflow import load_workflow require_phase_or_exit(cfg, session, load_workflow().decision_min_phase(type)) - sdir = session_dir(cfg, session) + snapshot = load_snapshot_or_exit(cfg, session) if type == "cte_approved": # Un CTE si approva solo se appartiene al piano persistito. Senza piano # (o con subject fuori piano) l'approvazione e' priva di significato: @@ -90,7 +90,7 @@ def add_cmd( # dove fu registrato un cte_approved:cte_plan senza alcun cte_plan.json. from tht.phase import cte_plan - plan = cte_plan(sdir) + plan = cte_plan(snapshot) if not plan: typer.secho( "ERRORE: nessun piano CTE (cte_plan.json) in sessione. Persisti prima " @@ -105,10 +105,10 @@ def add_cmd( fg=typer.colors.RED, err=True, ) raise typer.Exit(code=5) - record = append_decision( - sdir, type=type, subject=subject, - detail=detail, rationale=rationale, retracts=retracts, - ) + record = session_repository(cfg).append_decisions(session, [{ + "type": type, "subject": subject, "detail": detail, + "rationale": rationale, "retracts": retracts, + }])[0] typer.secho(f"OK: decisione [{record.seq}] {record.type}: {record.subject}", fg=typer.colors.GREEN) @@ -123,19 +123,18 @@ def retract_cmd( Granularita' (a) del rollback §4.8: 'rispondi di nuovo a questa domanda'. Scrive un marker decision_retracted (append-only, l'audit resta) che effective_decisions onora; il widget corrente puo' essere riproposto. Non cambia la fase.""" - from tht.decisions import append_decision from tht.phase import effective_decisions cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session) - sdir = session_dir(cfg, session) + snapshot = load_snapshot_or_exit(cfg, session) # Ultima decisione NON-meta della vista effective = quella associata al widget corrente. meta = { "phase_approved", "phase_auto_approved", "phase_reopened", "phase_skipped", "decision_retracted", } - substantive = [d for d in effective_decisions(sdir) if d.type not in meta] + substantive = [d for d in effective_decisions(snapshot) if d.type not in meta] if not substantive: typer.secho( "Nessuna decisione sostanziale da ritirare nella fase corrente.", @@ -143,10 +142,10 @@ def retract_cmd( ) raise typer.Exit(code=6) target = substantive[-1] - record = append_decision( - sdir, type="decision_retracted", subject=target.subject, - rationale=f"ritira [{target.seq}] {target.type}", retracts=target.seq, - ) + record = session_repository(cfg).append_decisions(session, [{ + "type": "decision_retracted", "subject": target.subject, + "rationale": f"ritira [{target.seq}] {target.type}", "retracts": target.seq, + }])[0] typer.secho( f"OK: ritirata decisione [{target.seq}] {target.type}: {target.subject} " f"(marker #{record.seq}).", @@ -160,11 +159,8 @@ def list_cmd( config: Path = CONFIG_OPT, ) -> None: """Elenca le decisioni della sessione.""" - from tht.decisions import list_decisions - cfg = _load_config_or_exit(config) - load_session_or_exit(cfg, session) - decisions = list_decisions(session_dir(cfg, session)) + decisions = load_snapshot_or_exit(cfg, session).decisions if not decisions: typer.echo("Nessuna decisione registrata.") return diff --git a/harness/tht/cli/memory_cmd.py b/harness/tht/cli/memory_cmd.py index 9b9df4ce..ea0d1fe9 100644 --- a/harness/tht/cli/memory_cmd.py +++ b/harness/tht/cli/memory_cmd.py @@ -17,7 +17,7 @@ from tht.cli._guards import ( require_server_profile, require_vector_write_allowed, ) -from tht.cli.session_cmd import load_session_or_exit, session_dir +from tht.cli.session_cmd import load_snapshot_or_exit from tht.cli.vector_cmd import require_vector_cfg memory_app = typer.Typer(help="Review memory (registro canonico + indice pgvector)") @@ -48,18 +48,17 @@ def promote_cmd( """Promuove le decisioni SCELTE nel registro globale. Usa --preview per vedere i candidati.""" import json as _json - from tht.memory import promote + from tht.memory import promote_snapshot cfg = _load_config_or_exit(config) - manifest = load_session_or_exit(cfg, session) + snapshot = load_snapshot_or_exit(cfg, session) if preview: from tht.memory import ( - MAX_PROMOTION_CANDIDATES, preview_promotions, reusable_promotions, + MAX_PROMOTION_CANDIDATES, preview_promotions_snapshot, reusable_promotions_snapshot, ) - sdir = session_dir(cfg, session) - cand = preview_promotions(sdir, manifest, registry_path(cfg)) - extra = len(reusable_promotions(sdir, manifest, registry_path(cfg))) - len(cand) + cand = preview_promotions_snapshot(snapshot, registry_path(cfg)) + extra = len(reusable_promotions_snapshot(snapshot, registry_path(cfg))) - len(cand) payload = [ {"decision_seq": c.decision_seq, "type": c.type, "subject": c.subject, "detail": c.detail, "rationale": c.rationale, @@ -89,10 +88,7 @@ def promote_cmd( require_server_profile(cfg, "memory promote") require_vector_cfg(cfg) - promoted = promote( - session_dir(cfg, session), manifest, - seqs=list(decision), registry_path=registry_path(cfg), - ) + promoted = promote_snapshot(snapshot, seqs=list(decision), registry_path=registry_path(cfg)) if not promoted: msg = "Nessuna nuova promozione (gia' presenti o seq inesistenti)." if json_out: @@ -155,18 +151,17 @@ def save_one_cmd( from tht.adapters.factory import build_vector_store from tht.cli.vector_cmd import make_embedder - from tht.memory import load_registry, promote, save_one_memory + from tht.memory import load_registry, promote_snapshot, save_one_memory cfg = _load_config_or_exit(config) - manifest = load_session_or_exit(cfg, session) + snapshot = load_snapshot_or_exit(cfg, session) require_vector_write_allowed(cfg, "memory save-one") store = build_vector_store(cfg, require_write=True) - sdir = session_dir(cfg, session) # Promuove la decisione scelta nel registro locale (idempotente: salta se gia' presente # o se stale post-rollback, perche' _compute_promotions usa la vista effective). - promote(sdir, manifest, seqs=[decision], registry_path=registry_path(cfg)) - records = [r for r in load_registry(registry_path(cfg)) if r.session_id == manifest.id] + promote_snapshot(snapshot, seqs=[decision], registry_path=registry_path(cfg)) + records = [r for r in load_registry(registry_path(cfg)) if r.session_id == snapshot.manifest.id] embedder = make_embedder(cfg.embeddings) count = save_one_memory(records, decision, store=store, embedder=embedder) @@ -388,16 +383,14 @@ def search_cmd( from rich.console import Console from rich.table import Table - from tht.cli.session_cmd import session_dir from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg from tht.memory import decided_memory_ids, load_registry - from tht.decisions import list_decisions cfg = _load_config_or_exit(config) require_vector_cfg(cfg) excluded: set[str] = set() if session is not None: - excluded = decided_memory_ids(list_decisions(session_dir(cfg, session))) + excluded = decided_memory_ids(load_snapshot_or_exit(cfg, session).decisions) searcher = open_searcher(cfg) embedder = make_embedder(cfg.embeddings) hits = searcher.search(embedder.embed_query(question), top_n=top, kinds=["memory"]) @@ -445,7 +438,7 @@ def index_solved_session(cfg, session_id: str) -> int: from tht.adapters.factory import build_vector_store from tht.cli.sql_cmd import promoted_tables_for from tht.cli.vector_cmd import make_embedder - from tht.solved import build_solved_record, save_solved_question + from tht.solved import build_solved_snapshot, save_solved_question if not has_vector_write_rest(cfg): raise RuntimeError( @@ -453,10 +446,7 @@ def index_solved_session(cfg, session_id: str) -> int: "writer key configurata nel workspace yaml" ) store = build_vector_store(cfg, require_write=True) - manifest = load_session_or_exit(cfg, session_id) - record = build_solved_record( - session_dir(cfg, session_id), manifest, promoted_tables_for(cfg, session_id) - ) + record = build_solved_snapshot(load_snapshot_or_exit(cfg, session_id), promoted_tables_for(cfg, session_id)) return save_solved_question( record, store=store, diff --git a/harness/tht/cli/phase_cmd.py b/harness/tht/cli/phase_cmd.py index 85b5051a..eb70ff23 100644 --- a/harness/tht/cli/phase_cmd.py +++ b/harness/tht/cli/phase_cmd.py @@ -8,7 +8,6 @@ commands. All read workflow facts from load_workflow() (no mirrored constants). from __future__ import annotations import json -from pathlib import Path import typer @@ -22,16 +21,12 @@ from tht.workflow import load_workflow phase_app = typer.Typer(help="Fase del workflow HITL (gate di avanzamento/ritorno)") -def session_dir(cfg, session_id: str) -> Path: - """Where a session's artifacts live. Mirror of session_cmd.session_dir (kept - here to avoid a circular import: session_cmd imports phase helpers too).""" - return cfg.paths.sessions / session_id - - def require_phase_or_exit(cfg, session: str, min_phase: int) -> None: """Refuse with exit 1 if the session hasn't reached min_phase yet. The phase name in the message comes from workflow.yaml (load_workflow), not a constant.""" - cur = current_phase(session_dir(cfg, session)) + from tht.cli.session_cmd import load_snapshot_or_exit + + cur = current_phase(load_snapshot_or_exit(cfg, session)) if cur < min_phase: wf = load_workflow() nome = wf.phase_name(min_phase) @@ -89,21 +84,24 @@ def advance_cmd( ), ) -> None: """Approva la fase corrente e passa alla successiva (persiste phase_approved).""" - from tht.decisions import append_decision + from tht.cli.session_cmd import load_snapshot_or_exit, session_repository - sdir = session_dir(_cfg(), session) - cur = current_phase(sdir) + cfg = _cfg() + snapshot = load_snapshot_or_exit(cfg, session) + cur = current_phase(snapshot) wf = load_workflow() if cur > wf.max_phase: typer.secho("Sessione già alla fase terminale.", fg=typer.colors.YELLOW) raise typer.Exit(0) if auto: - if not auto_advance_eligible(sdir): - problems = advance_problems(sdir, cur) + if not auto_advance_eligible(snapshot): + problems = advance_problems(snapshot, cur) for p in problems: typer.echo(p) raise typer.Exit(6) # needs human confirmation (gate contract) - append_decision(sdir, type="phase_approved", subject=f"phase:{cur}") + session_repository(cfg).append_decisions( + session, [{"type": "phase_approved", "subject": f"phase:{cur}"}] + ) typer.echo(f"Fase {cur} ({wf.phase_name(cur)}) approvata → Fase {cur + 1}.") @@ -113,16 +111,21 @@ def reopen_cmd( phase: int = typer.Option(..., "--phase", help="Fase a cui tornare (1..fase corrente -1)."), ) -> None: """Torna a una fase precedente (persiste phase_reopened + teardown artefatti).""" - from tht.decisions import append_decision - from tht.teardown import teardown_to_phase + from tht.cli.session_cmd import load_snapshot_or_exit, session_repository - sdir = session_dir(_cfg(), session) - cur = current_phase(sdir) + cfg = _cfg() + snapshot = load_snapshot_or_exit(cfg, session) + cur = current_phase(snapshot) if phase < 1 or phase >= cur: typer.secho(f"Target non valido (fase corrente {cur}).", fg=typer.colors.RED, err=True) raise typer.Exit(1) - report = teardown_to_phase(sdir, target_phase=phase) - append_decision(sdir, type="phase_reopened", subject=f"phase:{phase}") + from tht.teardown import teardown_snapshot + + repository = session_repository(cfg) + report = teardown_snapshot(repository, snapshot, phase) + repository.append_decisions( + session, [{"type": "phase_reopened", "subject": f"phase:{phase}"}] + ) for f in report.deleted_files: typer.echo(f" eliminato artefatto: {f}") typer.echo(f"Tornati alla Fase {phase} ({load_workflow().phase_name(phase)}).") @@ -133,13 +136,13 @@ def show_cmd( session: str = typer.Option(..., "--session"), ) -> None: """Mostra stato, fase corrente e ultime decisioni della sessione.""" - from tht.decisions import list_decisions + from tht.cli.session_cmd import load_snapshot_or_exit - sdir = session_dir(_cfg(), session) - cur = current_phase(sdir) + snapshot = load_snapshot_or_exit(_cfg(), session) + cur = current_phase(snapshot) wf = load_workflow() typer.echo(f"Fase corrente: {cur}/{wf.max_phase} ({wf.phase_name(min(cur, wf.max_phase))})") - decisions = list_decisions(sdir) + decisions = snapshot.decisions if decisions: typer.echo(f"Decisioni registrate: {len(decisions)}") for d in decisions[-5:]: diff --git a/harness/tht/cli/search_cmd.py b/harness/tht/cli/search_cmd.py index dcf91f32..b9bef1d5 100644 --- a/harness/tht/cli/search_cmd.py +++ b/harness/tht/cli/search_cmd.py @@ -371,14 +371,13 @@ def pack_cmd( md = "\n".join(md_lines) + "\n" if session: - from tht.cli.session_cmd import load_session_or_exit, session_dir + from tht.cli.session_cmd import load_session_or_exit, session_repository load_session_or_exit(cfg, session) - out = session_dir(cfg, session) / "retrieval_pack.md" - out.write_text(md) + session_repository(cfg).write_artifact(session, "retrieval_pack", md) if not json_out: typer.secho( - f"OK: retrieval pack scritto in {out} " + "OK: retrieval pack scritto " f"({len(tables)} tabelle, {len(evidence)} evidence, {len(solved)} solved).", fg=typer.colors.GREEN, ) diff --git a/harness/tht/cli/session_cmd.py b/harness/tht/cli/session_cmd.py index 05cd0814..e0a55cb4 100644 --- a/harness/tht/cli/session_cmd.py +++ b/harness/tht/cli/session_cmd.py @@ -10,6 +10,22 @@ from tht.cli.schema_cmd import _load_config_or_exit session_app = typer.Typer(help="Sessioni (directory artefatti)") +def session_repository(cfg): + """Configured persistence boundary for every workflow command.""" + from tht.session.repository import build_session_repository + + return build_session_repository(cfg) + + +def session_dir(cfg, session_id: str) -> Path: + """Legacy path bridge for out-of-scope datamart/memory compatibility only. + + Workflow session commands in this module use ``session_repository``; this + helper remains until those non-workflow consumers lose their path APIs. + """ + return cfg.paths.sessions / session_id + + @session_app.command("migrate") def migrate_cmd( database_url: str = typer.Option( @@ -43,15 +59,21 @@ def migrate_cmd( ) -def session_dir(cfg, session_id: str) -> Path: - return cfg.paths.sessions / session_id - - def load_session_or_exit(cfg, session_id: str): - from tht.session.store import SessionError, load_session + from tht.session.store import SessionError try: - return load_session(session_id, cfg.paths.sessions) + return session_repository(cfg).get(session_id).manifest + except SessionError as e: + typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + + +def load_snapshot_or_exit(cfg, session_id: str): + from tht.session.store import SessionError + + try: + return session_repository(cfg).get(session_id) except SessionError as e: typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) @@ -93,7 +115,14 @@ def list_cmd( ) -> None: """Elenca le sessioni esistenti.""" cfg = _load_config_or_exit(config) - rows = _list_sessions(cfg.paths.sessions) + rows = [ + {"id": s.manifest.id, "status": s.manifest.status, "question": s.manifest.question, + "summary": s.manifest.summary, "created_at": s.manifest.created_at.isoformat(), + "updated_at": s.manifest.updated_at.isoformat() if s.manifest.updated_at else None, + "author": s.manifest.author, "name": s.manifest.name, "group": s.manifest.group, + "archived": s.manifest.archived} + for s in session_repository(cfg).list() + ] if json_out: typer.echo(json.dumps(rows, ensure_ascii=False, indent=2)) return @@ -112,16 +141,19 @@ def new_cmd( config: Path = CONFIG_OPT, ) -> None: """Crea una sessione e stampa il suo id (ultima riga dell'output).""" - from tht.session.store import _extract_name, create_session + from tht.session.store import _extract_name, new_session_manifest, render_question_md cfg = _load_config_or_exit(config) - manifest = create_session(question, cfg.database, cfg.paths.sessions, + manifest = new_session_manifest(question, cfg.database, provider=provider, model=model, thinking=thinking, name=name or _extract_name(question)) + repository = session_repository(cfg) + repository.create(manifest) + repository.write_artifact(manifest.id, "question", render_question_md(question)) if json_out: typer.echo(json.dumps({"id": manifest.id}, ensure_ascii=False)) return - typer.secho(f"OK: sessione creata in {session_dir(cfg, manifest.id)}", fg=typer.colors.GREEN) + typer.secho(f"OK: sessione creata ({manifest.id})", fg=typer.colors.GREEN) typer.echo(manifest.id) @@ -136,12 +168,12 @@ def set_question_cmd( config: Path = CONFIG_OPT, ) -> None: """Scrive question.md (domanda riscritta + assunzioni) in modo deterministico.""" - from tht.session.store import set_question + from tht.session.store import render_question_md cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - path = set_question(session_id, question, assumption or [], cfg.paths.sessions) - typer.secho(f"OK: question.md aggiornato ({path}).", fg=typer.colors.GREEN) + session_repository(cfg).write_artifact(session_id, "question", render_question_md(question, assumption or [])) + typer.secho("OK: question.md aggiornato.", fg=typer.colors.GREEN) @session_app.command("set-schema-linking") @@ -157,7 +189,7 @@ def set_schema_linking_cmd( from pydantic import ValidationError - from tht.session.store import set_schema_linking + from tht.session.models import SchemaLinking cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) @@ -168,11 +200,12 @@ def set_schema_linking_cmd( typer.secho(f"ERRORE: JSON non valido: {e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=5) try: - path = set_schema_linking(session_id, data, cfg.paths.sessions) + model = SchemaLinking.model_validate(data) + session_repository(cfg).write_artifact(session_id, "schema_linking", json.dumps(model.model_dump(by_alias=True), indent=2, ensure_ascii=False)) except ValidationError as e: typer.secho(f"ERRORE: schema_linking non valido:\n{e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=5) - typer.secho(f"OK: schema_linking.json aggiornato ({path}).", fg=typer.colors.GREEN) + typer.secho("OK: schema_linking.json aggiornato.", fg=typer.colors.GREEN) @session_app.command("sync-schema-linking") @@ -183,16 +216,17 @@ def sync_schema_linking_cmd( """Riproietta schema_linking.json dalle decisioni F4 del ledger (deterministico).""" from pydantic import ValidationError - from tht.session.store import sync_schema_linking + from tht.session.store import sync_schema_linking_snapshot cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) try: - path = sync_schema_linking(session_id, cfg.paths.sessions) + content = sync_schema_linking_snapshot(load_snapshot_or_exit(cfg, session_id)) + session_repository(cfg).write_artifact(session_id, "schema_linking", content) except ValidationError as e: typer.secho(f"ERRORE: schema_linking non valido:\n{e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=5) - typer.secho(f"OK: schema_linking.json riproiettato ({path}).", fg=typer.colors.GREEN) + typer.secho("OK: schema_linking.json riproiettato.", fg=typer.colors.GREEN) @session_app.command("show") @@ -202,28 +236,27 @@ def show_cmd( config: Path = CONFIG_OPT, ) -> None: """Stato della sessione: manifest + decisioni registrate (per la ripresa).""" - from tht.decisions import list_decisions from tht.phase import current_phase cfg = _load_config_or_exit(config) manifest = load_session_or_exit(cfg, session_id) - sdir = session_dir(cfg, session_id) + snapshot = load_snapshot_or_exit(cfg, session_id) if json_out: - has_schema_linking = (sdir / "schema_linking.json").exists() - phase = current_phase(sdir) + has_schema_linking = "schema_linking" in snapshot.artifacts + phase = current_phase(snapshot) data = manifest.model_dump(mode="json", by_alias=True) data["phase"] = phase data["has_schema_linking"] = has_schema_linking # Ledger integrale: il gate lo usa per costruire deterministicamente il # recap delle decisioni nei riepiloghi di fase (v2). data["decisions"] = [ - d.model_dump(mode="json") for d in list_decisions(sdir) + d.model_dump(mode="json") for d in snapshot.decisions ] typer.echo(json.dumps(data, ensure_ascii=False, indent=2)) return - decisions = list_decisions(sdir) + decisions = snapshot.decisions typer.echo(f"id : {manifest.id}") typer.echo(f"stato : {manifest.status}") typer.echo(f"domanda : {manifest.question}") @@ -231,8 +264,7 @@ def show_cmd( typer.echo(f"decisioni: {len(decisions)}") for d in decisions[-10:]: typer.echo(f" [{d.seq}] {d.type}: {d.subject}" + (f" — {d.detail}" if d.detail else "")) - linking = sdir / "schema_linking.json" - typer.echo(f"schema_linking.json: {'presente' if linking.exists() else 'assente'}") + typer.echo(f"schema_linking.json: {'presente' if 'schema_linking' in snapshot.artifacts else 'assente'}") @session_app.command("retrieval-pack") @@ -243,12 +275,10 @@ def retrieval_pack_cmd( """Emette su stdout il retrieval pack persistito della sessione.""" cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - path = session_dir(cfg, session_id) / "retrieval_pack.md" - try: - content = path.read_text() - except OSError as exc: + content = load_snapshot_or_exit(cfg, session_id).artifacts.get("retrieval_pack") + if content is None: typer.secho( - f"ERRORE: retrieval pack non disponibile per la sessione {session_id}: {exc}", + f"ERRORE: retrieval pack non disponibile per la sessione {session_id}", fg=typer.colors.RED, err=True, ) @@ -259,32 +289,32 @@ def retrieval_pack_cmd( @session_app.command("close") def close_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: """Chiude la sessione (status=closed).""" - from tht.session.store import close_session - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - close_session(session_id, cfg.paths.sessions) + snapshot = load_snapshot_or_exit(cfg, session_id) + snapshot.manifest.status = "closed" + session_repository(cfg).save_manifest(snapshot.manifest) typer.secho(f"OK: sessione {session_id} chiusa.", fg=typer.colors.GREEN) @session_app.command("fail") def fail_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: """Registra un arresto del sistema non recuperabile automaticamente.""" - from tht.session.store import fail_session - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - fail_session(session_id, cfg.paths.sessions) + snapshot = load_snapshot_or_exit(cfg, session_id) + snapshot.manifest.status = "failed" + session_repository(cfg).save_manifest(snapshot.manifest) typer.secho(f"OK: sessione {session_id} marcata failed.", fg=typer.colors.RED) @session_app.command("reopen") def reopen_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: """Riapre una sessione per una ripresa manuale.""" - from tht.session.store import reopen_session - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - reopen_session(session_id, cfg.paths.sessions) + snapshot = load_snapshot_or_exit(cfg, session_id) + snapshot.manifest.status = "open" + session_repository(cfg).save_manifest(snapshot.manifest) typer.secho(f"OK: sessione {session_id} riaperta.", fg=typer.colors.GREEN) @@ -295,11 +325,11 @@ def set_name_cmd( config: Path = CONFIG_OPT, ) -> None: """Imposta il nome descrittivo della sessione.""" - from tht.session.store import set_name - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - set_name(session_id, name, cfg.paths.sessions) + snapshot = load_snapshot_or_exit(cfg, session_id) + snapshot.manifest.name = name.strip() or None + session_repository(cfg).save_manifest(snapshot.manifest) typer.secho(f"OK: nome aggiornato per {session_id}.", fg=typer.colors.GREEN) @@ -310,44 +340,42 @@ def set_group_cmd( config: Path = CONFIG_OPT, ) -> None: """Sposta la sessione in un gruppo (o la toglie da ogni gruppo).""" - from tht.session.store import set_group - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - set_group(session_id, group, cfg.paths.sessions) + snapshot = load_snapshot_or_exit(cfg, session_id) + snapshot.manifest.group = group.strip() or None + session_repository(cfg).save_manifest(snapshot.manifest) typer.secho(f"OK: gruppo aggiornato per {session_id}.", fg=typer.colors.GREEN) @session_app.command("archive") def archive_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: """Archivia la sessione (la toglie dalla lista attiva, sola lettura).""" - from tht.session.store import set_archived - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - set_archived(session_id, True, cfg.paths.sessions) + snapshot = load_snapshot_or_exit(cfg, session_id) + snapshot.manifest.archived = True + session_repository(cfg).save_manifest(snapshot.manifest) typer.secho(f"OK: sessione {session_id} archiviata.", fg=typer.colors.GREEN) @session_app.command("unarchive") def unarchive_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: """Ripristina la sessione dall'archivio (non ne cambia la ripristinabilità).""" - from tht.session.store import set_archived - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - set_archived(session_id, False, cfg.paths.sessions) + snapshot = load_snapshot_or_exit(cfg, session_id) + snapshot.manifest.archived = False + session_repository(cfg).save_manifest(snapshot.manifest) typer.secho(f"OK: sessione {session_id} ripristinata.", fg=typer.colors.GREEN) @session_app.command("delete") def delete_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) -> None: """Elimina definitivamente la cartella di sessione.""" - from tht.session.store import delete_session - cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - delete_session(session_id, cfg.paths.sessions) + session_repository(cfg).delete(session_id) typer.secho(f"OK: sessione {session_id} eliminata.", fg=typer.colors.GREEN) @@ -358,11 +386,10 @@ def documents_cmd( config: Path = CONFIG_OPT, ) -> None: """Documenti di sola lettura della sessione (domanda, rivista, schema, SQL, report, decisioni).""" - from tht.session.store import build_documents + from tht.session.store import build_snapshot_documents cfg = _load_config_or_exit(config) - manifest = load_session_or_exit(cfg, session_id) - docs = build_documents(manifest, session_dir(cfg, session_id)) + docs = build_snapshot_documents(load_snapshot_or_exit(cfg, session_id)) if json_out: typer.echo(json.dumps(docs, ensure_ascii=False, indent=2)) return @@ -376,19 +403,18 @@ def session_problems(cfg, session_id: str) -> list[str]: from pydantic import ValidationError - from tht.decisions import list_decisions from tht.session.models import SchemaLinking - sdir = session_dir(cfg, session_id) + snapshot = load_snapshot_or_exit(cfg, session_id) problems: list[str] = [] - if not list_decisions(sdir): + if not snapshot.decisions: problems.append("nessuna decisione registrata (review_decisions.jsonl vuoto o assente)") - linking_path = sdir / "schema_linking.json" - if not linking_path.exists(): + raw_linking = snapshot.artifacts.get("schema_linking") + if raw_linking is None: problems.append("schema_linking.json assente") else: try: - SchemaLinking.model_validate(json.loads(linking_path.read_text())) + SchemaLinking.model_validate(json.loads(raw_linking)) except (json.JSONDecodeError, ValidationError) as e: problems.append(f"schema_linking.json non valido: {e}") return problems @@ -402,7 +428,7 @@ def check_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OPT) cfg = _load_config_or_exit(config) load_session_or_exit(cfg, session_id) - cur = current_phase(session_dir(cfg, session_id)) + cur = current_phase(load_snapshot_or_exit(cfg, session_id)) wf = load_workflow() schema_linking_phase = wf.schema_linking_phase() if cur < schema_linking_phase: @@ -443,7 +469,7 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP do_run, promoted_tables_for, ) - from tht.ctetest import CteError, load_cte_tests + from tht.ctetest import CteError, CteTestRecord, _iter_json_objects from tht.execute import ExecutionError from tht.execute.warnings import plan_warnings, runtime_warnings, static_warnings from tht.report import extract_reviewer_notes, render_validation_report @@ -454,13 +480,13 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP from tht.sqlcheck import validate_sql cfg = _load_config_or_exit(config) - manifest = load_session_or_exit(cfg, session_id) - sdir = session_dir(cfg, session_id) + repository = session_repository(cfg) + snapshot = load_snapshot_or_exit(cfg, session_id) from tht.phase import current_phase from tht.workflow import load_workflow - cur = current_phase(sdir) + cur = current_phase(snapshot) wf = load_workflow() if cur <= wf.max_phase: typer.secho( @@ -474,9 +500,9 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP # --- gate di ingresso --- # Vista effective (D15): un sql_approved/cte ritirato o stale post-rollback non conta. problems = session_problems(cfg, session_id) - decisions = effective_decisions(sdir) - sql_file = sdir / "sql_final.sql" - if not sql_file.exists(): + decisions = effective_decisions(snapshot) + sql = snapshot.artifacts.get("sql_final") + if sql is None: problems.append("sql_final.sql assente") if not any(d.type == "sql_approved" for d in decisions): problems.append( @@ -485,10 +511,11 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP ) # I CTE richiesti sono quelli del PIANO effettivo, non i file glob su disco: un # ctes/*.sql orfano lasciato da un teardown incompleto non deve bloccare il finalize. - plan = effective_cte_plan(sdir) + plan = effective_cte_plan(snapshot) if plan: try: - tested = {r.name for r in load_cte_tests(sdir)} + raw_tests = snapshot.artifacts.get("cte_tests", "") + tested = {r.name for r in map(CteTestRecord.model_validate, _iter_json_objects(raw_tests))} except CteError as e: typer.secho( f"Finalize rifiutato: impossibile leggere cte_tests.json: {e}", @@ -505,7 +532,7 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP raise typer.Exit(code=3) # --- batteria di validazione su sql_final.sql --- - sql = sql_file.read_text() + assert sql is not None check = validate_sql( sql, physical=_load_physical_or_exit(cfg), @@ -527,34 +554,29 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP raise typer.Exit(code=3) # --- validation_report.md (preservando le note del reviewer) --- - report_path = sdir / "validation_report.md" - existing_notes = ( - extract_reviewer_notes(report_path.read_text()) if report_path.exists() else "" - ) - report_path.write_text(render_validation_report( + existing_report = snapshot.artifacts.get("validation_report", "") + existing_notes = extract_reviewer_notes(existing_report) if existing_report else "" + report = render_validation_report( session_id=session_id, check=check, plan=plan, plan_warnings=plan_warnings(plan, cfg.execution), result=result, runtime_warnings=runtime_warnings(result, cfg.execution), static_warnings=static_warnings(check.ast), limit=limit, reviewer_notes=existing_notes, - )) + ) # --- evidence.json --- linking = SchemaLinking.model_validate( - json.loads((sdir / "schema_linking.json").read_text()) + json.loads(snapshot.artifacts["schema_linking"]) ) entries = build_evidence_entries(decisions, linking, cfg.paths.artifacts / "evidence") - (sdir / "evidence.json").write_text(json.dumps(entries, ensure_ascii=False, indent=2)) + evidence = json.dumps(entries, ensure_ascii=False, indent=2) # --- manifest + riepilogo --- - from datetime import UTC, datetime + from tht.session.store import persist_verified_finalization - from tht.session.store import MANIFEST, current_author - - manifest.status = "finalized" - manifest.updated_at = datetime.now(UTC) - manifest.updated_by = current_author() - manifest.to_yaml(sdir / MANIFEST) + persist_verified_finalization( + repository, session_id, validation_report=report, evidence=evidence + ) # --- memoria attiva (parte B): indicizza la coppia domanda->SQL, best-effort --- # Import lazy: memory_cmd importa da session_cmd (un import top-level qui sarebbe # circolare). Qualunque errore (writer key assente, VPN giu', Ollama spento) NON @@ -581,5 +603,5 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP ) typer.secho(f"OK: sessione {session_id} finalizzata. Artefatti:", fg=typer.colors.GREEN) for name in ARTIFACT_FILES: - state = "presente" if (sdir / name).exists() else "assente" + state = "presente" if name not in {"session_manifest.yaml", "review_decisions.jsonl"} else "persistito" typer.echo(f" - {name}: {state}") diff --git a/harness/tht/cli/sql_cmd.py b/harness/tht/cli/sql_cmd.py index e1dc609b..e9a2d247 100644 --- a/harness/tht/cli/sql_cmd.py +++ b/harness/tht/cli/sql_cmd.py @@ -42,12 +42,14 @@ def require_action(cfg, action: str) -> None: def promoted_tables_for(cfg, session_id: str | None) -> set[str] | None: if session_id is None: return None - linking_path = cfg.paths.sessions / session_id / "schema_linking.json" - if not linking_path.exists(): + from tht.cli.session_cmd import load_snapshot_or_exit + + raw = load_snapshot_or_exit(cfg, session_id).artifacts.get("schema_linking") + if raw is None: return None from tht.session.models import SchemaLinking - linking = SchemaLinking.model_validate(json.loads(linking_path.read_text())) + linking = SchemaLinking.model_validate(json.loads(raw)) return { c.name for c in linking.candidates if c.kind == "table" and c.decision == "promoted" @@ -59,12 +61,14 @@ def promoted_tables_for(cfg, session_id: str | None) -> set[str] | None: def promoted_columns_for(cfg, session_id: str | None) -> set[str] | None: if session_id is None: return None - linking_path = cfg.paths.sessions / session_id / "schema_linking.json" - if not linking_path.exists(): + from tht.cli.session_cmd import load_snapshot_or_exit + + raw = load_snapshot_or_exit(cfg, session_id).artifacts.get("schema_linking") + if raw is None: return None from tht.session.models import SchemaLinking - linking = SchemaLinking.model_validate(json.loads(linking_path.read_text())) + linking = SchemaLinking.model_validate(json.loads(raw)) return { c.name for c in linking.candidates if c.kind == "column" and c.decision == "promoted" @@ -182,7 +186,7 @@ def preview_cmd( """Esecuzione controllata con LIMIT iniettato; aggregati mostrati per interi. Se FILE è omesso e --session è fornito, il file viene risolto automaticamente - come /sessions//sql_final.sql (tramite _session_sql_file). + dall'artefatto `sql_final` del repository della sessione. """ from tht.execute import ExecutionError @@ -195,8 +199,7 @@ def preview_cmd( fg=typer.colors.RED, err=True, ) raise typer.Exit(code=1) - resolved = _session_sql_file(cfg, session) - sql = resolved.read_text() + sql = _session_sql(cfg, session) else: sql = _read_sql(file) check = validate_or_exit(cfg, sql, session) @@ -247,15 +250,32 @@ def preview_cmd( typer.secho(f" warning: {w}", fg=typer.colors.YELLOW) -def _session_sql_file(cfg, session_id: str) -> Path: - from tht.cli.session_cmd import load_session_or_exit, session_dir +def _session_sql(cfg, session_id: str) -> str: + from tht.cli.session_cmd import load_snapshot_or_exit - load_session_or_exit(cfg, session_id) - sql_file = session_dir(cfg, session_id) / "sql_final.sql" - if not sql_file.exists(): - typer.secho(f"ERRORE: {sql_file} non trovato.", fg=typer.colors.RED, err=True) + sql = load_snapshot_or_exit(cfg, session_id).artifacts.get("sql_final") + if sql is None: + typer.secho("ERRORE: sql_final.sql non trovato.", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) - return sql_file + return sql + + +@sql_app.command("set-final") +def set_final_cmd( + session: str = typer.Option(..., "--session"), + file: str = typer.Option(..., "--file", help="File SQL, oppure '-' per stdin."), + config: Path = CONFIG_OPT, +) -> None: + """Persist clean final SQL through the configured session repository.""" + import sys + + from tht.cli.session_cmd import load_session_or_exit, session_repository + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session) + sql = sys.stdin.read() if file == "-" else _read_sql(Path(file)) + session_repository(cfg).write_artifact(session, "sql_final", sql) + typer.secho("OK: SQL finale salvato.", fg=typer.colors.GREEN) @sql_app.command("save") @@ -266,9 +286,8 @@ def save_cmd( ) -> None: """Salva una copia di sql_final.sql nel percorso indicato (su richiesta esplicita).""" cfg = _load_config_or_exit(config) - sql_file = _session_sql_file(cfg, session) dest.parent.mkdir(parents=True, exist_ok=True) - dest.write_text(sql_file.read_text()) + dest.write_text(_session_sql(cfg, session)) typer.secho(f"OK: SQL salvato in {dest}", fg=typer.colors.GREEN) @@ -286,8 +305,7 @@ def export_cmd( cfg = _load_config_or_exit(config) require_action(cfg, "export") - sql_file = _session_sql_file(cfg, session) - sql = sql_file.read_text() + sql = _session_sql(cfg, session) validate_or_exit(cfg, sql, session) try: result = do_run(cfg, sql, limit=cfg.execution.max_export_rows) diff --git a/harness/tht/ctetest.py b/harness/tht/ctetest.py index 1b82fb22..fdba06c4 100644 --- a/harness/tht/ctetest.py +++ b/harness/tht/ctetest.py @@ -99,7 +99,10 @@ def load_cte_tests(session_dir: Path) -> list[CteTestRecord]: path = session_dir / CTE_TESTS_FILE if not path.exists(): return [] - text = path.read_text() + return load_cte_tests_text(path.read_text()) + + +def load_cte_tests_text(text: str) -> list[CteTestRecord]: try: objs = list(_iter_json_objects(text)) except json.JSONDecodeError as e: @@ -107,6 +110,11 @@ def load_cte_tests(session_dir: Path) -> list[CteTestRecord]: return [CteTestRecord.model_validate(o) for o in objs] +def append_cte_test_snapshot(repository, snapshot, record: CteTestRecord) -> None: + previous = snapshot.artifacts.get("cte_tests", "") + repository.write_artifact(snapshot.manifest.id, "cte_tests", previous + record.model_dump_json() + "\n") + + def append_cte_test(session_dir: Path, record: CteTestRecord) -> None: """Append atomico (una riga JSON per esito): scritture concorrenti sulla stessa sessione non si sovrascrivono, a differenza del rewrite dell'array.""" diff --git a/harness/tht/memory.py b/harness/tht/memory.py index 8a6acc78..cdd2ed2d 100644 --- a/harness/tht/memory.py +++ b/harness/tht/memory.py @@ -186,6 +186,14 @@ def promote( return promoted +def promote_snapshot(snapshot, *, seqs: list[int] | None, registry_path: Path) -> list[MemoryRecord]: + existing = load_registry(registry_path) + promoted = _compute_promotions(snapshot, snapshot.manifest, seqs=seqs, existing=existing) + if promoted: + save_registry(existing + promoted, registry_path) + return promoted + + def reusable_promotions( session_dir: Path, manifest: SessionManifest, registry_path: Path ) -> list[MemoryRecord]: @@ -203,6 +211,14 @@ def reusable_promotions( ] +def reusable_promotions_snapshot(snapshot, registry_path: Path) -> list[MemoryRecord]: + from tht.phase import effective_decisions + + cand = _compute_promotions(snapshot, snapshot.manifest, seqs=None, existing=load_registry(registry_path)) + declined = declined_promotion_seqs(effective_decisions(snapshot)) + return [c for c in cand if c.type in REUSABLE_TYPES and c.decision_seq not in declined] + + def preview_promotions( session_dir: Path, manifest: SessionManifest, registry_path: Path ) -> list[MemoryRecord]: @@ -212,6 +228,10 @@ def preview_promotions( ] +def preview_promotions_snapshot(snapshot, registry_path: Path) -> list[MemoryRecord]: + return reusable_promotions_snapshot(snapshot, registry_path)[:MAX_PROMOTION_CANDIDATES] + + def memory_vector_records(records: list[MemoryRecord]) -> list[VectorRecord]: out: list[VectorRecord] = [] for r in records: diff --git a/harness/tht/phase.py b/harness/tht/phase.py index dbcc8d7a..41df97e4 100644 --- a/harness/tht/phase.py +++ b/harness/tht/phase.py @@ -17,7 +17,7 @@ from pathlib import Path from pydantic import ValidationError from tht.decisions import DecisionRecord, list_decisions -from tht.session.models import SchemaLinking +from tht.session.models import SchemaLinking, SessionSnapshot from tht.workflow import load_workflow @@ -31,11 +31,22 @@ def _phase_num(subject: str) -> int | None: return None -def _audit_excluding_retracted(session_dir: Path) -> list[DecisionRecord]: +def _decisions(source: Path | SessionSnapshot) -> list[DecisionRecord]: + return list(source.decisions) if isinstance(source, SessionSnapshot) else list_decisions(source) + + +def _artifact(source: Path | SessionSnapshot, key: str, filename: str) -> str | None: + if isinstance(source, SessionSnapshot): + return source.artifacts.get(key) + path = source / filename + return path.read_text() if path.exists() else None + + +def _audit_excluding_retracted(source: Path | SessionSnapshot) -> list[DecisionRecord]: """Tutto il ledger (append-only) tranne le decisioni ritirate e i marker di ritrazione. Base per il fold di current_phase: il guard 'n == cur' del fold e' gia' reopen-aware (una phase_approved:N dopo un reopen a M list[DecisionRecord]: ] -def current_phase(session_dir: Path) -> int: +def current_phase(source: Path | SessionSnapshot) -> int: """Fase corrente come fold cronologico sull'audit (con ritirate escluse). cur parte da 1; ogni phase_approved/phase_auto_approved per la fase CORRENTE avanza @@ -57,7 +68,7 @@ def current_phase(session_dir: Path) -> int: wf = load_workflow() max_plus_one = wf.max_phase + 1 cur = 1 - for d in _audit_excluding_retracted(session_dir): + for d in _audit_excluding_retracted(source): n = _phase_num(d.subject) if n is None: continue @@ -68,7 +79,7 @@ def current_phase(session_dir: Path) -> int: return cur -def effective_decisions(session_dir: Path) -> list[DecisionRecord]: +def effective_decisions(source: Path | SessionSnapshot) -> list[DecisionRecord]: """La vista canonica 'effective as of pointer'. TUTTI gli helper non-fold devono usare questa. Semantica: una decisione e' effective se appartiene a una fase <= current_phase. @@ -81,9 +92,9 @@ def effective_decisions(session_dir: Path) -> list[DecisionRecord]: Inoltre esclude le decisioni ritirate (decision_retracted) e i marker stessi. """ - cur = current_phase(session_dir) + cur = current_phase(source) out: list[DecisionRecord] = [] - for d in _audit_excluding_retracted(session_dir): + for d in _audit_excluding_retracted(source): n = _phase_num(d.subject) # Per le decisioni con subject "a nome" (es. cte_approved -> nome CTE, evidence_* # -> id evidence) il subject non porta la fase: si usa la fase emittente registrata @@ -106,9 +117,9 @@ _META_TYPES = frozenset( ) -def substantive_count_current_phase(session_dir: Path) -> int: +def substantive_count_current_phase(source: Path | SessionSnapshot) -> int: """Numero di decisioni sostanziali dall'ultimo confine di fase (vista effective).""" - decs = effective_decisions(session_dir) + decs = effective_decisions(source) start = 0 for i, d in enumerate(decs): if d.type in _BOUNDARY_TYPES: @@ -116,14 +127,14 @@ def substantive_count_current_phase(session_dir: Path) -> int: return sum(1 for d in decs[start:] if d.type not in _META_TYPES) -def auto_advance_eligible(session_dir: Path) -> bool: +def auto_advance_eligible(source: Path | SessionSnapshot) -> bool: """Vero sse la fase corrente puo' auto-avanzare (zero decisioni sostanziali + prereq ok).""" - cur = current_phase(session_dir) + cur = current_phase(source) if cur not in _AUTO_ADVANCE_PHASES: return False - if substantive_count_current_phase(session_dir) > 0: + if substantive_count_current_phase(source) > 0: return False - return not advance_problems(session_dir, cur) + return not advance_problems(source, cur) # --- CTE helpers (consultano effective_decisions) --------------------------- @@ -131,21 +142,21 @@ def auto_advance_eligible(session_dir: Path) -> bool: CTE_PLAN_FILE = "cte_plan.json" -def cte_plan(session_dir: Path) -> list[str]: - path = session_dir / CTE_PLAN_FILE - if not path.exists(): +def cte_plan(source: Path | SessionSnapshot) -> list[str]: + raw = _artifact(source, "cte_plan", CTE_PLAN_FILE) + if raw is None: return [] - return json.loads(path.read_text()) + return json.loads(raw) -def approved_ctes(session_dir: Path) -> set[str]: +def approved_ctes(source: Path | SessionSnapshot) -> set[str]: """Insieme dei CTE approvati, dalla vista effective (esclude stale post-reopen).""" - return {d.subject for d in effective_decisions(session_dir) if d.type == "cte_approved"} + return {d.subject for d in effective_decisions(source) if d.type == "cte_approved"} -def next_cte(session_dir: Path) -> str | None: - approved = approved_ctes(session_dir) - for name in cte_plan(session_dir): +def next_cte(source: Path | SessionSnapshot) -> str | None: + approved = approved_ctes(source) + for name in cte_plan(source): if name not in approved: return name return None @@ -153,18 +164,18 @@ def next_cte(session_dir: Path) -> str | None: # --- advance_problems (ladder if-phase-N che consulta effective_decisions) --- -def _has_decision(session_dir: Path, type_: str) -> bool: - return any(d.type == type_ for d in effective_decisions(session_dir)) +def _has_decision(source: Path | SessionSnapshot, type_: str) -> bool: + return any(d.type == type_ for d in effective_decisions(source)) -def _has_decision_subject(session_dir: Path, type_: str, subject: str) -> bool: +def _has_decision_subject(source: Path | SessionSnapshot, type_: str, subject: str) -> bool: return any( d.type == type_ and d.subject == subject - for d in effective_decisions(session_dir) + for d in effective_decisions(source) ) -def advance_problems(session_dir: Path, phase: int) -> list[str]: +def advance_problems(source: Path | SessionSnapshot, phase: int) -> list[str]: """Prerequisiti minimi per chiudere `phase` (lista vuota = ok). Ladder if-phase-N (Strada 2): la logica specifica resta, ma ogni lettura passa per @@ -172,19 +183,19 @@ def advance_problems(session_dir: Path, phase: int) -> list[str]: l'evaluator generico (F2 pieno) entra in un secondo momento. """ problems: list[str] = [] - if phase == 3 and not _has_decision(session_dir, "question_rewritten"): + if phase == 3 and not _has_decision(source, "question_rewritten"): problems.append("manca la decisione question_rewritten (Fase 3)") if phase == 5: - path = session_dir / "schema_linking.json" - if not path.exists(): + raw = _artifact(source, "schema_linking", "schema_linking.json") + if raw is None: problems.append("schema_linking.json assente (Fase 5)") else: try: - SchemaLinking.model_validate(json.loads(path.read_text())) + SchemaLinking.model_validate(json.loads(raw)) except (json.JSONDecodeError, ValidationError) as e: problems.append(f"schema_linking.json non valido (Fase 5): {e}") - if phase == 6 and not _has_decision_subject(session_dir, "phase_skipped", "phase:6"): - plan = cte_plan(session_dir) + if phase == 6 and not _has_decision_subject(source, "phase_skipped", "phase:6"): + plan = cte_plan(source) if not plan: problems.append( "Fase 6: nessun piano CTE (cte_plan.json) e nessun salto esplicito. " @@ -192,14 +203,14 @@ def advance_problems(session_dir: Path, phase: int) -> list[str]: "registrando una decisione phase_skipped subject phase:6." ) else: - nc = next_cte(session_dir) + nc = next_cte(source) if nc is not None: problems.append(f"CTE non ancora approvato: {nc} (Fase 6)") - if phase == 7 and not _has_decision(session_dir, "sql_approved"): + if phase == 7 and not _has_decision(source, "sql_approved"): problems.append("manca la decisione sql_approved (Fase 7)") if phase == 8 and not any( d.type in ("datamart_requested", "datamart_declined") - for d in effective_decisions(session_dir) + for d in effective_decisions(source) ): problems.append( "Fase 8: nessuna risposta sulla generazione dbt del datamart " diff --git a/harness/tht/session/filesystem_repository.py b/harness/tht/session/filesystem_repository.py index c30f95f7..98794e80 100644 --- a/harness/tht/session/filesystem_repository.py +++ b/harness/tht/session/filesystem_repository.py @@ -27,19 +27,24 @@ _ARTIFACT_FILES = { "validation_report": "validation_report.md", "retrieval_pack": "retrieval_pack.md", "cte_tests": "cte_tests.json", + "cte_plan": "cte_plan.json", + "cte_plan_doc": "cte_plan_doc.json", } _ARTIFACT_KEYS = {filename: key for key, filename in _ARTIFACT_FILES.items()} _SAFE_CTE_NAME = re.compile(r"[A-Za-z0-9_-]+\Z") +_LEGACY_SESSION_ID = re.compile(r"[A-Za-z0-9][A-Za-z0-9._-]{0,127}\Z") class FilesystemSessionRepository: """Readable phase documents under one private local user's ThothII home.""" - def __init__(self, home: Path, workspace: str, principal: PrincipalContext): + def __init__( + self, home: Path, workspace: str, principal: PrincipalContext, *, root: Path | None = None + ): self.home = Path(home).expanduser() self.workspace = workspace self.principal = principal - self.root = self.home / "workspaces" / workspace / "sessions" + self.root = Path(root) if root is not None else self.home / "workspaces" / workspace / "sessions" @property def preferences_path(self) -> Path: @@ -49,6 +54,8 @@ class FilesystemSessionRepository: return self.home / "principals" / digest / "preferences.json" def create(self, manifest: SessionManifest) -> SessionSnapshot: + # New records are opaque UUIDv4 only. Historical timestamp ids remain + # readable/mutable through _session_dir but are never created here. self._require_uuid4(manifest.id) session_dir = self.root / manifest.id with self._lock(session_dir): @@ -72,6 +79,18 @@ class FilesystemSessionRepository: decisions=list_decisions(session_dir), ) + def list(self) -> list[SessionSnapshot]: + if not self.root.exists(): + return [] + snapshots = [] + for path in self.root.iterdir(): + if path.is_dir() and (path / MANIFEST).exists(): + try: + snapshots.append(self.get(path.name)) + except SessionError: + continue + return sorted(snapshots, key=lambda item: item.manifest.created_at, reverse=True) + def save_manifest(self, manifest: SessionManifest) -> SessionSnapshot: session_dir = self._session_dir(manifest.id) with self._lock(session_dir): @@ -91,6 +110,12 @@ class FilesystemSessionRepository: self._require_existing(session_dir, session_id) self._write_private_text(self._artifact_path(session_dir, key), content) + def delete_artifact(self, session_id: str, key: str) -> None: + session_dir = self._session_dir(session_id) + with self._lock(session_dir): + self._require_existing(session_dir, session_id) + self._artifact_path(session_dir, key).unlink(missing_ok=True) + def append_decisions( self, session_id: str, decisions: Sequence[DecisionInput | dict] ) -> list[DecisionRecord]: @@ -101,6 +126,21 @@ class FilesystemSessionRepository: # separate lock also keeps direct workflow callers safe during Task 3. return append_decisions(session_dir, list(decisions)) + def finalize(self, manifest: SessionManifest, artifacts: dict[str, str]) -> SessionSnapshot: + """Commit verified final artifacts before publishing finalized status. + + The finalized manifest is the filesystem commit marker: a crash can leave + an open session with already-rendered artifacts, but never a finalized + session without its report and evidence. + """ + session_dir = self._session_dir(manifest.id) + with self._lock(session_dir): + self._require_existing(session_dir, manifest.id) + for key, content in artifacts.items(): + self._write_private_text(self._artifact_path(session_dir, key), content) + self._write_manifest(session_dir, manifest) + return self.get(manifest.id) + def get_preferences(self) -> dict: path = self.preferences_path if not path.exists(): @@ -125,9 +165,18 @@ class FilesystemSessionRepository: shutil.rmtree(session_dir) def _session_dir(self, session_id: str) -> Path: - self._require_uuid4(session_id) + self._require_session_id(session_id) return self.root / session_id + @classmethod + def _require_session_id(cls, session_id: str) -> None: + try: + cls._require_uuid4(session_id) + except SessionError: + if session_id not in {".", ".."} and _LEGACY_SESSION_ID.fullmatch(session_id): + return + raise + @staticmethod def _require_uuid4(session_id: str) -> None: try: diff --git a/harness/tht/session/postgres_repository.py b/harness/tht/session/postgres_repository.py index 6febba4a..65ad689b 100644 --- a/harness/tht/session/postgres_repository.py +++ b/harness/tht/session/postgres_repository.py @@ -33,6 +33,8 @@ _ARTIFACT_KEYS = { "validation_report", "retrieval_pack", "cte_tests", + "cte_plan", + "cte_plan_doc", } _LOCK_KEY = 8_420_613_069_444_020_731 _RUNTIME_ROLE = "thoth_sessions_runtime" @@ -278,6 +280,13 @@ class PostgresSessionRepository: decisions=decisions, ) + def list(self) -> list[SessionSnapshot]: + with self._transaction() as connection: + ids = [row[0] for row in connection.execute( + text("SELECT id FROM thoth_sessions.sessions ORDER BY created_at DESC") + ).all()] + return [self.get(session_id) for session_id in ids] + def save_manifest(self, manifest: SessionManifest) -> SessionSnapshot: self._require_uuid4(manifest.id) with self._transaction() as connection: @@ -322,6 +331,17 @@ class PostgresSessionRepository: {"id": session_id, "key": key, "content": content}, ) + def delete_artifact(self, session_id: str, key: str) -> None: + self._require_uuid4(session_id) + self._require_artifact_key(key) + with self._transaction() as connection: + self._lock_session(connection, session_id) + self._require_session(connection, session_id) + connection.execute( + text("DELETE FROM thoth_sessions.session_artifacts WHERE session_id = :id AND artifact_key = :key"), + {"id": session_id, "key": key}, + ) + def append_decisions( self, session_id: str, decisions: Sequence[DecisionInput | dict] ) -> list[DecisionRecord]: @@ -357,6 +377,36 @@ class PostgresSessionRepository: ) return records + def finalize(self, manifest: SessionManifest, artifacts: dict[str, str]) -> SessionSnapshot: + """Atomically publish DWH-verified artifacts and the final manifest.""" + self._require_uuid4(manifest.id) + for key in artifacts: + self._require_artifact_key(key) + with self._transaction() as connection: + self._lock_session(connection, manifest.id) + self._require_session(connection, manifest.id) + for key, content in artifacts.items(): + connection.execute( + text( + "INSERT INTO thoth_sessions.session_artifacts (session_id, artifact_key, content) " + "VALUES (:id, :key, :content) " + "ON CONFLICT (session_id, artifact_key) DO UPDATE " + "SET content = EXCLUDED.content, updated_at = pg_catalog.now()" + ), + {"id": manifest.id, "key": key, "content": content}, + ) + result = connection.execute( + text( + "UPDATE thoth_sessions.sessions " + "SET manifest = CAST(:manifest AS jsonb), updated_at = pg_catalog.now() " + "WHERE id = :id" + ), + {"id": manifest.id, "manifest": json.dumps(manifest.model_dump(by_alias=True, mode="json"))}, + ) + if result.rowcount != 1: + raise SessionError(f"Sessione non trovata: {manifest.id}") + return self.get(manifest.id) + def get_preferences(self) -> dict: with self._transaction() as connection: principal_id = self._upsert_principal(connection) diff --git a/harness/tht/session/repository.py b/harness/tht/session/repository.py index 8e588720..5357a27a 100644 --- a/harness/tht/session/repository.py +++ b/harness/tht/session/repository.py @@ -2,10 +2,12 @@ from __future__ import annotations +import os from typing import Protocol, Sequence from tht.decisions import DecisionInput, DecisionRecord from tht.session.models import PrincipalContext, SessionManifest, SessionSnapshot +from tht.session.store import SessionError class SessionRepository(Protocol): @@ -17,16 +19,22 @@ class SessionRepository(Protocol): def get(self, session_id: str) -> SessionSnapshot: ... + def list(self) -> list[SessionSnapshot]: ... + def save_manifest(self, manifest: SessionManifest) -> SessionSnapshot: ... def read_artifact(self, session_id: str, key: str) -> str | None: ... def write_artifact(self, session_id: str, key: str, content: str) -> None: ... + def delete_artifact(self, session_id: str, key: str) -> None: ... + def append_decisions( self, session_id: str, decisions: Sequence[DecisionInput | dict] ) -> list[DecisionRecord]: ... + def finalize(self, manifest: SessionManifest, artifacts: dict[str, str]) -> SessionSnapshot: ... + def get_preferences(self) -> dict: ... def set_preferences(self, preferences: dict) -> None: ... @@ -34,8 +42,35 @@ class SessionRepository(Protocol): def delete(self, session_id: str) -> None: ... -def build_session_repository(config, principal: PrincipalContext, *, home=None) -> SessionRepository: +def resolve_principal(config) -> PrincipalContext: + """Resolve the only principal source permitted for this CLI invocation. + + The backend injects these values from its authenticated request context before + spawning ``tht``. A server-storage command without them must fail closed; + substituting a workstation identity would cross user ownership boundaries. + """ + if getattr(config, "session_storage", None) is None: + from tht.session.models import local_principal + + return local_principal() + issuer = os.environ.get("THT_PRINCIPAL_ISSUER", "").strip() + subject = os.environ.get("THT_PRINCIPAL_SUBJECT", "").strip() + if not issuer or not subject: + raise SessionError( + "THT_PRINCIPAL_ISSUER e THT_PRINCIPAL_SUBJECT sono obbligatori per session storage server" + ) + display_name = os.environ.get("THT_PRINCIPAL_DISPLAY_NAME", "").strip() or None + is_admin = os.environ.get("THT_PRINCIPAL_IS_ADMIN", "").strip().lower() in {"1", "true"} + return PrincipalContext( + issuer=issuer, subject=subject, display_name=display_name, is_admin=is_admin + ) + + +def build_session_repository( + config, principal: PrincipalContext | None = None, *, home=None +) -> SessionRepository: """Build the configured private session persistence adapter.""" + principal = principal or resolve_principal(config) session_storage = getattr(config, "session_storage", None) if session_storage is not None: from tht.session.postgres_repository import PostgresSessionRepository @@ -46,4 +81,7 @@ def build_session_repository(config, principal: PrincipalContext, *, home=None) from tht.session.filesystem_repository import FilesystemSessionRepository workspace = getattr(config, "_workspace_id", "default") - return FilesystemSessionRepository(home or local_tht_home(), workspace, principal) + return FilesystemSessionRepository( + home or local_tht_home(), workspace, principal, + root=None if home is not None else getattr(config.paths, "sessions", None), + ) diff --git a/harness/tht/session/store.py b/harness/tht/session/store.py index b8e2c834..01427051 100644 --- a/harness/tht/session/store.py +++ b/harness/tht/session/store.py @@ -130,6 +130,35 @@ def create_session( return manifest +def new_session_manifest( + question: str, db: DatabaseConfig, *, provider=None, model=None, thinking=None, name=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, name=name, + ) + + +def build_snapshot_documents(snapshot) -> list[dict]: + docs = [{"phase": "—", "key": "question", "title": "Original question", "format": "text", "content": snapshot.manifest.question}] + spec = [ + ("question", "F3", "revised_question", "Revised question", "markdown"), + ("schema_linking", "F4", "schema_linking", "Schema linking", "schema-linking"), + ("sql_final", "F7", "sql", "Final SQL", "sql"), + ("validation_report", "finalize", "validation_report", "Validation report", "markdown"), + ] + for artifact, phase, key, title, fmt in spec: + if artifact in snapshot.artifacts: + docs.append({"phase": phase, "key": key, "title": title, "format": fmt, "content": snapshot.artifacts[artifact]}) + if snapshot.decisions: + docs.append({"phase": "—", "key": "decisions", "title": "Decisions", "format": "decisions", "content": "\n".join(d.model_dump_json() for d in snapshot.decisions) + "\n"}) + return docs + + def touch_manifest( session_id: str, sessions_root: Path, *, updated_by: str | None = None ) -> SessionManifest: @@ -227,6 +256,33 @@ def sync_schema_linking(session_id: str, sessions_root: Path) -> Path: 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(): @@ -298,6 +354,30 @@ def delete_session(session_id: str, sessions_root: Path) -> None: 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]: """Ordered, read-only document bundle for the UI panel. Only documents that exist on disk are returned. CTE artifacts (F6) are intentionally excluded (intermediate).""" diff --git a/harness/tht/solved.py b/harness/tht/solved.py index e62910de..4bff7521 100644 --- a/harness/tht/solved.py +++ b/harness/tht/solved.py @@ -84,3 +84,20 @@ def build_solved_record(session_dir, manifest, promoted_tables) -> VectorRecord: sql=sql_file.read_text().strip(), tables=sorted(promoted_tables or set()), ) + + +def build_solved_snapshot(snapshot, promoted_tables) -> VectorRecord: + from tht.memory import question_context + from tht.phase import effective_decisions + + sql = snapshot.artifacts.get("sql_final") + if sql is None: + raise SolvedIndexError("sql_final.sql assente") + decisions = effective_decisions(snapshot) + if not any(d.type == "sql_approved" for d in decisions): + raise SolvedIndexError("decisione sql_approved assente") + return solved_question_record( + session_id=snapshot.manifest.id, + question=question_context(decisions, snapshot.manifest), + sql=sql.strip(), tables=sorted(promoted_tables or set()), + ) diff --git a/harness/tht/taskdoc.py b/harness/tht/taskdoc.py index d5a6f32b..2aeea1e5 100644 --- a/harness/tht/taskdoc.py +++ b/harness/tht/taskdoc.py @@ -15,6 +15,7 @@ from dataclasses import dataclass from pathlib import Path from tht.phase import effective_decisions +from tht.session.models import SessionSnapshot from tht.workflow import load_workflow MAX_BODY_BYTES = 80_000 # ~20k token (target per task document di una fase) @@ -50,7 +51,7 @@ def _slice_schema_linking(raw: str, promoted_tables: list[str] | None) -> str: def generate_task_doc( - session_dir: Path | str, + session_dir: Path | str | SessionSnapshot, phase: int, promoted_tables: list[str] | None = None, ) -> TaskDoc: @@ -62,20 +63,21 @@ def generate_task_doc( - Brief delle decisioni effective (esclude stale post-rollback, esclude ritirate). - Header del task con il numero/nome della fase. """ - session_dir = Path(session_dir) + snapshot = session_dir if isinstance(session_dir, SessionSnapshot) else None + session_dir = None if snapshot is not None else Path(session_dir) parts: list[str] = [] - q = session_dir / "question.md" - if q.exists(): - parts.append("## Domanda\n" + q.read_text()) + question = snapshot.artifacts.get("question") if snapshot else (session_dir / "question.md").read_text() if (session_dir / "question.md").exists() else None + if question: + parts.append("## Domanda\n" + question) - sl = session_dir / "schema_linking.json" - if sl.exists() and phase >= 4: - sliced = _slice_schema_linking(sl.read_text(), promoted_tables) + linking = snapshot.artifacts.get("schema_linking") if snapshot else (session_dir / "schema_linking.json").read_text() if (session_dir / "schema_linking.json").exists() else None + if linking and phase >= 4: + sliced = _slice_schema_linking(linking, promoted_tables) parts.append("## Schema linking (deciso)\n```json\n" + sliced + "\n```") # Brief decisioni effective (D15-aware) - eff = effective_decisions(session_dir) + eff = effective_decisions(snapshot or session_dir) if eff: lines = [f"- {d.type} | {d.subject} | {d.detail}" for d in eff] parts.append("## Decisioni effettive (effective)\n" + "\n".join(lines)) diff --git a/harness/tht/teardown.py b/harness/tht/teardown.py index 0072c549..cbd01ded 100644 --- a/harness/tht/teardown.py +++ b/harness/tht/teardown.py @@ -52,3 +52,31 @@ def teardown_to_phase(session_dir: str | Path, target_phase: int) -> TeardownRep target.unlink() report.deleted_files.append(artifact) return report + + +def teardown_snapshot(repository, snapshot, target_phase: int) -> TeardownReport: + """Repository equivalent of teardown_to_phase, including orphaned CTE blobs.""" + wf = load_workflow() + report = TeardownReport(target_phase=target_phase) + files = { + "question.md": "question", "schema_linking.json": "schema_linking", + "evidence.json": "evidence", "cte_tests.json": "cte_tests", + "sql_final.sql": "sql_final", "validation_report.md": "validation_report", + "retrieval_pack.md": "retrieval_pack", "cte_plan.json": "cte_plan", + "cte_plan_doc.json": "cte_plan_doc", + } + for phase in wf.phases: + if phase.num <= target_phase: + continue + for artifact in phase.artifacts_out: + if artifact.endswith("/"): + for key in list(snapshot.artifacts): + if key.startswith("cte_sql:"): + repository.delete_artifact(snapshot.manifest.id, key) + report.deleted_files.append(key.removeprefix("cte_sql:") + ".sql") + else: + key = files.get(artifact) + if key and key in snapshot.artifacts: + repository.delete_artifact(snapshot.manifest.id, key) + report.deleted_files.append(artifact) + return report