mirror of
https://github.com/tiennm99/DocsGPT.git
synced 2026-10-11 03:12:55 +00:00
fix(tasks): a deferred duplicate stands down instead of failing the task
A task that runs longer than the broker's visibility timeout is redelivered while its first run is still going. The idempotency lease correctly stops the duplicate from doing the work, but the duplicate then re-queued itself once per LEASE_TTL until celery ran out of retries and raised MaxRetriesExceededError — so a perfectly healthy long task (a large graph extraction is the one that found this) reported a task failure, with a traceback, while the real run was still making progress next to it. Catch the exhaustion and return a "deferred" result instead. A normal deferral still re-queues: only the give-up path changes, and the lease holder's dedup row is left untouched so its own completion still records.
This commit is contained in:
1 parent
e3d819d9fd
commit
6ea737e57b
2 files changed
+93
-4
No files matched your search
@@ -9,6 +9,8 @@ import threading
|
||||
import uuid
|
||||
from typing import Any, Callable, Optional
|
||||
|
||||
from celery.exceptions import MaxRetriesExceededError
|
||||
|
||||
from docsgpt.storage.db.repositories.idempotency import IdempotencyRepository
|
||||
from docsgpt.storage.db.session import db_readonly, db_session
|
||||
|
||||
@@ -81,10 +83,28 @@ def with_idempotency(
|
||||
"idempotency: live lease held; deferring task=%s key=%s",
|
||||
task_name, key,
|
||||
)
|
||||
raise self.retry(
|
||||
countdown=LEASE_TTL_SECONDS,
|
||||
max_retries=LEASE_RETRY_MAX,
|
||||
)
|
||||
try:
|
||||
raise self.retry(
|
||||
countdown=LEASE_TTL_SECONDS,
|
||||
max_retries=LEASE_RETRY_MAX,
|
||||
)
|
||||
except MaxRetriesExceededError:
|
||||
# The holder is simply slower than LEASE_RETRY_MAX
|
||||
# deferrals — a task that outruns the broker's visibility
|
||||
# timeout is redelivered while its first run is still
|
||||
# going. Standing down is the correct end state for the
|
||||
# duplicate; raising here would report a failure for a
|
||||
# task that is running normally somewhere else.
|
||||
logger.info(
|
||||
"idempotency: lease still held after %s deferrals; "
|
||||
"leaving task=%s key=%s to its holder",
|
||||
LEASE_RETRY_MAX, task_name, key,
|
||||
)
|
||||
return {
|
||||
"status": "deferred",
|
||||
"reason": "another worker holds the lease",
|
||||
"idempotency_key": key,
|
||||
}
|
||||
|
||||
if attempt > MAX_TASK_ATTEMPTS:
|
||||
logger.error(
|
||||
|
||||
@@ -335,6 +335,75 @@ class TestLiveLeaseDefersConcurrentRun:
|
||||
assert row[2] == "completed"
|
||||
|
||||
|
||||
@pytest.mark.unit
|
||||
class TestLeaseDeferralGivesUpQuietly:
|
||||
"""Deferral is bookkeeping, not failure.
|
||||
|
||||
A task that outruns the broker's visibility timeout is redelivered while
|
||||
the first worker is still running it. The lease keeps the duplicate from
|
||||
doing the work, but the duplicate kept re-queueing itself until celery
|
||||
exhausted ``LEASE_RETRY_MAX`` and raised ``MaxRetriesExceededError``, so a
|
||||
healthy long task logged a task failure. The duplicate should stand down
|
||||
instead and leave the run to the worker that holds the lease.
|
||||
"""
|
||||
|
||||
def _hold_lease(self, pg_conn, key):
|
||||
from docsgpt.storage.db.repositories.idempotency import (
|
||||
IdempotencyRepository,
|
||||
)
|
||||
|
||||
IdempotencyRepository(pg_conn).try_claim_lease(
|
||||
key=key, task_name="thing",
|
||||
task_id="t-worker-1", owner_id="worker-1",
|
||||
)
|
||||
|
||||
def test_exhausted_retries_return_deferred_instead_of_raising(self, pg_conn):
|
||||
from celery.exceptions import MaxRetriesExceededError
|
||||
|
||||
from docsgpt.api.user.idempotency import with_idempotency
|
||||
|
||||
self._hold_lease(pg_conn, "k-long-run")
|
||||
invocations = {"count": 0}
|
||||
|
||||
@with_idempotency(task_name="thing")
|
||||
def task(self, idempotency_key=None):
|
||||
invocations["count"] += 1
|
||||
return {"ran": True}
|
||||
|
||||
# Celery raises this from ``self.retry`` once max_retries is hit.
|
||||
worker2 = _fake_celery_self("t-worker-2")
|
||||
worker2.retry.side_effect = MaxRetriesExceededError("out of retries")
|
||||
|
||||
with _patch_decorator_db(pg_conn):
|
||||
result = task(worker2, idempotency_key="k-long-run")
|
||||
|
||||
assert result["status"] == "deferred"
|
||||
# The lease holder is still running it; the duplicate did not.
|
||||
assert invocations["count"] == 0
|
||||
# The holder's row is untouched — not failed, not completed.
|
||||
row = _row_for(pg_conn, "k-long-run")
|
||||
assert row[2] == "pending"
|
||||
|
||||
def test_a_normal_retry_still_propagates(self, pg_conn):
|
||||
"""Only exhaustion stands down; the first deferrals must re-queue."""
|
||||
from docsgpt.api.user.idempotency import with_idempotency
|
||||
|
||||
self._hold_lease(pg_conn, "k-busy-once")
|
||||
|
||||
@with_idempotency(task_name="thing")
|
||||
def task(self, idempotency_key=None):
|
||||
return {"ran": True}
|
||||
|
||||
class _RetrySignal(Exception):
|
||||
pass
|
||||
|
||||
worker2 = _fake_celery_self("t-worker-2")
|
||||
worker2.retry.side_effect = _RetrySignal("retry scheduled")
|
||||
|
||||
with _patch_decorator_db(pg_conn), pytest.raises(_RetrySignal):
|
||||
task(worker2, idempotency_key="k-busy-once")
|
||||
|
||||
|
||||
@pytest.mark.unit
|
||||
class TestExceptionPathReleasesLease:
|
||||
"""When ``fn`` raises, the lease is dropped so the next attempt
|
||||
|
||||
Reference in new issue
Block a user