Files
arrr-erp/alembic/versions/20260611_phase_7s2_imap_incoming_email_reading.py
T
2026-06-20 15:01:44 +05:30

114 lines
6.3 KiB
Python

"""Phase 7S.2 IMAP incoming email reading
Revision ID: 20260611_phase_7s2_imap_incoming
Revises: 20260610_phase_7s1e_email_queue
Create Date: 2026-06-11
"""
from __future__ import annotations
from alembic import op
import sqlalchemy as sa
revision = "20260611_phase_7s2_imap_incoming"
down_revision = "20260610_phase_7s1e_email_queue"
branch_labels = None
depends_on = None
def _has_table(table_name: str) -> bool:
return table_name in sa.inspect(op.get_bind()).get_table_names()
def _create_index_if_missing(table_name: str, index_name: str, columns: list[str]) -> None:
existing = {idx["name"] for idx in sa.inspect(op.get_bind()).get_indexes(table_name)} if _has_table(table_name) else set()
if index_name not in existing:
op.create_index(index_name, table_name, columns)
def upgrade() -> None:
if not _has_table("email_incoming_messages"):
op.create_table(
"email_incoming_messages",
sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True),
sa.Column("tenant_id", sa.Integer(), sa.ForeignKey("tenants.id", ondelete="CASCADE"), nullable=True),
sa.Column("branch_id", sa.Integer(), sa.ForeignKey("branches.id", ondelete="SET NULL"), nullable=True),
sa.Column("mailbox_email", sa.String(length=255), nullable=False),
sa.Column("folder_name", sa.String(length=120), nullable=False, server_default="INBOX"),
sa.Column("provider_uid", sa.String(length=120), nullable=False),
sa.Column("provider_message_id", sa.String(length=500), nullable=True),
sa.Column("sender_email", sa.String(length=255), nullable=True),
sa.Column("sender_name", sa.String(length=255), nullable=True),
sa.Column("recipient_emails", sa.Text(), nullable=True),
sa.Column("cc_emails", sa.Text(), nullable=True),
sa.Column("subject", sa.String(length=500), nullable=True),
sa.Column("body_text", sa.Text(), nullable=True),
sa.Column("body_html", sa.Text(), nullable=True),
sa.Column("raw_headers", sa.Text(), nullable=True),
sa.Column("received_at_utc", sa.DateTime(timezone=True), nullable=True),
sa.Column("status", sa.String(length=30), nullable=False, server_default="NEW"),
sa.Column("matched_user_id", sa.Integer(), sa.ForeignKey("users.id", ondelete="SET NULL"), nullable=True),
sa.Column("matched_client_id", sa.Integer(), sa.ForeignKey("clients.id", ondelete="SET NULL"), nullable=True),
sa.Column("matched_consultant_id", sa.Integer(), sa.ForeignKey("consultant_profiles.id", ondelete="SET NULL"), nullable=True),
sa.Column("related_module", sa.String(length=80), nullable=True),
sa.Column("related_id", sa.Integer(), nullable=True),
sa.Column("has_attachments", sa.Boolean(), nullable=False, server_default=sa.false()),
sa.Column("attachment_count", sa.Integer(), nullable=False, server_default="0"),
sa.Column("error_message", sa.Text(), nullable=True),
sa.Column("fetched_at_utc", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("created_at_utc", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.Column("updated_at_utc", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
sa.UniqueConstraint("tenant_id", "branch_id", "mailbox_email", "folder_name", "provider_uid", name="uq_email_incoming_scope_mailbox_folder_uid"),
)
for name, cols in {
"ix_email_incoming_messages_tenant_id": ["tenant_id"],
"ix_email_incoming_messages_branch_id": ["branch_id"],
"ix_email_incoming_messages_mailbox_email": ["mailbox_email"],
"ix_email_incoming_messages_folder_name": ["folder_name"],
"ix_email_incoming_messages_provider_uid": ["provider_uid"],
"ix_email_incoming_messages_provider_message_id": ["provider_message_id"],
"ix_email_incoming_messages_sender_email": ["sender_email"],
"ix_email_incoming_messages_subject": ["subject"],
"ix_email_incoming_messages_received_at_utc": ["received_at_utc"],
"ix_email_incoming_messages_status": ["status"],
"ix_email_incoming_messages_matched_user_id": ["matched_user_id"],
"ix_email_incoming_messages_matched_client_id": ["matched_client_id"],
"ix_email_incoming_messages_matched_consultant_id": ["matched_consultant_id"],
"ix_email_incoming_messages_related_module": ["related_module"],
"ix_email_incoming_messages_related_id": ["related_id"],
"ix_email_incoming_messages_has_attachments": ["has_attachments"],
"ix_email_incoming_messages_fetched_at_utc": ["fetched_at_utc"],
}.items():
_create_index_if_missing("email_incoming_messages", name, cols)
if not _has_table("email_incoming_attachments"):
op.create_table(
"email_incoming_attachments",
sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True),
sa.Column("incoming_message_id", sa.Integer(), sa.ForeignKey("email_incoming_messages.id", ondelete="CASCADE"), nullable=False),
sa.Column("tenant_id", sa.Integer(), sa.ForeignKey("tenants.id", ondelete="CASCADE"), nullable=True),
sa.Column("branch_id", sa.Integer(), sa.ForeignKey("branches.id", ondelete="SET NULL"), nullable=True),
sa.Column("filename", sa.String(length=255), nullable=True),
sa.Column("content_type", sa.String(length=120), nullable=True),
sa.Column("size_bytes", sa.Integer(), nullable=False, server_default="0"),
sa.Column("storage_path", sa.String(length=1000), nullable=True),
sa.Column("created_at_utc", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
)
for name, cols in {
"ix_email_incoming_attachments_incoming_message_id": ["incoming_message_id"],
"ix_email_incoming_attachments_tenant_id": ["tenant_id"],
"ix_email_incoming_attachments_branch_id": ["branch_id"],
}.items():
_create_index_if_missing("email_incoming_attachments", name, cols)
def downgrade() -> None:
try:
op.drop_table("email_incoming_attachments")
except Exception:
pass
try:
op.drop_table("email_incoming_messages")
except Exception:
pass