Skip to content

Commit 170ec43

Browse files
UN-3056 [FEAT] Bound buffer redelivery + drop dead provider cluster
- NotificationBuffer.dispatch_attempts + NOTIFICATION_MAX_DISPATCH_ATTEMPTS: _dispatch_group dead-letters rows past the cap and increments on each SENDING claim, bounding the reaper reclaim loop so a lost terminal callback can't redeliver forever (self-review #3). - Delete the orphaned synchronous notification_v2/provider/ cluster — zero callers after the batched dispatch_notifications path replaced it (#2). - Fold dispatch_attempts into 0002_notification_batching; refresh lifecycle db_comments + BufferStatus docstring. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 8b05977 commit 170ec43

12 files changed

Lines changed: 56 additions & 254 deletions

File tree

‎backend/backend/settings/base.py‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -238,6 +238,14 @@ def get_required_setting(setting_key: str, default: str | None = None) -> str |
238238
NOTIFICATION_DISPATCH_LEASE_SECONDS = int(
239239
os.environ.get("NOTIFICATION_DISPATCH_LEASE_SECONDS", "900")
240240
)
241+
# Hard ceiling on how many times a buffer row may be claimed for dispatch. Each
242+
# SENDING claim increments NotificationBuffer.dispatch_attempts; once it reaches
243+
# this cap the row is dead-lettered instead of re-dispatched. Bounds the reaper
244+
# reclaim loop so a row whose terminal callback never fires (e.g. a crash that
245+
# recurs in the dispatch->callback window) cannot be redelivered forever.
246+
NOTIFICATION_MAX_DISPATCH_ATTEMPTS = int(
247+
os.environ.get("NOTIFICATION_MAX_DISPATCH_ATTEMPTS", "5")
248+
)
241249
ATOMIC_REQUESTS = CommonUtils.str_to_bool(
242250
os.environ.get("DJANGO_ATOMIC_REQUESTS", "False")
243251
)

