feat(feishu): 添加飞书入站事件inbox和混合数据库协调功能 - 实现飞书入站事件持久化inbox机制,支持状态管理、租约锁定和重试退避 - 添加混合数据库基线协调工具,确保平台PostgreSQL结构安全对齐 - 增加运行组件心跳检测和readiness就绪检查机制 - 实现app_ticket事件的安全轮换和验证处理 - 添加生产环境运行编排和fail-closed安全机制 - 支持webhook快速确认和长连接独立进程处理 - 完善个人数据擦除时的待处理事件清理功能 ```
92 lines
2.8 KiB
Python
92 lines
2.8 KiB
Python
from logging.config import fileConfig
|
|
|
|
from alembic import context
|
|
from sqlalchemy import engine_from_config, pool
|
|
|
|
from app.core.config import get_settings
|
|
from app.core.database import Base
|
|
from app.core.database.safety import validate_platform_migration_target
|
|
from app.modules.ai_memory import models as ai_memory_models
|
|
from app.modules.audit import models as audit_models
|
|
from app.modules.business import models as business_models
|
|
from app.modules.events import models as event_models
|
|
from app.modules.feishu import models as feishu_models
|
|
from app.modules.feishu_users import models as feishu_user_models
|
|
from app.modules.observability import models as observability_models
|
|
from app.modules.personalization import models as personalization_models
|
|
from app.modules.subscriptions import models as subscription_models
|
|
from app.modules.workflows import models as workflow_models
|
|
|
|
config = context.config
|
|
|
|
if config.config_file_name is not None:
|
|
fileConfig(config.config_file_name)
|
|
|
|
target_metadata = Base.metadata
|
|
settings = get_settings()
|
|
validate_platform_migration_target(
|
|
settings.database_url,
|
|
settings.legacy_database_url,
|
|
)
|
|
|
|
# Keep imports referenced so SQLAlchemy model classes register with Base.metadata.
|
|
_REGISTERED_MODEL_MODULES = (
|
|
ai_memory_models,
|
|
audit_models,
|
|
business_models,
|
|
event_models,
|
|
feishu_models,
|
|
feishu_user_models,
|
|
observability_models,
|
|
personalization_models,
|
|
subscription_models,
|
|
workflow_models,
|
|
)
|
|
|
|
|
|
def run_migrations_offline() -> None:
|
|
context.configure(
|
|
url=settings.database_url,
|
|
target_metadata=target_metadata,
|
|
literal_binds=True,
|
|
dialect_opts={"paramstyle": "named"},
|
|
)
|
|
|
|
with context.begin_transaction():
|
|
context.run_migrations()
|
|
|
|
|
|
def run_migrations_online() -> None:
|
|
supplied_connection = config.attributes.get("connection")
|
|
if supplied_connection is not None:
|
|
validate_platform_migration_target(
|
|
supplied_connection.engine.url,
|
|
settings.legacy_database_url,
|
|
)
|
|
context.configure(
|
|
connection=supplied_connection,
|
|
target_metadata=target_metadata,
|
|
)
|
|
with context.begin_transaction():
|
|
context.run_migrations()
|
|
return
|
|
|
|
configuration = config.get_section(config.config_ini_section, {})
|
|
configuration["sqlalchemy.url"] = settings.database_url
|
|
connectable = engine_from_config(
|
|
configuration,
|
|
prefix="sqlalchemy.",
|
|
poolclass=pool.NullPool,
|
|
)
|
|
|
|
with connectable.connect() as connection:
|
|
context.configure(connection=connection, target_metadata=target_metadata)
|
|
with context.begin_transaction():
|
|
context.run_migrations()
|
|
|
|
|
|
if context.is_offline_mode():
|
|
run_migrations_offline()
|
|
else:
|
|
run_migrations_online()
|