From f3357b50319cb0fa62cbe3430e986d070aa9ef18 Mon Sep 17 00:00:00 2001 From: yoyodine-industries <311904754+yoyodine-industries@users.noreply.github.com> Date: Sat, 12 Sep 2026 11:42:36 -0400 Subject: [PATCH] fix(kanban): stage review-bound handoff artifacts in request_review request_review ignored artifacts entirely, so a review-bound card lost every file its handoff named: the reviewer's complete_task is what runs _cleanup_workspace over the managed scratch workspace. Stage declared (explicit artifacts argument or metadata["artifacts"]) and prose-referenced files into the task's durable attachments dir at the review handoff, exactly as complete_task already does, carry the staged paths in the review_requested event payload, and let the gateway notifier upload them (its guard widens from completed to review_requested). ArtifactPreservationError still rolls the whole transition back: the task stays running and retryable with no attachments and no event. --- gateway/kanban_watchers_notifier.py | 9 ++- hermes_cli/kanban_db.py | 92 ++++++++++++++++++++------ tests/hermes_cli/test_kanban_db.py | 78 ++++++++++++++++++++++ tests/hermes_cli/test_kanban_notify.py | 74 +++++++++++++++++++++ tools/kanban_tools.py | 19 +++++- tools/kanban_tools_schemas.py | 18 +++++ 6 files changed, 262 insertions(+), 28 deletions(-) diff --git a/gateway/kanban_watchers_notifier.py b/gateway/kanban_watchers_notifier.py index 0d94f7b4d3..9e168538d6 100644 --- a/gateway/kanban_watchers_notifier.py +++ b/gateway/kanban_watchers_notifier.py @@ -548,9 +548,12 @@ class _KanbanNotification: raise RuntimeError(f"adapter send() reported failure: {getattr(_send_res, 'error', None) or 'unknown error'}") logger.debug("kanban notifier: delivered %s event for %s to %s/%s on board %s", ev.kind, self.task_id, self.platform_str, sub["chat_id"], self.board_slug) - # Upload artifact paths from the completion payload / legacy result as - # native files. Only on ``completed`` so retries never spam attachments. - if ev.kind == "completed": + # Upload artifact paths from the handoff payload / legacy result as + # native files. Both handoff kinds stage files for exactly this: a + # review-bound card's files exist precisely so the human sees them at + # handoff time. Retry exposure matches ``completed`` (the sub cursor is + # rewound only when a send failed). + if ev.kind in ("completed", "review_requested"): try: await self.runner._deliver_kanban_artifacts( adapter=adapter, chat_id=sub["chat_id"], metadata=metadata, diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 17e9c7b77c..bd5834543d 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -2649,17 +2649,30 @@ def _gate_created_cards( return verified_cards -def _stage_completion_artifacts(conn: sqlite3.Connection, task_id: str, metadata: dict, now: int) -> None: +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.""" _persist_scratch_completion_artifacts(conn, task_id, metadata) for stored_path in metadata.pop("_staged_artifacts", []): path = Path(stored_path) _insert_completion_attachment( conn, task_id, filename=path.name, stored_path=str(path), - size=path.stat().st_size, created_at=now, + size=path.stat().st_size, created_at=now, uploaded_by=uploaded_by, ) +def _cleaned_artifact_paths(metadata: Any) -> list[str]: + """Non-blank string paths declared in ``metadata["artifacts"]``.""" + if not isinstance(metadata, dict): + return [] + raw = metadata.get("artifacts") + if not isinstance(raw, (list, tuple)): + return [] + return [str(p).strip() for p in raw if isinstance(p, str) and str(p).strip()] + + def _completed_event_payload( result: Optional[str], event_summary: Optional[str], verified_cards: list[str], metadata: Any, ) -> dict: @@ -2679,11 +2692,9 @@ def _completed_event_payload( if verified_cards: payload["verified_cards"] = verified_cards if isinstance(metadata, dict): - md_artifacts = metadata.get("artifacts") - if isinstance(md_artifacts, (list, tuple)): - cleaned = [str(p).strip() for p in md_artifacts if isinstance(p, str) and str(p).strip()] - if cleaned: - payload["artifacts"] = cleaned + cleaned = _cleaned_artifact_paths(metadata) + if cleaned: + payload["artifacts"] = cleaned return payload @@ -2842,16 +2853,16 @@ def _copy_capped(src: Path, dest: Path, artifact: str) -> None: def _insert_completion_attachment( conn: sqlite3.Connection, task_id: str, *, filename: str, stored_path: str, size: int, - created_at: int, + created_at: int, uploaded_by: str = "kanban_complete", ) -> None: """Record a worker-produced artifact in the existing attachment table.""" conn.execute( "INSERT INTO task_attachments " "(task_id, filename, stored_path, content_type, size, uploaded_by, created_at) " - "VALUES (?, ?, ?, NULL, ?, 'kanban_complete', ?)", - (task_id, filename, stored_path, size, created_at), + "VALUES (?, ?, ?, NULL, ?, ?, ?)", + (task_id, filename, stored_path, size, uploaded_by, created_at), ) - _append_event(conn, task_id, "attached", {"filename": filename, "size": size, "by": "kanban_complete"}) + _append_event(conn, task_id, "attached", {"filename": filename, "size": size, "by": uploaded_by}) def _unique_attachment_path(directory: Path, filename: str, used: set[Path]) -> Path: @@ -3001,10 +3012,30 @@ def redact_review_value(value: Any) -> Any: return value +def _declare_handoff_artifacts( + metadata: Optional[dict], artifacts: Optional[Iterable[str]], +) -> Optional[dict]: + """Fold an explicit ``artifacts`` argument into ``metadata["artifacts"]`` + (order-preserving, deduped). Returns ``metadata`` untouched when there is + nothing to add, so callers can pass ``None`` through.""" + if not artifacts: + return metadata + items = [str(item).strip() for item in artifacts if item is not None and str(item).strip()] + if not items: + return metadata + updated = dict(metadata) if isinstance(metadata, dict) else {} + existing = updated.get("artifacts") + merged = list(existing) if isinstance(existing, (list, tuple)) else [] + merged.extend(items) + updated["artifacts"] = list(dict.fromkeys(str(p).strip() for p in merged if str(p).strip())) + return updated + + def request_review( conn: sqlite3.Connection, task_id: str, *, summary: Optional[str] = None, metadata: Optional[dict] = None, reviewer: Optional[str] = None, expected_run_id: Optional[int] = None, force: bool = False, with_reason: bool = False, + artifacts: Optional[Iterable[str]] = None, ): """``running``/``ready`` -> ``review``; never touches block recurrence accounting. @@ -3013,6 +3044,15 @@ def request_review( re-review defaults to the latest ``changes_requested`` provenance. A live claim is only cleared with proof of ownership (``expected_run_id``) or ``force=True``. Returns ``bool``, or ``(ok, reason)`` with ``with_reason``. + + ``artifacts`` (or ``metadata["artifacts"]``) names the handoff's deliverable + files; a review handoff is the last implementer transition, and the + *reviewer's* completion is what cleans the managed scratch workspace up, so + the files are staged into the task's durable attachments dir here and the + staged paths ride the ``review_requested`` payload for the notifier to + upload. A declared artifact that cannot be preserved raises + :class:`ArtifactPreservationError`, rolling the whole transition back: the + task stays ``running`` and retryable, with no attachments and no event. """ def _ret(ok: bool, reason: Optional[str] = None): @@ -3020,6 +3060,12 @@ def request_review( summary = redact_review_value(summary) metadata = redact_review_value(metadata) + # Declared (explicit arg or metadata["artifacts"]) and prose-referenced files + # must be durable BEFORE anything can clean the scratch workspace up: for a + # review-bound card the reviewer's completion is the cleanup trigger. + metadata = _declare_handoff_artifacts(metadata, artifacts) + 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") @@ -3075,21 +3121,23 @@ def request_review( return _ret( False, "task is not in running/ready (or expected_run_id did not match the current run)", ) + if isinstance(metadata, dict): + _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), ) - _append_event( - conn, - task_id, - "review_requested", - { - "summary": _first_line(summary, 400) or None, - "implementer": implementer, - "reviewer": reviewer, - }, - run_id=run_id, - ) + 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) return _ret(True) diff --git a/tests/hermes_cli/test_kanban_db.py b/tests/hermes_cli/test_kanban_db.py index 59f7619f8b..51c02d14c3 100644 --- a/tests/hermes_cli/test_kanban_db.py +++ b/tests/hermes_cli/test_kanban_db.py @@ -617,6 +617,84 @@ def test_complete_task_persists_scratch_artifacts_before_cleanup(kanban_home): ] +def test_review_bound_handoff_preserves_declared_artifacts(kanban_home): + """A review-bound card's declared files must outlive the reviewer's + completion — that completion is what cleans the scratch workspace up.""" + with kbc.connect() as conn: + t = kb.create_task(conn, title="review bound") + task = kb.get_task(conn, t) + ws = kbw.resolve_workspace(task) + kbw.set_workspace_path(conn, t, ws) + artifact = ws / "evidence.json" + artifact.write_bytes(b'{"ok": true}') + kb.claim_task(conn, t) + run_id = kb.get_task(conn, t).current_run_id + assert run_id is not None + assert kb.request_review( + conn, t, summary="ready for review", + artifacts=[str(artifact)], expected_run_id=run_id) + handoff = [e for e in kb.list_events(conn, t) if e.kind == "review_requested"][-1] + assert kb.complete_task(conn, t, summary="approved") + attachments = kb.list_attachments(conn, t) + persisted = Path(handoff.payload["artifacts"][0]) + assert not ws.exists(), "scratch workspace should still be cleaned up" + assert persisted.exists(), "staged copy must survive scratch cleanup" + assert persisted.parent == kb.task_attachments_dir(t) + assert persisted.read_bytes() == b'{"ok": true}' + assert [(a.filename, a.stored_path) for a in attachments] == [ + ("evidence.json", str(persisted.resolve())) + ] + + +def test_review_bound_handoff_preserves_prose_referenced_artifacts(kanban_home): + """Legacy workers name deliverables only by absolute scratch path in prose; + the review handoff must stage those too, before the reviewer completes.""" + with kbc.connect() as conn: + t = kb.create_task(conn, title="review bound prose") + task = kb.get_task(conn, t) + ws = kbw.resolve_workspace(task) + kbw.set_workspace_path(conn, t, ws) + artifact = ws / "notes.md" + artifact.write_bytes(b"# notes\n") + kb.claim_task(conn, t) + run_id = kb.get_task(conn, t).current_run_id + assert run_id is not None + assert kb.request_review( + conn, t, summary=f"ready for review, deliverable at {artifact}", + expected_run_id=run_id) + handoff = [e for e in kb.list_events(conn, t) if e.kind == "review_requested"][-1] + assert kb.complete_task(conn, t, summary="approved") + attachments = kb.list_attachments(conn, t) + persisted = Path(handoff.payload["artifacts"][0]) + assert not ws.exists(), "scratch workspace should still be cleaned up" + assert persisted.exists(), "staged copy must survive scratch cleanup" + assert persisted.parent == kb.task_attachments_dir(t) + assert persisted.read_bytes() == b"# notes\n" + assert [(a.filename, a.stored_path) for a in attachments] == [ + ("notes.md", str(persisted.resolve())) + ] + + +def test_review_bound_handoff_rolls_back_when_declared_artifact_missing(kanban_home): + """Fail-closed: an unresolvable declared artifact aborts the whole review + transition — task stays running/retryable, nothing staged, no event.""" + with kbc.connect() as conn: + t = kb.create_task(conn, title="review bound broken") + task = kb.get_task(conn, t) + ws = kbw.resolve_workspace(task) + kbw.set_workspace_path(conn, t, ws) + missing = ws / "missing.png" + kb.claim_task(conn, t) + run_id = kb.get_task(conn, t).current_run_id + assert run_id is not None + with pytest.raises(kb.ArtifactPreservationError): + kb.request_review( + conn, t, summary="ready for review", artifacts=[str(missing)], + expected_run_id=run_id) + assert kb.get_task(conn, t).status == "running" + assert kb.list_attachments(conn, t) == [] + assert [e for e in kb.list_events(conn, t) if e.kind == "review_requested"] == [] + assert not missing.exists(), "declared path must be left untouched" # --------------------------------------------------------------------------- diff --git a/tests/hermes_cli/test_kanban_notify.py b/tests/hermes_cli/test_kanban_notify.py index 9650b18ce7..472739de6e 100644 --- a/tests/hermes_cli/test_kanban_notify.py +++ b/tests/hermes_cli/test_kanban_notify.py @@ -897,6 +897,80 @@ async def test_notifier_artifact_delivery_skips_missing_files(kanban_home, tmp_p assert "real.pdf" in documents_uploaded[0] +@pytest.mark.asyncio +async def test_notifier_uploads_review_handoff_artifacts(kanban_home, tmp_path, monkeypatch): + """A review handoff's files are uploaded from the durable staged copy — + not the scratch original the reviewer's completion is about to delete.""" + import hermes_cli.kanban_db as kb + from hermes_cli import kanban_db_connect as kbc + from hermes_cli import kanban_db_notify as kbn + from hermes_cli import kanban_db_workspace as kbw + from gateway.run import GatewayRunner + from gateway.config import Platform + + monkeypatch.setenv("HERMES_MEDIA_ALLOW_DIRS", str(tmp_path)) + + conn = kbc.connect() + try: + tid = kb.create_task(conn, title="review handoff", assignee="worker1") + kbn.add_notify_sub(conn, task_id=tid, platform="telegram", chat_id="chat1") + ws = kbw.resolve_workspace(kb.get_task(conn, tid)) + kbw.set_workspace_path(conn, tid, ws) + scratch = ws / "report.pdf" + scratch.write_bytes(b"%PDF-fake") + kb.claim_task(conn, tid) + run_id = kb.get_task(conn, tid).current_run_id + assert kb.request_review( + conn, tid, summary="ready for review", 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 staged_path == attachments[0].stored_path + + runner = object.__new__(GatewayRunner) + runner._owns_kanban_dispatcher_lock = lambda: True + runner._running = True + runner._kanban_sub_fail_counts = {} + runner._kanban_dispatcher_lock_handle = object() + + fake_adapter = MagicMock() + fake_adapter.name = "telegram" + + documents_uploaded: list = [] + + async def _send(chat_id, msg, metadata=None): + runner._running = False + + async def _send_document(chat_id, file_path, metadata=None, **_kw): + documents_uploaded.append(file_path) + + fake_adapter.send = AsyncMock(side_effect=_send) + fake_adapter.send_document = AsyncMock(side_effect=_send_document) + fake_adapter.send_multiple_images = AsyncMock() + from gateway.platforms.base import BasePlatformAdapter + fake_adapter.extract_local_files = BasePlatformAdapter.extract_local_files + + runner.adapters = {Platform.TELEGRAM: fake_adapter} + + _orig_sleep = asyncio.sleep + + async def _fast_sleep(_): + await _orig_sleep(0) + + with patch("gateway.run.asyncio.sleep", side_effect=_fast_sleep): + await asyncio.wait_for( + runner._kanban_notifier_watcher(interval=1), + timeout=10.0, + ) + + assert documents_uploaded == [staged_path] + assert str(scratch) not in documents_uploaded + + # --------------------------------------------------------------------------- # Migration backfill: pre-delivery_mode gateway subscriptions keep active wake. # diff --git a/tools/kanban_tools.py b/tools/kanban_tools.py index f616a4d53a..8e90184437 100644 --- a/tools/kanban_tools.py +++ b/tools/kanban_tools.py @@ -658,6 +658,9 @@ def _handle_request_review(args: dict, **kw) -> str: if metadata is not None: metadata = _redact_metadata(metadata) _check(metadata is not None, "metadata could not be safely serialized") + artifacts = _coerce_str_list(args.get("artifacts"), "artifacts", "file paths", strip=True) + if artifacts: + metadata = _merge_artifacts(metadata, artifacts) metadata = _stamp_worker_session_metadata(tid, metadata) # Reviewer is model-supplied free text stored durably on the event payload. reviewer = _redact_opt(args.get("reviewer") or None) @@ -671,9 +674,19 @@ def _handle_request_review(args: dict, **kw) -> str: f"Installed profiles: {', '.join(list_profile_names())}") with _board(args.get("board")) as (kb, conn): _goal_gate("kanban_request_review", kb.get_task(conn, tid), tid, summary) - ok, fail_reason = kb.request_review( - conn, tid, summary=summary, metadata=metadata, reviewer=reviewer, - expected_run_id=_worker_run_id(tid), with_reason=True) + try: + ok, fail_reason = kb.request_review( + conn, tid, summary=summary, metadata=metadata, reviewer=reviewer, + artifacts=artifacts, expected_run_id=_worker_run_id(tid), with_reason=True) + except kb.ArtifactPreservationError as artifact_err: + # Same contract as kanban_complete (#22923): the transition rolled + # back, the task is untouched and retryable — say so explicitly or + # the model treats the tool_error as terminal. + return tool_error( + f"kanban_request_review could not preserve the declared artifacts: {artifact_err}. " + f"Your task is still in-flight (no state change) and its scratch workspace was " + f"kept. Fix the artifact path or storage error, then retry " + f"kanban_request_review with the same handoff.") _check(ok, f"could not request review for {tid}: " f"{fail_reason or 'unknown id or not in running/ready'}") return _ok_landed(kb, conn, tid, "review") diff --git a/tools/kanban_tools_schemas.py b/tools/kanban_tools_schemas.py index adcc30263c..b6f4be7fb6 100644 --- a/tools/kanban_tools_schemas.py +++ b/tools/kanban_tools_schemas.py @@ -229,6 +229,24 @@ KANBAN_REQUEST_REVIEW_SCHEMA = _schema( ), "additionalProperties": True, }, + "artifacts": { + "type": "array", + "items": {"type": "string"}, + "description": ( + "Optional list of absolute paths to deliverable " + "files this handoff names — generated charts, " + "PDFs, spreadsheets, images, archives. Examples: " + "['/tmp/q3-revenue.png', '/tmp/report.pdf']. " + "A review handoff is the last implementer " + "transition, so the kernel copies these into the " + "task's durable attachments before the reviewer's " + "completion cleans the scratch workspace up, and " + "the gateway notifier uploads them as native " + "attachments to the subscribed chat. A missing " + "declared scratch artifact keeps the task in place " + "so you can fix the path and retry." + ), + }, }, ["summary"], )