Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
57 changes: 56 additions & 1 deletion cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -15582,7 +15582,7 @@ def _block(reason: str) -> None:
except Exception:
pass

_run_loop(
loop_result = _run_loop(
task_id=task_id,
goal_text=goal_text,
run_turn=_run_turn,
Expand All @@ -15592,6 +15592,61 @@ def _block(reason: str) -> None:
first_response=first_response or "",
log=lambda m: logger.info("%s", m),
)
if not isinstance(loop_result, dict):
return
if loop_result.get("outcome") != "finalization_recovery":
return

recovery_contract = str(loop_result.get("recovery_contract") or "").strip()
recovery_reason = str(loop_result.get("reason") or "").strip()
run_id_raw = (_os.environ.get("HERMES_KANBAN_RUN_ID") or "").strip()
expected_run_id = int(run_id_raw) if run_id_raw.isdigit() else None
c = _kb.connect()
try:
saved = _kb.record_finalization_recovery(
c,
task_id,
reason=recovery_reason,
contract=recovery_contract,
expected_run_id=expected_run_id,
)
Comment on lines +15600 to +15612

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔒 Security & Privacy | 🟠 Major | ⚡ Quick win

Recovery contract text isn't redacted before being persisted.

recovery_contract embeds the goal judge's freeform reason (via KANBAN_GOAL_FINALIZATION_RECOVERY_TEMPLATE, which can echo content derived from the worker's response) and flows straight into record_finalization_recovery, which writes it verbatim into task_comments and the finalization_recovery_requested event payload. Every sibling lifecycle path (_handle_complete, _handle_block in tools/kanban_tools.py) redacts comparable user/LLM-influenced text with redact_sensitive_text(..., force=True) before persistence; this path skips that step.

🛡️ Proposed fix
     recovery_contract = str(loop_result.get("recovery_contract") or "").strip()
     recovery_reason = str(loop_result.get("reason") or "").strip()
