mirror of
https://github.com/tiennm99/DocsGPT.git
synced 2026-10-11 03:12:55 +00:00
fix(tasks): a deferred duplicate records no result, rather than success
Returning a deferred marker traded one wrong signal for a worse one. A redelivery reuses the original task id — Context.as_execution_options carries task_id into the retry — so the duplicate's return marked the very id the client polls as SUCCESS. /api/task_status reports celery's state verbatim and the UI maps SUCCESS to "done", so the GraphRAG enable modal would announce a finished build, rendered from a payload with no counts, while the run holding the lease was still extracting. Raise Ignore instead: celery records no state for the duplicate, so the task id keeps whatever the holder sets and the poller keeps waiting. The autoretry wrapper re-raises Ignore ahead of autoretry_for, so the wider autoretry_for=(Exception,) on these tasks cannot turn it back into a retry.
This commit is contained in:
1 parent
3f774d813c
commit
b7e7872bf7
2 files changed
+21
-15
No files matched your search
@@ -9,7 +9,7 @@ import threading
|
||||
import uuid
|
||||
from typing import Any, Callable, Optional
|
||||
|
||||
from celery.exceptions import MaxRetriesExceededError
|
||||
from celery.exceptions import Ignore, MaxRetriesExceededError
|
||||
|
||||
from docsgpt.storage.db.repositories.idempotency import IdempotencyRepository
|
||||
from docsgpt.storage.db.session import db_readonly, db_session
|
||||
@@ -90,21 +90,23 @@ def with_idempotency(
|
||||
)
|
||||
except MaxRetriesExceededError:
|
||||
# The holder is simply slower than LEASE_RETRY_MAX
|
||||
# deferrals — a task that outruns the broker's visibility
|
||||
# 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.
|
||||
# going. Letting the exhaustion propagate would report a
|
||||
# failure for a task that is running normally — but so
|
||||
# would returning a value, only less visibly. A redelivery
|
||||
# reuses the original task id (``Context`` carries
|
||||
# ``task_id`` into the retry), so a return marks the very
|
||||
# id the client polls SUCCESS, and ``/api/task_status``
|
||||
# hands that to the UI as a finished build. ``Ignore``
|
||||
# records no state at all, leaving the outcome to the run
|
||||
# that actually holds the lease.
|
||||
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,
|
||||
}
|
||||
raise Ignore() from None
|
||||
|
||||
if attempt > MAX_TASK_ATTEMPTS:
|
||||
logger.error(
|
||||
|
||||
@@ -357,8 +357,8 @@ class TestLeaseDeferralGivesUpQuietly:
|
||||
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
|
||||
def test_exhausted_retries_stand_down_without_recording_a_result(self, pg_conn):
|
||||
from celery.exceptions import Ignore, MaxRetriesExceededError
|
||||
|
||||
from docsgpt.api.user.idempotency import with_idempotency
|
||||
|
||||
@@ -374,10 +374,14 @@ class TestLeaseDeferralGivesUpQuietly:
|
||||
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")
|
||||
# Ignore rather than a return value: a redelivery reuses the original
|
||||
# task id, so returning would mark the id the client is polling
|
||||
# SUCCESS — /api/task_status hands that straight to the UI, which
|
||||
# would announce a finished (empty) build while the holder is still
|
||||
# working. Ignore leaves the id's state to the holder.
|
||||
with _patch_decorator_db(pg_conn), pytest.raises(Ignore):
|
||||
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.
|
||||
|
||||
Reference in new issue
Block a user