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
68 changes: 67 additions & 1 deletion agent/delegation_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,18 +5,33 @@
same Python process, but they are not dispatcher-owned Kanban workers. This
module lets code paths that resolve tool schemas or spawn subprocesses fail
closed for delegated children without mutating global os.environ for the parent.

Cron jobs need the same treatment for the same reason: ``cronjob(action="run")``
executes ``run_job()`` in-process, so a cron agent fired from inside a Kanban
worker would otherwise inherit that worker's dispatcher identity.
``non_dispatcher_owned_context()`` covers both cases.
"""
from __future__ import annotations

from contextlib import contextmanager
from contextvars import ContextVar
from contextvars import ContextVar, Token
from typing import Iterator, Mapping, MutableMapping

_DELEGATED_CHILD_CONTEXT: ContextVar[bool] = ContextVar(
"hermes_delegated_child_context",
default=False,
)

# Set for any in-process execution that is NOT the dispatcher-owned worker even
# though the worker's HERMES_KANBAN_* vars are legitimately in os.environ (cron
# jobs fired via the `cronjob` tool). Kept separate from
# _DELEGATED_CHILD_CONTEXT so the delegate_task-specific behaviour attached to
# that flag (subprocess env scrubbing, its own error strings) is unchanged.
_NON_DISPATCHER_OWNED_CONTEXT: ContextVar[bool] = ContextVar(
"hermes_non_dispatcher_owned_context",
default=False,
)

DELEGATED_CHILD_ENV_MARKER = "HERMES_DELEGATED_CHILD_CONTEXT"

KANBAN_ENV_KEYS: tuple[str, ...] = (
Expand Down Expand Up @@ -55,6 +70,57 @@ def is_delegated_child_context() -> bool:
return bool(_DELEGATED_CHILD_CONTEXT.get())


@contextmanager
def non_dispatcher_owned_context() -> Iterator[None]:
"""Mark in-process execution that does NOT own the dispatcher's Kanban task.

A Kanban worker is a normal CLI agent whose default toolset includes
``cronjob``; ``cronjob(action="run")`` runs ``run_job()`` inside the worker's
own process, where ``HERMES_KANBAN_TASK`` is legitimately set. Without this
marker the cron agent is misread as that worker: the kanban toolset is
force-added, the worker protocol is injected into its system prompt, and
``kanban_complete`` defaults ``task_id`` to ``$HERMES_KANBAN_TASK`` — letting
an unrelated cron job close the worker's task and overwrite real results.

Scoped via ContextVar rather than by clearing ``os.environ``: the env is
process-global and shared with the worker's own claim heartbeat, the
gateway's Kanban watchers, and concurrent cron jobs on the parallel pool, so
mutating it would starve the worker's claim and race those readers.
"""
token = _NON_DISPATCHER_OWNED_CONTEXT.set(True)
try:
yield
finally:
_NON_DISPATCHER_OWNED_CONTEXT.reset(token)


def is_dispatcher_owned_worker_context() -> bool:
"""Return True only when this execution owns the dispatcher's Kanban task.

The single predicate every ``HERMES_KANBAN_*`` identity gate should use
before trusting those vars. False for delegate_task children and for cron
jobs fired in-process from a worker.
"""
if _DELEGATED_CHILD_CONTEXT.get():
return False
return not _NON_DISPATCHER_OWNED_CONTEXT.get()


def enter_non_dispatcher_owned_context() -> Token[bool]:
"""Token-based form of :func:`non_dispatcher_owned_context`.

For callers whose scope is a long ``try`` with a matching ``finally`` rather
than a ``with`` block (``cron.scheduler.run_job``). Pair with
:func:`exit_non_dispatcher_owned_context`.
"""
return _NON_DISPATCHER_OWNED_CONTEXT.set(True)


def exit_non_dispatcher_owned_context(token: Token[bool]) -> None:
"""Restore the flag saved by :func:`enter_non_dispatcher_owned_context`."""
_NON_DISPATCHER_OWNED_CONTEXT.reset(token)


def is_delegated_child_process_context() -> bool:
"""Return True in this process or a subprocess spawned by a child."""
import os
Expand Down
22 changes: 19 additions & 3 deletions agent/skill_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -289,10 +289,12 @@ def skill_matches_platform(frontmatter: Dict[str, Any]) -> bool:
def _detect_environment(env: str) -> bool:
"""Return True when the named runtime environment is currently active.

