From efeb866cdf45f011b86e3ef822040ae99504b7d6 Mon Sep 17 00:00:00 2001 From: Matheus-TecDev Date: Fri, 11 Sep 2026 23:29:42 -0300 Subject: [PATCH] fix: align DLQ reprocessing contract with outbox flow --- backend/app/api/routes/dead_letter_events.py | 6 ------ backend/app/services/dead_letter_service.py | 4 ---- docs/api.md | 9 ++++++--- docs/architecture.md | 8 +++++--- 4 files changed, 11 insertions(+), 16 deletions(-) diff --git a/backend/app/api/routes/dead_letter_events.py b/backend/app/api/routes/dead_letter_events.py index 0155053..e041820 100644 --- a/backend/app/api/routes/dead_letter_events.py +++ b/backend/app/api/routes/dead_letter_events.py @@ -10,7 +10,6 @@ ) from app.services.dead_letter_service import ( DeadLetterEventNotFoundError, - DeadLetterPublishError, DeadLetterService, UnsafeReprocessError, ) @@ -81,11 +80,6 @@ def reprocess_dead_letter_event( raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Dead letter event not found.") from exc except UnsafeReprocessError as exc: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(exc)) from exc - except DeadLetterPublishError as exc: - raise HTTPException( - status_code=status.HTTP_503_SERVICE_UNAVAILABLE, - detail="Dead letter event was not republished.", - ) from exc event = dead_letter_event.event routing_key = dead_letter_event.original_routing_key or event.routing_key or "" diff --git a/backend/app/services/dead_letter_service.py b/backend/app/services/dead_letter_service.py index 6858078..23b0964 100644 --- a/backend/app/services/dead_letter_service.py +++ b/backend/app/services/dead_letter_service.py @@ -27,10 +27,6 @@ class UnsafeReprocessError(Exception): pass -class DeadLetterPublishError(Exception): - pass - - class DeadLetterService: def __init__(self, db: Session) -> None: self.db = db diff --git a/docs/api.md b/docs/api.md index fe88c5b..4b56ff5 100644 --- a/docs/api.md +++ b/docs/api.md @@ -211,7 +211,9 @@ Possible errors: ### `POST /api/dead-letter-events/{id}/reprocess` -Republishes the original event to the main exchange with its original routing key. Relay preserves `correlation_id` and `trace_id`, moves the original event back to `queued`, and records the operational action. +Queues the original event for reprocessing through the Transactional Outbox. Relay reuses the existing `Event`, preserves `correlation_id` and `trace_id`, resets the processing state to `pending`, moves the event back to `queued`, creates or resets its `OutboxMessage` to `pending`, and records the operational action. + +The request returns after the database transaction is persisted. RabbitMQ publication is asynchronous and performed later by the Outbox publisher, so `200` means the event was queued in the Outbox, not that broker delivery was already confirmed. Later publication failures remain retryable by the Outbox retry flow. ```json { @@ -228,8 +230,7 @@ Possible errors: - `401`: missing, invalid, or expired token; - `404`: DLQ event not found; -- `409`: reprocessing blocked by an operational safety rule; -- `503`: event could not be republished. +- `409`: reprocessing blocked by an operational safety rule. ## Reprocessing Semantics @@ -237,4 +238,6 @@ Possible errors: - No new `Event` is created. - The routing key comes from `original_routing_key`, falling back to `event.routing_key`. - `correlation_id` and `trace_id` are preserved. +- `EventProcessingState` is reset to `pending`. +- `OutboxMessage` is created or reset to `pending`. - Reprocessing does not replace real handler idempotency. diff --git a/docs/architecture.md b/docs/architecture.md index 13864ca..6802527 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -145,7 +145,9 @@ Shared retry queues return messages through `retry.*`. Workers inspect `x-origin ## Retry and Backoff -Workers avoid `basic_nack(requeue=True)`, preventing tight loops in the main queue. On failure, the original message is republished to the DLX with `x-retry-count` and the appropriate retry routing key, then acknowledged. +On normal handler failure, workers avoid `basic_nack(requeue=True)` to prevent tight loops in the main queue. The original message is published to the DLX with `x-retry-count` and the appropriate retry routing key. The original message is acknowledged only after the recovery message is published successfully and the corresponding state is persisted. + +If publishing the retry or DLQ message fails, the worker does not persist the state as sent. It issues `basic_nack(requeue=True)`, calls `stop_consuming()`, and leaves the original message available for redelivery after the worker or messaging infrastructure recovers. Failure progression: @@ -162,9 +164,9 @@ The operational DLQ stores messages that exhausted automated retries. PostgreSQL - `GET /api/dead-letter-events`: list DLQ events with failure and correlation data; - `GET /api/dead-letter-events/{id}`: inspect payload, original event, attempts, and logs; -- `POST /api/dead-letter-events/{id}/reprocess`: republish the existing event. +- `POST /api/dead-letter-events/{id}/reprocess`: queue the existing event for Outbox-backed reprocessing. -Manual reprocessing preserves `correlation_id`, `trace_id`, and `original_routing_key`, moves the original event to `queued`, and creates an operational `EventLog`. It does not create a second `Event`. +Manual reprocessing preserves `correlation_id`, `trace_id`, and `original_routing_key`, resets the processing state, moves the original event to `queued`, creates or resets an `OutboxMessage` to `pending`, and creates an operational `EventLog`. It does not create a second `Event` or publish directly inside the HTTP request. Recommended operation: