refactor(cron): pack hanging bracket spans in scheduler.py (AST-identical)
This commit is contained in:
+98
-261
@@ -39,19 +39,13 @@ sys.path.insert(0, str(Path(__file__).parent.parent))
|
||||
from hermes_constants import get_hermes_home
|
||||
from hermes_cli._subprocess_compat import windows_hide_flags # noqa: F401 (tests patch cron.scheduler.windows_hide_flags)
|
||||
from hermes_cli.config import (
|
||||
_expand_env_vars,
|
||||
cron_model_drift_axes,
|
||||
cron_model_drift_guard_enabled,
|
||||
load_config,
|
||||
resolve_cron_model_drift_defaults,
|
||||
)
|
||||
_expand_env_vars, cron_model_drift_axes, cron_model_drift_guard_enabled, load_config,
|
||||
resolve_cron_model_drift_defaults)
|
||||
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,
|
||||
)
|
||||
enter_non_dispatcher_owned_context, exit_non_dispatcher_owned_context)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -213,8 +207,7 @@ def _log_tick_yield_once(reason: str) -> None:
|
||||
"Cron tick yielded: this process is running stale code (%s) and a "
|
||||
"fresher gateway owns the runtime lock — jobs will fire from that "
|
||||
"process. Restart this one to reclaim its ticks.",
|
||||
reason,
|
||||
)
|
||||
reason)
|
||||
_last_yield_log = {"reason": reason, "at": now}
|
||||
|
||||
|
||||
@@ -341,19 +334,14 @@ def _upsert_incident_for_failure(
|
||||
from cron.incidents import get_incident, upsert_incident
|
||||
|
||||
incident_id, _is_new = upsert_incident(
|
||||
job["id"],
|
||||
str(error or ""),
|
||||
job_name=job.get("name"),
|
||||
output_file=output_file,
|
||||
)
|
||||
job["id"], str(error or ""), job_name=job.get("name"), output_file=output_file)
|
||||
incident = get_incident(incident_id)
|
||||
acked = bool(incident and incident.get("state") == "closed")
|
||||
return acked, incident_id
|
||||
except Exception as exc:
|
||||
logger.debug(
|
||||
"Incident store unavailable for job %s (delivery unaffected): %s",
|
||||
job["id"], exc,
|
||||
)
|
||||
job["id"], exc)
|
||||
return False, None
|
||||
|
||||
|
||||
@@ -437,8 +425,7 @@ def _resolve_cron_enabled_toolsets(job: dict, cfg: dict) -> list[str] | None:
|
||||
except Exception as exc:
|
||||
logger.warning(
|
||||
"Cron toolset resolution failed, falling back to full default toolset: %s",
|
||||
exc,
|
||||
)
|
||||
exc)
|
||||
return None
|
||||
|
||||
|
||||
@@ -464,25 +451,14 @@ def _resolve_job_reasoning_config(job: dict, cfg: dict, model: str) -> dict | No
|
||||
"minimal, low, medium, high, xhigh, max, ultra).",
|
||||
job.get("id", "?"),
|
||||
pinned,
|
||||
job.get("id", "?"),
|
||||
)
|
||||
job.get("id", "?"))
|
||||
return resolve_reasoning_config(cfg if isinstance(cfg, dict) else {}, str(model))
|
||||
|
||||
|
||||
from cron.jobs import (
|
||||
_ensure_cron_dir,
|
||||
advance_next_runs,
|
||||
claim_dispatch,
|
||||
claim_job_for_fire,
|
||||
fire_claim_fence,
|
||||
clear_run_claim,
|
||||
get_due_jobs,
|
||||
heartbeat_fire_claim,
|
||||
heartbeat_run_claim,
|
||||
mark_job_run,
|
||||
save_job_output,
|
||||
use_cron_store,
|
||||
)
|
||||
_ensure_cron_dir, advance_next_runs, claim_dispatch, claim_job_for_fire, fire_claim_fence,
|
||||
clear_run_claim, get_due_jobs, heartbeat_fire_claim, heartbeat_run_claim, mark_job_run,
|
||||
save_job_output, use_cron_store)
|
||||
from cron.executions import create_execution, finish_execution, mark_execution_running
|
||||
|
||||
# Response marker that suppresses delivery (output is still saved locally for audit).
|
||||
@@ -609,8 +585,7 @@ def _inflight_min_allowance_minutes() -> float:
|
||||
logger.warning(
|
||||
"Invalid HERMES_CRON_INFLIGHT_MAX_MINUTES=%r; using default %s",
|
||||
raw,
|
||||
_INFLIGHT_MIN_ALLOWANCE_MINUTES,
|
||||
)
|
||||
_INFLIGHT_MIN_ALLOWANCE_MINUTES)
|
||||
return _INFLIGHT_MIN_ALLOWANCE_MINUTES
|
||||
|
||||
|
||||
@@ -670,8 +645,7 @@ def get_inflight_guard_stats() -> dict:
|
||||
for jid, started in _running_since.items()
|
||||
},
|
||||
"forced_releases": _forced_release_count,
|
||||
"recent_forced_releases": list(_forced_releases),
|
||||
}
|
||||
"recent_forced_releases": list(_forced_releases)}
|
||||
|
||||
|
||||
def _record_forced_release(job_id: str, name: str, age_seconds: float, allowance_seconds: float) -> None:
|
||||
@@ -681,8 +655,7 @@ def _record_forced_release(job_id: str, name: str, age_seconds: float, allowance
|
||||
"name": name,
|
||||
"age_seconds": round(age_seconds, 1),
|
||||
"allowance_seconds": round(allowance_seconds, 1),
|
||||
"at": _hermes_now().isoformat(),
|
||||
}
|
||||
"at": _hermes_now().isoformat()}
|
||||
with _running_lock:
|
||||
_forced_releases.append(entry)
|
||||
del _forced_releases[:-_FORCED_RELEASE_HISTORY]
|
||||
@@ -806,8 +779,7 @@ def sweep_stale_inflight(due_jobs: Optional[list] = None) -> list:
|
||||
job_id,
|
||||
age,
|
||||
allowance,
|
||||
future_state,
|
||||
)
|
||||
future_state)
|
||||
_record_forced_release(job_id, name, age, allowance)
|
||||
# Ledger already records how the run ended: mark_job_run here would clobber an honest
|
||||
# ok status with a synthetic failure or double-write a failure.
|
||||
@@ -822,8 +794,7 @@ def sweep_stale_inflight(due_jobs: Optional[list] = None) -> list:
|
||||
"finite-repeat job released without mark_job_run (repeat budget "
|
||||
"preserved); row left in place so it re-fires normally",
|
||||
name,
|
||||
job_id,
|
||||
)
|
||||
job_id)
|
||||
continue
|
||||
try:
|
||||
mark_job_run(
|
||||
@@ -831,8 +802,7 @@ def sweep_stale_inflight(due_jobs: Optional[list] = None) -> list:
|
||||
False,
|
||||
f"Stale in-flight claim force-released after {age / 60:.1f}m "
|
||||
f"(allowance {allowance / 60:.1f}m); previous run never released "
|
||||
f"the scheduler in-flight guard",
|
||||
)
|
||||
f"the scheduler in-flight guard")
|
||||
except Exception as e:
|
||||
logger.warning("Could not record forced release for job %s: %s", job_id, e)
|
||||
|
||||
@@ -840,9 +810,7 @@ def sweep_stale_inflight(due_jobs: Optional[list] = None) -> list:
|
||||
|
||||
|
||||
def mark_running_jobs_interrupted(
|
||||
reason: str,
|
||||
*,
|
||||
only_owners: Optional[set] = None,
|
||||
reason: str, *, only_owners: Optional[set] = None,
|
||||
) -> list:
|
||||
"""Best-effort: mark every in-flight cron job interrupted; returns the job IDs marked.
|
||||
|
||||
@@ -875,8 +843,7 @@ def mark_running_jobs_interrupted(
|
||||
logger.warning(
|
||||
"Job '%s' interrupted before its durable fire owner was registered; "
|
||||
"leaving persisted state untouched",
|
||||
job_id,
|
||||
)
|
||||
job_id)
|
||||
# Still report it: shutdown uses the returned IDs for the interrupted-cron notice. The
|
||||
# in-memory flag WAS recorded above; only the persisted last_status write is skipped.
|
||||
marked.append(job_id)
|
||||
@@ -884,11 +851,7 @@ def mark_running_jobs_interrupted(
|
||||
try:
|
||||
with use_cron_store(profile_home):
|
||||
if mark_job_run(
|
||||
job_id,
|
||||
False,
|
||||
reason,
|
||||
expected_fire_owner=fire_owner,
|
||||
):
|
||||
job_id, False, reason, expected_fire_owner=fire_owner):
|
||||
marked.append(job_id)
|
||||
except Exception as e:
|
||||
logger.warning("Failed to mark job %s interrupted: %s", job_id, e)
|
||||
@@ -920,11 +883,7 @@ def _consume_interrupted_flag(job_id: str, token: Optional[object] = None) -> bo
|
||||
|
||||
|
||||
def _inactivity_watchdog_loop(
|
||||
*,
|
||||
get_idle_seconds: Callable[[], float],
|
||||
limit_s: float,
|
||||
poll_s: float,
|
||||
stop: threading.Event,
|
||||
*, get_idle_seconds: Callable[[], float], limit_s: float, poll_s: float, stop: threading.Event,
|
||||
future_done: Callable[[], bool],
|
||||
) -> bool:
|
||||
"""Poll idle time until limit (-> True), stop, or the future completes (-> False). Uses
|
||||
@@ -964,9 +923,7 @@ def _get_parallel_pool(max_workers: Optional[int]) -> concurrent.futures.ThreadP
|
||||
if _parallel_pool is not None:
|
||||
_parallel_pool.shutdown(wait=False, cancel_futures=False)
|
||||
_parallel_pool = concurrent.futures.ThreadPoolExecutor(
|
||||
max_workers=max_workers,
|
||||
thread_name_prefix="cron-parallel",
|
||||
)
|
||||
max_workers=max_workers, thread_name_prefix="cron-parallel")
|
||||
_parallel_pool_max_workers = max_workers
|
||||
return _parallel_pool
|
||||
|
||||
@@ -1107,11 +1064,7 @@ def _cron_cleanup_timeout_seconds() -> float:
|
||||
|
||||
|
||||
def _run_cron_cleanup_with_timeout(
|
||||
cleanup,
|
||||
*,
|
||||
job_id: str,
|
||||
label: str,
|
||||
timeout_seconds: Optional[float] = None,
|
||||
cleanup, *, job_id: str, label: str, timeout_seconds: Optional[float] = None,
|
||||
) -> bool:
|
||||
"""Run fallible post-run cleanup without permanently wedging a cron ID."""
|
||||
timeout = (_cron_cleanup_timeout_seconds() if timeout_seconds is None else float(timeout_seconds))
|
||||
@@ -1137,18 +1090,14 @@ def _run_cron_cleanup_with_timeout(
|
||||
# Daemon thread is deliberate: unlike ThreadPoolExecutor workers it is not joined at interpreter
|
||||
# exit if cleanup never returns, so the gateway can still shut down.
|
||||
worker = threading.Thread(
|
||||
target=_runner,
|
||||
name=f"cron-cleanup-{job_id}",
|
||||
daemon=True,
|
||||
)
|
||||
target=_runner, name=f"cron-cleanup-{job_id}", daemon=True)
|
||||
worker.start()
|
||||
if not done.wait(timeout):
|
||||
logger.error(
|
||||
"Job '%s': %s exceeded %.1fs; abandoning cleanup so future runs remain dispatchable",
|
||||
job_id,
|
||||
label,
|
||||
timeout,
|
||||
)
|
||||
timeout)
|
||||
return False
|
||||
if error:
|
||||
logger.debug("Job '%s': %s failed: %s", job_id, label, error[0])
|
||||
@@ -1184,10 +1133,7 @@ class _BoundedCronSessionDB:
|
||||
raise
|
||||
|
||||
ok = _run_cron_cleanup_with_timeout(
|
||||
_call,
|
||||
job_id=self._job_id,
|
||||
label=f"session finalization ({name})",
|
||||
)
|
||||
_call, job_id=self._job_id, label=f"session finalization ({name})")
|
||||
if not ok:
|
||||
error = result.get("error")
|
||||
if error is not None:
|
||||
@@ -1216,8 +1162,7 @@ def _resolve_job_workdir(job: dict, job_id: str) -> Optional[str]:
|
||||
if workdir and not Path(workdir).is_dir():
|
||||
logger.warning(
|
||||
"Job '%s': configured workdir %r no longer exists — running without it",
|
||||
job_id, workdir,
|
||||
)
|
||||
job_id, workdir)
|
||||
return None
|
||||
return workdir
|
||||
|
||||
@@ -1248,8 +1193,7 @@ def _run_no_agent_job(
|
||||
_job_workdir = _resolve_job_workdir(job, job_id)
|
||||
try:
|
||||
ok, output = _run_job_script_with_claim_heartbeat(
|
||||
job, script_path, workdir=_job_workdir, cancel_event=cancel_event,
|
||||
)
|
||||
job, script_path, workdir=_job_workdir, cancel_event=cancel_event)
|
||||
except Exception as exc:
|
||||
logger.exception("Job '%s': script execution raised unexpectedly", job_id)
|
||||
ok, output = False, f"Script execution failed: {exc}"
|
||||
@@ -1436,8 +1380,7 @@ def _preflight_or_block(job: dict, job_id: str, job_name: str, cfg: dict) -> Opt
|
||||
logger.warning(
|
||||
"Job '%s' (ID: %s): BLOCKED by pre-dispatch config "
|
||||
"validation — %s (no LLM call was made)",
|
||||
job_name, job_id, _pf_reason,
|
||||
)
|
||||
job_name, job_id, _pf_reason)
|
||||
already_alerted = False
|
||||
try:
|
||||
from cron.jobs import mark_preflight_alerted
|
||||
@@ -1470,9 +1413,7 @@ def _resolve_job_runtime(
|
||||
swap only the provider while keeping a paid primary model).
|
||||
"""
|
||||
from hermes_cli.runtime_provider import (
|
||||
resolve_runtime_provider,
|
||||
format_runtime_provider_error,
|
||||
)
|
||||
resolve_runtime_provider, format_runtime_provider_error)
|
||||
from hermes_cli.auth import AuthError
|
||||
|
||||
model = jc.model
|
||||
@@ -1515,8 +1456,7 @@ def _resolve_job_runtime(
|
||||
)
|
||||
logger.warning(
|
||||
"Job '%s': primary provider resolve failed (%s: %s), trying fallback",
|
||||
job_id, "auth" if is_auth else "transient network", resolve_exc,
|
||||
)
|
||||
job_id, "auth" if is_auth else "transient network", resolve_exc)
|
||||
for entry in get_fallback_chain(jc.cfg):
|
||||
if not isinstance(entry, dict):
|
||||
continue
|
||||
@@ -1536,8 +1476,7 @@ def _resolve_job_runtime(
|
||||
runtime = resolve_runtime_provider(**fb_kwargs)
|
||||
logger.info(
|
||||
"Job '%s': fallback resolved to %s model %s",
|
||||
job_id, runtime.get("provider"), fb_model,
|
||||
)
|
||||
job_id, runtime.get("provider"), fb_model)
|
||||
return runtime, fb_model, primary_provider_for_drift
|
||||
except Exception as fb_exc:
|
||||
logger.debug("Job '%s': fallback %s failed: %s", job_id, fb_provider, fb_exc)
|
||||
@@ -1563,8 +1502,7 @@ def _check_model_drift(
|
||||
_current_model = str(primary_model_for_drift or "").strip().lower()
|
||||
_drift: list[str] = []
|
||||
for _axis in cron_model_drift_axes(
|
||||
job, current_provider=_current_provider, current_model=_current_model, config=cfg,
|
||||
):
|
||||
job, current_provider=_current_provider, current_model=_current_model, config=cfg):
|
||||
_snapshot = str(job.get(f"{_axis}_snapshot") or "").strip().lower()
|
||||
_current = _current_provider if _axis == "provider" else _current_model
|
||||
_drift.append(f"{_axis} '{_snapshot}' -> '{_current}'")
|
||||
@@ -1596,8 +1534,7 @@ def _check_model_drift(
|
||||
"Job '%s': SKIPPED — global inference config drifted since "
|
||||
"creation (%s) and this job is unpinned. Skipped to prevent "
|
||||
"unintended spend. %s",
|
||||
job_id, _changes, _remediation,
|
||||
)
|
||||
job_id, _changes, _remediation)
|
||||
# Alert-once via drift_alerted bit (silent marker suppresses delivery); a successful run
|
||||
# clears it and re-arms the alert.
|
||||
_drift_already_alerted = False
|
||||
@@ -1626,8 +1563,7 @@ def _load_credential_pool(runtime: dict, job_id: str):
|
||||
if pool.has_credentials():
|
||||
logger.info(
|
||||
"Job '%s': loaded credential pool for provider %s with %d entries",
|
||||
job_id, runtime_provider, len(pool.entries()),
|
||||
)
|
||||
job_id, runtime_provider, len(pool.entries()))
|
||||
return pool
|
||||
except Exception as e:
|
||||
logger.debug("Job '%s': failed to load credential pool for %s: %s", job_id, runtime_provider, e)
|
||||
@@ -1675,8 +1611,7 @@ def _open_cron_session_db(job: dict):
|
||||
"Job '%s': SessionDB init did not return within %.0fs — proceeding "
|
||||
"without a session store for this run instead of blocking it "
|
||||
"forever",
|
||||
job.get("id", "?"), _session_db_timeout,
|
||||
)
|
||||
job.get("id", "?"), _session_db_timeout)
|
||||
except Exception as e:
|
||||
logger.debug("Job '%s': SQLite session store not available: %s", job.get("id", "?"), e)
|
||||
return None
|
||||
@@ -1725,8 +1660,7 @@ def _run_agent_with_watchdog(
|
||||
# Carry scheduler-scoped ContextVar state (e.g. env passthrough) into the worker thread.
|
||||
_cron_context = contextvars.copy_context()
|
||||
_cron_future = _cron_pool.submit(
|
||||
_cron_context.run, agent.run_conversation, prompt, task_id=task_id,
|
||||
)
|
||||
_cron_context.run, agent.run_conversation, prompt, task_id=task_id)
|
||||
_inactivity_timeout = False
|
||||
_watch_stop = threading.Event()
|
||||
|
||||
@@ -1744,19 +1678,12 @@ def _run_agent_with_watchdog(
|
||||
if _cron_inactivity_limit is None:
|
||||
return
|
||||
if _inactivity_watchdog_loop(
|
||||
get_idle_seconds=_idle_seconds,
|
||||
limit_s=_cron_inactivity_limit,
|
||||
poll_s=_POLL_INTERVAL,
|
||||
stop=_watch_stop,
|
||||
future_done=_cron_future.done,
|
||||
):
|
||||
get_idle_seconds=_idle_seconds, limit_s=_cron_inactivity_limit, poll_s=_POLL_INTERVAL,
|
||||
stop=_watch_stop, future_done=_cron_future.done):
|
||||
_inactivity_timeout = True
|
||||
|
||||
_watch_thread = threading.Thread(
|
||||
target=_watch_inactivity,
|
||||
name=f"cron-inactivity-{str(job_id)[:8]}",
|
||||
daemon=True,
|
||||
)
|
||||
target=_watch_inactivity, name=f"cron-inactivity-{str(job_id)[:8]}", daemon=True)
|
||||
try:
|
||||
if _cron_inactivity_limit is not None:
|
||||
# Separate daemon thread so a hung get_activity_summary can't stop the limit firing.
|
||||
@@ -1794,8 +1721,7 @@ def _run_agent_with_watchdog(
|
||||
"| last_activity=%s | iteration=%s/%s | tool=%s",
|
||||
job_name, _secs_ago, _cron_inactivity_limit,
|
||||
_last_desc, _activity.get("api_call_count", 0), _activity.get("max_iterations", 0),
|
||||
_activity.get("current_tool") or "none",
|
||||
)
|
||||
_activity.get("current_tool") or "none")
|
||||
request_hard_interrupt(agent, "Cron job timed out (inactivity)")
|
||||
raise TimeoutError(
|
||||
f"Cron job '{job_name}' idle for "
|
||||
@@ -1830,8 +1756,7 @@ def _final_response_from_result(result: dict, job_id: str, job_name: str, AIAgen
|
||||
logger.warning(
|
||||
"Job '%s' reached the iteration limit but produced a final fallback response; "
|
||||
"delivering the response instead of failing the cron run",
|
||||
job_name,
|
||||
)
|
||||
job_name)
|
||||
|
||||
final_response = result.get("final_response", "") or ""
|
||||
# Repair model-mangled computer_use media paths before delivery (fail-open, as in gateway).
|
||||
@@ -1839,8 +1764,7 @@ def _final_response_from_result(result: dict, job_id: str, job_name: str, AIAgen
|
||||
from gateway.media_repair import repair_explicit_computer_use_media_paths
|
||||
|
||||
final_response = repair_explicit_computer_use_media_paths(
|
||||
final_response, result.get("messages", []),
|
||||
)
|
||||
final_response, result.get("messages", []))
|
||||
if final_response.strip() == "(No response generated)":
|
||||
final_response = ""
|
||||
# The "⚠️ No reply" turn-completion explainer would be delivered as a cron warning; detect it
|
||||
@@ -1867,8 +1791,7 @@ def _final_response_from_result(result: dict, job_id: str, job_name: str, AIAgen
|
||||
if final_response.strip() in _explainer_variants:
|
||||
logger.info(
|
||||
"Job '%s': abnormal empty turn (%s) — suppressing explainer for cron delivery",
|
||||
job_id, turn_exit_reason,
|
||||
)
|
||||
job_id, turn_exit_reason)
|
||||
final_response = ""
|
||||
return final_response
|
||||
|
||||
@@ -1902,8 +1825,7 @@ def _finalize_cron_session(session_db, agent, job_id: str, job_name: str, cron_s
|
||||
# Never leave the session untitled.
|
||||
for _fallback in (
|
||||
getattr(_session_db, "get_next_title_in_lineage", lambda b: b)(f"cron {job_id}"),
|
||||
f"cron {job_id} {_final_cron_session_id[-6:]}",
|
||||
):
|
||||
f"cron {job_id} {_final_cron_session_id[-6:]}"):
|
||||
try:
|
||||
if _set_cron_session_title(_session_db, _final_cron_session_id, _fallback):
|
||||
break
|
||||
@@ -1921,8 +1843,7 @@ def _finalize_cron_session(session_db, agent, job_id: str, job_name: str, cron_s
|
||||
logger.warning(
|
||||
"Job '%s': session ended without a final assistant "
|
||||
"message (lifecycle=%s) — booking run as %s",
|
||||
job_id, _lifecycle, _end_reason,
|
||||
)
|
||||
job_id, _lifecycle, _end_reason)
|
||||
except (Exception, KeyboardInterrupt) as e:
|
||||
logger.debug("Job '%s': session lifecycle classification failed: %s", job_id, e)
|
||||
try:
|
||||
@@ -2028,8 +1949,7 @@ def _prepare_job_prompt(
|
||||
_CRON_DELIVERY_VARS = (
|
||||
"HERMES_CRON_AUTO_DELIVER_PLATFORM",
|
||||
"HERMES_CRON_AUTO_DELIVER_CHAT_ID",
|
||||
"HERMES_CRON_AUTO_DELIVER_THREAD_ID",
|
||||
)
|
||||
"HERMES_CRON_AUTO_DELIVER_THREAD_ID")
|
||||
|
||||
|
||||
class _CronRunScope:
|
||||
@@ -2153,8 +2073,7 @@ def _resolve_cron_agent_setup(job: dict, job_id: str, job_name: str, jc) -> _Cro
|
||||
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,
|
||||
)
|
||||
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.
|
||||
@@ -2219,18 +2138,13 @@ class _FireAudit:
|
||||
"deliver_target": self.job.get("deliver"),
|
||||
"model": self.model or None,
|
||||
"duration_ms": int((time.monotonic() - self.t_start) * 1000),
|
||||
"error": error,
|
||||
})
|
||||
"error": error})
|
||||
|
||||
|
||||
|
||||
def run_job(
|
||||
job: dict,
|
||||
*,
|
||||
defer_agent_teardown: Optional[list] = None,
|
||||
extra_prompt: Optional[str] = None,
|
||||
cancel_event: Optional[_CancelEventLike] = None,
|
||||
execution_id: Optional[str] = None,
|
||||
job: dict, *, defer_agent_teardown: Optional[list] = None, extra_prompt: Optional[str] = None,
|
||||
cancel_event: Optional[_CancelEventLike] = None, execution_id: Optional[str] = None,
|
||||
) -> tuple[bool, str, str, Optional[str]]:
|
||||
"""Execute a single cron job. Returns (success, full_output_doc, final_response, error).
|
||||
|
||||
@@ -2273,13 +2187,11 @@ def run_job(
|
||||
_session_db = _open_cron_session_db(job)
|
||||
agent = _construct_cron_agent(
|
||||
AIAgent, job, _cfg, setup, workdir=scope.workdir, session_id=_cron_session_id,
|
||||
session_db=_session_db,
|
||||
)
|
||||
session_db=_session_db)
|
||||
_audit = _FireAudit(job, job_id, model)
|
||||
|
||||
result = _run_agent_with_watchdog(
|
||||
agent, prompt, job, job_id, job_name, scope.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)"
|
||||
@@ -2335,11 +2247,8 @@ def _teardown_cron_agent(
|
||||
logger.debug("Job '%s': failed to reap stale auxiliary clients: %s", job_id, e)
|
||||
|
||||
_run_cron_cleanup_with_timeout(
|
||||
_cleanup_agent,
|
||||
job_id=job_id,
|
||||
label="agent resource teardown",
|
||||
timeout_seconds=timeout_seconds,
|
||||
)
|
||||
_cleanup_agent, job_id=job_id, label="agent resource teardown",
|
||||
timeout_seconds=timeout_seconds)
|
||||
|
||||
|
||||
def _run_with_fire_claim_heartbeat(job: dict, run) -> bool:
|
||||
@@ -2363,8 +2272,7 @@ def _run_with_fire_claim_heartbeat(job: dict, run) -> bool:
|
||||
logger.warning(
|
||||
"Job '%s': failed to close unstarted execution ledger row",
|
||||
job_id,
|
||||
exc_info=True,
|
||||
)
|
||||
exc_info=True)
|
||||
|
||||
try:
|
||||
owns_fire_claim = heartbeat_fire_claim(job_id, expected_owner=owner)
|
||||
@@ -2386,8 +2294,7 @@ def _run_with_fire_claim_heartbeat(job: dict, run) -> bool:
|
||||
lost_ownership.set()
|
||||
logger.warning(
|
||||
"Job '%s': fire claim ownership lost; interrupting stale run",
|
||||
job_id,
|
||||
)
|
||||
job_id)
|
||||
return
|
||||
last_confirmed = time.monotonic()
|
||||
except Exception:
|
||||
@@ -2401,16 +2308,13 @@ def _run_with_fire_claim_heartbeat(job: dict, run) -> bool:
|
||||
"Job '%s': fire_claim could not be renewed within %.1fs; "
|
||||
"interrupting uncertain run",
|
||||
job_id,
|
||||
_FIRE_CLAIM_HEARTBEAT_GRACE_SECONDS,
|
||||
)
|
||||
_FIRE_CLAIM_HEARTBEAT_GRACE_SECONDS)
|
||||
return
|
||||
|
||||
heartbeat_thread = _start_heartbeat_thread(
|
||||
_heartbeat_loop, "cron-fire-claim-heartbeat",
|
||||
lambda: logger.warning(
|
||||
"Job '%s': could not start fire_claim heartbeat", job_id, exc_info=True,
|
||||
),
|
||||
)
|
||||
"Job '%s': could not start fire_claim heartbeat", job_id, exc_info=True))
|
||||
if heartbeat_thread is None:
|
||||
_finish_unstarted("Fire claim heartbeat could not be started; execution was not run.")
|
||||
return True
|
||||
@@ -2423,13 +2327,8 @@ def _run_with_fire_claim_heartbeat(job: dict, run) -> bool:
|
||||
|
||||
|
||||
def run_one_job(
|
||||
job: dict,
|
||||
*,
|
||||
adapters=None,
|
||||
loop=None,
|
||||
verbose: bool = False,
|
||||
extra_prompt: Optional[str] = None,
|
||||
cancel_event: Optional[_CancelEventLike] = None,
|
||||
job: dict, *, adapters=None, loop=None, verbose: bool = False,
|
||||
extra_prompt: Optional[str] = None, cancel_event: Optional[_CancelEventLike] = None,
|
||||
) -> bool:
|
||||
"""Run ONE due job end-to-end: execute → save output → deliver → mark.
|
||||
|
||||
@@ -2450,9 +2349,7 @@ def run_one_job(
|
||||
profile_home = _get_hermes_home().resolve()
|
||||
with _running_lock:
|
||||
_running_fire_owners.setdefault(job["id"], {})[execution_token] = (
|
||||
fire_owner or None,
|
||||
profile_home,
|
||||
)
|
||||
fire_owner or None, profile_home)
|
||||
try:
|
||||
return _run_with_fire_claim_heartbeat(
|
||||
job,
|
||||
@@ -2467,9 +2364,7 @@ def run_one_job(
|
||||
if cancel_event is not None
|
||||
else lost_ownership
|
||||
),
|
||||
execution_token=execution_token,
|
||||
),
|
||||
)
|
||||
execution_token=execution_token))
|
||||
finally:
|
||||
with _running_lock:
|
||||
executions = _running_fire_owners.get(job["id"])
|
||||
@@ -2491,10 +2386,8 @@ def _record_fire_ownership_lost(job_id: str, fire_owner: Optional[str], executio
|
||||
finish_execution(execution_id, success=False, error=_OWNERSHIP_LOST_INTERRUPTED)
|
||||
else:
|
||||
finish_execution(
|
||||
execution_id,
|
||||
success=False,
|
||||
error="Fire claim ownership lost; stale result was discarded.",
|
||||
)
|
||||
execution_id, success=False,
|
||||
error="Fire claim ownership lost; stale result was discarded.")
|
||||
|
||||
|
||||
def _classify_delivery_outcome(
|
||||
@@ -2560,8 +2453,7 @@ def _compose_run_delivery(
|
||||
deliver_content = f"⚠️ Cron '{job.get('name') or job['id']}' skipped: {_drift_text}"
|
||||
return (
|
||||
deliver_content, blocked_config, blocked_config_silent or drift_skip_silent,
|
||||
incident_acked, failure_incident_id,
|
||||
)
|
||||
incident_acked, failure_incident_id)
|
||||
|
||||
|
||||
class _FireClaimLostDuringSideEffect(Exception):
|
||||
@@ -2592,8 +2484,7 @@ class _FireOwnership:
|
||||
return False
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"Job '%s': fire_claim ownership validation failed", self.job["id"], exc_info=True,
|
||||
)
|
||||
"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()
|
||||
@@ -2643,8 +2534,7 @@ def _save_compose_deliver(
|
||||
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,
|
||||
)
|
||||
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
|
||||
@@ -2693,13 +2583,10 @@ def _finish_interrupted_run(job: dict, execution_id: str, delivery_error: Option
|
||||
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,
|
||||
)
|
||||
"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.",
|
||||
)
|
||||
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:
|
||||
@@ -2713,10 +2600,8 @@ def _finish_completed_run(d: _RunDelivery, fire_owner: Optional[str], execution_
|
||||
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.",
|
||||
)
|
||||
execution_id, success=False,
|
||||
error="Fire claim ownership lost before terminal completion.")
|
||||
return True
|
||||
delivery_outcome = _classify_delivery_outcome(
|
||||
delivery_error=d.delivery_error,
|
||||
@@ -2731,11 +2616,7 @@ def _finish_completed_run(d: _RunDelivery, fire_owner: Optional[str], execution_
|
||||
# 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,
|
||||
)
|
||||
execution_id, success=d.success, error=d.error, delivery_outcome=delivery_outcome)
|
||||
return True
|
||||
|
||||
|
||||
@@ -2768,13 +2649,8 @@ def _deliver_crash_failure(
|
||||
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,
|
||||
)
|
||||
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
|
||||
@@ -2782,13 +2658,8 @@ def _deliver_crash_failure(
|
||||
|
||||
|
||||
def _run_one_job_body(
|
||||
job: dict,
|
||||
*,
|
||||
adapters=None,
|
||||
loop=None,
|
||||
verbose: bool = False,
|
||||
extra_prompt: Optional[str] = None,
|
||||
fire_claim_lost: Optional[_CancelEventLike] = None,
|
||||
job: dict, *, adapters=None, loop=None, verbose: bool = False,
|
||||
extra_prompt: Optional[str] = None, fire_claim_lost: Optional[_CancelEventLike] = None,
|
||||
execution_token: Optional[object] = None,
|
||||
) -> bool:
|
||||
fence = _FireOwnership(job, fire_claim_lost)
|
||||
@@ -2802,10 +2673,7 @@ def _run_one_job_body(
|
||||
delivery_attempted = False
|
||||
delivery_error = None
|
||||
from agent.secret_scope import (
|
||||
build_profile_secret_scope,
|
||||
reset_secret_scope,
|
||||
set_secret_scope,
|
||||
)
|
||||
build_profile_secret_scope, reset_secret_scope, set_secret_scope)
|
||||
|
||||
_scope_token = None
|
||||
_terminal_scope_token = None
|
||||
@@ -2815,13 +2683,10 @@ def _run_one_job_body(
|
||||
if not claim_dispatch(job["id"]):
|
||||
logger.info(
|
||||
"Job '%s': one-shot dispatch limit reached — skipping",
|
||||
job.get("name", job["id"]),
|
||||
)
|
||||
job.get("name", job["id"]))
|
||||
finish_execution(
|
||||
execution_id,
|
||||
success=False,
|
||||
error="Dispatch claim rejected; execution was not started.",
|
||||
)
|
||||
execution_id, success=False,
|
||||
error="Dispatch claim rejected; execution was not started.")
|
||||
return True # not an error — already handled/removed
|
||||
|
||||
mark_execution_running(execution_id)
|
||||
@@ -2833,8 +2698,7 @@ def _run_one_job_body(
|
||||
# process-global TERMINAL_* env a concurrent profile pinned. Resolution failure installs a
|
||||
# refusal scope — terminal execution raises instead of using the launch process's policy.
|
||||
from tools.terminal_scope import (
|
||||
install_profile_terminal_scope,
|
||||
)
|
||||
install_profile_terminal_scope)
|
||||
|
||||
_terminal_scope_token = install_profile_terminal_scope(_get_hermes_home())
|
||||
# Defer agent teardown until AFTER delivery; closing first races the live send against a
|
||||
@@ -2848,8 +2712,7 @@ def _run_one_job_body(
|
||||
_run_kwargs = {
|
||||
"defer_agent_teardown": _deferred_agents,
|
||||
"extra_prompt": extra_prompt,
|
||||
"execution_id": execution_id,
|
||||
}
|
||||
"execution_id": execution_id}
|
||||
if fire_claim_lost is not None:
|
||||
_run_kwargs["cancel_event"] = fire_claim_lost
|
||||
try:
|
||||
@@ -2871,8 +2734,7 @@ def _run_one_job_body(
|
||||
try:
|
||||
_save_compose_deliver(
|
||||
d, fence, final_response, output, adapters=adapters, loop=loop, verbose=verbose,
|
||||
execution_token=execution_token,
|
||||
)
|
||||
execution_token=execution_token)
|
||||
except _FireClaimLostDuringSideEffect:
|
||||
d.side_effect_ownership_lost = True
|
||||
finally:
|
||||
@@ -2905,8 +2767,7 @@ def _run_one_job_body(
|
||||
"Error processing job %s: %s",
|
||||
job["id"],
|
||||
_err_text,
|
||||
exc_info=(type(e), e, e.__traceback__),
|
||||
)
|
||||
exc_info=(type(e), e, e.__traceback__))
|
||||
delivery_outcome = "suppressed"
|
||||
# Owner fencing: a stale worker whose claim was taken over (or transport-cancelled) must not
|
||||
# send a failure alert on top of the replacement run's; fall through to fenced bookkeeping.
|
||||
@@ -2917,8 +2778,7 @@ def _run_one_job_body(
|
||||
and not _fire_claim_ownership_lost()
|
||||
):
|
||||
delivery_error, delivery_outcome = _deliver_crash_failure(
|
||||
job, _err_text, adapters=adapters, loop=loop,
|
||||
)
|
||||
job, _err_text, adapters=adapters, loop=loop)
|
||||
try:
|
||||
if not _consume_interrupted_flag(job["id"], execution_token):
|
||||
mark_kwargs = {}
|
||||
@@ -2932,11 +2792,7 @@ def _run_one_job_body(
|
||||
logger.error("Failed to record interrupted run for job %s: %s", job["id"], record_err)
|
||||
try:
|
||||
finish_execution(
|
||||
execution_id,
|
||||
success=False,
|
||||
error=_err_text,
|
||||
delivery_outcome=delivery_outcome,
|
||||
)
|
||||
execution_id, success=False, error=_err_text, delivery_outcome=delivery_outcome)
|
||||
except Exception as record_err:
|
||||
logger.error("Failed to finish execution record for job %s: %s", job["id"], record_err)
|
||||
if not isinstance(e, Exception):
|
||||
@@ -2992,8 +2848,7 @@ class CronSchedulerRegistrationError(RuntimeError):
|
||||
"job_id": self.job["id"],
|
||||
"job_saved": True,
|
||||
"scheduler_registered": False,
|
||||
"retry_create": False,
|
||||
}
|
||||
"retry_create": False}
|
||||
|
||||
|
||||
def create_job_with_scheduler_registration(**kwargs) -> dict:
|
||||
@@ -3042,8 +2897,7 @@ def _worktree_maintenance_repos() -> List[str]:
|
||||
probe = subprocess.run(
|
||||
["git", "rev-parse", "--show-toplevel"],
|
||||
capture_output=True, text=True, encoding="utf-8",
|
||||
errors="replace", timeout=5, cwd=workdir,
|
||||
)
|
||||
errors="replace", timeout=5, cwd=workdir)
|
||||
if probe.returncode == 0 and probe.stdout.strip():
|
||||
repos.add(probe.stdout.strip())
|
||||
except Exception:
|
||||
@@ -3112,8 +2966,7 @@ def _acquire_tick_lock(lock_file):
|
||||
logger.error(
|
||||
"Cron tick could not acquire tick lock: %s — scheduler will "
|
||||
"attempt fd reclamation and retry with backoff",
|
||||
exc,
|
||||
)
|
||||
exc)
|
||||
else:
|
||||
logger.error("Cron tick could not acquire tick lock: %s", exc)
|
||||
raise
|
||||
@@ -3148,8 +3001,7 @@ def _maybe_reap_dead_owners() -> None:
|
||||
logger.warning(
|
||||
"Reclaimed %d cron execution(s) whose owner process died "
|
||||
"before reaching a terminal state (marked unknown)",
|
||||
_reclaimed,
|
||||
)
|
||||
_reclaimed)
|
||||
except Exception as _reap_exc:
|
||||
logger.debug("Dead-owner execution reclaim failed: %s", _reap_exc)
|
||||
|
||||
@@ -3205,10 +3057,7 @@ def _process_due_job(job: dict, adapters, loop, verbose: bool) -> bool:
|
||||
claimed = claim_job_for_fire(job["id"], return_job=True)
|
||||
if not claimed:
|
||||
finish_execution(
|
||||
job["execution_id"],
|
||||
success=False,
|
||||
error="Fire claim lost; execution was not started.",
|
||||
)
|
||||
job["execution_id"], success=False, error="Fire claim lost; execution was not started.")
|
||||
return True
|
||||
# CAS returns the persisted record; bool fallback only for older test doubles.
|
||||
claimed_job = dict(claimed) if isinstance(claimed, dict) else dict(job)
|
||||
@@ -3235,8 +3084,7 @@ def _submit_with_guard(job: dict, pool: concurrent.futures.ThreadPoolExecutor, p
|
||||
logger.warning(
|
||||
"Could not clear run_claim for job '%s' after dispatch "
|
||||
"failure: %s (claim will expire at TTL)",
|
||||
job_label, claim_err,
|
||||
)
|
||||
job_label, claim_err)
|
||||
|
||||
def _not_dispatched_shutdown() -> None:
|
||||
logger.warning("Job '%s' not dispatched — interpreter is shutting down", job_label)
|
||||
@@ -3259,8 +3107,7 @@ def _submit_with_guard(job: dict, pool: concurrent.futures.ThreadPoolExecutor, p
|
||||
release_running_job(job_id)
|
||||
_clear_run_claim_best_effort()
|
||||
logger.exception(
|
||||
"Job '%s' not dispatched: execution creation failed: %s", job_label, execution_err,
|
||||
)
|
||||
"Job '%s' not dispatched: execution creation failed: %s", job_label, execution_err)
|
||||
return None
|
||||
|
||||
def _run_and_release(j=dispatched_job, ctx=_ctx):
|
||||
@@ -3275,10 +3122,7 @@ def _submit_with_guard(job: dict, pool: concurrent.futures.ThreadPoolExecutor, p
|
||||
release_running_job(job_id)
|
||||
_clear_run_claim_best_effort()
|
||||
finish_execution(
|
||||
execution["id"],
|
||||
success=False,
|
||||
error=f"Executor dispatch failed: {submit_err}",
|
||||
)
|
||||
execution["id"], success=False, error=f"Executor dispatch failed: {submit_err}")
|
||||
if isinstance(submit_err, RuntimeError) and _interpreter_shutting_down(submit_err):
|
||||
_not_dispatched_shutdown()
|
||||
else:
|
||||
@@ -3292,13 +3136,7 @@ def _submit_with_guard(job: dict, pool: concurrent.futures.ThreadPoolExecutor, p
|
||||
|
||||
|
||||
def tick(
|
||||
verbose: bool = True,
|
||||
adapters=None,
|
||||
loop=None,
|
||||
sync: bool = True,
|
||||
*,
|
||||
can_dispatch=None,
|
||||
):
|
||||
verbose: bool = True, adapters=None, loop=None, sync: bool = True, *, can_dispatch=None):
|
||||
"""Check and run all due jobs. File-locked so only one tick runs at a time (gateway ticker vs
|
||||
standalone daemon / manual tick). ``can_dispatch``: optional gate; false leaves due jobs for the
|
||||
next allowed tick. Returns the number of jobs executed (0 if another tick holds the lock)."""
|
||||
@@ -3357,8 +3195,7 @@ def tick(
|
||||
logger.info(
|
||||
"Running %d job(s) in parallel (max_workers=%s)",
|
||||
len(due_jobs),
|
||||
_max_workers if _max_workers else "unbounded",
|
||||
)
|
||||
_max_workers if _max_workers else "unbounded")
|
||||
|
||||
def _process_job(job: dict) -> bool:
|
||||
return _process_due_job(job, adapters, loop, verbose)
|
||||
|
||||
Reference in New Issue
Block a user