13fb87af92
cron/scheduler.py (8535 -> 6353): - _deliver_result split into per-target helpers: _resolve_target_transport (live/relay/standalone transport + enablement), _deliver_via_live_adapter (_live_route_metadata for Telegram DM-topic vs forum routing, _live_send_text with the cancel()-based timeout disambiguation, _live_send_media, _seed_live_delivery_sessions), _standalone_send/_deliver_standalone. The three interpreter-shutdown skip branches and the repeated log+append+continue pattern collapse into one _note_target_error / one shutdown message; _TargetDelivery carries per-target state. - Thread/channel session seeding unified into _seed_cron_session (was two near-identical functions). - run_job decomposed into _run_no_agent_job, _apply_monitor_gate, _load_cron_job_config (_CronJobConfig), _resolve_job_runtime, _check_model_drift, _open_cron_session_db, _run_agent_with_watchdog, _finalize_cron_session, plus one _run_doc_header and one _audit closure for the success/failure paths. - _run_one_job_body: ownership-lost bookkeeping, delivery composition and outcome classification extracted; tick: _acquire_tick_lock/_release_tick_lock, _maybe_reap_dead_owners, _sweep_stale_inflight_for_tick, _process_due_job, _submit_with_guard. - _build_job_prompt: context_from injection and skill loading extracted; one _prepend_context_block for the four fenced-data blocks. - One _start_heartbeat_thread for the script- and fire-claim heartbeat threads. - Dropped unreachable return in SharedRouteAdapters.get; boolean-return and nested-if shapes collapsed. - Comments/docstrings compacted by hand, rationale kept (fd-leak reason for the late SessionDB close callback, title-persistence rules, no_agent classification gate, inactivity-vs-provider timeout ordering, stale-claim force-release, interruption token keying).
122 lines
4.7 KiB
Python
122 lines
4.7 KiB
Python
"""Regression coverage for #58720 / #55924 — cron scheduling races
|
|
interpreter finalization.
|
|
|
|
When the gateway tears down (SIGTERM from ``hermes update`` /
|
|
``hermes gateway stop`` / systemd restart, or an OOM-kill), a cron tick can
|
|
still fire. Once the Python interpreter is finalizing, ``concurrent.futures``
|
|
refuses new work with ``RuntimeError: cannot schedule new futures after
|
|
interpreter shutdown`` and asyncio's default executor is gone. The cron
|
|
delivery + dispatch paths used to hit that unguarded, crashing the tick and
|
|
spraying a traceback into ``errors.log`` on every restart-race.
|
|
|
|
The fix adds ``_interpreter_shutting_down()`` and guards the scheduling
|
|
sites so they skip gracefully with a warning instead of raising.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import sys
|
|
from unittest.mock import AsyncMock, MagicMock, patch
|
|
|
|
import pytest
|
|
|
|
|
|
class TestInterpreterShuttingDownHelper:
|
|
def test_true_when_finalizing(self):
|
|
from cron.scheduler import _interpreter_shutting_down
|
|
|
|
with patch("sys.is_finalizing", return_value=True):
|
|
assert _interpreter_shutting_down() is True
|
|
|
|
def test_false_when_not_finalizing_and_no_exc(self):
|
|
from cron.scheduler import _interpreter_shutting_down
|
|
|
|
with patch("sys.is_finalizing", return_value=False):
|
|
assert _interpreter_shutting_down() is False
|
|
|
|
def test_matches_shutdown_error_text_as_fallback(self):
|
|
"""The concurrent.futures module-global flag can be set a hair before
|
|
``sys.is_finalizing()`` flips — matching the error text catches that
|
|
race so a shutdown RuntimeError isn't misread as a real failure."""
|
|
from cron.scheduler import _interpreter_shutting_down
|
|
|
|
exc = RuntimeError("cannot schedule new futures after interpreter shutdown")
|
|
with patch("sys.is_finalizing", return_value=False):
|
|
assert _interpreter_shutting_down(exc) is True
|
|
|
|
def test_unrelated_error_is_not_shutdown(self):
|
|
from cron.scheduler import _interpreter_shutting_down
|
|
|
|
exc = RuntimeError("some other problem")
|
|
with patch("sys.is_finalizing", return_value=False):
|
|
assert _interpreter_shutting_down(exc) is False
|
|
|
|
|
|
class TestStandaloneDeliverySkipsDuringShutdown:
|
|
def _telegram_cfg(self):
|
|
from gateway.config import Platform
|
|
|
|
pconfig = MagicMock()
|
|
pconfig.enabled = True
|
|
mock_cfg = MagicMock()
|
|
mock_cfg.platforms = {Platform.TELEGRAM: pconfig}
|
|
return mock_cfg
|
|
|
|
def test_standalone_path_skips_without_scheduling(self):
|
|
"""With the interpreter finalizing, the standalone delivery path must
|
|
skip BEFORE attempting to schedule the send — no ``_send_to_platform``
|
|
call, a graceful warning-level skip, and an error string returned
|
|
(not a raised exception)."""
|
|
from cron.scheduler import _deliver_result
|
|
|
|
job = {
|
|
"id": "gov-job",
|
|
"name": "model-governor",
|
|
"deliver": "origin",
|
|
"origin": {"platform": "telegram", "chat_id": "123"},
|
|
}
|
|
send_mock = AsyncMock(return_value={"success": True})
|
|
with patch("gateway.config.load_gateway_config", return_value=self._telegram_cfg()), \
|
|
patch("tools.send_message_tool._send_to_platform", new=send_mock), \
|
|
patch("sys.is_finalizing", return_value=True):
|
|
result = _deliver_result(job, "daily report body")
|
|
|
|
send_mock.assert_not_called()
|
|
assert result is not None
|
|
assert "shutting down" in result
|
|
|
|
def test_normal_delivery_still_works_when_not_finalizing(self):
|
|
"""Guard must not regress the happy path: a normal (non-finalizing)
|
|
run still delivers via the standalone send."""
|
|
from cron.scheduler import _deliver_result
|
|
|
|
job = {
|
|
"id": "gov-job",
|
|
"name": "model-governor",
|
|
"deliver": "origin",
|
|
"origin": {"platform": "telegram", "chat_id": "123"},
|
|
}
|
|
send_mock = AsyncMock(return_value={"success": True})
|
|
with patch("gateway.config.load_gateway_config", return_value=self._telegram_cfg()), \
|
|
patch("tools.send_message_tool._send_to_platform", new=send_mock), \
|
|
patch("sys.is_finalizing", return_value=False):
|
|
result = _deliver_result(job, "daily report body")
|
|
|
|
send_mock.assert_called_once()
|
|
assert result is None
|
|
|
|
|
|
class TestSourceGuardrail:
|
|
@pytest.fixture
|
|
def source(self) -> str:
|
|
from pathlib import Path
|
|
|
|
return (
|
|
Path(__file__).resolve().parents[2] / "cron" / "scheduler.py"
|
|
).read_text(encoding="utf-8")
|
|
|
|
def test_helper_defined(self, source):
|
|
assert "def _interpreter_shutting_down(" in source
|
|
|
|
|