From 333874a755a372bf403901a880d24ec02ecb1110 Mon Sep 17 00:00:00 2001 From: User Date: Tue, 14 Jul 2026 15:24:06 +0200 Subject: [PATCH] fix: make join approval atomic --- backend/test/routes-sessions.test.ts | 27 ++++++ .../2026-07-14-workflow-ui-regressions.md | 11 ++- ...26-07-14-workflow-ui-regressions-design.md | 8 +- .../gate/__tests__/gate_join_review.test.js | 53 +++++++++-- harness/.pi/extensions/tht-gate.js | 44 +++++++-- harness/.pi/skills/tht-sessione/SKILL.md | 5 +- harness/tests/test_decision_join_set_cli.py | 93 +++++++++++++++++++ harness/tht/cli/decision_cmd.py | 37 ++++++++ harness/tht/decisions.py | 70 +++++++++++--- 9 files changed, 314 insertions(+), 34 deletions(-) create mode 100644 harness/tests/test_decision_join_set_cli.py diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 4d4f8dab..07f90b36 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -106,6 +106,33 @@ test("POST /sessions/:id/resume configura Pi con il thinking persistito", async expect(configured.thinking).toBe("medium"); }); +test("POST /sessions/:id/resume usa il thinking globale se manca nel manifest", async () => { + let configured: any; + const bridge = { onClientEvent: () => {}, emitClientEvent: () => {} }; + const runtime = { bridge } as any; + const mgr = { + get: () => undefined, + createFor: () => runtime, + configure: async (_rt: any, options: any) => { configured = options; }, + start: () => {}, + teardown: () => {}, + } as any; + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + mgr, + thtRunner: { + ollamaEnsure: async () => ({ ok: true }), + sessionShow: async () => ({ status: "open", archived: false }), + reopenSession: async () => {}, + } as any, + getSettings: () => ({ workspace: "psd", thinking: "low" }) as any, + }); + + await app.inject({ method: "POST", url: "/sessions/s-thinking/resume" }); + await new Promise((resolve) => setImmediate(resolve)); + + expect(configured.thinking).toBe("low"); +}); + test("POST /sessions/:id/response inoltra al bridge (no error)", async () => { const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { diff --git a/docs/superpowers/plans/2026-07-14-workflow-ui-regressions.md b/docs/superpowers/plans/2026-07-14-workflow-ui-regressions.md index dc395ebc..af90f5b2 100644 --- a/docs/superpowers/plans/2026-07-14-workflow-ui-regressions.md +++ b/docs/superpowers/plans/2026-07-14-workflow-ui-regressions.md @@ -93,6 +93,9 @@ Expected: PASS. - Modify: `harness/.pi/extensions/tht-gate.js` - Create: `harness/.pi/extensions/gate/__tests__/gate_join_review.test.js` - Modify: `harness/.pi/skills/tht-sessione/SKILL.md` +- Modify: `harness/tht/decisions.py` +- Modify: `harness/tht/cli/decision_cmd.py` +- Create: `harness/tests/test_decision_join_set_cli.py` - Create: `frontend/src/widgets/JoinReviewWidget.tsx` - Create: `frontend/src/widgets/JoinReviewWidget.test.tsx` - Modify: `frontend/src/widgets/index.ts` @@ -108,7 +111,8 @@ Expected: PASS. Assert the builder emits read-only option details, `confirm_label: "Continue"`, and reserved controls. Exercise `reviewer_decide` with only `join_modified` decisions; queue a Continue response containing -all ids and assert every `tht decision add` call occurs. Queue `control:"freetext"` and assert no +all ids and assert one atomic `tht decision add-join-set` call occurs. Queue malformed and partial +responses and assert the widget is presented again. Queue `control:"freetext"` and assert no decision is written. - [x] **Step 2: Run harness tests and verify RED** @@ -125,7 +129,8 @@ Add `buildJoinReviewRequest`. In `reviewer_decide`, detect a non-empty, join-onl const joinOnly = opts.length > 0 && opts.every((o) => o.decision.type === "join_modified"); ``` -Emit `join-review` with `detail` and `rationale`; after Continue persist all original decisions. +Emit `join-review` with `detail` and `rationale`; accept only a response carrying the exact complete +id set. After Continue persist all original decisions through atomic `decision add-join-set`. On free text, return feedback without persistence. Keep all other decisions on `multiselect`. - [x] **Step 4: Write failing frontend widget tests** @@ -269,7 +274,7 @@ Expected: PASS. - Running Compose services must use the newly built `thothii-core:local` and `thothii-frontend:local` image ids. -- [ ] **Step 1: Run all regression suites** +- [x] **Step 1: Run all regression suites** Run: `cd backend && npx vitest run && npx tsc --noEmit -p . && npm run build` diff --git a/docs/superpowers/specs/2026-07-14-workflow-ui-regressions-design.md b/docs/superpowers/specs/2026-07-14-workflow-ui-regressions-design.md index 1d011f0b..67890ba5 100644 --- a/docs/superpowers/specs/2026-07-14-workflow-ui-regressions-design.md +++ b/docs/superpowers/specs/2026-07-14-workflow-ui-regressions-design.md @@ -35,9 +35,11 @@ the configured model activity. F4 calls to `reviewer_decide` whose merit decisions are all `join_modified` become a read-only `join-review` widget. Each proposed join is rendered as an informational card with its name, -join expression, and rationale. `Continue` returns every join id; the gate persists all decisions. -`Other — specify` returns textual feedback without persisting the current proposal, so the model -must revise and present the complete join set again. Mixed join/non-join calls remain regular +join expression, and rationale. `Continue` returns every join id; the gate accepts only that exact +complete response and persists all decisions through one atomic ledger replacement. Malformed or +partial responses re-present the widget. `Other — specify` returns textual feedback without +persisting the current proposal, so the model must revise and present the complete join set again. +Mixed join/non-join calls remain regular multiselects, and the skill instructs the model to keep joins in a separate call. ### CTE presentation diff --git a/harness/.pi/extensions/gate/__tests__/gate_join_review.test.js b/harness/.pi/extensions/gate/__tests__/gate_join_review.test.js index 6f0b0907..617d1af0 100644 --- a/harness/.pi/extensions/gate/__tests__/gate_join_review.test.js +++ b/harness/.pi/extensions/gate/__tests__/gate_join_review.test.js @@ -33,9 +33,10 @@ const JOIN_OPTIONS = [ }, ]; -function shellStub(calls) { - return (_file, args) => { +function shellStub(calls, inputs = []) { + return (_file, args, options = {}) => { calls.push(args.join(" ")); + if (options.input) inputs.push(options.input); if (args[0] === "phase" && args[1] === "meta") { return JSON.stringify({ phases: [{ num: 4, id: "F4" }] }); } @@ -46,8 +47,9 @@ function shellStub(calls) { test("join-only reviewer_decide is read-only and Continue persists every join", async () => { const calls = []; + const inputs = []; const original = cp.execFileSync; - cp.execFileSync = shellStub(calls); + cp.execFileSync = shellStub(calls, inputs); try { const gate = require(GATE); const { createFakePi } = require("./fake_pi_runtime.js"); @@ -77,9 +79,48 @@ test("join-only reviewer_decide is read-only and Continue persists every join", assert.equal(descriptor.widget, "join-review"); assert.equal(descriptor.options[0].detail, JOIN_OPTIONS[0].decision.detail); assert.equal(descriptor.options[0].rationale, JOIN_OPTIONS[0].decision.rationale); - const decisions = calls.filter((call) => call.startsWith("decision add")); - assert.equal(decisions.length, 2); - assert.ok(decisions.every((call) => call.includes("--type join_modified"))); + const decisions = calls.filter((call) => call.startsWith("decision add-join-set")); + assert.equal(decisions.length, 1); + assert.deepEqual(JSON.parse(inputs[0]), JOIN_OPTIONS.map((option) => option.decision)); + } finally { + cp.execFileSync = original; + } +}); + +test("partial join-review responses are rejected and the widget is presented again", async () => { + const calls = []; + const original = cp.execFileSync; + cp.execFileSync = shellStub(calls); + 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); + + let attempts = 0; + ctx.ui.input = async (title) => { + const descriptor = JSON.parse(title); + attempts += 1; + return JSON.stringify({ + id: descriptor.id, + kind: attempts === 1 ? "multiselect" : "join-review", + choices: attempts === 2 + ? [descriptor.options[0].id] + : descriptor.options.map((option) => option.id), + }); + }; + + const tool = tools.get("reviewer_decide"); + await tool.def.execute( + "call-partial", + { session: "s1", title: "Review joins", options: JOIN_OPTIONS, advance: false }, + null, + null, + ctx, + ); + + assert.equal(attempts, 3); } finally { cp.execFileSync = original; } diff --git a/harness/.pi/extensions/tht-gate.js b/harness/.pi/extensions/tht-gate.js index 80268823..eac2c836 100644 --- a/harness/.pi/extensions/tht-gate.js +++ b/harness/.pi/extensions/tht-gate.js @@ -246,9 +246,9 @@ function tht(ctx, args, input) { // Runs a privileged tht call; converts any failure into an actionable textResult // (never propagates a raw "Command failed" to the model). -function relayIfThtFails(ctx, args, recovery) { +function relayIfThtFails(ctx, args, recovery, input) { try { - tht(ctx, args); + tht(ctx, args, input); return null; } catch (e) { const cliMsg = (e.stderr || e.message || String(e)).toString().trim(); @@ -380,6 +380,16 @@ export function resolveConfirmOutcome(resp) { return { kind: "unknown", choice }; } +export function isJoinReviewApproval(resp, optionIds) { + if (resp?.kind !== "join-review" || !Array.isArray(resp.choices)) return false; + const expected = new Set(optionIds); + const chosen = new Set(resp.choices); + return expected.size === optionIds.length && + chosen.size === resp.choices.length && + chosen.size === expected.size && + [...chosen].every((id) => expected.has(id)); +} + // True when a reviewer_decide has no merito options but is allowed to close empty // and advance (the empty memory phase F2). The gate then shows an info notice and // auto-advances instead of presenting an empty checklist. @@ -764,7 +774,15 @@ export default function (pi) { options: meritOptions, ...(phase === "F2" ? memorySelectionWidgetProps(opts) : {}), }); - const resp = await emitAndWait(ctx, widget); + let resp; + for (;;) { + resp = await emitAndWait(ctx, widget); + if (resp.control === "freetext" || resp.control === "back" || resp.control === "exit") + break; + if (!joinOnly || isJoinReviewApproval(resp, meritOptions.map((option) => option.id))) + break; + await reLoop(ctx); + } if (resp.control === "freetext") { return textResult( `Altro (reviewer): ${resp.text}. Riformula la proposta tenendo conto.`, @@ -777,11 +795,23 @@ export default function (pi) { const chosen = joinOnly ? opts.filter((o) => !isReserved(o.label)) : opts.filter((o) => (resp.choices ?? []).includes(o.id)); - for (const c of chosen) { - const d = c.decision; - const err = relayIfThtFails(ctx, decisionAddArgs(session, d), ""); + if (joinOnly) { + const decisions = chosen.map((choice) => choice.decision); + const err = relayIfThtFails( + ctx, + ["decision", "add-join-set", "--session", session, "--doc", "-"], + "Il set di join non è stato registrato; correggi l'errore e riprova.", + JSON.stringify(decisions), + ); if (err) return err; - toAdd.push(d); + toAdd.push(...decisions); + } else { + for (const c of chosen) { + const d = c.decision; + const err = relayIfThtFails(ctx, decisionAddArgs(session, d), ""); + if (err) return err; + toAdd.push(d); + } } if (advance) advanceIfReady(ctx, session); return textResult( diff --git a/harness/.pi/skills/tht-sessione/SKILL.md b/harness/.pi/skills/tht-sessione/SKILL.md index 1bc9625b..a7d86b70 100644 --- a/harness/.pi/skills/tht-sessione/SKILL.md +++ b/harness/.pi/skills/tht-sessione/SKILL.md @@ -284,8 +284,9 @@ Prerequisite: Phase 3 closed. `reviewer_decide(advance:false)`, registering `join_modified`. Do not mix `join_modified` with other decision types in that call. The gate renders this proposal as read-only information: **Continue records every proposed join**; the reviewer cannot - remove individual joins (which could create an accidental Cartesian product). If the - reviewer uses **Other — specify**, none of the current joins is recorded: incorporate + remove individual joins (which could create an accidental Cartesian product). The complete + set is persisted atomically: an invalid response or write failure records none of it. If + the reviewer uses **Other — specify**, none of the current joins is recorded: incorporate the textual correction and present the complete revised join set again. Ground joins in the `【Foreign keys】` section of the mschema-text render: it lists the curated logical FKs of the workspace (e.g. diff --git a/harness/tests/test_decision_join_set_cli.py b/harness/tests/test_decision_join_set_cli.py new file mode 100644 index 00000000..5d9b1a8e --- /dev/null +++ b/harness/tests/test_decision_join_set_cli.py @@ -0,0 +1,93 @@ +import json + +import pytest +from typer.testing import CliRunner + +from tht.cli.decision_cmd import decision_app +from tht.decisions import append_decision, append_decisions, list_decisions +from tht.phase import current_phase + + +def _walk_to_phase(session, target): + while current_phase(session) < target: + append_decision(session, type="phase_approved", subject=f"phase:{current_phase(session)}") + + +def _configure_command(monkeypatch, sessions): + import tht.cli.decision_cmd as mod + + class _Cfg: + class paths: + pass + + _Cfg.paths.sessions = sessions + 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) + + +def test_add_join_set_rejects_the_whole_batch_when_one_item_is_invalid(tmp_path, monkeypatch): + session_id = "s1" + session = tmp_path / session_id + session.mkdir() + _walk_to_phase(session, 4) + _configure_command(monkeypatch, tmp_path) + + payload = [ + {"type": "join_modified", "subject": "a-b", "detail": "a.id = b.a_id"}, + {"type": "not_a_decision", "subject": "b-c", "detail": "b.id = c.b_id"}, + ] + result = CliRunner().invoke( + decision_app, + ["add-join-set", "--session", session_id, "--doc", "-"], + input=json.dumps(payload), + ) + + assert result.exit_code != 0 + assert not [decision for decision in list_decisions(session) if decision.type == "join_modified"] + + +def test_add_join_set_appends_the_complete_valid_batch(tmp_path, monkeypatch): + session_id = "s1" + session = tmp_path / session_id + session.mkdir() + _walk_to_phase(session, 4) + _configure_command(monkeypatch, tmp_path) + + payload = [ + {"type": "join_modified", "subject": "a-b", "detail": "a.id = b.a_id"}, + {"type": "join_modified", "subject": "b-c", "detail": "b.id = c.b_id"}, + ] + result = CliRunner().invoke( + decision_app, + ["add-join-set", "--session", session_id, "--doc", "-"], + input=json.dumps(payload), + ) + + assert result.exit_code == 0, result.output + joins = [decision for decision in list_decisions(session) if decision.type == "join_modified"] + assert [decision.subject for decision in joins] == ["a-b", "b-c"] + + +def test_atomic_replace_failure_keeps_the_original_ledger(tmp_path, monkeypatch): + session = tmp_path / "s1" + session.mkdir() + append_decision(session, type="phase_approved", subject="phase:1") + original = (session / "review_decisions.jsonl").read_text() + + import tht.decisions as mod + + def fail_replace(_source, _target): + raise OSError("disk failure") + + monkeypatch.setattr(mod.os, "replace", fail_replace) + with pytest.raises(OSError, match="disk failure"): + append_decisions( + session, + [ + {"type": "join_modified", "subject": "a-b"}, + {"type": "join_modified", "subject": "b-c"}, + ], + ) + + assert (session / "review_decisions.jsonl").read_text() == original diff --git a/harness/tht/cli/decision_cmd.py b/harness/tht/cli/decision_cmd.py index de7ea37d..0f13d9bf 100644 --- a/harness/tht/cli/decision_cmd.py +++ b/harness/tht/cli/decision_cmd.py @@ -1,7 +1,10 @@ +import json +import sys from pathlib import Path from typing import get_args import typer +from pydantic import ValidationError from tht.cli.config_cmd import CONFIG_OPT from tht.cli.schema_cmd import _load_config_or_exit @@ -10,6 +13,40 @@ from tht.cli.session_cmd import load_session_or_exit, session_dir decision_app = typer.Typer(help="Decisioni del reviewer (per sessione, append-only)") +@decision_app.command("add-join-set") +def add_join_set_cmd( + session: str = typer.Option(..., "--session", help="Id della sessione."), + doc: str = typer.Option(..., "--doc", help="Array JSON di join; usa '-' per stdin."), + 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 + + cfg = _load_config_or_exit(config) + load_session_or_exit(cfg, session) + try: + raw = sys.stdin.read() if doc == "-" else Path(doc).read_text() + payload = json.loads(raw) + if not isinstance(payload, list) or not payload: + raise ValueError("il documento deve essere un array JSON non vuoto") + decisions = [DecisionInput.model_validate(item) for item in payload] + if any(decision.type != "join_modified" for decision in decisions): + raise ValueError("tutte le decisioni devono avere type=join_modified") + except (OSError, json.JSONDecodeError, ValidationError, ValueError) as error: + typer.secho(f"ERRORE: set di join non valido: {error}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from error + + from tht.cli.phase_cmd import require_phase_or_exit + 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) + typer.secho( + f"OK: registrato set atomico di {len(records)} join.", + fg=typer.colors.GREEN, + ) + + @decision_app.command("add") def add_cmd( session: str = typer.Option(..., "--session", help="Id della sessione."), diff --git a/harness/tht/decisions.py b/harness/tht/decisions.py index 35916e69..f432c3eb 100644 --- a/harness/tht/decisions.py +++ b/harness/tht/decisions.py @@ -1,3 +1,5 @@ +import os +import tempfile from datetime import UTC, datetime from pathlib import Path from typing import Literal @@ -70,6 +72,14 @@ class DecisionRecord(BaseModel): phase: int | None = None +class DecisionInput(BaseModel): + type: DecisionType + subject: str + detail: str = "" + rationale: str = "" + retracts: int | None = None + + def list_decisions(session_dir: Path) -> list[DecisionRecord]: path = session_dir / DECISIONS_FILE if not path.exists(): @@ -90,24 +100,58 @@ def append_decision( rationale: str = "", retracts: int | None = None, ) -> DecisionRecord: + return append_decisions( + session_dir, + [{ + "type": type, + "subject": subject, + "detail": detail, + "rationale": rationale, + "retracts": retracts, + }], + )[0] + + +def append_decisions( + session_dir: Path, + decisions: list[DecisionInput | dict], +) -> list[DecisionRecord]: + """Validate and append a decision set through one atomic file replacement.""" + inputs = [DecisionInput.model_validate(decision) for decision in decisions] + if not inputs: + return [] + # Fase corrente PRIMA dell'append (high-water-mark D15). Import lazy: phase.py # importa decisions.py (ciclo). Per i marker di fase (subject "phase:N") il valore # e' ridondante col subject; per le decisioni sostanziali con subject "a nome" # (cte_approved, ...) e' l'unico modo per filtrarle dopo un reopen. from tht.phase import current_phase - record = DecisionRecord( - seq=len(list_decisions(session_dir)) + 1, - ts=datetime.now(UTC), - type=type, - subject=subject, - detail=detail, - rationale=rationale, - retracts=retracts, - phase=current_phase(session_dir), - ) + existing = list_decisions(session_dir) + phase = current_phase(session_dir) + records = [ + DecisionRecord( + seq=len(existing) + index, + ts=datetime.now(UTC), + phase=phase, + **decision.model_dump(), + ) + for index, decision in enumerate(inputs, start=1) + ] path = session_dir / DECISIONS_FILE path.parent.mkdir(parents=True, exist_ok=True) - with path.open("a") as f: - f.write(record.model_dump_json() + "\n") - return record + previous = path.read_text() if path.exists() else "" + separator = "" if not previous or previous.endswith("\n") else "\n" + content = previous + separator + "".join(record.model_dump_json() + "\n" for record in records) + + fd, temporary_name = tempfile.mkstemp(prefix=f".{DECISIONS_FILE}.", dir=path.parent) + temporary = Path(temporary_name) + try: + with os.fdopen(fd, "w") as handle: + handle.write(content) + handle.flush() + os.fsync(handle.fileno()) + os.replace(temporary, path) + finally: + temporary.unlink(missing_ok=True) + return records