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
11 changes: 11 additions & 0 deletions hermes_cli/goals.py
Original file line number Diff line number Diff line change
Expand Up @@ -1386,6 +1386,17 @@ def kanban_handoff_rejection(
if worker_run_id is not None:
blocked = kb.block_task(conn, task_id, reason=reason, kind="transient", expected_run_id=worker_run_id)
return f"{reason}; task {'blocked transient after judge retry' if blocked else 'not blocked (run ownership changed)'}"
# A goal-mode worker handing off through the CLI runs it as a CHILD of
# its terminal tool: it inherits the card's task env but is never the run
# owner, so the resolver above returns None. That is the worker's lineage,
# not an operator -- a judge error must not let its handoff through
# (FleetReview #960). Refuse without blocking (run ownership is unproven).
if (
worker_run_id is None
and task_id
and (os.environ.get("HERMES_KANBAN_TASK") or "").strip() == task_id
):
return f"{reason}; handoff refused (caller inherits this card's worker env; judge errors fail closed)"
return None


Expand Down
18 changes: 13 additions & 5 deletions hermes_cli/kanban.py
Original file line number Diff line number Diff line change
Expand Up @@ -1971,16 +1971,24 @@ def _caller_session_id() -> Optional[str]:
explicit = (_SLASH_SESSION_ID.get() or "").strip()
if explicit:
return explicit
in_gateway = os.environ.get("_HERMES_GATEWAY") == "1"
try:
from gateway.session_context import resolve_current_session_id

from gateway.session_context import _SESSION_ID, resolve_current_session_id

if in_gateway:
# In-process gateway: ONLY a per-turn contextvar bound in this
# context is ours. A plain slash command runs before that bind,
# and the resolver's os.environ fallback is process-global --
# another chat's session. Sessionless means None, never a
# borrowed identity (FleetReview #951).
bound = _SESSION_ID.get()
return (bound.strip() or None) if isinstance(bound, str) else None
resolved = (resolve_current_session_id() or "").strip()
if resolved:
return resolved
if os.environ.get("_HERMES_GATEWAY") == "1":
return None # in-process: env is another session's, never ours
except Exception:
pass
if in_gateway:
return None
return (os.environ.get("HERMES_SESSION_ID") or "").strip() or None


Expand Down
61 changes: 43 additions & 18 deletions hermes_cli/kanban_db.py
Original file line number Diff line number Diff line change
Expand Up @@ -5347,14 +5347,22 @@ def session_owner_profile(session_id: Optional[str]) -> Optional[str]:


def _actor_profiles(actor: MutationActor) -> frozenset[str]:
"""Every profile identity the actor legitimately holds: the bound
profile plus the owner profile of each caller session."""
names = {actor.profile} if actor.profile else set()
for sid in actor.session_ids:
owner = session_owner_profile(sid)
if owner:
names.add(owner)
return frozenset(names)
"""Every profile identity the actor legitimately holds.

The owner profile of each caller session is authoritative. The bound
(env-derived) profile counts only when no caller session resolves to an
owner: a root-home repoint makes it read ``default`` -- an operator
profile -- from inside another profile's session, so keeping both let
that caller pass ``--operator`` or mutate a ``default``-assigned card
(FleetReview #1074).
"""
owners = {
owner for owner in (session_owner_profile(sid) for sid in actor.session_ids)
if owner
}
if owners:
return frozenset(owners)
return frozenset({actor.profile} if actor.profile else ())


def _valid_operator_reason(reason: str) -> bool:
Expand Down Expand Up @@ -5605,10 +5613,9 @@ def check_home_session(
return None
if home in actor.session_ids:
return None
# Execution lane: the assignee works its card wherever it was born, and a
# dispatched worker always owns the card it was spawned for.
if actor.profile and (row["assignee"] or "") == actor.profile:
return None
# Execution lane: a dispatched worker always owns the card it was spawned
# for. The assignee match lives below, on :func:`_actor_profiles` (session
# owner over env profile), not on the raw env profile (FleetReview #1074).
if (os.environ.get("HERMES_KANBAN_TASK") or "").strip() == task_id:
return None
# ...and the cards it fanned out: item 1 stamps a worker's children with
Expand Down Expand Up @@ -8378,8 +8385,11 @@ def _prior_worker_still_alive(
# An outcome cannot certify exit: operators can write the same outcomes as
# worker tools, and a newer synthetic row can hide an older live owner.
runs = conn.execute(
"SELECT r.id, r.outcome, r.ended_at, r.started_at, t.max_runtime_seconds "
"FROM task_runs r JOIN tasks t ON t.id = r.task_id "
# The run's OWN runtime cap, snapshotted at claim: the card's current
# cap can be shortened after release while that worker still runs
# (FleetReview #956). A run with no snapshot is probed, never skipped.
"SELECT r.id, r.outcome, r.ended_at, r.started_at, r.max_runtime_seconds "
"FROM task_runs r "
"WHERE r.task_id = ? ORDER BY r.id DESC",
(task_id,),
).fetchall()
Expand Down Expand Up @@ -8564,8 +8574,9 @@ def _real_pid_started_in_claim(pid, claimed_at, spawned_at,
window ``[claimed_at - 1 s, spawned_at + 2 s]``.

``False``: provably not the recorded worker. ``None``: the needed reading
is unreadable. A missing bound is simply not applied, so missing evidence
never proves a PID recycled.
is unreadable, or the ``spawned`` upper bound is missing (identity
unproven). Missing evidence never proves a PID recycled, and never proves
it is the worker either.
"""
if start_token is not None:
try:
Expand All @@ -8584,6 +8595,13 @@ def _real_pid_started_in_claim(pid, claimed_at, spawned_at,
return False
if spawned_at is not None and created > float(spawned_at) + _OWNER_CREATE_LAG_SECONDS:
return False
if spawned_at is None:
# Only the claim lower bound is known (no run-scoped ``spawned``
# evidence: legacy/migrated runs). Any process created after the
# claim -- including one that reused the worker's PID -- fits, so
# this is UNPROVEN, not proven: termination never signals it, while
# every liveness caller still treats it as alive (FleetReview #1021).
return None
return True


Expand Down Expand Up @@ -9875,9 +9893,16 @@ def complete_task(
# non-milestone PR through fleet-merge.sh.
from hermes_cli import kanban_pr_freshness as _fresh
try:
# Arm only the card's OWN handoff PR(s), never a PR the prose
# merely mentions (FleetReview #1234 finding).
freshness = _fresh.check(
still_open, task_id=task_id,
allow_arm=not is_milestone_card(conn, task_id),
armable={
f"{r.repo}#{r.number}" for r in _open_pr.extract_pr_refs(
metadata=metadata, survivor_pr=survivor_pr,
)
},
)
except _fresh.DraftPrError as draft_err:
with write_txn(conn):
Expand Down Expand Up @@ -12727,8 +12752,8 @@ def _ret(ok: bool, reason: Optional[str] = None):
_REVIEW_NA_INABILITY = re.compile(
r"\b(?:"
r"skip(?:ped|ping|s)?|"
r"(?:could|can)\s*(?:not|n't)|cannot|unable|"
r"(?:did|was|were|does|do)\s*(?:not|n't)\s+(?:run|ran|execute|executed|attempt|attempted|try|tried|finish|finished|complete|completed|get|reach)|"
r"(?:could|can)\s*(?:not|n['\u2019]t)|can['\u2019]t|cannot|unable|"
r"(?:did|was|were|does|do)\s*(?:not|n['\u2019]t)\s+(?:run|ran|execute|executed|attempt|attempted|try|tried|finish|finished|complete|completed|get|reach)|"
r"not\s+(?:run|ran|executed|attempted|tried|finished|completed|reached)|"
r"ran\s+out|out\s+of\s+time|no\s+time\b|timed?\s*out|"
r"fail(?:ed|s)?\s+to\b|errored|crashed|blocked\s+(?:by|on)\b|"
Expand Down
14 changes: 12 additions & 2 deletions hermes_cli/kanban_pr_freshness.py
Original file line number Diff line number Diff line change
Expand Up @@ -152,9 +152,16 @@ def spawn_arm(repo: str, number: int, sha: str, task_id: str) -> Optional[str]:


def check(refs, *, task_id: str, allow_arm: bool, gh: Optional[GhFn] = None,
arm: Optional[ArmFn] = None, behind_max: int = STALE_BEHIND_MAX) -> dict:
arm: Optional[ArmFn] = None, behind_max: int = STALE_BEHIND_MAX,
armable=None) -> dict:
"""Apply (a)/(b)/(c) to OPEN PR ``refs``. Raises :class:`DraftPrError`.

``armable`` is the set of ``"o/r#n"`` keys this card actually handed off
(``metadata.pr_url``/``pr``/``pr_urls`` or ``--survivor-pr``). Only those
are ever armed for merge: a PR merely MENTIONED in summary/result prose is
context, not merge authorization. ``complete_task`` always passes it;
``None`` (direct callers) keeps the legacy every-ref behavior.

Returns a report ``{"prs": {"o/r#n": {...}}}`` for the handoff metadata.
All draft checks run before any mutation, so a refused handoff has changed
nothing on GitHub.
Expand All @@ -171,6 +178,7 @@ def check(refs, *, task_id: str, allow_arm: bool, gh: Optional[GhFn] = None,
if gh is None:
return report
arm = arm or spawn_arm
armable_keys = None if armable is None else {str(k).lower() for k in armable}

views = []
drafts = []
Expand Down Expand Up @@ -208,7 +216,9 @@ def check(refs, *, task_id: str, allow_arm: bool, gh: Optional[GhFn] = None,
entry["update_branch"] = "requested" if upd is not None else "failed (fail-open)"
# The new head has fresh CI; arming is fleet-merge's job once it is green.
continue
if allow_arm and automerge_enabled() and _ci_green(gh, ref.repo, head):
if allow_arm and armable_keys is not None and key.lower() not in armable_keys:
entry["automerge"] = "not armed: PR only mentioned, not this card's handoff PR"
elif allow_arm and automerge_enabled() and _ci_green(gh, ref.repo, head):
log = arm(ref.repo, ref.number, head, task_id)
entry["automerge"] = f"fleet-merge spawned ({log})" if log else "not spawned"
elif allow_arm and automerge_enabled():
Expand Down
16 changes: 15 additions & 1 deletion hermes_cli/kanban_survivor.py
Original file line number Diff line number Diff line change
Expand Up @@ -2515,7 +2515,21 @@ def preserve(conn, task_id, metadata=None, *, cleanup=False, workspace=None,
survivor_ref=None, survivor_pr=None, survivor_unbound=False, evidence=(),
survivor_none=False, survivor_reason=None):
"""Return a verified survivor or None for non-code work; fail closed on doubt."""
bases, held, previous = _state(conn, task_id)
try:
bases, held, previous = _state(conn, task_id)
if not isinstance(bases, dict) or not (previous is None or isinstance(previous, dict)):
raise TypeError("survivor row is not a JSON object")
except (ValueError, TypeError) as exc:
# The row is read BEFORE the main try below, so its malformed-record
# backstop cannot see a partially written `bases`/`survivor` value
# (JSONDecodeError) or a non-object one (AttributeError on `.get`):
# the workspace was retained but with no `held_reason` and no
# `workspace_held` event (FleetReview #1034). Same HOLD, same reason.
_log.exception("kanban survivor: unreadable recovery state for %s", task_id)
reason = ("survivor_unavailable: recorded survivor state is malformed "
f"({type(exc).__name__}); repair the row or recover the workspace by hand")
_hold(conn, task_id, reason)
raise _refusal(reason) from exc
if cleanup and held:
raise SurvivorUnavailable(held)
if cleanup and previous and previous.get("kind") == "none":
Expand Down
28 changes: 27 additions & 1 deletion plugins/kanban/dashboard/plugin_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -996,6 +996,27 @@ def _review_exit_refused(current_status: Optional[str], new_status: Optional[str
)


def _claimed_review_exit_refused(
conn: sqlite3.Connection, task_id: str, new_status: Optional[str],
) -> bool:
"""A ``running`` card whose current run was claimed FROM ``review`` is a
review in progress: moving it to todo/triage/scheduled would close the
reviewer run without ``request_changes`` and hand unreviewed work back
(FleetReview #999). ``ready`` is allowed -- it resumes to ``review``."""
if new_status is None or new_status in ("ready", "running", "review"):
return False
if new_status in _REVIEW_EXIT_STATUSES:
return False
row = conn.execute(
"SELECT status, current_run_id FROM tasks WHERE id = ?", (task_id,),
).fetchone()
if row is None or row["status"] != "running" or row["current_run_id"] is None:
return False
return kanban_db._retry_status_for_run(
conn, task_id, row["current_run_id"],
) == "review"


@router.patch("/tasks/{task_id}")
def update_task(task_id: str, payload: UpdateTaskBody, board: Optional[str] = Query(None)):
board = _resolve_board(board)
Expand All @@ -1004,7 +1025,9 @@ def update_task(task_id: str, payload: UpdateTaskBody, board: Optional[str] = Qu
task = kanban_db.get_task(conn, task_id)
if task is None:
raise HTTPException(status_code=404, detail=f"task {task_id} not found")
if _review_exit_refused(task.status, payload.status):
if _review_exit_refused(task.status, payload.status) or (
_claimed_review_exit_refused(conn, task_id, payload.status)
):
raise HTTPException(status_code=409, detail=_REVIEW_EXIT_REFUSAL)

review_assignee_deferred = (
Expand Down Expand Up @@ -1286,6 +1309,9 @@ def _set_status_direct(
).fetchone()
if held is None:
return False
if _claimed_review_exit_refused(conn, task_id, new_status):
# Refused BEFORE terminating the reviewer worker.
return False
released_lock = held["claim_lock"]
if held["status"] == "running" and new_status != "running":
termination = kanban_db._terminate_reclaimed_worker(
Expand Down
Loading
Loading