126 lines
4.3 KiB
Python
126 lines
4.3 KiB
Python
"""SQLAlchemy model for the transactional event outbox table."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import uuid
|
|
from datetime import datetime
|
|
|
|
from sqlalchemy import DateTime, ForeignKey, Integer, String, Text, UniqueConstraint, 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,
|
|
)
|
|
|
|
|
|
class OutboxDelivery(Base):
|
|
"""Per-consumer delivery status for an event_outbox row (Migration 0075).
|
|
|
|
Tracks whether each consumer successfully processed an event; an event is
|
|
only 'published' when all mandatory deliveries succeed.
|
|
"""
|
|
|
|
__tablename__ = "outbox_deliveries"
|
|
__table_args__ = (
|
|
UniqueConstraint(
|
|
"event_id", "consumer_name",
|
|
name="uq_outbox_deliveries_event_consumer",
|
|
),
|
|
)
|
|
|
|
id: Mapped[uuid.UUID] = mapped_column(
|
|
PGUUID(as_uuid=True), primary_key=True,
|
|
server_default=func.gen_random_uuid(),
|
|
)
|
|
event_id: Mapped[uuid.UUID] = mapped_column(
|
|
PGUUID(as_uuid=True),
|
|
ForeignKey("event_outbox.id", ondelete="CASCADE"),
|
|
nullable=False,
|
|
)
|
|
consumer_name: Mapped[str] = mapped_column(String(150), nullable=False)
|
|
status: Mapped[str] = mapped_column(
|
|
String(30), nullable=False, server_default="pending",
|
|
)
|
|
attempt_count: Mapped[int] = mapped_column(
|
|
Integer, nullable=False, server_default="0",
|
|
)
|
|
next_attempt_at: Mapped[datetime | None] = mapped_column(
|
|
DateTime(timezone=True), nullable=True,
|
|
)
|
|
last_error: Mapped[str | None] = mapped_column(Text, nullable=True)
|
|
processed_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(),
|
|
)
|