Skip to content

Commit 266ce5f

Browse files
authored
Merge pull request #456 from anumukul/feat/issues-301-302-308-309
Add draft persistence, rate-limit headers, webhook dead letters, and …
2 parents e5f57ea + 3f87c53 commit 266ce5f

18 files changed

Lines changed: 887 additions & 145 deletions
Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,26 @@
1+
"""add webhook delivery status tracking
2+
3+
Revision ID: 005
4+
Revises: 004
5+
Create Date: 2026-05-30
6+
"""
7+
from typing import Sequence, Union
8+
9+
from alembic import op
10+
import sqlalchemy as sa
11+
12+
revision: str = "005"
13+
down_revision: Union[str, None] = "004"
14+
branch_labels: Union[str, Sequence[str], None] = None
15+
depends_on: Union[str, Sequence[str], None] = None
16+
17+
18+
def upgrade() -> None:
19+
op.add_column(
20+
"webhook_deliveries",
21+
sa.Column("delivery_status", sa.String(length=32), nullable=False, server_default="pending"),
22+
)
23+
24+
25+
def downgrade() -> None:
26+
op.drop_column("webhook_deliveries", "delivery_status")

backend/src/models.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -184,6 +184,7 @@ class WebhookDelivery(Base):
184184
response_body = Column(Text, nullable=True)
185185
success = Column(Boolean, default=False, nullable=False)
186186
attempts = Column(Integer, default=0, nullable=False)
187+
delivery_status = Column(String(32), default="pending", nullable=False)
187188
last_attempt_at = Column(DateTime, nullable=True)
188189
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
189190

@@ -192,6 +193,7 @@ class WebhookDelivery(Base):
192193
def __init__(self, **kwargs):
193194
kwargs.setdefault("success", False)
194195
kwargs.setdefault("attempts", 0)
196+
kwargs.setdefault("delivery_status", "pending")
195197
super().__init__(**kwargs)
196198

197199
def __repr__(self):

backend/src/rate_limiter.py

