From 5b2c417db858a09eca066165a49e4171b783dd3d Mon Sep 17 00:00:00 2001 From: liuzikaii <2319582736@qq.com> Date: Sun, 6 Sep 2026 10:51:53 +0800 Subject: [PATCH] fix(cron): serialize delivery deduplication with terminal retention --- cron/delivery_queue.py | 3 + .../test_delivery_queue_retention_race.py | 62 +++++++++++++++++++ 2 files changed, 65 insertions(+) create mode 100644 tests/cron/test_delivery_queue_retention_race.py diff --git a/cron/delivery_queue.py b/cron/delivery_queue.py index 689fd620d6..612b626047 100644 --- a/cron/delivery_queue.py +++ b/cron/delivery_queue.py @@ -140,6 +140,9 @@ def enqueue( ) -> dict: """Persist one idempotent delivery request before the worker waits.""" with _transaction() as conn: + # Serialize the tombstone check and insert with retention in other + # processes, which can move a terminal delivery into the tombstone table. + conn.execute("BEGIN IMMEDIATE") tombstone = conn.execute( "SELECT terminal_status, finished_at FROM delivery_tombstones " "WHERE execution_id=?", diff --git a/tests/cron/test_delivery_queue_retention_race.py b/tests/cron/test_delivery_queue_retention_race.py new file mode 100644 index 0000000000..c1ff028397 --- /dev/null +++ b/tests/cron/test_delivery_queue_retention_race.py @@ -0,0 +1,62 @@ +"""A delivery must stay terminal when retention races a repeated enqueue.""" + +import sqlite3 +import subprocess +import sys + + +def test_enqueue_during_retention_never_recreates_delivered_request(tmp_path, monkeypatch): + from cron import delivery_queue as queue + + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + db = tmp_path / "deliveries.db" + monkeypatch.setattr(queue, "DELIVERY_DB", db) + queue.enqueue("execution", {"id": "job"}, "original") + sends = [] + assert queue.drain(lambda *args: sends.append(args)) == 1 + connect = sqlite3.connect + prunes = [] + + class InterleavedConnection(sqlite3.Connection): + def execute(self, sql, parameters=()): + cursor = super().execute(sql, parameters) + if sql.startswith("SELECT terminal_status, finished_at FROM delivery_tombstones"): + # Pause after the real tombstone lookup, before enqueue can + # INSERT. A separate process exercises SQLite's actual lock + # boundary rather than the module's process-local RLock. + result = subprocess.run( + [sys.executable, "-c", """ +import sqlite3 +import sys +from cron import delivery_queue as queue +queue.MAX_TERMINAL_DELIVERIES = 0 +try: + with sqlite3.connect(sys.argv[1], timeout=0) as conn: + queue._prune_terminal_unlocked(conn) +except sqlite3.OperationalError as exc: + if 'locked' not in str(exc): + raise + sys.exit(75) +""", str(db)], capture_output=True, text=True, timeout=30, + ) + assert result.returncode in (0, 75), result.stderr + prunes.append(result.returncode) + return cursor + + with monkeypatch.context() as patch: + patch.setattr( + queue.sqlite3, "connect", + lambda *args, **kwargs: connect(*args, **kwargs, factory=InterleavedConnection), + ) + replay = queue.enqueue("execution", {"id": "job"}, "duplicate") + + assert prunes, "the concurrent retention path must actually be attempted" + # If enqueue held the SQLite writer lock, complete the deferred pruning now. + monkeypatch.setattr(queue, "MAX_TERMINAL_DELIVERIES", 0) + with queue._transaction() as conn: + queue._prune_terminal_unlocked(conn) + + assert replay["status"] == "delivered" + assert queue.get_status("execution")["status"] == "delivered" + assert queue.drain(lambda *args: sends.append(args)) == 0 + assert len(sends) == 1