diff --git a/cron/scheduler.py b/cron/scheduler.py index f1cd08125c..ab9d2905f6 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -1947,6 +1947,283 @@ def _run_doc_header(job: dict, title: str, job_id: str, prompt: str) -> str: ) +_RunResult = tuple[bool, str, str, Optional[str]] + + +def _prepare_job_prompt( + job: dict, job_id: str, job_name: str, extra_prompt: Optional[str], cancel_event, +) -> tuple[Optional[_RunResult], Optional[str]]: + """Run every pre-agent gate and build the prompt. Returns ``(early_result, prompt)``: an early + result short-circuits ``run_job`` (no_agent job, empty payload, monitor gate, wake gate, + injection block, empty prompt); otherwise ``prompt`` is set.""" + # Fail closed on a corrupt config.yaml: defaults would let auto-detection bill a provider the + # user never chose. no_agent jobs are exempt. Escape hatch: HERMES_IGNORE_USER_CONFIG=1. + if not job.get("no_agent"): + from hermes_cli.config import InvalidUserConfigError, require_parseable_user_config + + try: + require_parseable_user_config() + except InvalidUserConfigError as exc: + logger.error("Job '%s': refusing to run — %s", job_id, exc) + return (False, f"# Cron Job: {job_name}\n\nError: {exc}\n", "", str(exc)), None + + # no_agent short-circuits BEFORE importing run_agent / opening SessionDB. + if job.get("no_agent"): + return _run_no_agent_job(job, job_id, job_name, cancel_event), None + + # Legacy / hand-edited job with nothing to run: pause it instead of waking the LLM every fire. + from cron.jobs import EMPTY_PAYLOAD_ERROR, job_payload_is_empty + + if job_payload_is_empty(job): + return _block_and_pause_job(job_id, job_name, EMPTY_PAYLOAD_ERROR), None + + _early, extra_prompt = _apply_monitor_gate(job, job_id, job_name, extra_prompt) + if _early is not None: + return _early, None + + # Wake-gate: run the pre-check script BEFORE building the prompt; its result is passed into + # _build_job_prompt so the script runs only once. + prerun_script = None + script_path = job.get("script") + if script_path: + prerun_script = _run_job_script_with_claim_heartbeat(job, script_path, cancel_event=cancel_event) + _ran_ok, _script_output = prerun_script + if _ran_ok and not _parse_wake_gate(_script_output): + logger.info("Job '%s' (ID: %s): wakeAgent=false, skipping agent run", job_name, job_id) + silent_doc = ( + f"# Cron Job: {job_name}\n\n" + f"**Job ID:** {job_id}\n" + f"**Run Time:** {_hermes_now().strftime('%Y-%m-%d %H:%M:%S')}\n\n" + "Script gate returned `wakeAgent=false` — agent skipped.\n" + ) + return (True, silent_doc, SILENT_MARKER, None), None + + try: + prompt = _build_job_prompt(job, prerun_script=prerun_script, extra_prompt=extra_prompt) + except CronPromptInjectionBlocked as block_exc: + # Injection scanner tripped: refuse this tick and tell the operator WHY. + logger.warning( + "Job '%s' (ID: %s): blocked by prompt-injection scanner — %s", job_name, job_id, block_exc, + ) + blocked_doc = ( + f"# Cron Job: {job_name}\n\n" + f"**Job ID:** {job_id}\n" + f"**Run Time:** {_hermes_now().strftime('%Y-%m-%d %H:%M:%S')}\n" + f"**Status:** BLOCKED\n\n" + "The assembled prompt (user prompt + loaded skill content) tripped " + "the cron injection scanner and the agent was NOT run.\n\n" + f"**Scanner result:** {block_exc}\n\n" + "Audit the skill(s) attached to this job for prompt-injection " + "payloads or invisible-unicode markers. If the skill is legitimate " + "and the match is a false positive, rephrase the content to avoid " + "the threat pattern (`tools/cronjob_tools.py::_CRON_THREAT_PATTERNS`)." + ) + return (False, blocked_doc, "", str(block_exc)), None + if prompt is None: + logger.info("Job '%s': script produced no output, skipping AI call.", job_name) + return (True, "", SILENT_MARKER, None), None + return None, prompt + + +_CRON_DELIVERY_VARS = ( + "HERMES_CRON_AUTO_DELIVER_PLATFORM", + "HERMES_CRON_AUTO_DELIVER_CHAT_ID", + "HERMES_CRON_AUTO_DELIVER_THREAD_ID", +) + + +class _CronRunScope: + """Per-run ContextVar / tool-cwd scope for ``run_job`` (ContextVars, not os.environ, so + parallel jobs don't clobber each other). Construct before the try, ``enter()`` as its first + statement, ``exit()`` in the finally — every setter here has a matching reset there. + + HERMES_SESSION_* are deliberately NOT seeded from job["origin"]: it is delivery metadata, not + a sender, and terminal/tts/skills/send_message tools would act as if the origin user were + driving the agent. Delivery reads job["origin"] / HERMES_CRON_AUTO_DELIVER_* directly. + """ + + def __init__(self, job: dict, job_id: str, execution_id: Optional[str]): + from gateway.session_context import set_session_vars, _VAR_MAP + from tools.terminal_tool import record_session_cwd + + self._var_map = _VAR_MAP + # Resolve workdir BEFORE set_session_vars so it owns the _SESSION_CWD set/clear. + self.workdir = _resolve_job_workdir(job, job_id) + self._ctx_tokens = set_session_vars( + platform="", + chat_id="", + chat_name="", + # Cron can't receive completions after its turn; async delegation output could + # otherwise route to an unrelated chat via the ambient session key => inline delegation. + async_delivery=False, + cwd=self.workdir or "", + ) + for name in _CRON_DELIVERY_VARS: + _VAR_MAP[name].set("") + # Workdir binds to the per-run task id (tool-layer cwd authority) instead of mutating + # global TERMINAL_CWD; _SESSION_CWD above remains the prompt/context-file authority. + self.task_id = f"cron:{job_id}:{execution_id or job.get('execution_id') or uuid.uuid4().hex}" + if self.workdir: + record_session_cwd(self.task_id, self.workdir) + self._cron_session_var = _VAR_MAP["HERMES_CRON_SESSION"] + self._cron_session_token = None + self._non_dispatcher_token = None + + def enter(self) -> None: + # Scope cron approval policy; exit() RESETS via token (pinning "" would suppress the legacy + # os.environ fallback used by standalone entrypoints/tests). + self._cron_session_token = self._cron_session_var.set("1") + # Mark NOT the kanban worker: a worker's cronjob(action="run") lands here with + # HERMES_KANBAN_TASK in env, and an unrelated job could close the worker's task. Must be a + # ContextVar, NOT an os.environ clear (env is shared with the worker heartbeat and + # concurrent jobs); copy_context() carries it into the agent thread. + self._non_dispatcher_token = enter_non_dispatcher_owned_context() + + def exit(self) -> None: + from gateway.session_context import clear_session_vars + from tools.terminal_tool import clear_session_cwd + + clear_session_cwd(self.task_id) + clear_session_vars(self._ctx_tokens) # also clears _SESSION_CWD + if self._cron_session_token is not None: + self._cron_session_var.reset(self._cron_session_token) + if self._non_dispatcher_token is not None: + exit_non_dispatcher_owned_context(self._non_dispatcher_token) + for name in _CRON_DELIVERY_VARS: + self._var_map[name].set("") + + +def _reload_dotenv_and_publish_delivery_target(job: dict) -> None: + """Re-read .env for this run and publish the auto-deliver target into the session ContextVars.""" + # Reset the secret-source cache FIRST or a Bitwarden/BSM-backed secret is never re-resolved + # (only the placeholder reloads -> 401s). + from hermes_cli.env_loader import load_hermes_dotenv, reset_secret_source_cache + from gateway.session_context import _VAR_MAP + + reset_secret_source_cache() + load_hermes_dotenv(hermes_home=_get_hermes_home()) + + delivery_target = _resolve_delivery_target(job) + if delivery_target: + _VAR_MAP["HERMES_CRON_AUTO_DELIVER_PLATFORM"].set(delivery_target["platform"]) + _VAR_MAP["HERMES_CRON_AUTO_DELIVER_CHAT_ID"].set(str(delivery_target["chat_id"])) + _VAR_MAP["HERMES_CRON_AUTO_DELIVER_THREAD_ID"].set( + "" if delivery_target.get("thread_id") is None else str(delivery_target["thread_id"]) + ) + + +@dataclass +class _CronAgentSetup: + """Everything ``AIAgent(...)`` needs that is resolved from job + config (or a preflight block).""" + blocked: Optional[_RunResult] = None + model: str = "" + runtime: dict = None + prefill_messages: Any = None + max_iterations: Any = None + reasoning_config: Any = None + fallback_model: Any = None + credential_pool: Any = None + + +def _resolve_cron_agent_setup(job: dict, job_id: str, job_name: str, jc) -> _CronAgentSetup: + """Resolve model/runtime/reasoning/pool for the run, in the original gate order: exfil guard -> + preflight (may block) -> runtime -> drift check -> fallback chain -> credential pool -> MCP.""" + _cfg = jc.cfg + setup = _CronAgentSetup(model=jc.model) + setup.prefill_messages = _load_prefill_messages(_cfg, job_id) + + # resolve_turn_limit() honors none/unlimited (sys.maxsize) and explicit 0 / null. + from hermes_cli.config import resolve_turn_limit as _resolve_turn_limit + _mt = _cfg.get("agent", {}).get("max_turns") + if _mt is None: + _mt = _cfg.get("max_turns") + setup.max_iterations = _resolve_turn_limit(_mt) + + # Runtime backstop (CWE-200/522): fail closed BEFORE resolution on a provider/base_url pair + # that would ship a stored credential off-host; hand-written jobs bypass create-time checks. + _guard_job_credential_exfil(job) + + setup.blocked = _preflight_or_block(job, job_id, job_name, _cfg) + if setup.blocked is not None: + return setup + + primary_model_for_drift = setup.model + setup.runtime, setup.model, primary_provider_for_drift = _resolve_job_runtime(job, job_id, jc) + setup.reasoning_config = _resolve_job_reasoning_config( + job, _cfg if isinstance(_cfg, dict) else {}, str(setup.model) + ) + _check_model_drift( + job, job_id, _cfg, setup.runtime, primary_provider_for_drift, primary_model_for_drift, + ) + setup.fallback_model = get_fallback_chain(_cfg) or None + setup.credential_pool = _load_credential_pool(setup.runtime, job_id) + # MCP servers must be registered before AIAgent is constructed. + _init_cron_mcp_tools(job_id) + return setup + + +def _construct_cron_agent(AIAgent, job: dict, _cfg: dict, setup: _CronAgentSetup, *, workdir, session_id, session_db): + runtime = setup.runtime + pr = _cfg.get("provider_routing") or {} + return AIAgent( + model=setup.model, + api_key=runtime.get("api_key"), + base_url=runtime.get("base_url"), + provider=runtime.get("provider"), + requested_provider=runtime.get("requested_provider"), + api_mode=runtime.get("api_mode"), + request_overrides=runtime.get("request_overrides"), + acp_command=runtime.get("command"), + acp_args=runtime.get("args"), + max_iterations=setup.max_iterations, + reasoning_config=setup.reasoning_config, + prefill_messages=setup.prefill_messages, + fallback_model=setup.fallback_model, + credential_pool=setup.credential_pool, + providers_allowed=pr.get("only"), + providers_ignored=pr.get("ignore"), + providers_order=pr.get("order"), + provider_sort=pr.get("sort"), + openrouter_min_coding_score=(_cfg.get("openrouter") or {}).get("min_coding_score"), + enabled_toolsets=_resolve_cron_enabled_toolsets(job, _cfg), + disabled_toolsets=_resolve_cron_disabled_toolsets(_cfg), + quiet_mode=True, + # Project context files only with a configured workdir; SOUL.md always. + skip_context_files=not bool(workdir), + load_soul_identity=True, + skip_memory=False, + skip_background_review=True, # Cron has no human-in-the-loop need for skill/memory review forks (~30K tok/event) + platform="cron", + session_id=session_id, + session_db=session_db, + ) + + +class _FireAudit: + """One usage_audit.jsonl line per fire (created once the agent exists; fire id + start clock).""" + + def __init__(self, job: dict, job_id: str, model: str): + self.job, self.job_id, self.model = job, job_id, model + self.fire_id = uuid.uuid4().hex + self.t_start = time.monotonic() + + def write(self, result: dict, error: Optional[str]) -> None: + _write_usage_audit({ + "ts": _utcnow_iso_ms(), + "job_id": self.job_id, + "fire_id": self.fire_id, + "prompt_tokens": result.get("prompt_tokens"), + "completion_tokens": result.get("completion_tokens"), + "total_tokens": result.get("total_tokens"), + "response_silent": bool(result.get("response_silent")), + "deliver_target": self.job.get("deliver"), + "model": self.model or None, + "duration_ms": int((time.monotonic() - self.t_start) * 1000), + "error": error, + }) + + + def run_job( job: dict, *, @@ -1964,271 +2241,59 @@ def run_job( job_id = job["id"] job_name = str(job.get("name") or job.get("prompt") or job_id or "cron job") - # Fail closed on a corrupt config.yaml: defaults would let auto-detection bill a provider the - # user never chose. no_agent jobs are exempt. Escape hatch: HERMES_IGNORE_USER_CONFIG=1. - if not job.get("no_agent"): - from hermes_cli.config import ( - InvalidUserConfigError, - require_parseable_user_config, - ) - - try: - require_parseable_user_config() - except InvalidUserConfigError as exc: - logger.error("Job '%s': refusing to run — %s", job_id, exc) - return (False, f"# Cron Job: {job_name}\n\nError: {exc}\n", "", str(exc)) - - # no_agent short-circuits BEFORE importing run_agent / opening SessionDB. - if job.get("no_agent"): - return _run_no_agent_job(job, job_id, job_name, cancel_event) - - # Legacy / hand-edited job with nothing to run: pause it instead of waking the LLM every fire. - from cron.jobs import EMPTY_PAYLOAD_ERROR, job_payload_is_empty - - if job_payload_is_empty(job): - return _block_and_pause_job(job_id, job_name, EMPTY_PAYLOAD_ERROR) - - _early, extra_prompt = _apply_monitor_gate(job, job_id, job_name, extra_prompt) - if _early is not None: - return _early - + early, prompt = _prepare_job_prompt(job, job_id, job_name, extra_prompt, cancel_event) + if early is not None: + return early from run_agent import AIAgent - # Wake-gate: run the pre-check script BEFORE building the prompt; its result is passed into - # _build_job_prompt so the script runs only once. - prerun_script = None - script_path = job.get("script") - if script_path: - prerun_script = _run_job_script_with_claim_heartbeat( - job, script_path, cancel_event=cancel_event, - ) - _ran_ok, _script_output = prerun_script - if _ran_ok and not _parse_wake_gate(_script_output): - logger.info("Job '%s' (ID: %s): wakeAgent=false, skipping agent run", job_name, job_id) - silent_doc = ( - f"# Cron Job: {job_name}\n\n" - f"**Job ID:** {job_id}\n" - f"**Run Time:** {_hermes_now().strftime('%Y-%m-%d %H:%M:%S')}\n\n" - "Script gate returned `wakeAgent=false` — agent skipped.\n" - ) - return True, silent_doc, SILENT_MARKER, None - - try: - prompt = _build_job_prompt(job, prerun_script=prerun_script, extra_prompt=extra_prompt) - except CronPromptInjectionBlocked as block_exc: - # Injection scanner tripped: refuse this tick and tell the operator WHY. - logger.warning( - "Job '%s' (ID: %s): blocked by prompt-injection scanner — %s", - job_name, job_id, block_exc, - ) - blocked_doc = ( - f"# Cron Job: {job_name}\n\n" - f"**Job ID:** {job_id}\n" - f"**Run Time:** {_hermes_now().strftime('%Y-%m-%d %H:%M:%S')}\n" - f"**Status:** BLOCKED\n\n" - "The assembled prompt (user prompt + loaded skill content) tripped " - "the cron injection scanner and the agent was NOT run.\n\n" - f"**Scanner result:** {block_exc}\n\n" - "Audit the skill(s) attached to this job for prompt-injection " - "payloads or invisible-unicode markers. If the skill is legitimate " - "and the match is a false positive, rephrase the content to avoid " - "the threat pattern (`tools/cronjob_tools.py::_CRON_THREAT_PATTERNS`)." - ) - return False, blocked_doc, "", str(block_exc) - if prompt is None: - logger.info("Job '%s': script produced no output, skipping AI call.", job_name) - return True, "", SILENT_MARKER, None _cron_session_id = f"cron_{job_id}_{_hermes_now().strftime('%Y%m%d_%H%M%S')}" - logger.info("Running job '%s' (ID: %s)", job_name, job_id) logger.info("Prompt: %s", prompt[:100]) agent = None model = "" - - # ContextVars, not os.environ (process-global), so parallel jobs don't clobber each other. - from gateway.session_context import set_session_vars, clear_session_vars, _VAR_MAP - - # Do NOT seed HERMES_SESSION_* from job["origin"]: it is delivery metadata, not a sender, and - # terminal/tts/skills/send_message tools would act as if the origin user were driving the - # agent. Delivery reads job["origin"] and HERMES_CRON_AUTO_DELIVER_* directly, so blanking is - # safe. Resolve workdir BEFORE set_session_vars so it owns the _SESSION_CWD set/clear. - _job_workdir = _resolve_job_workdir(job, job_id) - - _ctx_tokens = set_session_vars( - platform="", - chat_id="", - chat_name="", - # Cron can't receive completions after its turn; async delegation output could otherwise - # route to an unrelated chat via the ambient session key. Stateless => inline delegation. - async_delivery=False, - cwd=_job_workdir or "", - ) - _cron_delivery_vars = ( - "HERMES_CRON_AUTO_DELIVER_PLATFORM", - "HERMES_CRON_AUTO_DELIVER_CHAT_ID", - "HERMES_CRON_AUTO_DELIVER_THREAD_ID", - ) - for _var_name in _cron_delivery_vars: - _VAR_MAP[_var_name].set("") - - # Bind workdir to the per-run task id (tool-layer cwd authority) instead of mutating global - # TERMINAL_CWD; _SESSION_CWD above remains the prompt/context-file authority. - _cron_task_id = ( - f"cron:{job_id}:" - f"{execution_id or job.get('execution_id') or uuid.uuid4().hex}" - ) - from tools.terminal_tool import clear_session_cwd as _clear_tool_session_cwd - from tools.terminal_tool import record_session_cwd as _record_tool_session_cwd - if _job_workdir: - _record_tool_session_cwd(_cron_task_id, _job_workdir) - _cron_session_var = _VAR_MAP["HERMES_CRON_SESSION"] - _cron_session_token = None - _non_dispatcher_token = None _session_db = None + _audit: Optional[_FireAudit] = None + scope = _CronRunScope(job, job_id, execution_id) try: - - # Scope cron approval policy; the finally RESETS via token (pinning "" would suppress the - # legacy os.environ fallback used by standalone entrypoints/tests). - _cron_session_token = _cron_session_var.set("1") - - # Mark NOT the kanban worker: a worker's cronjob(action="run") lands here with - # HERMES_KANBAN_TASK in env, and an unrelated job could close the worker's task. Must be a - # ContextVar, NOT an os.environ clear (env is shared with the worker heartbeat and - # concurrent jobs); copy_context() carries it into the agent thread. - _non_dispatcher_token = enter_non_dispatcher_owned_context() - if _job_workdir: - logger.info("Job '%s': using task-scoped workdir %s", job_id, _job_workdir) - - # Re-read .env every run; reset the secret-source cache FIRST or a Bitwarden/BSM-backed - # secret is never re-resolved (only the placeholder reloads -> 401s). - from hermes_cli.env_loader import ( - load_hermes_dotenv, - reset_secret_source_cache, - ) - reset_secret_source_cache() - load_hermes_dotenv(hermes_home=_get_hermes_home()) - - delivery_target = _resolve_delivery_target(job) - if delivery_target: - _VAR_MAP["HERMES_CRON_AUTO_DELIVER_PLATFORM"].set(delivery_target["platform"]) - _VAR_MAP["HERMES_CRON_AUTO_DELIVER_CHAT_ID"].set(str(delivery_target["chat_id"])) - _VAR_MAP["HERMES_CRON_AUTO_DELIVER_THREAD_ID"].set( - "" - if delivery_target.get("thread_id") is None - else str(delivery_target["thread_id"]) - ) + scope.enter() + if scope.workdir: + logger.info("Job '%s': using task-scoped workdir %s", job_id, scope.workdir) + _reload_dotenv_and_publish_delivery_target(job) jc = _load_cron_job_config(job, job_id, job_name) _cfg = jc.cfg model = jc.model - - prefill_messages = _load_prefill_messages(_cfg, job_id) - - # resolve_turn_limit() honors none/unlimited (sys.maxsize) and explicit 0 / null. - from hermes_cli.config import resolve_turn_limit as _resolve_turn_limit - _mt = _cfg.get("agent", {}).get("max_turns") - if _mt is None: - _mt = _cfg.get("max_turns") - max_iterations = _resolve_turn_limit(_mt) - - pr = _cfg.get("provider_routing") or {} - - # Runtime backstop (CWE-200/522): fail closed BEFORE resolution on a provider/base_url pair - # that would ship a stored credential off-host; hand-written jobs bypass create-time checks. - _guard_job_credential_exfil(job) - - _blocked = _preflight_or_block(job, job_id, job_name, _cfg) - if _blocked is not None: - return _blocked - - primary_model_for_drift = model - runtime, model, primary_provider_for_drift = _resolve_job_runtime(job, job_id, jc) - - reasoning_config = _resolve_job_reasoning_config( - job, _cfg if isinstance(_cfg, dict) else {}, str(model) - ) - _check_model_drift( - job, job_id, _cfg, runtime, primary_provider_for_drift, primary_model_for_drift, - ) - - fallback_model = get_fallback_chain(_cfg) or None - credential_pool = _load_credential_pool(runtime, job_id) - # MCP servers must be registered before AIAgent is constructed. - _init_cron_mcp_tools(job_id) + setup = _resolve_cron_agent_setup(job, job_id, job_name, jc) + if setup.blocked is not None: + return setup.blocked + model = setup.model # Open state.db only after every early-return gate has passed. _session_db = _open_cron_session_db(job) - - agent = AIAgent( - model=model, - api_key=runtime.get("api_key"), - base_url=runtime.get("base_url"), - provider=runtime.get("provider"), - requested_provider=runtime.get("requested_provider"), - api_mode=runtime.get("api_mode"), - request_overrides=runtime.get("request_overrides"), - acp_command=runtime.get("command"), - acp_args=runtime.get("args"), - max_iterations=max_iterations, - reasoning_config=reasoning_config, - prefill_messages=prefill_messages, - fallback_model=fallback_model, - credential_pool=credential_pool, - providers_allowed=pr.get("only"), - providers_ignored=pr.get("ignore"), - providers_order=pr.get("order"), - provider_sort=pr.get("sort"), - openrouter_min_coding_score=(_cfg.get("openrouter") or {}).get("min_coding_score"), - enabled_toolsets=_resolve_cron_enabled_toolsets(job, _cfg), - disabled_toolsets=_resolve_cron_disabled_toolsets(_cfg), - quiet_mode=True, - # Project context files only with a configured workdir; SOUL.md always. - skip_context_files=not bool(_job_workdir), - load_soul_identity=True, - skip_memory=False, - skip_background_review=True, # Cron has no human-in-the-loop need for skill/memory review forks (~30K tok/event) - platform="cron", - session_id=_cron_session_id, + agent = _construct_cron_agent( + AIAgent, job, _cfg, setup, workdir=scope.workdir, session_id=_cron_session_id, session_db=_session_db, ) - - _audit_fire_id = uuid.uuid4().hex - _audit_t_start = time.monotonic() - - def _audit(result: dict, error: Optional[str]) -> None: - """One usage_audit.jsonl line per fire.""" - _write_usage_audit({ - "ts": _utcnow_iso_ms(), - "job_id": job_id, - "fire_id": _audit_fire_id, - "prompt_tokens": result.get("prompt_tokens"), - "completion_tokens": result.get("completion_tokens"), - "total_tokens": result.get("total_tokens"), - "response_silent": bool(result.get("response_silent")), - "deliver_target": job.get("deliver"), - "model": model or None, - "duration_ms": int((time.monotonic() - _audit_t_start) * 1000), - "error": error, - }) + _audit = _FireAudit(job, job_id, model) result = _run_agent_with_watchdog( - agent, prompt, job, job_id, job_name, _cron_task_id, cancel_event, + agent, prompt, job, job_id, job_name, scope.task_id, cancel_event, ) final_response = _final_response_from_result(result, job_id, job_name, AIAgent) # Keep final_response clean for delivery logic (empty = no delivery). logged_response = final_response if final_response else "(No response generated)" output = _run_doc_header(job, job_name, job_id, prompt) + f"## Response\n\n{logged_response}\n" logger.info("Job '%s' completed successfully", job_name) - _audit(dict(result, response_silent=_is_cron_silence_response(final_response or "")), None) + _audit.write(dict(result, response_silent=_is_cron_silence_response(final_response or "")), None) return True, output, final_response, None except Exception as e: error_msg = f"{type(e).__name__}: {str(e)}" logger.exception("Job '%s' failed: %s", job_name, error_msg) - # _audit is unbound if we failed before the agent ran; the audit write must never raise. - if "_audit" in locals(): - _audit({}, error_msg) + # No audit row when we failed before the agent existed; the audit write must never raise. + if _audit is not None: + _audit.write({}, error_msg) output = ( _run_doc_header(job, f"{job_name} (FAILED)", job_id, prompt) + f"## Error\n\n```\n{error_msg}\n```\n" @@ -2236,15 +2301,7 @@ def run_job( return False, output, "", error_msg finally: - _clear_tool_session_cwd(_cron_task_id) - # clear_session_vars also clears _SESSION_CWD. - 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("") + scope.exit() if _session_db: _finalize_cron_session(_session_db, agent, job_id, job_name, _cron_session_id) # Tear down the ephemeral agent or the gateway leaks fds per tick (EMFILE). With deferred @@ -2507,6 +2564,223 @@ def _compose_run_delivery( ) +class _FireClaimLostDuringSideEffect(Exception): + """Raised inside a side-effect fence when the durable fire claim is no longer ours.""" + + +class _FireOwnership: + """Fire-claim ownership checks for one run (``owner`` is None when the job carries no claim).""" + + def __init__(self, job: dict, fire_claim_lost: Optional[_CancelEventLike]): + self.job = job + self.fire_claim_lost = fire_claim_lost + claim = job.get("fire_claim") + self.owner = str(claim.get("by") or "") if isinstance(claim, dict) else None + + def side_effect_fence(self): + if self.owner is None: + return contextlib.nullcontext(True) + return fire_claim_fence(self.job["id"], expected_owner=self.owner) + + def lost(self) -> bool: + if self.fire_claim_lost is not None and self.fire_claim_lost.is_set(): + return True + if self.owner is None: + return False + try: + if heartbeat_fire_claim(self.job["id"], expected_owner=self.owner): + return False + except Exception: + logger.debug( + "Job '%s': fire_claim ownership validation failed", self.job["id"], exc_info=True, + ) + return False + if self.fire_claim_lost is not None: + self.fire_claim_lost.set() + return True + + +@dataclass +class _RunDelivery: + """Mutable outcome of the save/compose/deliver phase, read back by the bookkeeping tail.""" + job: dict + success: bool + error: Optional[str] + delivery_attempted: bool = False + delivery_error: Optional[str] = None + should_deliver: bool = False + unresolved_origin: bool = False + blocked_config: bool = False + incident_acked: bool = False + failure_incident_id: Optional[str] = None + side_effect_ownership_lost: bool = False + + +def _save_compose_deliver( + d: _RunDelivery, fence: _FireOwnership, final_response: str, output: str, *, + adapters, loop, verbose: bool, execution_token, +) -> None: + """Save output, compose the notice and deliver it (both side effects run under the fire-claim + fence; a lost claim raises ``_FireClaimLostDuringSideEffect`` for the caller).""" + job = d.job + with fence.side_effect_fence() as owns_output: + if not owns_output: + raise _FireClaimLostDuringSideEffect + output_file = save_job_output(job["id"], output) + if verbose: + logger.info("Output saved to: %s", output_file) + + # A shutdown-killed tool subprocess can leave a plausible final_response from truncated + # output; force the honest "interrupted" failure path. Peek-only (consumed later). + if d.success and _is_interrupted(job["id"], execution_token): + d.success = False + d.error = ( + "Interrupted by gateway shutdown before the run finished " + "(tool subprocess was killed mid-flight)." + ) + + ( + deliver_content, d.blocked_config, _silent_alert, d.incident_acked, d.failure_incident_id, + ) = _compose_run_delivery( + job, success=d.success, error=d.error, final_response=final_response, + output_file=output_file, + ) + # Whitespace-only == empty: skip delivery; the guard below marks it a soft failure. + d.should_deliver = bool(deliver_content.strip()) and not _silent_alert + # Not a substring check: bare "SILENT"/"NO_REPLY" or a report quoting "[SILENT]" must + # not be swallowed; bracketed-prefix / trailing-line tolerance is kept. + if d.should_deliver and d.success and _is_cron_silence_response(deliver_content): + logger.info("Job '%s': agent returned %s — skipping delivery", job["id"], SILENT_MARKER) + d.should_deliver = False + + if d.should_deliver and fence.lost(): + d.should_deliver = False + logger.warning("Job '%s': skipping delivery after fire claim ownership loss", job["id"]) + + if not d.should_deliver: + return + d.unresolved_origin = ( + _normalize_deliver_value(_delivery_lane_value(job, for_failure=not d.success)) == "origin" + and not _resolve_delivery_targets(job, for_failure=not d.success) + ) + try: + with fence.side_effect_fence() as owns_delivery: + if not owns_delivery: + raise _FireClaimLostDuringSideEffect + d.delivery_attempted = True + d.delivery_error = _deliver_result( + job, + deliver_content, + adapters=adapters, + loop=loop, + # Failure summaries (and drift/blocked-config alerts composed into deliver_content + # on the failure path) honor the job's failure_deliver override (NS-788). + for_failure=not d.success, + ) + except Exception as de: + if isinstance(de, _FireClaimLostDuringSideEffect): + raise + d.delivery_error = str(de) + logger.error("Delivery failed for job %s: %s", job["id"], de) + + +def _finish_interrupted_run(job: dict, execution_id: str, delivery_error: Optional[str]) -> None: + """Shutdown already wrote last_status, so mark_job_run is skipped (a second call would skip a + fire or auto-delete the job); an unsent notice is recorded via update_job instead.""" + if delivery_error: + try: + from cron.jobs import update_job + update_job(job["id"], {"last_delivery_error": delivery_error}) + except Exception as _rec_err: + logger.debug( + "Failed recording delivery_error for interrupted job %s: %s", job["id"], _rec_err, + ) + finish_execution( + execution_id, + success=False, + error="Interrupted by gateway shutdown before terminal completion.", + ) + + +def _finish_completed_run(d: _RunDelivery, fire_owner: Optional[str], execution_id: str) -> bool: + """mark_job_run (owner-fenced) + execution ledger row for a run that reached delivery.""" + job = d.job + mark_kwargs = {"delivery_error": d.delivery_error} + if fire_owner is not None: + mark_kwargs["expected_fire_owner"] = fire_owner + if d.blocked_config: + mark_kwargs["status"] = "blocked_config" + marked = mark_job_run(job["id"], d.success, d.error, **mark_kwargs) + if fire_owner is not None and not marked: + finish_execution( + execution_id, + success=False, + error="Fire claim ownership lost before terminal completion.", + ) + return True + delivery_outcome = _classify_delivery_outcome( + delivery_error=d.delivery_error, + should_deliver=d.should_deliver, + unresolved_origin=d.unresolved_origin, + # Read the lane the notice was actually routed through (failure_deliver on failure). + normalized_deliver=_normalize_deliver_value(_delivery_lane_value(job, for_failure=not d.success)), + incident_acked=d.incident_acked, + success=d.success, + ) + if delivery_outcome in ("delivered", "not_configured") and not d.success: + # Failure ping left the process (or had a configured target): mark the incident alerted. + _mark_incident_alerted(d.failure_incident_id) + finish_execution( + execution_id, + success=d.success, + error=d.error, + delivery_outcome=delivery_outcome, + ) + return True + + +def _deliver_crash_failure( + job: dict, err_text: str, *, adapters, loop, +) -> tuple[Optional[str], str]: + """Failure notice for a run that raised out of run_job. Returns (delivery_error, outcome).""" + normalized_deliver = _normalize_deliver_value(_delivery_lane_value(job, for_failure=True)) + # Same ack gate as the normal failure delivery: acked signatures stay silent here too. + incident_acked, failure_incident_id = _upsert_incident_for_failure(job, err_text) + if incident_acked: + return None, "suppressed_acked" + delivery_error = None + try: + delivery_error = _deliver_result( + job, + # Same text as the normal failure delivery: this run also counts toward + # failure_streak, so the nudge must leave through here too. + _summarize_cron_failure_for_delivery(job, err_text) + _failure_streak_nudge(job), + adapters=adapters, + loop=loop, + for_failure=True, + ) + except Exception as delivery_exc: + delivery_error = str(delivery_exc) + logger.error("Delivery failed for job %s: %s", job["id"], delivery_exc) + unresolved_origin = bool( + not delivery_error + and normalized_deliver == "origin" + and not _resolve_delivery_targets(job, for_failure=True) + ) + delivery_outcome = _classify_delivery_outcome( + delivery_error=delivery_error, + should_deliver=True, + unresolved_origin=unresolved_origin, + normalized_deliver=normalized_deliver, + incident_acked=False, + success=False, + ) + if delivery_outcome in ("delivered", "not_configured"): + _mark_incident_alerted(failure_incident_id) + return delivery_error, delivery_outcome + + + def _run_one_job_body( job: dict, *, @@ -2517,43 +2791,16 @@ def _run_one_job_body( fire_claim_lost: Optional[_CancelEventLike] = None, execution_token: Optional[object] = None, ) -> bool: - claim = job.get("fire_claim") - fire_owner = str(claim.get("by") or "") if isinstance(claim, dict) else None - - class _FireClaimLostDuringSideEffect(Exception): - pass - - def _side_effect_fence(): - if fire_owner is None: - return contextlib.nullcontext(True) - return fire_claim_fence(job["id"], expected_owner=fire_owner) - - def _fire_claim_ownership_lost() -> bool: - if fire_claim_lost is not None and fire_claim_lost.is_set(): - return True - if fire_owner is None: - return False - try: - if heartbeat_fire_claim(job["id"], expected_owner=fire_owner): - return False - except Exception: - logger.debug( - "Job '%s': fire_claim ownership validation failed", - job["id"], - exc_info=True, - ) - return False - if fire_claim_lost is not None: - fire_claim_lost.set() - return True + fence = _FireOwnership(job, fire_claim_lost) + fire_owner = fence.owner + _side_effect_fence = fence.side_effect_fence + _fire_claim_ownership_lost = fence.lost execution_id = job.get("execution_id") if not execution_id: execution_id = create_execution(job["id"], source="direct")["id"] delivery_attempted = False delivery_error = None - incident_acked = False - failure_incident_id = None from agent.secret_scope import ( build_profile_secret_scope, reset_secret_scope, @@ -2620,141 +2867,33 @@ def _run_one_job_body( # Agent is still live through delivery; wrap ALL of save/compose/deliver in try/finally so a # raise anywhere still tears the deferred agent down. - blocked_config = False - side_effect_ownership_lost = False + d = _RunDelivery(job=job, success=success, error=error) try: - with _side_effect_fence() as owns_output: - if not owns_output: - raise _FireClaimLostDuringSideEffect - output_file = save_job_output(job["id"], output) - if verbose: - logger.info("Output saved to: %s", output_file) - - # A shutdown-killed tool subprocess can leave a plausible final_response from truncated - # output; force the honest "interrupted" failure path. Peek-only (consumed later). - if success and _is_interrupted(job["id"], execution_token): - success = False - error = ( - "Interrupted by gateway shutdown before the run finished " - "(tool subprocess was killed mid-flight)." - ) - - ( - deliver_content, blocked_config, _silent_alert, - incident_acked, failure_incident_id, - ) = _compose_run_delivery( - job, success=success, error=error, final_response=final_response, - output_file=output_file, + _save_compose_deliver( + d, fence, final_response, output, adapters=adapters, loop=loop, verbose=verbose, + execution_token=execution_token, ) - # Whitespace-only == empty: skip delivery; the guard below marks it a soft failure. - should_deliver = bool(deliver_content.strip()) and not _silent_alert - unresolved_origin = False - # Not a substring check: bare "SILENT"/"NO_REPLY" or a report quoting "[SILENT]" must - # not be swallowed; bracketed-prefix / trailing-line tolerance is kept. - if should_deliver and success and _is_cron_silence_response(deliver_content): - logger.info("Job '%s': agent returned %s — skipping delivery", job["id"], SILENT_MARKER) - should_deliver = False - - if should_deliver and _fire_claim_ownership_lost(): - should_deliver = False - logger.warning( - "Job '%s': skipping delivery after fire claim ownership loss", - job["id"], - ) - - if should_deliver: - unresolved_origin = ( - _normalize_deliver_value(_delivery_lane_value(job, for_failure=not success)) - == "origin" - and not _resolve_delivery_targets(job, for_failure=not success) - ) - try: - with _side_effect_fence() as owns_delivery: - if not owns_delivery: - raise _FireClaimLostDuringSideEffect - delivery_attempted = True - delivery_error = _deliver_result( - job, - deliver_content, - adapters=adapters, - loop=loop, - # Failure summaries (and drift/blocked-config alerts - # composed into deliver_content on the failure path) - # honor the job's failure_deliver override (NS-788). - for_failure=not success, - ) - except Exception as de: - if isinstance(de, _FireClaimLostDuringSideEffect): - raise - delivery_error = str(de) - logger.error("Delivery failed for job %s: %s", job["id"], de) except _FireClaimLostDuringSideEffect: - side_effect_ownership_lost = True + d.side_effect_ownership_lost = True finally: + delivery_attempted, delivery_error = d.delivery_attempted, d.delivery_error # Every path must tear down deferred agent(s) so they never leak subprocesses/clients. _teardown_deferred() - if side_effect_ownership_lost or _fire_claim_ownership_lost(): + if d.side_effect_ownership_lost or _fire_claim_ownership_lost(): _record_fire_ownership_lost(job["id"], fire_owner, execution_id) return True # Empty final_response is a soft failure so last_status is not "ok". - if success and not final_response.strip(): - success = False - error = "Agent completed but produced empty response (model error, timeout, or misconfiguration)" + if d.success and not final_response.strip(): + d.success = False + d.error = "Agent completed but produced empty response (model error, timeout, or misconfiguration)" - interrupted = _consume_interrupted_flag(job["id"], execution_token) - if interrupted: - if delivery_error: - # Shutdown already wrote last_status so mark_job_run is skipped below (a second call - # would skip a fire or auto-delete the job); note the unsent notice via update_job. - try: - from cron.jobs import update_job - update_job(job["id"], {"last_delivery_error": delivery_error}) - except Exception as _rec_err: - logger.debug( - "Failed recording delivery_error for interrupted job %s: %s", - job["id"], _rec_err, - ) - finish_execution( - execution_id, - success=False, - error="Interrupted by gateway shutdown before terminal completion.", - ) + if _consume_interrupted_flag(job["id"], execution_token): + _finish_interrupted_run(job, execution_id, delivery_error) return True - mark_kwargs = {"delivery_error": delivery_error} - if fire_owner is not None: - mark_kwargs["expected_fire_owner"] = fire_owner - if blocked_config: - mark_kwargs["status"] = "blocked_config" - marked = mark_job_run(job["id"], success, error, **mark_kwargs) - if fire_owner is not None and not marked: - finish_execution( - execution_id, - success=False, - error="Fire claim ownership lost before terminal completion.", - ) - return True - delivery_outcome = _classify_delivery_outcome( - delivery_error=delivery_error, - should_deliver=should_deliver, - unresolved_origin=unresolved_origin, - # Read the lane the notice was actually routed through (failure_deliver on failure). - normalized_deliver=_normalize_deliver_value(_delivery_lane_value(job, for_failure=not success)), - incident_acked=incident_acked, - success=success, - ) - if delivery_outcome in ("delivered", "not_configured") and not success: - # Failure ping left the process (or had a configured target): mark the incident alerted. - _mark_incident_alerted(failure_incident_id) - finish_execution( - execution_id, - success=success, - error=error, - delivery_outcome=delivery_outcome, - ) - return True + return _finish_completed_run(d, fire_owner, execution_id) except BaseException as e: # noqa: BLE001 — deliberate: see below # BaseException, not Exception: CancelledError/KeyboardInterrupt/SystemExit propagate here. @@ -2777,42 +2916,9 @@ def _run_one_job_body( and not isinstance(e, _FireClaimLostDuringSideEffect) and not _fire_claim_ownership_lost() ): - normalized_deliver = _normalize_deliver_value(_delivery_lane_value(job, for_failure=True)) - # Same ack gate as the normal failure delivery: acked signatures stay silent here too. - incident_acked, failure_incident_id = _upsert_incident_for_failure(job, _err_text) - if incident_acked: - delivery_outcome = "suppressed_acked" - else: - try: - delivery_attempted = True - delivery_error = _deliver_result( - job, - # Same text as the normal failure delivery: this run also counts toward - # failure_streak, so the nudge must leave through here too. - _summarize_cron_failure_for_delivery(job, _err_text) - + _failure_streak_nudge(job), - adapters=adapters, - loop=loop, - for_failure=True, - ) - except Exception as delivery_exc: - delivery_error = str(delivery_exc) - logger.error("Delivery failed for job %s: %s", job["id"], delivery_exc) - unresolved_origin = bool( - not delivery_error - and normalized_deliver == "origin" - and not _resolve_delivery_targets(job, for_failure=True) - ) - delivery_outcome = _classify_delivery_outcome( - delivery_error=delivery_error, - should_deliver=True, - unresolved_origin=unresolved_origin, - normalized_deliver=normalized_deliver, - incident_acked=False, - success=False, - ) - if delivery_outcome in ("delivered", "not_configured"): - _mark_incident_alerted(failure_incident_id) + delivery_error, delivery_outcome = _deliver_crash_failure( + job, _err_text, adapters=adapters, loop=loop, + ) try: if not _consume_interrupted_flag(job["id"], execution_token): mark_kwargs = {} diff --git a/tests/gateway/test_cron_interrupt_notification.py b/tests/gateway/test_cron_interrupt_notification.py index f157e17478..ad65f56597 100644 --- a/tests/gateway/test_cron_interrupt_notification.py +++ b/tests/gateway/test_cron_interrupt_notification.py @@ -251,10 +251,12 @@ class TestDeliveryErrorIsRecordedWhenTheNoticeCannotBeSent: import cron.scheduler as sched - src = inspect.getsource(sched._run_one_job_body) + body_src = inspect.getsource(sched._run_one_job_body) + src = inspect.getsource(sched._finish_interrupted_run) assert 'update_job(job["id"], {"last_delivery_error": delivery_error})' in src, ( "interrupted runs must still persist the delivery failure" ) # The recovery branch hangs off the interrupted-flag short-circuit, # not off a second mark_job_run call. - assert "if interrupted:" in src and "if delivery_error:" in src + assert "_consume_interrupted_flag(" in body_src and "_finish_interrupted_run(" in body_src + assert "if delivery_error:" in src and "mark_job_run(" not in src