Lines changed: 89 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -3,26 +3,33 @@
33
Uses slowapi for IP-based rate limiting with authenticated user bypass support.
44
"""
55
import logging
6+
import re
7+
import time
8+
from typing import Optional
9+
610
from slowapi import Limiter
711
from slowapi.util import get_remote_address
812
from slowapi.errors import RateLimitExceeded
913
from slowapi.middleware import SlowAPIMiddleware
10-
from fastapi import Request, Response
14+
from fastapi import Request
1115
from fastapi.responses import JSONResponse
16+
from starlette.middleware.base import BaseHTTPMiddleware
17+
1218
from .config import get_settings
1319

1420
logger = logging.getLogger(__name__)
1521

1622
settings = get_settings()
1723

24+
RATE_LIMIT_HEADER_NAMES = (
25+
"Retry-After",
26+
"X-RateLimit-Limit",
27+
"X-RateLimit-Remaining",
28+
"X-RateLimit-Reset",
29+
)
30+
1831

1932
def _get_rate_limit_key(request: Request) -> str:
20-
"""
21-
Determine rate limit key based on request.
22-
If authenticated user bypass is enabled and a valid Bearer token is present,
23-
return a user-specific key that receives higher limits.
24-
Otherwise, fall back to IP-based limiting.
25-
"""
2633
if settings.rate_limit_auth_bypass:
2734
auth_header = request.headers.get("Authorization", "")
2835
if auth_header.startswith("Bearer ") and len(auth_header) > 7:
@@ -34,34 +41,99 @@ def _get_rate_limit_key(request: Request) -> str:
3441
return get_remote_address(request)
3542

3643

44+
def _parse_limit_value(limit_detail: str) -> int:
45+
match = re.match(r"(\d+)", limit_detail or "")
46+
return int(match.group(1)) if match else 60
47+
48+
49+
def _parse_retry_after_seconds(limit_detail: str) -> int:
50+
detail = (limit_detail or "").lower()
51+
count_match = re.search(r"(\d+)\s*per", detail)
52+
window_match = re.search(r"per\s+(\d+)\s*(second|minute|hour|day)", detail)
53+
54+
count = int(count_match.group(1)) if count_match else 60
55+
if not window_match:
56+
return max(count, 1)
57+
58+
amount = int(window_match.group(1))
59+
unit = window_match.group(2)
60+
multipliers = {"second": 1, "minute": 60, "hour": 3600, "day": 86400}
61+
window_seconds = amount * multipliers.get(unit, 60)
62+
return max(window_seconds // max(count, 1), 1)
63+
64+
65+
def build_rate_limit_headers(limit_detail: str, remaining: Optional[int] = None) -> dict[str, str]:
66+
limit_value = _parse_limit_value(limit_detail)
67+
retry_after = _parse_retry_after_seconds(limit_detail)
68+
reset_at = int(time.time()) + retry_after
69+
resolved_remaining = str(max(remaining, 0)) if remaining is not None else str(max(limit_value - 1, 0))
70+
71+
return {
72+
"Retry-After": str(retry_after),
73+
"X-RateLimit-Limit": str(limit_value),
74+
"X-RateLimit-Remaining": resolved_remaining,
75+
"X-RateLimit-Reset": str(reset_at),
76+
}
77+
78+
3779
limiter = Limiter(
3880
key_func=_get_rate_limit_key,
3981
default_limits=[settings.rate_limit_default],
4082
storage_uri=settings.redis_url if settings.redis_enabled else "memory://",
83+
headers_enabled=True,
4184
)
4285

4386

4487
def rate_limit_exceeded_handler(request: Request, exc: RateLimitExceeded) -> JSONResponse:
45-
"""Handle rate limit exceeded errors with proper retry information."""
46-
retry_after = exc.detail.split("per")[-1].strip() if "per" in exc.detail else "60"
47-
logger.warning("Rate limit exceeded for %s on %s", _get_rate_limit_key(request), request.url.path)
88+
limit_detail = str(exc.detail)
89+
headers = build_rate_limit_headers(limit_detail, remaining=0)
90+
retry_after = int(headers["Retry-After"])
91+
92+
logger.warning(
93+
"Rate limit exceeded for %s on %s",
94+
_get_rate_limit_key(request),
95+
request.url.path,
96+
)
97+
4898
return JSONResponse(
4999
status_code=429,
50100
content={
51101
"error_code": "RATE_001",
52-
"detail": f"Rate limit exceeded: {exc.detail}",
102+
"detail": f"Rate limit exceeded: {limit_detail}",
53103
"retry_after": retry_after,
54104
},
55-
headers={
56-
"Retry-After": retry_after,
57-
"X-RateLimit-Limit": str(exc.detail),
58-
},
105+
headers=headers,
59106
)
60107

61108

109+
class RateLimitHeaderMiddleware(BaseHTTPMiddleware):
110+
async def dispatch(self, request: Request, call_next):
111+
response = await call_next(request)
112+
113+
if response.status_code == 429:
114+
return response
115+
116+
limit_header = response.headers.get("X-RateLimit-Limit")
117+
if limit_header:
118+
for header_name in RATE_LIMIT_HEADER_NAMES:
119+
if header_name in response.headers:
120+
continue
121+
if "X-RateLimit-Remaining" not in response.headers:
122+
response.headers["X-RateLimit-Remaining"] = str(max(_parse_limit_value(limit_header) - 1, 0))
123+
if "X-RateLimit-Reset" not in response.headers:
124+
retry_after = _parse_retry_after_seconds(settings.rate_limit_default)
125+
response.headers["X-RateLimit-Reset"] = str(int(time.time()) + retry_after)
126+
127+
return response
128+
129+
62130
def setup_rate_limiting(app):
63-
"""Attach rate limiter and exception handler to the FastAPI app."""
64131
app.state.limiter = limiter
132+
app.add_middleware(RateLimitHeaderMiddleware)
65133
app.add_middleware(SlowAPIMiddleware)
66134
app.add_exception_handler(RateLimitExceeded, rate_limit_exceeded_handler)
67-
logger.info("Rate limiting enabled: default=%s, auth=%s", settings.rate_limit_default, settings.rate_limit_auth)
135+
logger.info(
136+
"Rate limiting enabled: default=%s, auth=%s",
137+
settings.rate_limit_default,
138+
settings.rate_limit_auth,
139+
)

backend/src/routes/webhooks.py

Lines changed: 66 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
WebhookUpdateRequest,
1717
WebhookResponse,
1818
WebhookDeliveryResponse,
19+
WebhookDeliveryDetailResponse,
1920
WebhookDeliveryListResponse,
2021
MessageResponse,
2122
)
@@ -38,6 +39,39 @@ def __init__(self, detail: str = "Webhook limit reached for this user"):
3839
super().__init__(status.HTTP_429_TOO_MANY_REQUESTS, detail, "WEBHOOK_003")
3940

4041

42+
class WebhookDeliveryNotFoundError(StellarInsureError):
43+
def __init__(self, detail: str = "Webhook delivery not found"):
44+
super().__init__(status.HTTP_404_NOT_FOUND, detail, "WEBHOOK_004")
45+
46+
47+
def _format_delivery(delivery: WebhookDelivery) -> WebhookDeliveryResponse:
48+
return WebhookDeliveryResponse(
49+
id=delivery.id,
50+
webhook_id=delivery.webhook_id,
51+
event_type=delivery.event_type,
52+
response_status=delivery.response_status,
53+
success=delivery.success,
54+
attempts=delivery.attempts,
55+
delivery_status=delivery.delivery_status,
56+
created_at=delivery.created_at,
57+
)
58+
59+
60+
def _format_delivery_detail(delivery: WebhookDelivery) -> WebhookDeliveryDetailResponse:
61+
return WebhookDeliveryDetailResponse(
62+
id=delivery.id,
63+
webhook_id=delivery.webhook_id,
64+
event_type=delivery.event_type,
65+
response_status=delivery.response_status,
66+
success=delivery.success,
67+
attempts=delivery.attempts,
68+
delivery_status=delivery.delivery_status,
69+
created_at=delivery.created_at,
70+
response_body=delivery.response_body,
71+
last_attempt_at=delivery.last_attempt_at,
72+
)
73+
74+
4175
def _format_webhook(webhook: Webhook) -> WebhookResponse:
4276
return WebhookResponse(
4377
id=webhook.id,
@@ -119,6 +153,37 @@ async def list_webhooks(
119153
return [_format_webhook(w) for w in webhooks]
120154

121155

156+
@router.get(
157+
"/deliveries/{delivery_id}",
158+
response_model=WebhookDeliveryDetailResponse,
159+
summary="Get webhook delivery status",
160+
description="Returns the current delivery state for operator review, including retry and dead-letter status.",
161+
responses={
162+
200: {"description": "Delivery status"},
163+
404: {"description": "Delivery not found"},
164+
401: {"description": "Not authenticated"},
165+
},
166+
)
167+
async def get_webhook_delivery_status(
168+
delivery_id: int,
169+
current_user: User = Depends(get_current_active_user),
170+
db: Session = Depends(get_db),
171+
):
172+
delivery = (
173+
db.query(WebhookDelivery)
174+
.join(Webhook, Webhook.id == WebhookDelivery.webhook_id)
175+
.filter(
176+
WebhookDelivery.id == delivery_id,
177+
Webhook.user_id == current_user.id,
178+
)
179+
.first()
180+
)
181+
if delivery is None:
182+
raise WebhookDeliveryNotFoundError()
183+
184+
return _format_delivery_detail(delivery)
185+
186+
122187
@router.get(
123188
"/{webhook_id}",
124189
response_model=WebhookResponse,
@@ -267,18 +332,7 @@ async def list_webhook_deliveries(
267332
)
268333

269334
return WebhookDeliveryListResponse(
270-
deliveries=[
271-
WebhookDeliveryResponse(
272-
id=d.id,
273-
webhook_id=d.webhook_id,
274-
event_type=d.event_type,
275-
response_status=d.response_status,
276-
success=d.success,
277-
attempts=d.attempts,
278-
created_at=d.created_at,
279-
)
280-
for d in deliveries
281-
],
335+
deliveries=[_format_delivery(d) for d in deliveries],
282336
total=total,
283337
page=page,
284338
per_page=per_page,

backend/src/schemas.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -366,12 +366,18 @@ class WebhookDeliveryResponse(BaseModel):
366366
response_status: Optional[int] = Field(None, description="HTTP response status code")
367367
success: bool = Field(..., description="Whether delivery was successful")
368368
attempts: int = Field(..., description="Number of delivery attempts")
369+
delivery_status: str = Field(..., description="Delivery lifecycle status")
369370
created_at: datetime = Field(..., description="Delivery creation timestamp")
370371

371372
class Config:
372373
from_attributes = True
373374

374375

376+
class WebhookDeliveryDetailResponse(WebhookDeliveryResponse):
377+
response_body: Optional[str] = Field(None, description="Last response body or error message")
378+
last_attempt_at: Optional[datetime] = Field(None, description="Timestamp of the last delivery attempt")
379+
380+
375381
class WebhookDeliveryListResponse(BaseModel):
376382
deliveries: List[WebhookDeliveryResponse] = Field(..., description="List of webhook deliveries")
377383
total: int = Field(..., description="Total number of deliveries", example=100)

0 commit comments

Comments
 (0)