diff --git a/gateway/kanban_watchers.py b/gateway/kanban_watchers.py index c637086db4..96303924ca 100644 --- a/gateway/kanban_watchers.py +++ b/gateway/kanban_watchers.py @@ -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 "" diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index cd87c2f4ab..91d6dc724f 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -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) diff --git a/tests/hermes_cli/test_kanban_db.py b/tests/hermes_cli/test_kanban_db.py index 774d545308..8a6053da39 100644 --- a/tests/hermes_cli/test_kanban_db.py +++ b/tests/hermes_cli/test_kanban_db.py @@ -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) # --------------------------------------------------------------------------- diff --git a/tests/hermes_cli/test_kanban_notify.py b/tests/hermes_cli/test_kanban_notify.py index b3542544c1..a06b34a1f7 100644 --- a/tests/hermes_cli/test_kanban_notify.py +++ b/tests/hermes_cli/test_kanban_notify.py @@ -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)