task-queue-system/app/models.py

75 lines
3.8 KiB
Python

import enum, uuid
from datetime import datetime
from sqlalchemy import String, Enum, ForeignKey, Integer, DateTime, JSON
from sqlalchemy import LargeBinary, Index
from sqlalchemy.orm import relationship, Mapped, mapped_column
from sqlalchemy.dialects.postgresql import UUID, TIMESTAMP
from sqlalchemy.sql import func
from app.db import Base
class JobStatus(str, enum.Enum):
PENDING = "pending"
QUEUED = "queued"
RUNNING = "running"
SUCCESS = "success"
FAILED = "failed"
class TaskStatus(str, enum.Enum):
PENDING = "pending"
QUEUED = "queued"
RUNNING = "running"
SUCCESS = "success"
FAILED = "failed"
class Job(Base):
__tablename__ = "jobs"
id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
job_type: Mapped[str] = mapped_column(String(64), index=True)
status: Mapped[JobStatus] = mapped_column(Enum(JobStatus), default=JobStatus.PENDING, index=True)
domain: Mapped[str] = mapped_column(String(255), index=True)
payload: Mapped[dict | None] = mapped_column(JSON, nullable=True)
created_at: Mapped[str] = mapped_column(DateTime(timezone=True), server_default=func.now())
updated_at: Mapped[str] = mapped_column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now())
tasks: Mapped[list["Task"]] = relationship("Task", back_populates="job", cascade="all, delete-orphan", lazy="selectin")
class Task(Base):
__tablename__ = "tasks"
id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
job_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("jobs.id", ondelete="CASCADE"), index=True)
name: Mapped[str] = mapped_column(String(64), index=True)
status: Mapped[TaskStatus] = mapped_column(Enum(TaskStatus), default=TaskStatus.PENDING, index=True)
routing_key: Mapped[str] = mapped_column(String(128), index=True)
retries: Mapped[int] = mapped_column(Integer, default=0)
max_retries: Mapped[int] = mapped_column(Integer, default=3)
worker_id: Mapped[str | None] = mapped_column(String(64), nullable=True)
result: Mapped[dict | None] = mapped_column(JSON, nullable=True)
error: Mapped[str | None] = mapped_column(String(1024), nullable=True)
created_at: Mapped[str] = mapped_column(DateTime(timezone=True), server_default=func.now())
started_at: Mapped[str | None] = mapped_column(DateTime(timezone=True), nullable=True)
finished_at: Mapped[str | None] = mapped_column(DateTime(timezone=True), nullable=True)
job: Mapped["Job"] = relationship("Job", back_populates="tasks")
class TaskDependency(Base):
__tablename__ = "task_dependencies"
task_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("tasks.id", ondelete="CASCADE"), primary_key=True)
depends_on_task_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("tasks.id", ondelete="CASCADE"), primary_key=True)
class Message(Base):
__tablename__ = "messages"
id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
mailbox_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), nullable=False, index=True)
message_id: Mapped[str | None]= mapped_column(String(512), index=True, nullable=True)
subject: Mapped[str | None] = mapped_column(String(2048))
from_addr: Mapped[str | None] = mapped_column(String(2048))
to_addr: Mapped[str | None] = mapped_column(String(4096))
date_header: Mapped[datetime | None] = mapped_column(TIMESTAMP(timezone=True), nullable=True)
received_at: Mapped[datetime] = mapped_column(TIMESTAMP(timezone=True), default=datetime.utcnow, nullable=False)
raw_eml: Mapped[bytes] = mapped_column(LargeBinary, nullable=False)
Index("ix_messages_date", Message.date_header)
Index("ix_messages_message_id", Message.message_id)