The cron SQLite stores write hermes_time.now().isoformat(), an ISO string that carries the local UTC offset. That offset changes at a DST transition and after a timezone config change, so the text order of these columns is not their time order. Across a fall-back, 01:10-05:00 is 20 minutes after 01:50-04:00 but sorts before it. latest_execution(), latest_executions() (which fills job["latest_execution"] for `hermes cron list` and the dashboard), list_executions() and its before_claimed_at cursor, the delivery queue's pending-claim order, delivery and execution pruning, and list_incidents() all returned the older row first. Every ORDER BY and the cursor comparison on these columns now use julianday(col), which SQLite resolves to the same instant for any offset. The raw text stays as the next sort key, so rows within the same millisecond (julianday's resolution) keep their microsecond order, and the cursor compares the same (instant, text) key on both sides. Stored values and the schema are unchanged, so existing rows with mixed offsets order correctly without a migration. latest_executions() moves from a correlated per-row subquery to one ROW_NUMBER() window. With julianday() in the ORDER BY the index on claimed_at no longer serves the per-row LIMIT 1, and the correlated form took about 90 ms at 1000 rows; the window form takes about 1 ms. Refs #86520 Thanks to #86889 (zhao0112) for the executions-ledger diagnosis.
237 lines
8.0 KiB
Python
237 lines
8.0 KiB
Python
"""Durable at-most-once delivery handoff for restart-safe cron workers."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import sqlite3
|
|
from unittest.mock import Mock
|
|
|
|
import pytest
|
|
|
|
|
|
def test_pending_delivery_is_claimed_and_sent_once(tmp_path, monkeypatch):
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
queue.enqueue("exec-1", {"id": "job-1"}, "brief")
|
|
send = Mock(return_value=None)
|
|
|
|
assert queue.drain(send) == 1
|
|
assert queue.drain(send) == 0
|
|
send.assert_called_once_with({"id": "job-1"}, "brief", False)
|
|
status = queue.get_status("exec-1")
|
|
assert status["status"] == "delivered"
|
|
assert status["job_json"] == "{}"
|
|
assert status["content"] == ""
|
|
|
|
|
|
def test_pending_deliveries_are_claimed_in_instant_order_across_dst_fall_back(
|
|
tmp_path, monkeypatch
|
|
):
|
|
"""01:10-05:00 is 20 minutes after 01:50-04:00 but sorts first as text."""
|
|
from datetime import datetime
|
|
from zoneinfo import ZoneInfo
|
|
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
new_york = ZoneInfo("America/New_York")
|
|
monkeypatch.setattr(
|
|
queue, "_hermes_now", lambda: datetime(2026, 11, 1, 1, 50, tzinfo=new_york, fold=0)
|
|
)
|
|
queue.enqueue("exec-z-earlier", {"id": "job-1"}, "first")
|
|
monkeypatch.setattr(
|
|
queue, "_hermes_now", lambda: datetime(2026, 11, 1, 1, 10, tzinfo=new_york, fold=1)
|
|
)
|
|
queue.enqueue("exec-a-later", {"id": "job-2"}, "second")
|
|
|
|
assert queue.claim_next()["execution_id"] == "exec-z-earlier"
|
|
assert queue.claim_next()["execution_id"] == "exec-a-later"
|
|
|
|
|
|
def test_terminal_delivery_retention_is_bounded(tmp_path, monkeypatch):
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
monkeypatch.setattr(queue, "MAX_TERMINAL_DELIVERIES", 2, raising=False)
|
|
for index in range(4):
|
|
execution_id = f"exec-{index}"
|
|
queue.enqueue(execution_id, {"id": f"job-{index}"}, f"brief-{index}")
|
|
assert queue.claim_next()["execution_id"] == execution_id
|
|
assert queue._finish(execution_id, error=None)
|
|
|
|
# Pruning may discard verbose outcome rows, but never the durable
|
|
# idempotency tombstone for an execution that could be replayed later.
|
|
pruned = queue.get_status("exec-0")
|
|
assert pruned is not None
|
|
assert pruned["status"] == "delivered"
|
|
assert queue.get_status("exec-1")["status"] == "delivered"
|
|
assert queue.get_status("exec-2")["status"] == "delivered"
|
|
assert queue.get_status("exec-3")["status"] == "delivered"
|
|
|
|
# Pruning may discard verbose outcome rows, but never the durable
|
|
# idempotency tombstone for an execution that could be replayed later.
|
|
queue.enqueue("exec-0", {"id": "job-replayed"}, "duplicate brief")
|
|
send = Mock(return_value=None)
|
|
assert queue.drain(send) == 0
|
|
send.assert_not_called()
|
|
|
|
|
|
def test_failure_delivery_lane_survives_durable_handoff(tmp_path, monkeypatch):
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
queue.enqueue(
|
|
"exec-failure",
|
|
{"id": "job-failure", "failure_deliver": "local"},
|
|
"failed",
|
|
for_failure=True,
|
|
)
|
|
send = Mock(return_value=None)
|
|
|
|
assert queue.drain(send) == 1
|
|
send.assert_called_once_with(
|
|
{"id": "job-failure", "failure_deliver": "local"},
|
|
"failed",
|
|
True,
|
|
)
|
|
|
|
|
|
def test_legacy_queue_schema_adds_failure_lane_before_enqueue(tmp_path, monkeypatch):
|
|
import cron.delivery_queue as queue
|
|
|
|
db = tmp_path / "deliveries.db"
|
|
with sqlite3.connect(db) as conn:
|
|
conn.execute(
|
|
"""CREATE TABLE deliveries (
|
|
execution_id TEXT PRIMARY KEY,
|
|
job_json TEXT NOT NULL,
|
|
content TEXT NOT NULL,
|
|
status TEXT NOT NULL,
|
|
owner_process_id TEXT,
|
|
owner_pid INTEGER,
|
|
owner_started_at INTEGER,
|
|
created_at TEXT NOT NULL,
|
|
finished_at TEXT,
|
|
error TEXT
|
|
)"""
|
|
)
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", db)
|
|
|
|
queue.enqueue(
|
|
"exec-migrated",
|
|
{"id": "job-migrated"},
|
|
"failed",
|
|
for_failure=True,
|
|
)
|
|
|
|
assert queue.get_status("exec-migrated")["for_failure"] == 1
|
|
|
|
|
|
def test_wait_timeout_marks_inflight_delivery_unknown_without_retry(
|
|
tmp_path, monkeypatch
|
|
):
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
queue.enqueue("exec-inflight", {"id": "job-inflight"}, "result")
|
|
assert queue.claim_next() is not None
|
|
|
|
error = queue.enqueue_and_wait(
|
|
"exec-inflight", {"id": "job-inflight"}, "result", timeout=0
|
|
)
|
|
|
|
assert error is not None
|
|
assert "outcome is unknown" in error
|
|
status = queue.get_status("exec-inflight")
|
|
assert status is not None
|
|
assert status["status"] == "unknown"
|
|
send = Mock(return_value=None)
|
|
assert queue.drain(send) == 0
|
|
send.assert_not_called()
|
|
|
|
|
|
def test_dead_delivery_owner_becomes_unknown_and_is_not_retried(
|
|
tmp_path, monkeypatch
|
|
):
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
queue.enqueue("exec-1", {"id": "job-1"}, "brief")
|
|
assert queue.claim_next() is not None
|
|
monkeypatch.setattr(queue, "_PROCESS_ID", "replacement-gateway")
|
|
monkeypatch.setattr(queue, "_owner_is_live", lambda _pid, _started: False)
|
|
|
|
assert queue.recover_abandoned() == 1
|
|
send = Mock()
|
|
assert queue.drain(send) == 0
|
|
send.assert_not_called()
|
|
assert queue.get_status("exec-1")["status"] == "unknown"
|
|
|
|
|
|
def test_delivery_failure_is_terminal_not_retried_and_redacted(
|
|
tmp_path, monkeypatch
|
|
):
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
queue.enqueue("exec-1", {"id": "job-1"}, "brief")
|
|
send = Mock(return_value="request failed: https://example.test/?token=TOKEN123")
|
|
|
|
assert queue.drain(send) == 1
|
|
assert queue.drain(send) == 0
|
|
assert send.call_count == 1
|
|
status = queue.get_status("exec-1")
|
|
assert status is not None
|
|
assert status["status"] == "failed"
|
|
assert "TOKEN123" not in status["error"]
|
|
assert "token=***" in status["error"]
|
|
|
|
|
|
def test_wait_timeout_leaves_unclaimed_delivery_queued_for_next_gateway(
|
|
tmp_path, monkeypatch
|
|
):
|
|
"""A row nobody claimed was never attempted: it is not uncertain, so a
|
|
gateway outage longer than the worker's wait budget must not lose it."""
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
job = {"id": "job-3", "deliver": "origin"}
|
|
|
|
error = queue.enqueue_and_wait("exec-3", job, "result", timeout=0)
|
|
|
|
# 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"
|
|
send = Mock(return_value=None)
|
|
assert queue.drain(send) == 1
|
|
send.assert_called_once_with(job, "result", False)
|
|
assert queue.get_status("exec-3")["status"] == "delivered"
|
|
|
|
|
|
def test_same_gateway_recovers_terminalization_failure_without_resending(
|
|
tmp_path, monkeypatch
|
|
):
|
|
import cron.delivery_queue as queue
|
|
|
|
monkeypatch.setattr(queue, "DELIVERY_DB", tmp_path / "deliveries.db")
|
|
queue.enqueue("exec-4", {"id": "job-4"}, "result")
|
|
send = Mock(return_value=None)
|
|
original_finish = queue._finish
|
|
monkeypatch.setattr(
|
|
queue,
|
|
"_finish",
|
|
Mock(side_effect=OSError("database temporarily unavailable")),
|
|
)
|
|
|
|
with pytest.raises(OSError, match="temporarily unavailable"):
|
|
queue.drain(send)
|
|
|
|
monkeypatch.setattr(queue, "_finish", original_finish)
|
|
assert queue.drain(send) == 0
|
|
send.assert_called_once()
|
|
status = queue.get_status("exec-4")
|
|
assert status["status"] == "unknown"
|
|
assert "not retried" in status["error"]
|