-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdlq_repository.py
More file actions
151 lines (135 loc) · 6.11 KB
/
Copy pathdlq_repository.py
File metadata and controls
151 lines (135 loc) · 6.11 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
from datetime import datetime, timezone
from typing import Any
import structlog
from beanie.odm.enums import SortDirection
from beanie.operators import Set
from monggregate import Pipeline, S
from app.db.docs import DLQMessageDocument
from app.dlq import (
DLQMessage,
DLQMessageListResult,
DLQMessageStatus,
DLQMessageUpdate,
)
from app.domain.enums import EventType
from app.schemas_pydantic.dlq import DLQTopicSummary
class DLQRepository:
def __init__(self, logger: structlog.stdlib.BoundLogger):
self.logger = logger
async def get_messages(
self,
status: DLQMessageStatus | None = None,
topic: str | None = None,
event_type: EventType | None = None,
limit: int = 50,
offset: int = 0,
) -> DLQMessageListResult:
conditions: list[Any] = [
DLQMessageDocument.status == status if status else None,
DLQMessageDocument.original_topic == topic if topic else None,
DLQMessageDocument.event.event_type == event_type if event_type else None,
]
conditions = [c for c in conditions if c is not None]
query = DLQMessageDocument.find(*conditions)
total_count = await query.count()
docs = await query.sort([("failed_at", SortDirection.DESCENDING)]).skip(offset).limit(limit).to_list()
return DLQMessageListResult(
messages=[DLQMessage.model_validate(d) for d in docs],
total=total_count,
offset=offset,
limit=limit,
)
async def get_message_by_id(self, event_id: str) -> DLQMessage | None:
doc = await DLQMessageDocument.find_one({"event.event_id": event_id})
return DLQMessage.model_validate(doc) if doc else None
async def get_topics_summary(self) -> list[DLQTopicSummary]:
# Two-stage aggregation: group by topic+status first, then by topic with $arrayToObject
# Note: compound keys need S.field() wrapper for monggregate to add $ prefix
pipeline = (
Pipeline()
.group(
by={"topic": S.field(DLQMessageDocument.original_topic), "status": S.field(DLQMessageDocument.status)},
query={
"count": S.sum(1),
"oldest": S.min(S.field(DLQMessageDocument.failed_at)),
"newest": S.max(S.field(DLQMessageDocument.failed_at)),
"sum_retry": S.sum(S.field(DLQMessageDocument.retry_count)),
"max_retry": S.max(S.field(DLQMessageDocument.retry_count)),
},
)
.group(
by="$_id.topic",
query={
"status_pairs": S.push({"k": "$_id.status", "v": "$count"}),
"total_messages": S.sum("$count"),
"oldest_message": S.min("$oldest"),
"newest_message": S.max("$newest"),
"total_retry": S.sum("$sum_retry"),
"doc_count": S.sum("$count"),
"max_retry_count": S.max("$max_retry"),
},
)
.sort(by="total_messages", descending=True)
.project(
_id=0,
topic="$_id",
total_messages=1,
status_breakdown={"$arrayToObject": "$status_pairs"},
oldest_message=1,
newest_message=1,
avg_retry_count={"$round": [{"$divide": ["$total_retry", "$doc_count"]}, 2]},
max_retry_count=1,
)
)
results = await DLQMessageDocument.aggregate(pipeline.export()).to_list()
return [DLQTopicSummary.model_validate(r) for r in results]
async def save_message(self, message: DLQMessage) -> None:
"""Upsert a DLQ message by event_id (atomic, no TOCTOU race)."""
payload = message.model_dump()
await DLQMessageDocument.find_one({"event.event_id": message.event.event_id}).upsert(
Set(payload), # type: ignore[no-untyped-call]
on_insert=DLQMessageDocument(**payload),
)
async def update_status(self, event_id: str, update: DLQMessageUpdate) -> None:
"""Apply a status update to a DLQ message."""
doc = await DLQMessageDocument.find_one({"event.event_id": event_id})
if not doc:
return
update_dict: dict[str, Any] = {"status": update.status, "last_updated": datetime.now(timezone.utc)}
if update.next_retry_at is not None:
update_dict["next_retry_at"] = update.next_retry_at
if update.retried_at is not None:
update_dict["retried_at"] = update.retried_at
if update.discarded_at is not None:
update_dict["discarded_at"] = update.discarded_at
if update.retry_count is not None:
update_dict["retry_count"] = update.retry_count
if update.discard_reason is not None:
update_dict["discard_reason"] = update.discard_reason
if update.last_error is not None:
update_dict["last_error"] = update.last_error
await doc.set(update_dict)
async def find_due_retries(self, limit: int = 100) -> list[DLQMessage]:
"""Find scheduled messages whose retry time has arrived."""
now = datetime.now(timezone.utc)
docs = (
await DLQMessageDocument.find(
{
"status": DLQMessageStatus.SCHEDULED,
"next_retry_at": {"$lte": now},
}
)
.limit(limit)
.to_list()
)
return [DLQMessage.model_validate(doc) for doc in docs]
async def get_queue_sizes_by_topic(self) -> dict[str, int]:
"""Get message counts per topic for active (pending/scheduled) messages."""
pipeline: list[dict[str, Any]] = [
{"$match": {"status": {"$in": [DLQMessageStatus.PENDING, DLQMessageStatus.SCHEDULED]}}},
{"$group": {"_id": "$original_topic", "count": {"$sum": 1}}},
]
result: dict[str, int] = {}
async for row in DLQMessageDocument.aggregate(pipeline):
result[row["_id"]] = row["count"]
return result