"""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