Files

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(),
)