fix(kanban): review handoff rollback discards staged copies; no double upload
Two review follow-ups on request_review's artifact staging. Staging copies files into attachments/<tid>/ inside the write txn, but the copy is a filesystem side effect the rollback cannot undo. When a later step in the same txn raised (anything other than ArtifactPreservationError, e.g. run bookkeeping), the task correctly stayed `running` but the copy leaked, so the retry staged `a_1.txt` beside an orphan `a.txt`. _stage_completion_artifacts now returns the copies and request_review discards them on any exception around the txn, reusing the same unlink/rmdir logic the staging helper already had. complete_task is left alone: its txn has a different shape (the early-return paths and acceptance recording) and its cleanup runs the scratch workspace anyway, so it was not the identical one-line change. The notifier unions payload['artifacts'] with paths parsed from the summary prose and dedupes by full path only. For `review_requested` the scratch original still exists (the reviewer's completion is what deletes it), so a summary naming the original uploaded the file twice: staged copy and original. Prose-parsed paths whose basename matches a staged artifact are now skipped; `completed` delivery is unaffected in practice because there the original is already gone by delivery time.
This commit is contained in:
@@ -140,15 +140,21 @@ class GatewayKanbanWatchersMixin:
|
||||
reference only), and upload errors are logged, never raised.
|
||||
"""
|
||||
raw_paths: list[str] = []
|
||||
prose_paths: list[str] = []
|
||||
if isinstance(event_payload, dict):
|
||||
raw = event_payload.get("artifacts")
|
||||
if isinstance(raw, (list, tuple)):
|
||||
raw_paths += [item for item in raw if isinstance(item, str)]
|
||||
summary = event_payload.get("summary")
|
||||
if isinstance(summary, str) and summary:
|
||||
raw_paths += adapter.extract_local_files(summary)[0]
|
||||
prose_paths += adapter.extract_local_files(summary)[0]
|
||||
if task is not None and getattr(task, "result", None):
|
||||
raw_paths += adapter.extract_local_files(str(task.result))[0]
|
||||
prose_paths += adapter.extract_local_files(str(task.result))[0]
|
||||
# A staged copy and the scratch original it was copied from are the
|
||||
# same deliverable; on a review handoff the original still exists, so
|
||||
# prose mentions of it must not upload the file a second time.
|
||||
staged_names = {os.path.basename(p) for p in raw_paths}
|
||||
raw_paths += [p for p in prose_paths if os.path.basename(p) not in staged_names]
|
||||
candidates: list[str] = []
|
||||
for path in raw_paths:
|
||||
expanded = os.path.expanduser(path) if path else ""
|
||||
|
||||
@@ -2652,15 +2652,18 @@ def _gate_created_cards(
|
||||
def _stage_completion_artifacts(
|
||||
conn: sqlite3.Connection, task_id: str, metadata: dict, now: int, *,
|
||||
uploaded_by: str = "kanban_complete",
|
||||
) -> None:
|
||||
"""Copy scratch artifacts to the attachments dir and record each as an attachment row."""
|
||||
) -> list[Path]:
|
||||
"""Copy scratch artifacts to the attachments dir and record each as an
|
||||
attachment row; returns the copies so the caller can discard them if its
|
||||
transaction rolls back."""
|
||||
_persist_scratch_completion_artifacts(conn, task_id, metadata)
|
||||
for stored_path in metadata.pop("_staged_artifacts", []):
|
||||
path = Path(stored_path)
|
||||
staged = [Path(stored_path) for stored_path in metadata.pop("_staged_artifacts", [])]
|
||||
for path in staged:
|
||||
_insert_completion_attachment(
|
||||
conn, task_id, filename=path.name, stored_path=str(path),
|
||||
size=path.stat().st_size, created_at=now, uploaded_by=uploaded_by,
|
||||
)
|
||||
return staged
|
||||
|
||||
|
||||
def _cleaned_artifact_paths(metadata: Any) -> list[str]:
|
||||
@@ -2779,11 +2782,7 @@ def _persist_scratch_completion_artifacts(
|
||||
changed = False
|
||||
|
||||
def _discard_copies() -> None:
|
||||
for copied in used_destinations:
|
||||
with contextlib.suppress(OSError):
|
||||
copied.unlink(missing_ok=True)
|
||||
with contextlib.suppress(OSError):
|
||||
attachment_dir.rmdir()
|
||||
_discard_staged_copies(used_destinations, attachment_dir)
|
||||
|
||||
for item in raw_artifacts:
|
||||
artifact = str(item).strip() if isinstance(item, str) else ""
|
||||
@@ -2838,6 +2837,16 @@ def _persist_scratch_completion_artifacts(
|
||||
]
|
||||
|
||||
|
||||
def _discard_staged_copies(copies: Iterable[Path], attachment_dir: Path) -> None:
|
||||
"""Remove staged attachment copies whose DB rows never committed; a leaked
|
||||
copy would make the retry stage ``name_1.ext`` next to an orphan."""
|
||||
for copied in copies:
|
||||
with contextlib.suppress(OSError):
|
||||
Path(copied).unlink(missing_ok=True)
|
||||
with contextlib.suppress(OSError):
|
||||
attachment_dir.rmdir()
|
||||
|
||||
|
||||
def _copy_capped(src: Path, dest: Path, artifact: str) -> None:
|
||||
"""Chunked copy that aborts if the file grows past the attachment cap mid-copy."""
|
||||
with src.open("rb") as source_file, dest.open("xb") as destination_file:
|
||||
@@ -3045,78 +3054,86 @@ def request_review(
|
||||
# review-bound card the reviewer's completion is the cleanup trigger.
|
||||
metadata = _merge_completion_prose_artifacts(conn, task_id, metadata, summary=summary, result=None)
|
||||
now = int(time.time())
|
||||
with write_txn(conn):
|
||||
if not _parents_satisfied(conn, task_id):
|
||||
return _ret(False, "parent dependencies are not satisfied")
|
||||
trow = conn.execute(
|
||||
"SELECT assignee, status, claim_lock, current_run_id "
|
||||
"FROM tasks WHERE id = ?", (task_id,),
|
||||
).fetchone()
|
||||
if trow is None:
|
||||
return _ret(False, "task not found")
|
||||
# Refuse to clear a live worker's claim without proof of ownership
|
||||
# (expected_run_id) or an explicit human override (force=True).
|
||||
if (
|
||||
expected_run_id is None
|
||||
and not force
|
||||
and trow["status"] == "running"
|
||||
and trow["claim_lock"] is not None
|
||||
):
|
||||
return _ret(
|
||||
False, "task is running under a live claim; pass expected_run_id "
|
||||
"(worker ownership) or force=True (explicit operator "
|
||||
"override) instead of clearing the live run's claim",
|
||||
)
|
||||
implementer = trow["assignee"]
|
||||
if reviewer is None:
|
||||
reviewer = _prior_reviewer(conn, task_id)
|
||||
if reviewer is False:
|
||||
# Staged copies live outside the txn: a rollback after staging must not
|
||||
# leave orphans that make the retry stage ``name_1.ext`` beside them.
|
||||
staged_copies: list[Path] = []
|
||||
try:
|
||||
with write_txn(conn):
|
||||
if not _parents_satisfied(conn, task_id):
|
||||
return _ret(False, "parent dependencies are not satisfied")
|
||||
trow = conn.execute(
|
||||
"SELECT assignee, status, claim_lock, current_run_id "
|
||||
"FROM tasks WHERE id = ?", (task_id,),
|
||||
).fetchone()
|
||||
if trow is None:
|
||||
return _ret(False, "task not found")
|
||||
# Refuse to clear a live worker's claim without proof of ownership
|
||||
# (expected_run_id) or an explicit human override (force=True).
|
||||
if (
|
||||
expected_run_id is None
|
||||
and not force
|
||||
and trow["status"] == "running"
|
||||
and trow["claim_lock"] is not None
|
||||
):
|
||||
return _ret(
|
||||
False, "re-review has no durable reviewer provenance (the "
|
||||
"latest changes_requested event is missing or "
|
||||
"malformed); pass reviewer= explicitly",
|
||||
False, "task is running under a live claim; pass expected_run_id "
|
||||
"(worker ownership) or force=True (explicit operator "
|
||||
"override) instead of clearing the live run's claim",
|
||||
)
|
||||
reviewer = _canonical_assignee(reviewer)
|
||||
assignee_sql = ", assignee = ?" if reviewer is not None else ""
|
||||
run_guard = "" if expected_run_id is None else " AND current_run_id = ?"
|
||||
params: tuple[Any, ...] = (
|
||||
*(() if reviewer is None else (reviewer,)), task_id,
|
||||
*(() if expected_run_id is None else (int(expected_run_id),)),
|
||||
)
|
||||
cur = conn.execute(
|
||||
"""
|
||||
UPDATE tasks
|
||||
SET status = 'review',
|
||||
claim_lock = NULL,
|
||||
claim_expires = NULL,
|
||||
worker_pid = NULL
|
||||
""" + assignee_sql + """
|
||||
WHERE id = ?
|
||||
AND status IN ('running', 'ready')
|
||||
""" + run_guard,
|
||||
params,
|
||||
)
|
||||
if cur.rowcount != 1:
|
||||
return _ret(
|
||||
False, "task is not in running/ready (or expected_run_id did not match the current run)",
|
||||
implementer = trow["assignee"]
|
||||
if reviewer is None:
|
||||
reviewer = _prior_reviewer(conn, task_id)
|
||||
if reviewer is False:
|
||||
return _ret(
|
||||
False, "re-review has no durable reviewer provenance (the "
|
||||
"latest changes_requested event is missing or "
|
||||
"malformed); pass reviewer= explicitly",
|
||||
)
|
||||
reviewer = _canonical_assignee(reviewer)
|
||||
assignee_sql = ", assignee = ?" if reviewer is not None else ""
|
||||
run_guard = "" if expected_run_id is None else " AND current_run_id = ?"
|
||||
params: tuple[Any, ...] = (
|
||||
*(() if reviewer is None else (reviewer,)), task_id,
|
||||
*(() if expected_run_id is None else (int(expected_run_id),)),
|
||||
)
|
||||
if isinstance(metadata, dict):
|
||||
_stage_completion_artifacts(
|
||||
conn, task_id, metadata, now, uploaded_by="kanban_request_review",
|
||||
cur = conn.execute(
|
||||
"""
|
||||
UPDATE tasks
|
||||
SET status = 'review',
|
||||
claim_lock = NULL,
|
||||
claim_expires = NULL,
|
||||
worker_pid = NULL
|
||||
""" + assignee_sql + """
|
||||
WHERE id = ?
|
||||
AND status IN ('running', 'ready')
|
||||
""" + run_guard,
|
||||
params,
|
||||
)
|
||||
run_id = _end_or_synthesize_run(
|
||||
conn, task_id, outcome="review_requested", status="review",
|
||||
summary=summary, metadata=metadata, synthesize=bool(summary or metadata),
|
||||
)
|
||||
payload: dict = {
|
||||
"summary": _first_line(summary, 400) or None,
|
||||
"implementer": implementer,
|
||||
"reviewer": reviewer,
|
||||
}
|
||||
staged = _cleaned_artifact_paths(metadata)
|
||||
if staged:
|
||||
payload["artifacts"] = staged
|
||||
_append_event(conn, task_id, "review_requested", payload, run_id=run_id)
|
||||
if cur.rowcount != 1:
|
||||
return _ret(
|
||||
False, "task is not in running/ready (or expected_run_id did not match the current run)",
|
||||
)
|
||||
if isinstance(metadata, dict):
|
||||
staged_copies = _stage_completion_artifacts(
|
||||
conn, task_id, metadata, now, uploaded_by="kanban_request_review",
|
||||
)
|
||||
run_id = _end_or_synthesize_run(
|
||||
conn, task_id, outcome="review_requested", status="review",
|
||||
summary=summary, metadata=metadata, synthesize=bool(summary or metadata),
|
||||
)
|
||||
payload: dict = {
|
||||
"summary": _first_line(summary, 400) or None,
|
||||
"implementer": implementer,
|
||||
"reviewer": reviewer,
|
||||
}
|
||||
staged = _cleaned_artifact_paths(metadata)
|
||||
if staged:
|
||||
payload["artifacts"] = staged
|
||||
_append_event(conn, task_id, "review_requested", payload, run_id=run_id)
|
||||
except Exception:
|
||||
if staged_copies:
|
||||
_discard_staged_copies(staged_copies, staged_copies[0].parent)
|
||||
raise
|
||||
return _ret(True)
|
||||
|
||||
|
||||
|
||||
@@ -646,6 +646,34 @@ def test_review_bound_handoff_preserves_declared_artifacts(kanban_home):
|
||||
]
|
||||
|
||||
|
||||
def test_request_review_rollback_discards_staged_copies(kanban_home):
|
||||
"""A failure after staging rolls the txn back; the copied file must go
|
||||
too, or the retry stages ``evidence_1.json`` next to an orphan."""
|
||||
with kbc.connect() as conn:
|
||||
t = kb.create_task(conn, title="review rollback")
|
||||
ws = kbw.resolve_workspace(kb.get_task(conn, t))
|
||||
kbw.set_workspace_path(conn, t, ws)
|
||||
artifact = ws / "evidence.json"
|
||||
artifact.write_bytes(b"{}")
|
||||
kb.claim_task(conn, t)
|
||||
run_id = kb.get_task(conn, t).current_run_id
|
||||
kwargs = dict(summary="ready", metadata={"artifacts": [str(artifact)]}, expected_run_id=run_id)
|
||||
|
||||
def _boom(*_a, **_k):
|
||||
raise RuntimeError("run bookkeeping failed")
|
||||
|
||||
with pytest.MonkeyPatch.context() as mp:
|
||||
mp.setattr(kb, "_end_or_synthesize_run", _boom)
|
||||
with pytest.raises(RuntimeError):
|
||||
kb.request_review(conn, t, **kwargs)
|
||||
attachment_dir = kb.task_attachments_dir(t)
|
||||
assert kb.get_task(conn, t).status == "running"
|
||||
assert not attachment_dir.exists() or not any(attachment_dir.iterdir())
|
||||
assert kb.request_review(conn, t, **kwargs)
|
||||
assert [a.filename for a in kb.list_attachments(conn, t)] == ["evidence.json"]
|
||||
assert sorted(p.name for p in attachment_dir.iterdir()) == ["evidence.json"]
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Deferred scratch cleanup for parent/child handoff (#33774)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@@ -920,15 +920,18 @@ async def test_notifier_uploads_review_handoff_artifacts(kanban_home, tmp_path,
|
||||
scratch.write_bytes(b"%PDF-fake")
|
||||
kb.claim_task(conn, tid)
|
||||
run_id = kb.get_task(conn, tid).current_run_id
|
||||
# The summary names the scratch original, which still exists at
|
||||
# handoff time: it must not ride along as a second upload.
|
||||
assert kb.request_review(
|
||||
conn, tid, summary="ready for review", metadata={"artifacts": [str(scratch)]},
|
||||
expected_run_id=run_id)
|
||||
conn, tid, summary=f"ready for review: {scratch}",
|
||||
metadata={"artifacts": [str(scratch)]}, expected_run_id=run_id)
|
||||
handoff = [e for e in kb.list_events(conn, tid) if e.kind == "review_requested"][-1]
|
||||
attachments = kb.list_attachments(conn, tid)
|
||||
finally:
|
||||
conn.close()
|
||||
staged_path = handoff.payload["artifacts"][0]
|
||||
assert staged_path != str(scratch), "handoff must name the staged copy, not the scratch original"
|
||||
assert scratch.exists(), "scratch original survives until the reviewer completes"
|
||||
assert staged_path == attachments[0].stored_path
|
||||
|
||||
runner = object.__new__(GatewayRunner)
|
||||
|
||||
Reference in New Issue
Block a user