Skip to content
Open
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
94 changes: 45 additions & 49 deletions src/forge/queue/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
99 changes: 99 additions & 0 deletions tests/integration/redis/test_retry_recovery.py
Original file line number Diff line number Diff line change
@@ -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
4 changes: 3 additions & 1 deletion tests/unit/queue/test_consumer_retry.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading