Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 25 additions & 0 deletions fxsharing/shares/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -113,13 +113,18 @@ 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,)
retry_backoff = True
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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down
50 changes: 50 additions & 0 deletions fxsharing/shares/tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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):
Expand Down
Loading