Cached per process. Unknown env names return True (fail-open: never hide a
skill because of a tag we don't understand).
Cached per process, EXCEPT ``kanban``: that verdict is context-dependent
(a delegate_task child or an in-process cron job sees the worker's
HERMES_KANBAN_* vars without owning them), so caching it process-wide would
freeze whichever context asked first and leak it to the others.
"""
if env in _ENV_DETECT_CACHE:
if env != "kanban" and env in _ENV_DETECT_CACHE:
return _ENV_DETECT_CACHE[env]

result = True
Expand All @@ -304,6 +306,20 @@ def _detect_environment(env: str) -> bool:
# gate on (``tools/kanban_tools.py``) so the offer filter agrees with
# tool availability.
if os.getenv("HERMES_KANBAN_TASK") or os.getenv("HERMES_KANBAN_BOARD"):
# ...but only when this execution actually owns the dispatcher's
# task. A delegate_task child or a cron job fired in-process from a
# worker sees the worker's vars without being that worker.
try:
from agent.delegation_context import (
is_dispatcher_owned_worker_context,
)

_owns_dispatcher_task = is_dispatcher_owned_worker_context()
except Exception:
_owns_dispatcher_task = True
else:
_owns_dispatcher_task = False
if _owns_dispatcher_task:
result = True
else:
try:
Expand Down
27 changes: 27 additions & 0 deletions cron/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,10 @@
from hermes_cli.fallback_config import get_fallback_chain
from hermes_time import now as _hermes_now
from agent.interrupt_compat import request_hard_interrupt
from agent.delegation_context import (
enter_non_dispatcher_owned_context,
exit_non_dispatcher_owned_context,
)

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -3123,12 +3127,33 @@ def run_job(
# future writers. Acquire itself can't leak (it either blocks or returns).
_cron_session_var = _VAR_MAP["HERMES_CRON_SESSION"]
_cron_session_token = None
_non_dispatcher_token = None
try:
# Scope cron approval policy to this job. Keep the token so the finally
# restores the pre-job state instead of pinning an explicit empty value,
# which would suppress the legacy os.environ fallback used by standalone
# cron entrypoints and tests.
_cron_session_token = _cron_session_var.set("1")

# Mark this job as NOT the dispatcher-owned kanban worker.
#
# A kanban worker is a normal `hermes chat -q` CLI agent whose default
# toolset includes `cronjob`, running with HERMES_KANBAN_TASK
# legitimately in its own env; `cronjob(action="run")` calls
# run_one_job() -> run_job() right here in that process. Without this
# marker the cron agent is misread as that worker: the kanban toolset is
# force-added, the worker protocol is injected into its system prompt,
# and kanban_complete defaults task_id to $HERMES_KANBAN_TASK -- letting
# an unrelated cron job close the worker's task and overwrite real
# results.
#
# A ContextVar, NOT an os.environ clear: the env is process-global and
# shared with the worker's own claim heartbeat (run_agent._touch_activity
# -> heartbeat_current_worker_from_env, which would starve and let the
# dispatcher reclaim a live task), the gateway's kanban watchers, and
# concurrent cron jobs on the parallel pool. contextvars.copy_context()
# at the run_conversation hop carries this into the agent thread.
_non_dispatcher_token = enter_non_dispatcher_owned_context()
if _job_workdir:
os.environ["TERMINAL_CWD"] = _job_workdir
logger.info("Job '%s': using workdir %s", job_id, _job_workdir)
Expand Down Expand Up @@ -3777,6 +3802,8 @@ def _heartbeat_run_claim_if_due():
clear_session_vars(_ctx_tokens)
if _cron_session_token is not None:
_cron_session_var.reset(_cron_session_token)
if _non_dispatcher_token is not None:
exit_non_dispatcher_owned_context(_non_dispatcher_token)
for _var_name in _cron_delivery_vars:
_VAR_MAP[_var_name].set("")
if _session_db:
Expand Down
13 changes: 13 additions & 0 deletions model_tools.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,17 @@ def _is_delegated_child_context() -> bool:
return False


def _is_dispatcher_owned_worker() -> bool:
"""False when HERMES_KANBAN_* is present but this execution does not own it
(delegate_task child, or a cron job fired in-process from a worker)."""
try:
from agent.delegation_context import is_dispatcher_owned_worker_context

return is_dispatcher_owned_worker_context()
except Exception:
return True


# =============================================================================
# Async Bridging (single source of truth -- used by registry.dispatch too)
# =============================================================================
Expand Down Expand Up @@ -342,6 +353,7 @@ def get_tool_definitions(
bool(os.environ.get("HERMES_KANBAN_TASK")),
bool(skip_tool_search_assembly),
_is_delegated_child_context(),
_is_dispatcher_owned_worker(),
profile_scope,
)
cached = _tool_defs_cache.get(cache_key) if cache_key is not None else None
Expand Down Expand Up @@ -391,6 +403,7 @@ def _compute_tool_definitions(
if (
os.environ.get("HERMES_KANBAN_TASK")
and not _is_delegated_child_context()
and _is_dispatcher_owned_worker()
and "kanban" not in effective_enabled_toolsets
):
# Dispatcher-spawned workers are scoped by HERMES_KANBAN_TASK and
Expand Down
Loading
Loading