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
13 changes: 9 additions & 4 deletions hermes_cli/kanban_budget.py
Original file line number Diff line number Diff line change
Expand Up @@ -324,10 +324,14 @@ def read_pause_marker(board_dir) -> Optional[dict]:
return None


def _run_notify(argv: list[str]) -> None:
"""Fire notify.py out-of-agent. Best-effort by contract."""
def _run_notify(argv: list[str]) -> bool:
"""Fire notify.py out-of-agent. Best-effort by contract.

Returns whether the send exited 0, so a caller that latches "already
announced" state can refuse to latch a page that never went out.
"""
try:
subprocess.run(
proc = subprocess.run(
argv,
check=False,
stdin=subprocess.DEVNULL,
Expand All @@ -336,7 +340,8 @@ def _run_notify(argv: list[str]) -> None:
timeout=30,
)
except Exception:
pass
return False
return proc.returncode == 0


def _notify_script_path(home=None) -> Optional[str]:
Expand Down
104 changes: 89 additions & 15 deletions hermes_cli/kanban_db.py
Original file line number Diff line number Diff line change
Expand Up @@ -5726,18 +5726,34 @@ def _can_adopt_home(session_id: str) -> bool:
return not os.environ.get("HERMES_DELEGATED_CHILD_CONTEXT")


_HOME_UNREAD = object()


def _read_home_session(conn: sqlite3.Connection, task_id: str) -> Optional[str]:
row = conn.execute(
"SELECT session_id FROM tasks WHERE id = ?", (task_id,)
).fetchone()
return row["session_id"] if row is not None else None


def record_foreign_action(
conn: sqlite3.Connection, task_id: str, action: str, actor: MutationActor
conn: sqlite3.Connection, task_id: str, action: str, actor: MutationActor,
*, home_before: Any = _HOME_UNREAD,
) -> None:
"""Append the audit comment for an overridden foreign-session mutation.

An ``--operator`` override records an ``operator_override`` event only:
no comment, so nothing pages the home session.

``home_before`` is the home read BEFORE the guarded mutation ran; the
mutation itself may have re-stamped ``tasks.session_id`` (``update
--session``), and the audit must name the displaced home (C7 k103).
"""
sess = ", ".join(actor.session_ids) or "no-session"
home_row = conn.execute(
"SELECT session_id FROM tasks WHERE id = ?", (task_id,)
).fetchone()
prev_home = (
_read_home_session(conn, task_id)
if home_before is _HOME_UNREAD else home_before
)
if actor.operator:
last = conn.execute(
"SELECT kind, payload FROM task_events WHERE task_id = ? "
Expand Down Expand Up @@ -5766,11 +5782,10 @@ def record_foreign_action(
"reason": actor.operator,
"by_sessions": list(actor.session_ids),
"by_profile": actor.profile,
"home": (home_row["session_id"] if home_row is not None else None),
"home": prev_home,
},
)
return
prev_home = home_row["session_id"] if home_row is not None else None
new_home = (
actor.session_ids[0]
if action in REHOME_ON_TAKEOVER_ACTIONS
Expand Down Expand Up @@ -5846,13 +5861,19 @@ def wrapper(conn, *args, **kwargs):
return fn(conn, *args, **kwargs)
task_id = sig.bind_partial(conn, *args, **kwargs).arguments.get(task_param)
override = check_home_session(conn, str(task_id), action)
home_before = (
_read_home_session(conn, str(task_id))
if override is not None else _HOME_UNREAD
)
token = _MUTATION_ACTOR.set(None)
try:
result = fn(conn, *args, **kwargs)
finally:
_MUTATION_ACTOR.reset(token)
if override is not None and _mutation_succeeded(result):
record_foreign_action(conn, str(task_id), action, override)
record_foreign_action(
conn, str(task_id), action, override, home_before=home_before
)
return result

wrapper.__home_session_action__ = action
Expand Down Expand Up @@ -8480,7 +8501,7 @@ def _spawned_owner_alive(conn, task_id, row, host_prefix):
pid = int(payload["pid"])
except (TypeError, ValueError, KeyError):
continue
token = payload.get("start_token") if isinstance(payload, dict) else None
token = _spawn_start_token(payload)
candidates.append((pid, spawned["run_id"] is None,
spawned["created_at"], token))
for pid, late, spawned_at, token in candidates:
Expand Down Expand Up @@ -8547,6 +8568,36 @@ def _pid_start_token(pid: int) -> Optional[float]:
return None


def _boot_id() -> Optional[str]:
"""This boot's identity, or None where the start token is not boot-relative.

A Linux start token is /proc starttime ticks SINCE BOOT, so after a reboot a
new process at a reused PID can carry the very token a pre-boot worker
recorded. The ``spawned`` event stamps this id next to the token so a token
from another boot is never read as identity (C7 k109).
"""
try:
with open("/proc/sys/kernel/random/boot_id", encoding="ascii") as fh:
return fh.read().strip() or None
except OSError:
return None


def _spawn_start_token(payload: Any) -> Any:
"""The ``start_token`` a ``spawned`` payload recorded, if still meaningful.

A token stamped under a different boot proves nothing about a live PID
now, so it is dropped and the causal-window check decides instead.
"""
if not isinstance(payload, dict):
return None
token = payload.get("start_token")
recorded_boot = payload.get("boot_id")
if token is not None and recorded_boot and recorded_boot != _boot_id():
return None
return token


def _pid_create_time(pid: int) -> Optional[float]:
"""Wall-clock epoch creation time of ``pid``, or None if unreadable."""
try:
Expand Down Expand Up @@ -8703,7 +8754,7 @@ def _worker_owner_window(
payload = json.loads(ev["payload"] or "{}")
if int(payload["pid"]) == int(pid):
spawned_at = ev["created_at"]
start_token = payload.get("start_token")
start_token = _spawn_start_token(payload)
break
except (TypeError, ValueError, KeyError, AttributeError):
continue
Expand Down Expand Up @@ -13192,7 +13243,10 @@ def _in_txn() -> tuple[bool, Optional[str]]:
try:
with write_txn(conn):
ok, detail = _in_txn()
if not ok and opened:
# A refusal rolls back everything this call wrote: the opened
# claim AND a posted coverage comment the gate just rejected
# (C7 k107) -- a refused record must not persist on the run.
if not ok and (opened or posted):
raise _SendBackRefused(detail)
except _SendBackRefused as exc:
return False, exc.detail
Expand Down Expand Up @@ -15957,12 +16011,17 @@ def _notify_rate_limit_circuit(
if same_episode and int(prior) >= until:
return
seen[pool] = max(int(prior), until) if same_episode else until
marker.parent.mkdir(parents=True, exist_ok=True)
marker.write_text(json.dumps(seen), encoding="utf-8")

def _latch() -> None:
marker.parent.mkdir(parents=True, exist_ok=True)
marker.write_text(json.dumps(seen), encoding="utf-8")

if same_episode:
_latch()
return # hold extended; already announced
script = _kbudget._notify_script_path()
if script is None:
_latch()
return
import sys as _sys
body = (
Expand All @@ -15971,10 +16030,14 @@ def _notify_rate_limit_circuit(
f"Holding {pool} spawns until "
f"{time.strftime('%H:%M:%S', time.localtime(until))}; other pools unaffected."
)
_kbudget._run_notify([
delivered = _kbudget._run_notify([
_sys.executable, script, "--channel", "discord",
"--target", _kbudget.RECOVERY_TARGET, "--send", body,
])
# Latch the episode only once the page went out: a failed send
# must be retried on the next tick, not silently swallowed (C7 k105).
if delivered is not False:
_latch()
except Exception as exc:
_log.warning("kanban rate-limit circuit notify failed (%s: %s)", type(exc).__name__, exc)

Expand Down Expand Up @@ -17681,7 +17744,8 @@ def end_orphaned_terminal_runs(
host_prefix = f"{_claimer_id().split(':', 1)[0]}:"
rows = conn.execute(
"SELECT r.id, r.task_id, r.worker_pid, r.claim_lock, r.started_at, "
" r.last_heartbeat_at, r.max_runtime_seconds, t.status AS task_status "
" r.last_heartbeat_at, r.max_runtime_seconds, r.metadata, "
" t.status AS task_status "
"FROM task_runs r JOIN tasks t ON t.id = r.task_id "
"WHERE r.ended_at IS NULL AND t.status IN ('done', 'archived')"
).fetchall()
Expand Down Expand Up @@ -17709,6 +17773,13 @@ def end_orphaned_terminal_runs(
"max_runtime_seconds": limit,
"now": now,
}
# Merge, never replace: the open run may already carry metadata
# (``pool`` from _stamp_run_pool) that the ledger reads (C7 k110).
try:
prior_meta = json.loads(row["metadata"]) if row["metadata"] else {}
except (json.JSONDecodeError, TypeError):
prior_meta = {}
run_meta = {**prior_meta, **payload} if isinstance(prior_meta, dict) else payload
with write_txn(conn):
cur = conn.execute(
"UPDATE task_runs SET status = 'reclaimed', outcome = ?, "
Expand All @@ -17718,7 +17789,7 @@ def end_orphaned_terminal_runs(
(
ORPHANED_TERMINAL_TASK_OUTCOME,
f"run left open on a {row['task_status']} card; ended by reaper",
json.dumps(payload, ensure_ascii=False),
json.dumps(run_meta, ensure_ascii=False),
now,
run_id,
),
Expand Down Expand Up @@ -18845,6 +18916,9 @@ def _set_worker_pid(
start_token = _pid_start_token(int(pid))
if start_token is not None:
spawn_payload["start_token"] = start_token
boot_id = _boot_id()
if boot_id:
spawn_payload["boot_id"] = boot_id
if pool is not None:
# The relay pool this spawn was charged to (t_38be6b10).
spawn_payload["pool"] = pool
Expand Down
10 changes: 10 additions & 0 deletions hermes_cli/model_switch.py
Original file line number Diff line number Diff line change
Expand Up @@ -2507,6 +2507,16 @@ def switch_model(
and not validation.get("corrected_model")
and not any(
isinstance(_cp, dict)
# Only a declaration on the CURRENT provider vouches for the id;
# another provider's declaration says nothing about this endpoint
# (C7 k112) -- same match as the override block above.
and (
target_provider.lower() in custom_provider_aliases(
str(_cp.get("name", "") or ""),
str(_cp.get("provider_key") or ""),
)
or _cp.get("base_url", "") == base_url
)
and (
_cp.get("model") == new_model
or new_model in _declared_model_ids(_cp.get("models", {}))
Expand Down
10 changes: 7 additions & 3 deletions hermes_cli/process_env_files.py
Original file line number Diff line number Diff line change
Expand Up @@ -137,10 +137,14 @@ def strip_overlay(env: MutableMapping[str, str]) -> MutableMapping[str, str]:
# PATH is routinely re-edited after start (venv/tool dirs), so match
# components rather than the whole value: drop only what we added.
before = set((old or "").split(os.pathsep))
after = set(new.split(os.pathsep))
added = {p for p in new.split(os.pathsep) if p and p not in before}
env["PATH"] = os.pathsep.join(
p for p in env["PATH"].split(os.pathsep) if p not in added
)
kept = [p for p in env["PATH"].split(os.pathsep) if p not in added]
# ...and put back what the overlay REMOVED (a file that replaced
# PATH rather than prepending to it), in original order (C7 k114).
kept += [p for p in (old or "").split(os.pathsep)
if p and p not in after and p not in kept]
env["PATH"] = os.pathsep.join(kept)
continue
if env.get(key) != new:
continue
Expand Down
32 changes: 32 additions & 0 deletions hermes_cli/web_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -267,6 +267,37 @@ def _parent_start_markers_match(actual: str, expected: str) -> bool:
# when the same module is used across TestClient instances or uvicorn reloads.
# ---------------------------------------------------------------------------

class _DesktopCronDispatchGate:
"""Serve-process cron dispatch gate: the shared-checkout admission fence.

The gateway ticker passes ``_CronDispatchGate``; the desktop in-process
ticker passed nothing, so ``cron.scheduler.tick`` never consulted the serve
admission hold and could launch a cron agent after the serve process had
acknowledged quiescence (C7 k115). The gate is resolved per tick, so a hold
configured after the ticker thread started is still honoured.
"""

def __call__(self) -> bool:
return True

def admit(self):
try:
from gateway.checkout_admission import AdmissionRefused, process_gate

gate = process_gate("serve")
except Exception:
# Same as the gateway's _checkout_admission_gate: an unresolvable
# gate is "disabled", never a permanent cron stop.
_log.exception("desktop cron: checkout admission gate unavailable")
return lambda: None
if gate is None:
return lambda: None
try:
return gate.admit("cron:tick", internal=False).release
except AdmissionRefused:
return None


def _start_desktop_cron_ticker(stop_event: "threading.Event", interval: int = 60) -> None:
"""Tick the cron scheduler from inside the desktop dashboard backend.

Expand Down Expand Up @@ -297,6 +328,7 @@ def _start_desktop_cron_ticker(stop_event: "threading.Event", interval: int = 60

start_kwargs: dict = {"interval": interval}
if isinstance(provider, InProcessCronScheduler):
start_kwargs["can_dispatch"] = _DesktopCronDispatchGate()
try:
from hermes_cli.profiles import profiles_to_serve

Expand Down
7 changes: 6 additions & 1 deletion plugins/kanban/dashboard/dist/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -991,10 +991,15 @@
}, [selectedIds, requestMoveConfirm, requestCompletionSummary, performMoveTask]);

const createTask = useCallback(function (body) {
// Home the card to the viewing session, so the default "this" facet of
// a /kanban?session=<id> link still shows the card just created.
const payload = viewerSession && !(body && body.session_id)
? Object.assign({}, body, { session_id: viewerSession })
: body;
return SDK.fetchJSON(withBoard(`${API}/tasks`, board), {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(body),
body: JSON.stringify(payload),
}).then(function (res) {
// Surface dispatcher-presence warnings (e.g. "no gateway is
// running") via the existing error banner channel. Not fatal —
Expand Down
5 changes: 5 additions & 0 deletions plugins/kanban/dashboard/plugin_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -737,6 +737,10 @@ class CreateTaskBody(BaseModel):
# Explicit project link; when omitted, create_task inherits the board's
# scoped project (if any) so a project-scoped board anchors every task.
project_id: Optional[str] = None
# Home session. The dashboard sends the viewer's ``?session=`` so a card
# created from a session link lands in that session's home and stays
# visible under the default "this" facet (C7 k121).
session_id: Optional[str] = None


@router.post("/tasks")
Expand Down Expand Up @@ -765,6 +769,7 @@ def create_task(payload: CreateTaskBody, board: Optional[str] = Query(None)):
provider_override=payload.provider_override,
reasoning_effort=payload.reasoning_effort,
project_id=payload.project_id,
session_id=(payload.session_id or "").strip() or None,
board=board,
)
task = kanban_db.get_task(conn, task_id)
Expand Down
8 changes: 6 additions & 2 deletions tests/hermes_cli/test_desktop_cron_ticker_profiles.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,9 @@ def test_single_profile_keeps_legacy_path(monkeypatch, _providers, tmp_path):

ws._start_desktop_cron_ticker(threading.Event(), interval=9)

assert builtin.start_kwargs == {"interval": 9}
# can_dispatch is the serve admission gate (C7 k115); the legacy
# single-store path means no profile_homes.
assert {k: v for k, v in builtin.start_kwargs.items() if k != "can_dispatch"} == {"interval": 9}


def test_enumeration_failure_fails_open(monkeypatch, _providers):
Expand All @@ -97,7 +99,9 @@ def _boom(**_kw):

ws._start_desktop_cron_ticker(threading.Event(), interval=11)

assert builtin.start_kwargs == {"interval": 11}
# can_dispatch is the serve admission gate (C7 k115); the legacy
# single-store path means no profile_homes.
assert {k: v for k, v in builtin.start_kwargs.items() if k != "can_dispatch"} == {"interval": 11}


def test_external_provider_never_gets_profile_homes(monkeypatch, tmp_path):
Expand Down
Loading
Loading