diff --git a/src/forge/queue/consumer.py b/src/forge/queue/consumer.py index f86e0cd1..cd1cb101 100644 --- a/src/forge/queue/consumer.py +++ b/src/forge/queue/consumer.py @@ -240,59 +240,55 @@ async def _consume_stream(self, stream: str, _source: EventSource) -> None: if self._active_tasks: await asyncio.gather(*self._active_tasks, return_exceptions=True) - async def _process_retry_queue(self) -> None: - """Poll the retry queue and re-dispatch due messages. + async def _process_due_retries_once(self) -> None: + """Dispatch one batch of due retries. - Runs as a background task alongside the stream consumers. Polls on a - fixed interval so retries are dispatched once their backoff window has - elapsed. + Kept separate from the polling loop so recovery and terminal DLQ + behavior can be exercised deterministically in integration tests. """ + entries = await self._retry_queue.get_due_messages() + for entry in entries: + retry_stream = ( + JIRA_STREAM if entry.message.source == EventSource.JIRA else GITHUB_STREAM + ) + try: + await self._process_message( + entry.message, retry_stream, raise_on_error=True, skip_ack=True + ) + except Exception as e: + logger.warning( + f"Retry attempt {entry.attempt} failed for " + f"{entry.message.ticket_key}:{entry.message.event_id}: {e}" + ) + # Preserve the counter across re-enqueues so exhaustion can + # eventually move the message to the dead-letter queue. + await self._retry_queue.remove_from_retry_without_counter_reset(entry) + queued = await self._retry_queue.enqueue_for_retry(entry.message, str(e)) + if not queued: + # Terminal failure: the DLQ now owns the message, so clear + # its original stream PEL entry. + await self._ack(retry_stream, entry.message.message_id) + continue + + # Success: clear retry state and acknowledge the original entry. + await self._retry_queue.remove_from_retry(entry) + logger.info( + f"Retry succeeded for {entry.message.ticket_key}:" + f"{entry.message.event_id} (attempt {entry.attempt})" + ) + try: + await self._ack(retry_stream, entry.message.message_id) + except Exception as xack_err: + logger.warning( + f"xack failed for {entry.message.event_id} after successful " + f"retry (PEL entry may linger): {xack_err}" + ) + + async def _process_retry_queue(self) -> None: + """Poll the retry queue and re-dispatch due messages.""" while self._running: try: - entries = await self._retry_queue.get_due_messages() - for entry in entries: - retry_stream = ( - JIRA_STREAM if entry.message.source == EventSource.JIRA else GITHUB_STREAM - ) - try: - await self._process_message( - entry.message, retry_stream, raise_on_error=True, skip_ack=True - ) - except Exception as e: - logger.warning( - f"Retry attempt {entry.attempt} failed for " - f"{entry.message.ticket_key}:{entry.message.event_id}: {e}" - ) - # Remove only the sorted-set entry — do NOT delete the attempt - # counter key. enqueue_for_retry will INCR the existing key so - # the counter keeps accumulating and the message can eventually - # reach the dead-letter queue. - await self._retry_queue.remove_from_retry_without_counter_reset(entry) - await self._retry_queue.enqueue_for_retry(entry.message, str(e)) - continue - - # Message processing succeeded — clean up retry state and - # acknowledge the original stream entry (best-effort). - await self._retry_queue.remove_from_retry(entry) - logger.info( - f"Retry succeeded for {entry.message.ticket_key}:" - f"{entry.message.event_id} (attempt {entry.attempt})" - ) - stream = ( - JIRA_STREAM if entry.message.source == EventSource.JIRA else GITHUB_STREAM - ) - try: - redis_client = await self._get_redis() - await redis_client.xack(stream, CONSUMER_GROUP, entry.message.message_id) - except Exception as xack_err: - # xack is best-effort: the message was already successfully - # processed and removed from the retry queue. A failure here - # only means the PEL entry lingers; it will not be reprocessed - # because the retry-queue entry is gone. - logger.warning( - f"xack failed for {entry.message.event_id} after successful " - f"retry (PEL entry may linger): {xack_err}" - ) + await self._process_due_retries_once() except asyncio.CancelledError: break except Exception as e: diff --git a/tests/integration/redis/test_retry_recovery.py b/tests/integration/redis/test_retry_recovery.py new file mode 100644 index 00000000..a221216c --- /dev/null +++ b/tests/integration/redis/test_retry_recovery.py @@ -0,0 +1,99 @@ +"""Real-Redis integration tests for consumer retry and dead-letter recovery.""" + +from forge.models.events import EventSource +from forge.queue.consumer import CONSUMER_GROUP, QueueConsumer +from forge.queue.models import QueueMessage +from forge.queue.producer import JIRA_STREAM, QueueProducer +from forge.queue.retry import RETRY_QUEUE_KEY + + +async def _claim_one(redis_client, consumer: QueueConsumer) -> QueueMessage: + await consumer._ensure_consumer_groups() + entries = await redis_client.xreadgroup( + CONSUMER_GROUP, + consumer.consumer_name, + {JIRA_STREAM: ">"}, + count=1, + block=1000, + ) + message_id, fields = entries[0][1][0] + return QueueMessage.from_redis(message_id, fields) + + +async def _make_retries_due(redis_client) -> None: + members = await redis_client.zrange(RETRY_QUEUE_KEY, 0, -1) + if members: + await redis_client.zadd(RETRY_QUEUE_KEY, dict.fromkeys(members, 0)) + + +async def _pending_count(redis_client) -> int: + pending = await redis_client.xpending(JIRA_STREAM, CONSUMER_GROUP) + return int(pending["pending"]) + + +async def test_failed_message_succeeds_on_retry_and_clears_pending_state(redis_client) -> None: + producer = QueueProducer(redis_client=redis_client) + consumer = QueueConsumer("retry-success", redis_client=redis_client) + consumer._retry_queue._redis = redis_client + calls = 0 + + async def fail_once(_message: QueueMessage) -> None: + nonlocal calls + calls += 1 + if calls == 1: + raise RuntimeError("temporary failure") + + consumer.register_handler(EventSource.JIRA, fail_once) + await producer.publish( + event_id="retry-success-1", + source=EventSource.JIRA, + event_type="issue_updated", + ticket_key="TEST-RETRY", + payload={}, + ) + message = await _claim_one(redis_client, consumer) + + await consumer._process_message(message, JIRA_STREAM) + assert await _pending_count(redis_client) == 1 + assert (await consumer._retry_queue.get_queue_stats())["retry_queue_depth"] == 1 + + await _make_retries_due(redis_client) + await consumer._process_due_retries_once() + + assert calls == 2 + assert await _pending_count(redis_client) == 0 + assert await consumer._retry_queue.get_queue_stats() == { + "retry_queue_depth": 0, + "dead_letter_depth": 0, + } + + +async def test_exhausted_message_moves_to_dlq_and_clears_pending_state(redis_client) -> None: + producer = QueueProducer(redis_client=redis_client) + consumer = QueueConsumer("retry-dlq", redis_client=redis_client) + consumer._retry_queue._redis = redis_client + + async def always_fail(_message: QueueMessage) -> None: + raise RuntimeError("persistent failure") + + consumer.register_handler(EventSource.JIRA, always_fail) + await producer.publish( + event_id="retry-dlq-1", + source=EventSource.JIRA, + event_type="issue_updated", + ticket_key="TEST-DLQ", + payload={}, + ) + message = await _claim_one(redis_client, consumer) + await consumer._process_message(message, JIRA_STREAM) + + for _ in range(3): + await _make_retries_due(redis_client) + await consumer._process_due_retries_once() + + stats = await consumer._retry_queue.get_queue_stats() + dead_letters = await consumer._retry_queue.get_dead_letter_entries() + assert stats == {"retry_queue_depth": 0, "dead_letter_depth": 1} + assert dead_letters[0]["message"]["event_id"] == "retry-dlq-1" + assert dead_letters[0]["attempts"] == 4 + assert await _pending_count(redis_client) == 0 diff --git a/tests/unit/queue/test_consumer_retry.py b/tests/unit/queue/test_consumer_retry.py index f854a921..ae50c727 100644 --- a/tests/unit/queue/test_consumer_retry.py +++ b/tests/unit/queue/test_consumer_retry.py @@ -31,7 +31,9 @@ def make_message( def make_consumer() -> QueueConsumer: """Return a QueueConsumer with a mocked RetryQueue.""" - consumer = QueueConsumer(consumer_name="test-worker") + redis_mock = MagicMock() + redis_mock.xack = AsyncMock(return_value=1) + consumer = QueueConsumer(consumer_name="test-worker", redis_client=redis_mock) consumer._running = True # Replace the real RetryQueue with a mock consumer._retry_queue = MagicMock(spec=RetryQueue)