fix: serialize decision ledger writes
This commit is contained in:
@@ -130,7 +130,8 @@ const joinOnly = opts.length > 0 && opts.every((o) => o.decision.type === "join_
|
||||
```
|
||||
|
||||
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`.
|
||||
id set. After Continue persist all original decisions through atomic `decision add-join-set`, with
|
||||
a per-session cross-process lock covering ledger read, sequence assignment, and replacement.
|
||||
On free text, return feedback without persistence. Keep all other decisions on `multiselect`.
|
||||
|
||||
- [x] **Step 4: Write failing frontend widget tests**
|
||||
|
||||
@@ -36,8 +36,8 @@ 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 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
|
||||
complete response and persists all decisions through one atomic ledger replacement serialized by
|
||||
a per-session cross-process lock. 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.
|
||||
|
||||
@@ -285,9 +285,10 @@ Prerequisite: Phase 3 closed.
|
||||
`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). 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.
|
||||
set is persisted atomically under a per-session writer lock: 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.
|
||||
`fact_x.cod_paz=dim_patient.cod_paz`, `*_time_key=dim_time.day_key`) — prefer
|
||||
|
||||
@@ -1,4 +1,7 @@
|
||||
import fcntl
|
||||
import json
|
||||
import multiprocessing
|
||||
import time
|
||||
|
||||
import pytest
|
||||
from typer.testing import CliRunner
|
||||
@@ -8,6 +11,16 @@ from tht.decisions import append_decision, append_decisions, list_decisions
|
||||
from tht.phase import current_phase
|
||||
|
||||
|
||||
def _concurrent_append_worker(session, subject, start, ready, done):
|
||||
ready.put(subject)
|
||||
start.wait()
|
||||
try:
|
||||
append_decision(session, type="concept_clarified", subject=subject)
|
||||
done.put((subject, None))
|
||||
except Exception as error: # pragma: no cover - surfaced through the parent assertion
|
||||
done.put((subject, repr(error)))
|
||||
|
||||
|
||||
def _walk_to_phase(session, target):
|
||||
while current_phase(session) < target:
|
||||
append_decision(session, type="phase_approved", subject=f"phase:{current_phase(session)}")
|
||||
@@ -91,3 +104,41 @@ def test_atomic_replace_failure_keeps_the_original_ledger(tmp_path, monkeypatch)
|
||||
)
|
||||
|
||||
assert (session / "review_decisions.jsonl").read_text() == original
|
||||
|
||||
|
||||
def test_concurrent_appends_wait_for_the_session_lock_and_keep_both_records(tmp_path):
|
||||
session = tmp_path / "s1"
|
||||
session.mkdir()
|
||||
context = multiprocessing.get_context("fork")
|
||||
start = context.Event()
|
||||
ready = context.Queue()
|
||||
done = context.Queue()
|
||||
processes = [
|
||||
context.Process(
|
||||
target=_concurrent_append_worker,
|
||||
args=(session, subject, start, ready, done),
|
||||
)
|
||||
for subject in ("first", "second")
|
||||
]
|
||||
for process in processes:
|
||||
process.start()
|
||||
assert {ready.get(timeout=5), ready.get(timeout=5)} == {"first", "second"}
|
||||
|
||||
lock_path = session / ".review_decisions.lock"
|
||||
with lock_path.open("a+") as lock:
|
||||
fcntl.flock(lock, fcntl.LOCK_EX)
|
||||
start.set()
|
||||
time.sleep(0.2)
|
||||
writers_waited = not (session / "review_decisions.jsonl").exists()
|
||||
fcntl.flock(lock, fcntl.LOCK_UN)
|
||||
|
||||
results = [done.get(timeout=5), done.get(timeout=5)]
|
||||
for process in processes:
|
||||
process.join(timeout=5)
|
||||
assert process.exitcode == 0
|
||||
|
||||
assert writers_waited
|
||||
assert all(error is None for _subject, error in results)
|
||||
records = list_decisions(session)
|
||||
assert {record.subject for record in records} == {"first", "second"}
|
||||
assert [record.seq for record in records] == [1, 2]
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import fcntl
|
||||
import os
|
||||
import tempfile
|
||||
from datetime import UTC, datetime
|
||||
@@ -7,6 +8,7 @@ from typing import Literal
|
||||
from pydantic import BaseModel
|
||||
|
||||
DECISIONS_FILE = "review_decisions.jsonl"
|
||||
DECISIONS_LOCK_FILE = ".review_decisions.lock"
|
||||
|
||||
# 22 tipi di the reference implementation (verified leggendo session/decisions.py) + 1 nuovo (D15):
|
||||
# `decision_retracted` per il rollback a granularità step (ritira una decisione
|
||||
@@ -121,6 +123,22 @@ def append_decisions(
|
||||
if not inputs:
|
||||
return []
|
||||
|
||||
session_dir.mkdir(parents=True, exist_ok=True)
|
||||
lock_path = session_dir / DECISIONS_LOCK_FILE
|
||||
with lock_path.open("a+") as lock:
|
||||
fcntl.flock(lock, fcntl.LOCK_EX)
|
||||
try:
|
||||
return _append_decisions_locked(session_dir, inputs)
|
||||
finally:
|
||||
fcntl.flock(lock, fcntl.LOCK_UN)
|
||||
|
||||
|
||||
def _append_decisions_locked(
|
||||
session_dir: Path,
|
||||
inputs: list[DecisionInput],
|
||||
) -> list[DecisionRecord]:
|
||||
"""Append while the caller holds the session's cross-process ledger lock."""
|
||||
|
||||
# Fase corrente PRIMA dell'append (high-water-mark D15). Import lazy: phase.py
|
||||
# importa decisions.py (ciclo). Per i marker di fase (subject "phase:N") il valore
|
||||
# e' ridondante col subject; per le decisioni sostanziali con subject "a nome"
|
||||
|
||||
Reference in New Issue
Block a user