fix(kanban): make write_txn nesting explicit opt-in
Plain write_txn raises loudly on nesting again (the historical main invariant); composition primitives (create_task, add_comment) opt in with allow_nested=True for savepoint semantics. create_swarm activates the swarm root with an inline blocked->done CAS flip + synthesized run + event instead of nesting complete_task, so complete_task's post-commit side effects (workspace cleanup, failure-counter clear, recompute_ready) can no longer fire under an open outer transaction; recompute_ready now runs after the outer commit. recompute_ready docstring corrected. Regression: plain nesting raises; allow_nested composes and an outer rollback discards inner work with no side effects fired.
This commit is contained in:
@@ -2798,17 +2798,23 @@ def _execute_boundary_with_retry(conn: sqlite3.Connection, sql: str) -> None:
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def write_txn(conn: sqlite3.Connection):
|
||||
def write_txn(conn: sqlite3.Connection, *, allow_nested: bool = False):
|
||||
"""Context manager for an IMMEDIATE write transaction.
|
||||
|
||||
Use for any multi-statement write (creating a task + link, claiming a
|
||||
task + recording an event, etc.). A claim CAS inside this context is
|
||||
atomic -- at most one concurrent writer can succeed.
|
||||
|
||||
Nested callers use SQLite savepoints rather than opening a second
|
||||
``BEGIN IMMEDIATE``. This lets graph builders compose the existing atomic
|
||||
task operations under one outer commit, so the dispatcher can never
|
||||
observe a partially constructed multi-card graph.
|
||||
Nesting is an explicit opt-in: a caller already inside a transaction
|
||||
gets a loud ``RuntimeError`` unless it passes ``allow_nested=True``,
|
||||
in which case a SQLite savepoint is used instead of a second
|
||||
``BEGIN IMMEDIATE``. Only composition primitives that graph builders
|
||||
deliberately run under one outer commit (``create_task``,
|
||||
``add_comment``) opt in — helpers with post-commit side effects
|
||||
(``complete_task`` & co.) must never run under an open outer
|
||||
transaction, because their side effects (workspace cleanup, ready
|
||||
recomputation, failure-counter clears) would fire while the outer
|
||||
transaction can still roll back.
|
||||
|
||||
The explicit ROLLBACK on exception is wrapped in try/except so that
|
||||
a SQLite auto-rollback (which leaves no active transaction) does not
|
||||
@@ -2816,6 +2822,13 @@ def write_txn(conn: sqlite3.Connection):
|
||||
"""
|
||||
_assert_not_delegated_child_mutation()
|
||||
if getattr(conn, "in_transaction", False):
|
||||
if not allow_nested:
|
||||
raise RuntimeError(
|
||||
"write_txn: already inside a transaction. Nested composition "
|
||||
"must opt in explicitly with write_txn(conn, allow_nested=True) "
|
||||
"(savepoint semantics; the inner RELEASE is not durable until "
|
||||
"the outer transaction commits)."
|
||||
)
|
||||
savepoint = f"hermes_nested_{secrets.token_hex(8)}"
|
||||
conn.execute(f"SAVEPOINT {savepoint}")
|
||||
try:
|
||||
@@ -3178,7 +3191,10 @@ def create_task(
|
||||
for attempt in range(2):
|
||||
task_id = _new_task_id()
|
||||
try:
|
||||
with write_txn(conn):
|
||||
# ``allow_nested=True``: graph builders (kanban_swarm.create_swarm)
|
||||
# compose create_task calls under one outer commit so the
|
||||
# dispatcher can never observe a partially constructed graph.
|
||||
with write_txn(conn, allow_nested=True):
|
||||
# Determine task status from parent status, unless the caller
|
||||
# parks it directly in blocked for human-ops review or in
|
||||
# triage for a specifier.
|
||||
@@ -3707,7 +3723,9 @@ def add_comment(
|
||||
if not author or not author.strip():
|
||||
raise ValueError("comment author is required")
|
||||
now = int(time.time())
|
||||
with write_txn(conn):
|
||||
# ``allow_nested=True``: graph builders (kanban_swarm blackboard seeding)
|
||||
# compose comment writes under one outer commit.
|
||||
with write_txn(conn, allow_nested=True):
|
||||
if not conn.execute(
|
||||
"SELECT 1 FROM tasks WHERE id = ?", (task_id,)
|
||||
).fetchone():
|
||||
@@ -4230,8 +4248,9 @@ def recompute_ready(
|
||||
) -> int:
|
||||
"""Promote ``todo`` tasks to ``ready`` when all parents are ``done`` or ``archived``.
|
||||
|
||||
Returns the number of tasks promoted. Safe to call inside or outside
|
||||
an existing transaction; it opens its own IMMEDIATE txn.
|
||||
Returns the number of tasks promoted. Opens its own IMMEDIATE txn, so it
|
||||
MUST be called OUTSIDE any open write transaction (plain ``write_txn``
|
||||
raises on nesting); call it after the enclosing txn commits.
|
||||
|
||||
``blocked`` tasks are also considered for promotion (so a task
|
||||
blocked purely by a parent dependency unblocks itself when the
|
||||
|
||||
@@ -74,6 +74,57 @@ def _swarm_context(root_id: str, goal: str) -> str:
|
||||
)
|
||||
|
||||
|
||||
def _activate_root_inline(
|
||||
conn: sqlite3.Connection,
|
||||
root_id: str,
|
||||
*,
|
||||
summary: str,
|
||||
metadata: dict[str, Any],
|
||||
) -> bool:
|
||||
"""Inline blocked→done CAS flip + event insert for the swarm root.
|
||||
|
||||
Runs INSIDE create_swarm's outer write_txn, so it must not call
|
||||
``kb.complete_task`` — that helper opens its own transaction and fires
|
||||
post-commit side effects (workspace cleanup, failure-counter clear,
|
||||
``recompute_ready``) that would execute while the outer transaction can
|
||||
still roll back. Instead we do the minimal durable writes here and let
|
||||
the caller run ``recompute_ready`` after the outer commit.
|
||||
"""
|
||||
import time as _time
|
||||
|
||||
now = int(_time.time())
|
||||
cur = conn.execute(
|
||||
"""
|
||||
UPDATE tasks
|
||||
SET status = 'done',
|
||||
completed_at = ?,
|
||||
claim_lock = NULL,
|
||||
claim_expires= NULL,
|
||||
worker_pid = NULL
|
||||
WHERE id = ?
|
||||
AND status = 'blocked'
|
||||
""",
|
||||
(now, root_id),
|
||||
)
|
||||
if cur.rowcount != 1:
|
||||
return False
|
||||
run_id = kb._synthesize_ended_run(
|
||||
conn,
|
||||
root_id,
|
||||
outcome="completed",
|
||||
summary=summary,
|
||||
metadata=metadata,
|
||||
)
|
||||
kb._append_event(
|
||||
conn,
|
||||
root_id,
|
||||
"completed",
|
||||
{"result_len": 0, "summary": summary[:400] or None},
|
||||
run_id=run_id,
|
||||
)
|
||||
return True
|
||||
|
||||
|
||||
def create_swarm(
|
||||
conn: sqlite3.Connection,
|
||||
*,
|
||||
@@ -115,7 +166,7 @@ def create_swarm(
|
||||
)
|
||||
root = kb.get_task(conn, created.root_id)
|
||||
if root is not None and root.status == "blocked":
|
||||
if not kb.complete_task(
|
||||
if not _activate_root_inline(
|
||||
conn,
|
||||
created.root_id,
|
||||
summary=activation_summary,
|
||||
@@ -124,11 +175,14 @@ def create_swarm(
|
||||
"goal": goal.strip(),
|
||||
"worker_count": len(created.worker_ids),
|
||||
},
|
||||
fire_lifecycle_hook=False,
|
||||
):
|
||||
raise RuntimeError("could not activate the completed swarm topology")
|
||||
activated = True
|
||||
if activated:
|
||||
# Outside the outer transaction: promote the root's children now
|
||||
# that its 'done' flip is durable (recompute_ready opens its own
|
||||
# txn and must never run under an open write_txn).
|
||||
kb.recompute_ready(conn)
|
||||
root = kb.get_task(conn, created.root_id)
|
||||
run = kb.latest_run(conn, created.root_id)
|
||||
kb._fire_kanban_lifecycle_hook(
|
||||
|
||||
@@ -88,7 +88,12 @@ def test_create_swarm_graph_is_atomic_and_rolls_back_partial_build(
|
||||
assert reader.execute("SELECT COUNT(*) FROM tasks").fetchone()[0] == 0
|
||||
|
||||
monkeypatch.setattr(kb, "create_task", original_create)
|
||||
monkeypatch.setattr(kb, "complete_task", lambda *args, **kwargs: False)
|
||||
import hermes_cli.kanban_swarm as ks
|
||||
|
||||
original_activate = ks._activate_root_inline
|
||||
monkeypatch.setattr(
|
||||
ks, "_activate_root_inline", lambda *args, **kwargs: False
|
||||
)
|
||||
with pytest.raises(RuntimeError, match="could not activate"):
|
||||
create_swarm(
|
||||
writer,
|
||||
@@ -103,7 +108,7 @@ def test_create_swarm_graph_is_atomic_and_rolls_back_partial_build(
|
||||
assert reader.execute("SELECT COUNT(*) FROM tasks").fetchone()[0] == 0
|
||||
|
||||
hooks: list[tuple[str, bool]] = []
|
||||
monkeypatch.setattr(kb, "complete_task", original_complete)
|
||||
monkeypatch.setattr(ks, "_activate_root_inline", original_activate)
|
||||
monkeypatch.setattr(
|
||||
kb,
|
||||
"_fire_kanban_lifecycle_hook",
|
||||
@@ -124,6 +129,54 @@ def test_create_swarm_graph_is_atomic_and_rolls_back_partial_build(
|
||||
writer.close()
|
||||
|
||||
|
||||
def test_plain_write_txn_nesting_raises_and_allow_nested_composes(tmp_path):
|
||||
"""B1 regression: nesting is explicit opt-in, never silent.
|
||||
|
||||
Plain ``write_txn`` inside an open transaction must raise loudly (the
|
||||
historical invariant). ``allow_nested=True`` composes via a savepoint,
|
||||
and an outer rollback discards the inner work without any post-commit
|
||||
side effects having fired (the workspace directory survives).
|
||||
"""
|
||||
conn = kb.connect(tmp_path / "kanban.db")
|
||||
try:
|
||||
workspace = tmp_path / "scratch-ws"
|
||||
workspace.mkdir()
|
||||
tid = kb.create_task(conn, title="ws task", assignee="worker")
|
||||
with kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE tasks SET workspace_path = ? WHERE id = ?",
|
||||
(str(workspace), tid),
|
||||
)
|
||||
|
||||
# 1) Plain nesting raises loudly.
|
||||
with pytest.raises(RuntimeError, match="already inside a transaction"):
|
||||
with kb.write_txn(conn):
|
||||
with kb.write_txn(conn):
|
||||
pass
|
||||
assert not conn.in_transaction
|
||||
|
||||
# 2) allow_nested composes; outer rollback discards inner work
|
||||
# and no side effects (workspace cleanup) fired meanwhile.
|
||||
with pytest.raises(RuntimeError, match="outer failure"):
|
||||
with kb.write_txn(conn):
|
||||
with kb.write_txn(conn, allow_nested=True):
|
||||
conn.execute(
|
||||
"UPDATE tasks SET status = 'done' WHERE id = ?", (tid,)
|
||||
)
|
||||
kb._append_event(conn, tid, "completed", {"result_len": 0})
|
||||
# Inner savepoint released, but the outer txn now fails.
|
||||
raise RuntimeError("outer failure")
|
||||
task = kb.get_task(conn, tid)
|
||||
assert task is not None
|
||||
assert task.status == "ready" # inner 'done' flip was discarded
|
||||
assert not any(
|
||||
e.kind == "completed" for e in kb.list_events(conn, tid)
|
||||
)
|
||||
assert workspace.is_dir() # no _cleanup_workspace side effect fired
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def test_swarm_blackboard_merges_structured_updates(tmp_path):
|
||||
conn = kb.connect(tmp_path / "kanban.db")
|
||||
try:
|
||||
|
||||
Reference in New Issue
Block a user