+    if recovery_contract:
+        from agent.redact import redact_sensitive_text
+        recovery_contract = redact_sensitive_text(recovery_contract, force=True)
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
recovery_contract = str(loop_result.get("recovery_contract") or "").strip()
recovery_reason = str(loop_result.get("reason") or "").strip()
run_id_raw = (_os.environ.get("HERMES_KANBAN_RUN_ID") or "").strip()
expected_run_id = int(run_id_raw) if run_id_raw.isdigit() else None
c = _kb.connect()
try:
saved = _kb.record_finalization_recovery(
c,
task_id,
reason=recovery_reason,
contract=recovery_contract,
expected_run_id=expected_run_id,
)
recovery_contract = str(loop_result.get("recovery_contract") or "").strip()
recovery_reason = str(loop_result.get("reason") or "").strip()
if recovery_contract:
from agent.redact import redact_sensitive_text
recovery_contract = redact_sensitive_text(recovery_contract, force=True)
run_id_raw = (_os.environ.get("HERMES_KANBAN_RUN_ID") or "").strip()
expected_run_id = int(run_id_raw) if run_id_raw.isdigit() else None
c = _kb.connect()
try:
saved = _kb.record_finalization_recovery(
c,
task_id,
reason=recovery_reason,
contract=recovery_contract,
expected_run_id=expected_run_id,
)
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@cli.py` around lines 15600 - 15612, Redact the worker/LLM-influenced
recovery_contract before passing it to _kb.record_finalization_recovery, using
the same redact_sensitive_text(..., force=True) pattern as _handle_complete and
_handle_block in tools/kanban_tools.py. Preserve the existing recovery_reason
handling and pass the redacted contract to both persisted task comments and
event payloads through record_finalization_recovery.

if saved:
logger.info(
"kanban goal loop: recorded finalization recovery contract for "
"task %s (run_id=%s)",
task_id,
expected_run_id,
)
return
task_after = _kb.get_task(c, task_id)
if task_after is not None and task_after.status in {
"done",
"blocked",
"review",
"human_review",
"triage",
"todo",
"scheduled",
}:
logger.info(
"kanban goal loop: finalization recovery contract became a "
"no-op for task %s (already in %s)",
task_id,
task_after.status,
)
else:
logger.info(
"kanban goal loop: skipped finalization recovery contract for "
"task %s due to ownership/state drift (status=%s, run_id=%s)",
task_id,
task_after.status if task_after is not None else None,
expected_run_id,
)
finally:
try:
c.close()
except Exception:
pass


def main(
Expand Down
88 changes: 73 additions & 15 deletions hermes_cli/goals.py
Original file line number Diff line number Diff line change
Expand Up @@ -1658,6 +1658,22 @@ def render_contract(self) -> str:
"kanban_block with the reason instead."
)

# Bounded recovery contract when a worker keeps producing "done" prose but
# still exits without a terminal lifecycle call. This does NOT block the task;
# callers should persist this contract as recovery context and let normal
# worker retry logic pick it up.
KANBAN_GOAL_FINALIZATION_RECOVERY_TEMPLATE = (
"[Recovery contract: finalize this kanban task cleanly]\n"
"Judge reason: {reason}\n\n"
"The previous run appears to have finished the work but did not end the "
"task lifecycle correctly. Re-open the existing workspace state, verify "
"artifacts, and end with exactly one terminal lifecycle call:\n"
"- call kanban_complete with an evidence-backed summary if done, OR\n"
"- call kanban_block(kind='dependency' or 'needs_input', reason=...) if "
"human/external input is still required.\n"
"Do not exit without one of those lifecycle calls."
)


def run_kanban_goal_loop(
*,
Expand Down Expand Up @@ -1693,9 +1709,14 @@ def run_kanban_goal_loop(
(str -> str), ``task_status_fn`` (() -> str|None), and ``block_fn``
(reason: str -> None).

Returns a decision dict: ``{"outcome", "turns_used", "reason"}`` where
outcome is one of ``"completed_by_worker"``, ``"blocked_budget"``,
``"blocked_by_worker"``, or ``"stopped"``.
Returns a decision dict with ``{"outcome", "turns_used", "reason"}``.
``outcome`` is one of:
- ``"completed_by_worker"``
- ``"blocked_budget"``
- ``"blocked_by_worker"``
- ``"handed_off_by_worker"`` (worker moved task to review/human_review)
- ``"finalization_recovery"`` (done prose, missing lifecycle finalizer)
- ``"stopped"``
"""

def _log(msg: str) -> None:
Expand Down Expand Up @@ -1728,6 +1749,16 @@ def _log(msg: str) -> None:
if status == "blocked":
_log(f"kanban goal loop: task {task_id} blocked by worker after {turns_used} turn(s)")
return {"outcome": "blocked_by_worker", "turns_used": turns_used, "reason": "worker blocked the task"}
if status in {"review", "human_review"}:
_log(
f"kanban goal loop: task {task_id} handed off by worker to "
f"{status} after {turns_used} turn(s)"
)
return {
"outcome": "handed_off_by_worker",
"turns_used": turns_used,
"reason": f"worker moved task to {status}",
}
if status not in ("running", "ready"):
# Reclaimed / archived / unexpected — let the dispatcher own it.
_log(f"kanban goal loop: task {task_id} status={status!r}; stopping")
Expand All @@ -1743,25 +1774,51 @@ def _log(msg: str) -> None:
_log(f"kanban goal loop: turn {turns_used}/{max_turns} verdict={verdict} reason={_truncate(reason, 120)}")

if verdict == "done":
recovery_contract = KANBAN_GOAL_FINALIZATION_RECOVERY_TEMPLATE.format(
reason=_truncate(reason, 500)
)
if nudged_to_finalize:
# Already asked once to call kanban_complete and it still
# didn't — block for review rather than spin.
_log(f"kanban goal loop: task {task_id} judged done but worker won't finalize; blocking")
try:
block_fn(
f"Goal-mode worker's output looked complete but it never "
f"called kanban_complete after a finalize nudge ({reason})."
)
except Exception as exc:
_log(f"kanban goal loop: block_fn failed ({exc})")
return {"outcome": "blocked_budget", "turns_used": turns_used, "reason": "judged done, never finalized"}
# Bounded recovery route: never hard-block solely because prose
# omitted a lifecycle tool call. The dispatcher's protocol-
# violation retry path can recover this in a fresh run.
_log(
f"kanban goal loop: task {task_id} judged done after finalize "
"nudge but still not finalized; returning recoverable "
"finalization contract (no block)"
)
return {
"outcome": "finalization_recovery",
"turns_used": turns_used,
"reason": (
"judged done after finalize nudge, but worker never "
"called kanban_complete/kanban_block"
),
"recovery_contract": recovery_contract,
}
if turns_used >= max_turns:
# Out of turn budget before we can spend a dedicated finalizer
# nudge turn: emit recovery contract rather than hard-block.
_log(
f"kanban goal loop: task {task_id} judged done at budget edge "
f"({turns_used}/{max_turns}) before finalize nudge; "
"returning recoverable finalization contract (no block)"
)
return {
"outcome": "finalization_recovery",
"turns_used": turns_used,
"reason": (
"judged done at turn-budget edge without a lifecycle "
"finalizer call"
),
"recovery_contract": recovery_contract,
}
prompt = KANBAN_GOAL_FINALIZE_TEMPLATE.format(reason=_truncate(reason, 400))
nudged_to_finalize = True
else:
prompt = KANBAN_GOAL_CONTINUATION_TEMPLATE.format(reason=_truncate(reason, 400))

# Budget check BEFORE spending another turn.
if turns_used >= max_turns:
if verdict != "done" and turns_used >= max_turns:
_log(f"kanban goal loop: task {task_id} exhausted {turns_used}/{max_turns} turns; blocking")
try:
block_fn(
Expand Down Expand Up @@ -1797,6 +1854,7 @@ def _log(msg: str) -> None:
"DRAFT_CONTRACT_SYSTEM_PROMPT",
"KANBAN_GOAL_CONTINUATION_TEMPLATE",
"KANBAN_GOAL_FINALIZE_TEMPLATE",
"KANBAN_GOAL_FINALIZATION_RECOVERY_TEMPLATE",
"DEFAULT_MAX_TURNS",
"load_goal",
"save_goal",
Expand Down
102 changes: 102 additions & 0 deletions hermes_cli/kanban_db.py
Original file line number Diff line number Diff line change
Expand Up @@ -3220,6 +3220,108 @@ def list_comments(conn: sqlite3.Connection, task_id: str) -> list[Comment]:
]


def record_finalization_recovery(
conn: sqlite3.Connection,
task_id: str,
*,
reason: Optional[str],
contract: Optional[str],
expected_run_id: Optional[int] = None,
author: str = "kanban-goal-loop",
) -> bool:
"""Persist a deterministic recovery contract for missed lifecycle finalizers.

Used by goal-mode workers when the judge says the work appears complete but
the run still did not call ``kanban_complete``/``kanban_block`` after the
finalize nudge. This is intentionally non-terminal: it records context,
closes the abandoned run, and moves the card back to ``ready`` so the
dispatcher deterministically gives a recovery worker the existing
workspace. It never forces ``blocked`` / ``done`` from outside the worker
lifecycle tools.

CAS guard: when ``expected_run_id`` is supplied, we only write if the task
still points at that run. If another writer already transitioned the task,
this becomes a no-op (False) and callers can reconcile based on the current
status.
"""
clean_reason = (
(reason or "").strip()
or "judge marked work as done but the lifecycle finalizer call was missing"
)
clean_contract = (contract or "").strip()
now = int(time.time())
with write_txn(conn):
row = conn.execute(
"SELECT status, current_run_id FROM tasks WHERE id = ?",
(task_id,),
).fetchone()
if row is None:
return False
cur_run_id = int(row["current_run_id"]) if row["current_run_id"] else None
if expected_run_id is not None and cur_run_id != int(expected_run_id):
return False
if row["status"] != "running":
return False

# Release this run before a recovery worker can be claimed. Leaving the
# task running here strands its current_run_id until the stale-claim
# reaper fires, then turns an ordinary missed finalizer into a timeout
# failure. ``current_run_id`` is the CAS ownership token, so this update
# and the subsequent _end_run are one serialized transition.
cur = conn.execute(
"""
UPDATE tasks
SET status = 'ready',
claim_lock = NULL,
claim_expires = NULL,
worker_pid = NULL,
last_heartbeat_at = NULL
WHERE id = ?
AND status = 'running'
""" + ("" if expected_run_id is None else " AND current_run_id = ?"),
(task_id,) if expected_run_id is None else (task_id, int(expected_run_id)),
)
if cur.rowcount != 1:
return False

contract_preview = clean_contract[:2000] if clean_contract else ""
comment_body = (
"Goal-mode finalization recovery requested.\n"
f"Judge reason: {clean_reason[:500]}\n\n"
f"{contract_preview}"
).strip()
conn.execute(
"INSERT INTO task_comments (task_id, author, body, created_at) "
"VALUES (?, ?, ?, ?)",
(task_id, author, comment_body, now),
)
ended_run_id = _end_run(
conn,
task_id,
outcome="finalization_recovery",
status="finalization_recovery",
summary=clean_reason[:500],
metadata={"recovery_contract": contract_preview or None},
)
_append_event(
conn,
task_id,
"finalization_recovery_requested",
{
"reason": clean_reason[:500],
"contract": contract_preview or None,
},
run_id=ended_run_id,
)
_log.info(
"record_finalization_recovery: task=%s run_id=%s reason=%s",
task_id,
cur_run_id,
clean_reason[:120],
)
return True


# ---------------------------------------------------------------------------
# Attachments
# ---------------------------------------------------------------------------
Expand Down
Loading
Loading