fix(cron): key worker deliveries by the job's own attempt; deferred send is not a failure
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 <other>`) 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.
This commit is contained in:
@@ -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"
|
||||
|
||||
@@ -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 <other>``) 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
|
||||
|
||||
|
||||
Reference in New Issue
Block a user