fix: make join approval atomic
This commit is contained in:
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
@@ -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."),
|
||||
|
||||
+57
-13
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user