from __future__ import annotations import threading import time from sqlalchemy import inspect, text from sqlalchemy.exc import OperationalError from services.shared.db import engine from services.shared.schema_migrations import schema_management_mode, validate_schema_migrations_applied from services.shared.sql_models import Base _SCHEMA_INIT_LOCK = threading.Lock() _BASE_SCHEMA_INITIALIZED = False def _table_columns(inspector, table_name: str) -> set[str]: return {item["name"] for item in inspector.get_columns(table_name)} def _table_indexes(inspector, table_name: str) -> set[str]: return {item["name"] for item in inspector.get_indexes(table_name)} def _add_column_if_missing(conn, columns: set[str], table_name: str, column_name: str, ddl: str) -> None: if column_name in columns: return conn.execute(text(f"ALTER TABLE {table_name} ADD COLUMN {column_name} {ddl}")) columns.add(column_name) def _apply_runtime_schema_compatibility() -> None: inspector = inspect(engine) table_names = set(inspector.get_table_names()) with engine.begin() as conn: if "telegram_messages" in table_names: columns = _table_columns(inspector, "telegram_messages") _add_column_if_missing(conn, columns, "telegram_messages", "thread_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "telegram_messages", "interaction_id", "VARCHAR(64)") _add_column_if_missing( conn, columns, "telegram_messages", "direction", "VARCHAR(16) DEFAULT 'inbound'", ) _add_column_if_missing( conn, columns, "telegram_messages", "telegram_message_id_external", "VARCHAR(128)", ) _add_column_if_missing(conn, columns, "telegram_messages", "operator_user", "VARCHAR(128)") _add_column_if_missing(conn, columns, "telegram_messages", "delivery_status", "VARCHAR(32)") _add_column_if_missing(conn, columns, "telegram_messages", "delivery_attempts", "INTEGER DEFAULT 0") _add_column_if_missing( conn, columns, "telegram_messages", "next_delivery_attempt_at", "VARCHAR(64)", ) _add_column_if_missing(conn, columns, "telegram_messages", "delivery_locked_until", "VARCHAR(64)") _add_column_if_missing(conn, columns, "telegram_messages", "last_delivery_error", "TEXT") if "author_type" not in columns: _add_column_if_missing( conn, columns, "telegram_messages", "author_type", "VARCHAR(32) DEFAULT 'customer'", ) if "direction" in columns: conn.execute( text( """ UPDATE telegram_messages SET author_type = CASE WHEN direction = 'outbound' THEN 'human' WHEN direction = 'system' THEN 'system' ELSE 'customer' END WHERE author_type IS NULL OR TRIM(author_type) = '' """ ) ) if "author_id" not in columns: _add_column_if_missing(conn, columns, "telegram_messages", "author_id", "VARCHAR(128)") if "operator_user" in columns: conn.execute( text( """ UPDATE telegram_messages SET author_id = operator_user WHERE operator_user IS NOT NULL AND ( direction = 'outbound' OR direction IS NULL OR TRIM(direction) = '' ) AND (author_id IS NULL OR TRIM(author_id) = '') """ ) ) indexes = _table_indexes(inspector, "telegram_messages") if "idx_telegram_messages_thread_id" not in indexes: conn.execute( text("CREATE INDEX IF NOT EXISTS idx_telegram_messages_thread_id ON telegram_messages(thread_id)") ) if "idx_telegram_messages_interaction_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_messages_interaction_id ON telegram_messages(interaction_id)" ) ) if "idx_telegram_messages_direction" not in indexes: conn.execute( text("CREATE INDEX IF NOT EXISTS idx_telegram_messages_direction ON telegram_messages(direction)") ) if "idx_telegram_messages_telegram_message_id_external" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_messages_telegram_message_id_external " "ON telegram_messages(telegram_message_id_external)" ) ) if "idx_telegram_messages_operator_user" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_messages_operator_user ON telegram_messages(operator_user)" ) ) if "idx_telegram_messages_delivery_status" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_messages_delivery_status ON telegram_messages(delivery_status)" ) ) if "idx_telegram_messages_next_delivery_attempt_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_messages_next_delivery_attempt_at " "ON telegram_messages(next_delivery_attempt_at)" ) ) if "idx_telegram_messages_delivery_locked_until" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_messages_delivery_locked_until " "ON telegram_messages(delivery_locked_until)" ) ) if "idx_telegram_messages_author_type" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_messages_author_type ON telegram_messages(author_type)" ) ) if "idx_telegram_messages_author_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_messages_author_id ON telegram_messages(author_id)" ) ) if "telegram_threads" in table_names: columns = _table_columns(inspector, "telegram_threads") _add_column_if_missing(conn, columns, "telegram_threads", "ai_session_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "telegram_threads", "ai_state", "VARCHAR(32)") _add_column_if_missing(conn, columns, "telegram_threads", "ai_handoff_reason", "TEXT") _add_column_if_missing(conn, columns, "telegram_threads", "ai_last_model_at", "VARCHAR(64)") indexes = _table_indexes(inspector, "telegram_threads") if "idx_telegram_threads_ai_session_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_threads_ai_session_id ON telegram_threads(ai_session_id)" ) ) if "idx_telegram_threads_ai_state" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_threads_ai_state ON telegram_threads(ai_state)" ) ) if "idx_telegram_threads_ai_last_model_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_telegram_threads_ai_last_model_at ON telegram_threads(ai_last_model_at)" ) ) if "whatsapp_messages" in table_names: columns = _table_columns(inspector, "whatsapp_messages") _add_column_if_missing(conn, columns, "whatsapp_messages", "chat_id", "VARCHAR(128)") _add_column_if_missing(conn, columns, "whatsapp_messages", "thread_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "whatsapp_messages", "interaction_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "whatsapp_messages", "customer_id", "VARCHAR(64)") _add_column_if_missing( conn, columns, "whatsapp_messages", "direction", "VARCHAR(16) DEFAULT 'inbound'", ) _add_column_if_missing( conn, columns, "whatsapp_messages", "whatsapp_message_id_external", "VARCHAR(128)", ) _add_column_if_missing(conn, columns, "whatsapp_messages", "operator_user", "VARCHAR(128)") _add_column_if_missing(conn, columns, "whatsapp_messages", "delivery_status", "VARCHAR(32)") _add_column_if_missing(conn, columns, "whatsapp_messages", "delivery_attempts", "INTEGER DEFAULT 0") _add_column_if_missing( conn, columns, "whatsapp_messages", "next_delivery_attempt_at", "VARCHAR(64)", ) _add_column_if_missing(conn, columns, "whatsapp_messages", "delivery_locked_until", "VARCHAR(64)") _add_column_if_missing(conn, columns, "whatsapp_messages", "last_delivery_error", "TEXT") if "author_type" not in columns: _add_column_if_missing( conn, columns, "whatsapp_messages", "author_type", "VARCHAR(32) DEFAULT 'customer'", ) if "direction" in columns: conn.execute( text( """ UPDATE whatsapp_messages SET author_type = CASE WHEN direction = 'outbound' THEN 'human' WHEN direction = 'system' THEN 'system' ELSE 'customer' END WHERE author_type IS NULL OR TRIM(author_type) = '' """ ) ) if "author_id" not in columns: _add_column_if_missing(conn, columns, "whatsapp_messages", "author_id", "VARCHAR(128)") if "operator_user" in columns: conn.execute( text( """ UPDATE whatsapp_messages SET author_id = operator_user WHERE operator_user IS NOT NULL AND ( direction = 'outbound' OR direction IS NULL OR TRIM(direction) = '' ) AND (author_id IS NULL OR TRIM(author_id) = '') """ ) ) indexes = _table_indexes(inspector, "whatsapp_messages") if "ix_whatsapp_messages_chat_external_unique" not in indexes: conn.execute( text( "CREATE UNIQUE INDEX IF NOT EXISTS ix_whatsapp_messages_chat_external_unique " "ON whatsapp_messages(chat_id, whatsapp_message_id_external)" ) ) if "idx_whatsapp_messages_thread_id" not in indexes: conn.execute( text("CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_thread_id ON whatsapp_messages(thread_id)") ) if "idx_whatsapp_messages_interaction_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_interaction_id ON whatsapp_messages(interaction_id)" ) ) if "idx_whatsapp_messages_customer_id" not in indexes: conn.execute( text("CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_customer_id ON whatsapp_messages(customer_id)") ) if "idx_whatsapp_messages_direction" not in indexes: conn.execute( text("CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_direction ON whatsapp_messages(direction)") ) if "idx_whatsapp_messages_whatsapp_message_id_external" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_whatsapp_message_id_external " "ON whatsapp_messages(whatsapp_message_id_external)" ) ) if "idx_whatsapp_messages_operator_user" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_operator_user ON whatsapp_messages(operator_user)" ) ) if "idx_whatsapp_messages_delivery_status" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_delivery_status ON whatsapp_messages(delivery_status)" ) ) if "idx_whatsapp_messages_next_delivery_attempt_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_next_delivery_attempt_at " "ON whatsapp_messages(next_delivery_attempt_at)" ) ) if "idx_whatsapp_messages_delivery_locked_until" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_delivery_locked_until " "ON whatsapp_messages(delivery_locked_until)" ) ) if "idx_whatsapp_messages_author_type" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_author_type ON whatsapp_messages(author_type)" ) ) if "idx_whatsapp_messages_author_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_messages_author_id ON whatsapp_messages(author_id)" ) ) if "whatsapp_threads" in table_names: columns = _table_columns(inspector, "whatsapp_threads") _add_column_if_missing(conn, columns, "whatsapp_threads", "phone_number", "VARCHAR(64)") _add_column_if_missing(conn, columns, "whatsapp_threads", "is_group", "BOOLEAN DEFAULT 0") _add_column_if_missing(conn, columns, "whatsapp_threads", "ai_session_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "whatsapp_threads", "ai_state", "VARCHAR(32)") _add_column_if_missing(conn, columns, "whatsapp_threads", "ai_handoff_reason", "TEXT") _add_column_if_missing(conn, columns, "whatsapp_threads", "ai_last_model_at", "VARCHAR(64)") indexes = _table_indexes(inspector, "whatsapp_threads") if "idx_whatsapp_threads_phone_number" not in indexes: conn.execute( text("CREATE INDEX IF NOT EXISTS idx_whatsapp_threads_phone_number ON whatsapp_threads(phone_number)") ) if "idx_whatsapp_threads_is_group" not in indexes: conn.execute( text("CREATE INDEX IF NOT EXISTS idx_whatsapp_threads_is_group ON whatsapp_threads(is_group)") ) if "idx_whatsapp_threads_ai_session_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_threads_ai_session_id ON whatsapp_threads(ai_session_id)" ) ) if "idx_whatsapp_threads_ai_state" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_threads_ai_state ON whatsapp_threads(ai_state)" ) ) if "idx_whatsapp_threads_ai_last_model_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_whatsapp_threads_ai_last_model_at ON whatsapp_threads(ai_last_model_at)" ) ) if "ai_sessions" in table_names: columns = _table_columns(inspector, "ai_sessions") _add_column_if_missing(conn, columns, "ai_sessions", "call_id", "VARCHAR(128)") _add_column_if_missing(conn, columns, "ai_sessions", "context_summary_json", "TEXT DEFAULT '{}'") _add_column_if_missing(conn, columns, "ai_sessions", "context_summary_updated_at", "VARCHAR(64)") indexes = _table_indexes(inspector, "ai_sessions") if "idx_ai_sessions_call_id" not in indexes: conn.execute( text("CREATE INDEX IF NOT EXISTS idx_ai_sessions_call_id ON ai_sessions(call_id)") ) if "idx_ai_sessions_context_summary_updated_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_ai_sessions_context_summary_updated_at " "ON ai_sessions(context_summary_updated_at)" ) ) if "voice_events" in table_names: columns = _table_columns(inspector, "voice_events") _add_column_if_missing(conn, columns, "voice_events", "source_event_id", "VARCHAR(64)") indexes = _table_indexes(inspector, "voice_events") if "ix_voice_events_source_event_id" not in indexes: conn.execute( text( "CREATE UNIQUE INDEX IF NOT EXISTS ix_voice_events_source_event_id " "ON voice_events(source_event_id)" ) ) if "asterisk_call_links" in table_names: columns = _table_columns(inspector, "asterisk_call_links") _add_column_if_missing(conn, columns, "asterisk_call_links", "voice_session_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "asterisk_call_links", "ai_session_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "asterisk_call_links", "ai_state", "VARCHAR(32)") _add_column_if_missing(conn, columns, "asterisk_call_links", "ai_handoff_reason", "TEXT") _add_column_if_missing(conn, columns, "asterisk_call_links", "ai_last_model_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "asterisk_call_links", "voice_start_language", "VARCHAR(16)") _add_column_if_missing(conn, columns, "asterisk_call_links", "customer_name_status", "VARCHAR(32)") _add_column_if_missing(conn, columns, "asterisk_call_links", "customer_name_value", "VARCHAR(256)") _add_column_if_missing(conn, columns, "asterisk_call_links", "customer_name_source", "VARCHAR(32)") _add_column_if_missing(conn, columns, "asterisk_call_links", "customer_name_resolved_at", "VARCHAR(64)") indexes = _table_indexes(inspector, "asterisk_call_links") if "idx_asterisk_call_links_voice_session_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_voice_session_id " "ON asterisk_call_links(voice_session_id)" ) ) if "idx_asterisk_call_links_ai_session_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_ai_session_id " "ON asterisk_call_links(ai_session_id)" ) ) if "idx_asterisk_call_links_ai_state" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_ai_state " "ON asterisk_call_links(ai_state)" ) ) if "idx_asterisk_call_links_ai_last_model_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_ai_last_model_at " "ON asterisk_call_links(ai_last_model_at)" ) ) if "idx_asterisk_call_links_voice_start_language" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_voice_start_language " "ON asterisk_call_links(voice_start_language)" ) ) if "idx_asterisk_call_links_customer_name_status" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_customer_name_status " "ON asterisk_call_links(customer_name_status)" ) ) if "idx_asterisk_call_links_customer_name_source" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_customer_name_source " "ON asterisk_call_links(customer_name_source)" ) ) if "idx_asterisk_call_links_customer_name_resolved_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_customer_name_resolved_at " "ON asterisk_call_links(customer_name_resolved_at)" ) ) _add_column_if_missing(conn, columns, "asterisk_call_links", "tenant_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "asterisk_call_links", "current_level", "VARCHAR(16)") _add_column_if_missing(conn, columns, "asterisk_call_links", "required_skills_json", "TEXT DEFAULT '[]'") _add_column_if_missing(conn, columns, "asterisk_call_links", "priority", "INTEGER DEFAULT 3") indexes = _table_indexes(inspector, "asterisk_call_links") if "idx_asterisk_call_links_tenant_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_tenant_id " "ON asterisk_call_links(tenant_id)" ) ) if "idx_asterisk_call_links_current_level" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_asterisk_call_links_current_level " "ON asterisk_call_links(current_level)" ) ) if "voice_ai_sessions" in table_names: columns = _table_columns(inspector, "voice_ai_sessions") _add_column_if_missing(conn, columns, "voice_ai_sessions", "voice_start_language", "VARCHAR(16)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "customer_name_status", "VARCHAR(32)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "customer_name_value", "VARCHAR(256)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "customer_name_source", "VARCHAR(32)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "customer_name_resolved_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "media_uuid", "VARCHAR(64)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "media_status", "VARCHAR(32)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "media_connected_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "media_ended_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "voice_ai_sessions", "last_media_frame_at", "VARCHAR(64)") indexes = _table_indexes(inspector, "voice_ai_sessions") if "idx_voice_ai_sessions_voice_start_language" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_voice_start_language " "ON voice_ai_sessions(voice_start_language)" ) ) if "idx_voice_ai_sessions_customer_name_status" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_customer_name_status " "ON voice_ai_sessions(customer_name_status)" ) ) if "idx_voice_ai_sessions_customer_name_source" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_customer_name_source " "ON voice_ai_sessions(customer_name_source)" ) ) if "idx_voice_ai_sessions_customer_name_resolved_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_customer_name_resolved_at " "ON voice_ai_sessions(customer_name_resolved_at)" ) ) if "idx_voice_ai_sessions_media_uuid" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_media_uuid " "ON voice_ai_sessions(media_uuid)" ) ) if "idx_voice_ai_sessions_media_status" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_media_status " "ON voice_ai_sessions(media_status)" ) ) if "idx_voice_ai_sessions_media_connected_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_media_connected_at " "ON voice_ai_sessions(media_connected_at)" ) ) if "idx_voice_ai_sessions_media_ended_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_media_ended_at " "ON voice_ai_sessions(media_ended_at)" ) ) if "idx_voice_ai_sessions_last_media_frame_at" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_voice_ai_sessions_last_media_frame_at " "ON voice_ai_sessions(last_media_frame_at)" ) ) if "ivr_sessions" in table_names: columns = _table_columns(inspector, "ivr_sessions") _add_column_if_missing(conn, columns, "ivr_sessions", "resolved_queue_code", "VARCHAR(64)") indexes = _table_indexes(inspector, "ivr_sessions") if "idx_ivr_sessions_resolved_queue_code" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_ivr_sessions_resolved_queue_code " "ON ivr_sessions(resolved_queue_code)" ) ) if "kb_articles" in table_names: columns = _table_columns(inspector, "kb_articles") _add_column_if_missing(conn, columns, "kb_articles", "article_group_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "kb_articles", "intent_code", "VARCHAR(64)") _add_column_if_missing(conn, columns, "kb_articles", "language", "VARCHAR(8) DEFAULT 'ru'") if "language" in columns: conn.execute( text( """ UPDATE kb_articles SET language = 'ru' WHERE language IS NULL OR TRIM(language) = '' """ ) ) if "article_group_id" in columns: conn.execute( text( """ UPDATE kb_articles SET article_group_id = article_id WHERE article_group_id IS NULL OR TRIM(article_group_id) = '' """ ) ) indexes = _table_indexes(inspector, "kb_articles") if "idx_kb_articles_article_group_id" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_kb_articles_article_group_id " "ON kb_articles(article_group_id)" ) ) if "idx_kb_articles_language" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_kb_articles_language " "ON kb_articles(language)" ) ) if "idx_kb_articles_intent_code" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_kb_articles_intent_code " "ON kb_articles(intent_code)" ) ) if "sales_automation_tasks" in table_names: columns = _table_columns(inspector, "sales_automation_tasks") _add_column_if_missing(conn, columns, "sales_automation_tasks", "max_retries", "INTEGER DEFAULT 3") _add_column_if_missing(conn, columns, "sales_automation_tasks", "locked_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_automation_tasks", "locked_by", "VARCHAR(128)") _add_column_if_missing(conn, columns, "sales_automation_tasks", "completed_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_automation_tasks", "failed_at", "VARCHAR(64)") conn.execute( text( """ UPDATE sales_automation_tasks SET max_retries = 3 WHERE max_retries IS NULL """ ) ) indexes = _table_indexes(inspector, "sales_automation_tasks") if "idx_sales_tasks_lock" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_sales_tasks_lock " "ON sales_automation_tasks(status, locked_at, locked_by)" ) ) if "sales_communication_sessions" in table_names: columns = _table_columns(inspector, "sales_communication_sessions") _add_column_if_missing(conn, columns, "sales_communication_sessions", "channel_provider", "VARCHAR(32)") indexes = _table_indexes(inspector, "sales_communication_sessions") if "idx_sales_comm_channel_provider" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_sales_comm_channel_provider " "ON sales_communication_sessions(channel_provider)" ) ) if "sales_notes" in table_names: columns = _table_columns(inspector, "sales_notes") _add_column_if_missing(conn, columns, "sales_notes", "source_type", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_notes", "source_id", "VARCHAR(128)") _add_column_if_missing(conn, columns, "sales_notes", "updated_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_notes", "archived_at", "VARCHAR(64)") conn.execute( text( """ UPDATE sales_notes SET updated_at = created_at WHERE updated_at IS NULL OR TRIM(updated_at) = '' """ ) ) if "sales_escalations" in table_names: columns = _table_columns(inspector, "sales_escalations") _add_column_if_missing(conn, columns, "sales_escalations", "assigned_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_escalations", "started_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_escalations", "canceled_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_escalations", "resolution_code", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_escalations", "resolution_summary", "TEXT") _add_column_if_missing(conn, columns, "sales_escalations", "sla_due_at", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_escalations", "source_channel", "VARCHAR(32)") _add_column_if_missing(conn, columns, "sales_escalations", "source_communication_session_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_escalations", "source_automation_task_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_escalations", "updated_at", "VARCHAR(64)") conn.execute( text( """ UPDATE sales_escalations SET updated_at = COALESCE(updated_at, resolved_at, created_at) WHERE updated_at IS NULL OR TRIM(updated_at) = '' """ ) ) if "sales_channel_switches" in table_names: columns = _table_columns(inspector, "sales_channel_switches") _add_column_if_missing(conn, columns, "sales_channel_switches", "communication_session_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_channel_switches", "previous_communication_session_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_channel_switches", "new_communication_session_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_channel_switches", "reason_code", "VARCHAR(64) DEFAULT 'human_decision'") _add_column_if_missing(conn, columns, "sales_channel_switches", "reason_text", "TEXT") _add_column_if_missing(conn, columns, "sales_channel_switches", "initiated_by_type", "VARCHAR(32) DEFAULT 'human'") _add_column_if_missing(conn, columns, "sales_channel_switches", "initiated_by_id", "VARCHAR(128)") _add_column_if_missing(conn, columns, "sales_channel_switches", "source_type", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_channel_switches", "source_id", "VARCHAR(128)") _add_column_if_missing(conn, columns, "sales_channel_switches", "recommended_by_task_id", "VARCHAR(64)") _add_column_if_missing(conn, columns, "sales_channel_switches", "created_at", "VARCHAR(64)") conn.execute( text( """ UPDATE sales_channel_switches SET communication_session_id = COALESCE(communication_session_id, communication_id), previous_communication_session_id = COALESCE(previous_communication_session_id, communication_id), reason_code = COALESCE(reason_code, 'human_decision'), reason_text = COALESCE(reason_text, reason_for_channel_switch), initiated_by_type = COALESCE(initiated_by_type, 'human'), created_at = COALESCE(created_at, switched_at) WHERE created_at IS NULL OR TRIM(created_at) = '' """ ) ) if "sales_payment_webhook_events" in table_names: indexes = _table_indexes(inspector, "sales_payment_webhook_events") if "idx_sales_payment_webhook_events_received" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_sales_payment_webhook_events_received " "ON sales_payment_webhook_events(tenant_id, received_at)" ) ) if "idx_sales_payment_webhook_events_payment" not in indexes: conn.execute( text( "CREATE INDEX IF NOT EXISTS idx_sales_payment_webhook_events_payment " "ON sales_payment_webhook_events(tenant_id, payment_provider, external_payment_id)" ) ) if "idx_sales_payment_webhook_events_external_event_unique" not in indexes: conn.execute( text( "CREATE UNIQUE INDEX IF NOT EXISTS idx_sales_payment_webhook_events_external_event_unique " "ON sales_payment_webhook_events(tenant_id, payment_provider, external_event_id) " "WHERE external_event_id IS NOT NULL" ) ) if "sales_payments" in table_names: indexes = _table_indexes(inspector, "sales_payments") if "idx_sales_payments_tenant_external_payment_unique" not in indexes: conn.execute( text( "CREATE UNIQUE INDEX IF NOT EXISTS idx_sales_payments_tenant_external_payment_unique " "ON sales_payments(tenant_id, payment_provider, external_payment_id) " "WHERE external_payment_id IS NOT NULL" ) ) def init_sql_schema() -> None: global _BASE_SCHEMA_INITIALIZED with _SCHEMA_INIT_LOCK: if schema_management_mode() == "migrations": validate_schema_migrations_applied() elif not _BASE_SCHEMA_INITIALIZED: for attempt in range(5): try: Base.metadata.create_all(bind=engine) _BASE_SCHEMA_INITIALIZED = True break except OperationalError as exc: message = str(exc).lower() if "already exists" in message: _BASE_SCHEMA_INITIALIZED = True break if "database is locked" not in message or attempt >= 4: raise time.sleep(0.25 * (attempt + 1)) for attempt in range(5): try: _apply_runtime_schema_compatibility() return except OperationalError as exc: if "database is locked" not in str(exc).lower() or attempt >= 4: raise time.sleep(0.25 * (attempt + 1))