Files

44 lines
1.5 KiB
Python

"""event_log 事件日志仓储:幂等消费(event_id 唯一)+ 补拉未消费事件。"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import select, update
from common_const import EVENT_STATUS_CONSUMED, EVENT_STATUS_PENDING
from model.event_log import EventLog
from repositories.base import BaseRepository
class EventLogRepo(BaseRepository):
model = EventLog
async def get_by_event_id(self, event_id: str) -> EventLog | None:
return await self.db.scalar(
select(EventLog).where(EventLog.event_id == event_id)
)
async def list_pending(
self, event_names: tuple[str, ...], limit: int = 100
) -> list[EventLog]:
"""补拉指定事件名中仍未消费的记录(Redis Pub/Sub 丢消息的兜底)。"""
stmt = (
select(EventLog)
.where(
EventLog.event_name.in_(event_names),
EventLog.status == EVENT_STATUS_PENDING,
)
.order_by(EventLog.id.asc())
.limit(limit)
)
return list((await self.db.scalars(stmt)).all())
async def mark_consumed(self, event_id: str) -> None:
"""标记事件已消费(含消费时间),并提交。"""
await self.db.execute(
update(EventLog)
.where(EventLog.event_id == event_id)
.values(status=EVENT_STATUS_CONSUMED, consume_time=datetime.now())
)
await self.db.commit()