diff --git a/fxsharing/shares/tasks.py b/fxsharing/shares/tasks.py index 6954ff6..4ddad64 100644 --- a/fxsharing/shares/tasks.py +++ b/fxsharing/shares/tasks.py @@ -113,6 +113,10 @@ class BaseTaskWithRetry(DjangoTask): - exponential backoff capped at `retry_backoff_max` seconds, with jitter - a structured log line on every retry - a `DeadLetterTask` row + structured log line when retries are exhausted + + A task can opt out of the DLQ row with `dead_letter=False` (e.g. when its + failures are not retryable and so a stored row would only be noise). The + failure is still logged and reaches Sentry via the Celery integration. """ autoretry_for = (Exception,) @@ -120,6 +124,7 @@ class BaseTaskWithRetry(DjangoTask): retry_backoff_max = 600 retry_jitter = True max_retries = 3 + dead_letter = True def on_retry(self, exc, task_id, args, kwargs, einfo): attempt = (self.request.retries or 0) + 1 @@ -147,6 +152,25 @@ def on_failure(self, exc, task_id, args, kwargs, einfo): traceback = einfo.traceback if einfo else "" queue = (self.request.delivery_info or {}).get("routing_key", "") or "" + + if not self.dead_letter: + # Failure is not retryable, so a DLQ row would only be noise. Still + # log it (and let `super().on_failure` report it to Sentry via the + # Celery integration) so the failure stays observable. + logger.error( + "celery task failed (not dead-lettered): %s exc=%s", + self.name, + exc, + extra={ + "task_name": self.name, + "task_id": task_id, + "exception_class": type(exc).__name__, + "queue": queue, + }, + ) + super().on_failure(exc, task_id, args, kwargs, einfo) + return + logger.error( "celery task moved to DLQ: %s exc=%s", self.name, @@ -177,6 +201,7 @@ def on_failure(self, exc, task_id, args, kwargs, einfo): @shared_task( base=BaseTaskWithRetry, autoretry_for=(requests.exceptions.RequestException,), + dead_letter=False, ) def fetch_link_preview(link_id): from .models import Link diff --git a/fxsharing/shares/tests.py b/fxsharing/shares/tests.py index d675841..c0c0c0a 100644 --- a/fxsharing/shares/tests.py +++ b/fxsharing/shares/tests.py @@ -56,6 +56,16 @@ def _always_failing_task(value): raise ValueError(f"boom: {value}") +@shared_task( + base=BaseTaskWithRetry, + max_retries=0, + retry_backoff=False, + dead_letter=False, +) +def _always_failing_no_dlq_task(value): + raise ValueError(f"boom: {value}") + + class TestShareModel(TestCase): @classmethod def setUpTestData(cls): @@ -984,6 +994,46 @@ def test_on_failure_increments_task_deadlettered(self): 1, {"task": _always_failing_task.name, "exception_class": "ValueError"} ) + def test_on_failure_skips_dlq_when_dead_letter_false(self): + einfo = MagicMock() + einfo.traceback = "ValueError: boom" + + with self.assertLogs("fxsharing.shares.tasks", level="ERROR") as cm: + _always_failing_no_dlq_task.on_failure( + exc=ValueError("boom"), + task_id="task-no-dlq", + args=("hi",), + kwargs={}, + einfo=einfo, + ) + + assert not DeadLetterTask.objects.filter(task_id="task-no-dlq").exists() + assert any("not dead-lettered" in msg for msg in cm.output) + assert any(_always_failing_no_dlq_task.name in msg for msg in cm.output) + + def test_on_failure_no_deadlettered_metric_when_dead_letter_false(self): + einfo = MagicMock() + einfo.traceback = "ValueError: boom" + with patch("fxsharing.shares.metrics.task_deadlettered") as counter: + _always_failing_no_dlq_task.on_failure( + exc=ValueError("boom"), + task_id="task-no-dlq", + args=("hi",), + kwargs={}, + einfo=einfo, + ) + counter.add.assert_not_called() + + def test_failing_no_dlq_task_apply_creates_no_dlq_row(self): + result = _always_failing_no_dlq_task.apply(args=("payload",)) + assert result.failed() + assert not DeadLetterTask.objects.filter( + task_name=_always_failing_no_dlq_task.name + ).exists() + + def test_fetch_link_preview_opts_out_of_dlq(self): + assert fetch_link_preview.dead_letter is False + @override_settings(DEBUG=True) class TestSeedCommand(TestCase):