‎backend/notification_v2/enums.py‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,9 @@ class BufferStatus(Enum):
4646
its success/failure callback. Reclaimed to PENDING by the
4747
reaper if it stays here past the dispatch lease (crash window).
4848
DISPATCHED — delivery succeeded.
49-
DEAD_LETTER — Celery exhausted retries; terminal, never re-picked.
49+
DEAD_LETTER — Celery exhausted retries, or the row hit
50+
NOTIFICATION_MAX_DISPATCH_ATTEMPTS reclaim attempts; terminal,
51+
never re-picked.
5052
"""
5153

5254
PENDING = "PENDING"

‎backend/notification_v2/internal_api_views.py‎

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
from api_v2.models import APIDeployment
1818
from django.conf import settings
1919
from django.db import transaction
20-
from django.db.models import Min, QuerySet
20+
from django.db.models import F, Min, QuerySet
2121
from django.http import HttpRequest, JsonResponse
2222
from django.shortcuts import get_object_or_404
2323
from django.utils import timezone
@@ -522,6 +522,29 @@ def _dispatch_group(
522522
# row-level lock. Either way: nothing to do here.
523523
return 0, 0
524524

525+
# Bound the reaper reclaim loop: a row reclaimed past its dispatch budget
526+
# (e.g. a crash that recurs in the dispatch->callback window keeps
527+
# returning it to PENDING) is dead-lettered here rather than re-dispatched
528+
# forever. The cap is checked at claim time so it covers both the reaper
529+
# path and the broker-failure revert path in a single chokepoint.
530+
cap = settings.NOTIFICATION_MAX_DISPATCH_ATTEMPTS
531+
exhausted_ids = [str(r.id) for r in rows if r.dispatch_attempts >= cap]
532+
if exhausted_ids:
533+
NotificationBuffer.objects.filter(id__in=exhausted_ids).update(
534+
status=BufferStatus.DEAD_LETTER.value,
535+
)
536+
logger.warning(
537+
"metric=notification_buffer_dispatch_exhausted_total rows=%d "
538+
"org_id=%s platform=%s cap=%d",
539+
len(exhausted_ids),
540+
org_id,
541+
platform,
542+
cap,
543+
)
544+
rows = [r for r in rows if r.dispatch_attempts < cap]
545+
if not rows:
546+
return 0, 0
547+
525548
# Live auth — read from the FIRST row's notification. If multiple
526549
# notifications collide on (url, auth_sig, platform) we have, by
527550
# definition, identical auth + format, so this is safe. Retry budget
@@ -546,6 +569,7 @@ def _dispatch_group(
546569
NotificationBuffer.objects.filter(id__in=buffer_ids).update(
547570
status=BufferStatus.SENDING.value,
548571
dispatched_at=now,
572+
dispatch_attempts=F("dispatch_attempts") + 1,
549573
)
550574
transaction.on_commit(
551575
lambda: _send_clubbed(

‎backend/notification_v2/migrations/0002_notification_batching.py‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,13 @@ class Migration(migrations.Migration):
7272
),
7373
),
7474
("dispatched_at", models.DateTimeField(blank=True, null=True)),
75+
(
76+
"dispatch_attempts",
77+
models.PositiveIntegerField(
78+
db_comment="Count of times this row has been claimed for dispatch (incremented on each PENDING -> SENDING transition). Bounds the reaper reclaim loop: at NOTIFICATION_MAX_DISPATCH_ATTEMPTS the row is dead-lettered instead of re-dispatched, so a lost terminal callback cannot redeliver forever.",
79+
default=0,
80+
),
81+
),
7582
(
7683
"status",
7784
models.CharField(
@@ -81,7 +88,7 @@ class Migration(migrations.Migration):
8188
("DISPATCHED", "Dispatched"),
8289
("DEAD_LETTER", "Dead letter"),
8390
],
84-
db_comment="Lifecycle: PENDING -> SENDING (claimed by a flush tick) -> DISPATCHED on success / DEAD_LETTER on retry exhaustion. A SENDING row whose lease expires is reclaimed back to PENDING by the reaper.",
91+
db_comment="Lifecycle: PENDING -> SENDING (claimed by a flush tick) -> DISPATCHED on success / DEAD_LETTER on retry exhaustion or once dispatch_attempts hits NOTIFICATION_MAX_DISPATCH_ATTEMPTS. A SENDING row whose lease expires is reclaimed back to PENDING by the reaper.",
8592
default="PENDING",
8693
max_length=16,
8794
),

‎backend/notification_v2/models.py‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -169,13 +169,24 @@ class NotificationBuffer(BaseModel):
169169
),
170170
)
171171
dispatched_at = models.DateTimeField(null=True, blank=True)
172+
dispatch_attempts = models.PositiveIntegerField(
173+
default=0,
174+
db_comment=(
175+
"Count of times this row has been claimed for dispatch (incremented "
176+
"on each PENDING -> SENDING transition). Bounds the reaper reclaim "
177+
"loop: at NOTIFICATION_MAX_DISPATCH_ATTEMPTS the row is dead-lettered "
178+
"instead of re-dispatched, so a lost terminal callback cannot redeliver "
179+
"forever."
180+
),
181+
)
172182
status = models.CharField(
173183
max_length=16,
174184
choices=BufferStatus.choices(),
175185
default=BufferStatus.PENDING.value,
176186
db_comment=(
177187
"Lifecycle: PENDING -> SENDING (claimed by a flush tick) -> "
178-
"DISPATCHED on success / DEAD_LETTER on retry exhaustion. A SENDING "
188+
"DISPATCHED on success / DEAD_LETTER on retry exhaustion or once "
189+
"dispatch_attempts hits NOTIFICATION_MAX_DISPATCH_ATTEMPTS. A SENDING "
179190
"row whose lease expires is reclaimed back to PENDING by the reaper."
180191
),
181192
)

‎backend/notification_v2/provider/__init__.py‎

Whitespace-only changes.

‎backend/notification_v2/provider/notification_provider.py‎

Lines changed: 0 additions & 30 deletions
This file was deleted.

‎backend/notification_v2/provider/registry.py‎

Lines changed: 0 additions & 54 deletions
This file was deleted.

‎backend/notification_v2/provider/webhook/__init__.py‎

Whitespace-only changes.

‎backend/notification_v2/provider/webhook/api_webhook.py‎

Lines changed: 0 additions & 30 deletions
This file was deleted.

0 commit comments

Comments
 (0)