From c3e9b28a4214fef7136d4b854beb1904941962bb Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 3 Sep 2026 12:29:24 +0530 Subject: [PATCH] fix(cron): key worker deliveries by the job's own attempt; deferred send is not a failure MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review fold-in on the salvage of #101877: - `_deliver_result` routed to the durable queue whenever the worker's `_HERMES_CRON_EXTERNAL_WORKER` marker was set, regardless of WHICH job was delivering. A worker whose script dispatches another job in-process (`hermes cron run `) inherits that env and would have queued the nested job's message under the outer execution id — `INSERT OR IGNORE` then drops it silently. Match the marker against the delivering job's own `execution_id`, as `run_one_job` already does. Regression test added (mutation-checked: fails with the guard removed). - A `pending` row left queued at the worker's wait timeout was still reported as a delivery error, so `mark_job_run` recorded `last_status=delivery_failed` for a message the next gateway's drain goes on to send, and nothing ever corrects the job record. Log and return success instead; the deliveries row is the authority for the send. - Reuse `cron.executions._TERMINAL_STATES` in the parent wait loop instead of a second hardcoded terminal set. --- cron/delivery_queue.py | 18 ++++++---- cron/scheduler.py | 14 +++++--- tests/cron/test_delivery_queue.py | 3 +- tests/cron/test_restart_safe_worker.py | 49 ++++++++++++++++++++++++++ 4 files changed, 73 insertions(+), 11 deletions(-) diff --git a/cron/delivery_queue.py b/cron/delivery_queue.py index f03413b727..635a295b2f 100644 --- a/cron/delivery_queue.py +++ b/cron/delivery_queue.py @@ -10,6 +10,7 @@ possibly-completed send. from __future__ import annotations import json +import logging import os import sqlite3 import threading @@ -25,6 +26,8 @@ from hermes_cli.sqlite_util import add_column_if_missing from hermes_constants import get_hermes_home from hermes_time import now as _hermes_now +logger = logging.getLogger(__name__) + DELIVERY_DB: Optional[Path] = None _PROCESS_ID = uuid.uuid4().hex _lock = threading.RLock() @@ -309,14 +312,12 @@ def _terminalize_wait_timeout(execution_id: str) -> str: A row still ``pending`` was provably never attempted, so it is left queued for whichever gateway comes up next (a restart that includes an update can - easily exceed the worker's wait budget). Only a row caught mid-send is + easily exceed the worker's wait budget). That is a deferral, not a + failure: report success so the job is not recorded ``delivery_failed`` for + a message the drain will still send. Only a row caught mid-send is uncertain and gets fenced ``unknown``. """ now = _hermes_now().isoformat() - pending_error = ( - "timed out waiting for a live gateway; delivery is still queued and " - "will be sent by the next gateway" - ) uncertain_error = ( "timed out while gateway delivery was in progress; outcome is unknown and " "was not retried" @@ -327,7 +328,12 @@ def _terminalize_wait_timeout(execution_id: str) -> str: (str(execution_id),), ).fetchone() if row is not None and row["status"] == "pending": - return pending_error + logger.warning( + "Cron delivery %s: no live gateway within the wait budget; " + "left queued for the next gateway", + execution_id, + ) + return "" conn.execute( """UPDATE deliveries SET status='unknown', finished_at=?, error=? WHERE execution_id=? AND status='delivering'""", diff --git a/cron/scheduler.py b/cron/scheduler.py index 65313c4a98..267805a8f1 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -726,6 +726,7 @@ from cron.jobs import ( use_cron_store, ) from cron.executions import ( + _TERMINAL_STATES, create_execution, finish_execution, get_execution, @@ -3222,8 +3223,15 @@ def _deliver_result( # Hand the send back through a durable queue so the current or replacement # gateway performs it with relay/E2EE parity. The execution id is the # idempotency key; the queue never retries an uncertain claimed send. + # Match on this job's own attempt: a worker's script may itself dispatch + # another job in-process (``hermes cron run``), and that nested delivery + # must not be keyed under the outer execution id. external_execution = os.environ.get("_HERMES_CRON_EXTERNAL_WORKER", "") - if external_execution and adapters is None: + if ( + external_execution + and adapters is None + and external_execution == str(job.get("execution_id") or "") + ): from cron.delivery_queue import enqueue_and_wait return enqueue_and_wait( @@ -8068,11 +8076,9 @@ def _wait_for_external_cron_worker_body( A gateway replacement may kill this waiter; it does not kill the scoped worker or change its ledger ownership. """ - terminal_states = {"completed", "failed", "unknown"} - def _is_terminal() -> bool: current = get_execution(execution_id) - return bool(current and current.get("status") in terminal_states) + return bool(current and current.get("status") in _TERMINAL_STATES) # The worker commits its terminal row before its process exits, so exit is # the correct wakeup. Each ledger read opens a connection and re-runs diff --git a/tests/cron/test_delivery_queue.py b/tests/cron/test_delivery_queue.py index 66ad30bb3f..b4c56edfaf 100644 --- a/tests/cron/test_delivery_queue.py +++ b/tests/cron/test_delivery_queue.py @@ -175,7 +175,8 @@ def test_wait_timeout_leaves_unclaimed_delivery_queued_for_next_gateway( error = queue.enqueue_and_wait("exec-3", job, "result", timeout=0) - assert "timed out" in error and "still queued" in error + # Deferred, not failed: the worker must not record delivery_failed. + assert error is None status = queue.get_status("exec-3") assert status is not None assert status["status"] == "pending" diff --git a/tests/cron/test_restart_safe_worker.py b/tests/cron/test_restart_safe_worker.py index 4820c3e899..1be7e30654 100644 --- a/tests/cron/test_restart_safe_worker.py +++ b/tests/cron/test_restart_safe_worker.py @@ -358,6 +358,55 @@ def test_shutdown_does_not_interrupt_restart_safe_waiter(): scheduler._interrupted_job_ids.discard(job_id) +def test_worker_delivery_queue_is_keyed_by_the_delivering_jobs_own_execution( + monkeypatch, tmp_path +): + """A nested in-process dispatch inside a worker (e.g. a script running + ``hermes cron run ``) must not queue under the OUTER execution id.""" + import cron.scheduler as scheduler + + queued = [] + monkeypatch.setattr( + "cron.delivery_queue.enqueue_and_wait", + lambda execution_id, job, content, **kw: ( + queued.append(execution_id) or "queued-marker" + ), + ) + monkeypatch.setattr( + scheduler, + "_resolve_delivery_targets", + lambda job, for_failure=False: [{"platform": "telegram", "chat_id": "123"}], + ) + + def _standalone(*_args, **_kwargs): + raise RuntimeError("standalone path reached") + + # First call the standalone (non-queue) path makes after the guard; the + # failure is reported as the delivery error string. + monkeypatch.setattr("gateway.config.load_gateway_config", _standalone) + monkeypatch.setenv("_HERMES_CRON_EXTERNAL_WORKER", "exec-outer") + + # Own attempt: routed through the durable queue. + assert scheduler._deliver_result( + {"id": "job-1", "execution_id": "exec-outer", "deliver": "telegram:123"}, + "done", + adapters=None, + loop=None, + ) == "queued-marker" + assert queued == ["exec-outer"] + + # A different job's attempt: must NOT be queued under exec-outer; it falls + # through to the standalone path. + error = scheduler._deliver_result( + {"id": "job-2", "execution_id": "exec-inner", "deliver": "telegram:123"}, + "done", + adapters=None, + loop=None, + ) + assert error == "failed to load gateway config: standalone path reached" + assert queued == ["exec-outer"] + + def test_gateway_tool_run_without_adapter_objects_hands_off(monkeypatch): import cron.scheduler as scheduler