Files
leocrm/app/models/outbox.py
T
Agent Zero 07a99975ec Phase 5: Outbox DLQ, Monitoring, Consumer-Registry
- Migration 0092: DLQ columns (error_message, failed_at) + consumer_inbox RLS fix
- outbox.py: DLQ logic, replay functions, stats, consumer registry
- app/routes/outbox.py: 5 API endpoints (stats, failed, replay, replay-all, consumer-registry)
- outbox_deliveries tracking per consumer handler
- 18/18 tests passing
2026-08-02 23:25:54 +02:00

80 lines
2.7 KiB
Python

"""SQLAlchemy model for the transactional event outbox table."""
from __future__ import annotations
import uuid
from datetime import datetime
from sqlalchemy import DateTime, Integer, String, Text, func
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.dialects.postgresql import UUID as PGUUID
from sqlalchemy.orm import Mapped, mapped_column
from app.core.db import Base
class EventOutbox(Base):
"""Row in the ``event_outbox`` table.
Each row represents a domain event that was written within a business
transaction and is waiting to be published to the in-process event bus
by the outbox worker.
"""
__tablename__ = "event_outbox"
id: Mapped[uuid.UUID] = mapped_column(
PGUUID(as_uuid=True),
primary_key=True,
server_default=func.gen_random_uuid(),
)
tenant_id: Mapped[uuid.UUID] = mapped_column(
PGUUID(as_uuid=True), nullable=False, index=True,
)
event_name: Mapped[str] = mapped_column(String(255), nullable=False)
payload: Mapped[dict] = mapped_column(JSONB, nullable=False)
status: Mapped[str] = mapped_column(
String(20), nullable=False, server_default="pending",
)
attempts: Mapped[int] = mapped_column(
Integer, nullable=False, server_default="0",
)
max_attempts: Mapped[int] = mapped_column(
Integer, nullable=False, server_default="5",
)
next_retry_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), nullable=True,
)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now(),
)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now(),
)
published_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), nullable=True,
)
# Envelope columns (migration 0075)
aggregate_type: Mapped[str | None] = mapped_column(
String(100), nullable=True,
)
aggregate_id: Mapped[uuid.UUID | None] = mapped_column(
PGUUID(as_uuid=True), nullable=True,
)
occurred_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, server_default=func.now(),
)
correlation_id: Mapped[uuid.UUID | None] = mapped_column(
PGUUID(as_uuid=True), nullable=True,
)
schema_version: Mapped[int] = mapped_column(
Integer, nullable=False, server_default="1",
)
# Phase 5: DLQ columns
error_message: Mapped[str | None] = mapped_column(
Text, nullable=True,
)
failed_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), nullable=True,
)