Files
xiaoxia-saas/packages/adapters/sqlalchemy_impl/session.py
T
xiaoxia 3976f21c31
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 4s
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m26s
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 1m16s
CI/CD Pipeline / Frontend Lint (pull_request) Successful in 1m27s
CI/CD Pipeline / PR Build Worker Image (pull_request) Successful in 1m1s
CI/CD Pipeline / Frontend Unit Tests (pull_request) Successful in 1m34s
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m20s
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 3m51s
CI/CD Pipeline / Validate - Style (pull_request) Successful in 5m41s
CI/CD Pipeline / Integration Tests (pull_request) Successful in 4m58s
AI Code Review / AI Code Review (pull_request) Successful in 6m53s
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Successful in 5m58s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (pull_request) Successful in 15m37s
CI/CD Pipeline / Validate - Security (pull_request) Successful in 16m28s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / CI Gate (pull_request) Successful in 2s
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
PR Automation / Auto Merge on CI Green + Approved (pull_request) Successful in 10m49s
ACR Cleanup / ACR Image Cleanup (pull_request_target) Successful in 29s
Preview Cleanup / Cleanup Preview Environment (pull_request) Successful in 1m15s
feat(viral-video): #2115 三阶段分步 pipeline(图片分析/生成文案/生成视频拆分)
- 新增状态 IMAGE_ANALYZED / COPY_GENERATED,最小改动状态机:PENDING→RUNNING→IMAGE_ANALYZED→RUNNING→COPY_GENERATED→RUNNING→COMPLETED
- 新增三个 Celery task:run_viral_video_analyze / run_viral_video_generate_copy / run_viral_video_render
- 共享 _run_render_pipeline 渲染函数,兼容旧 /generate + /confirm-intent 路径(storyboard 缺失时自动兜底重算)
- 新增三个端点:
  * POST /viral-video/analyze-images  阶段1:图片VLM分析+可选参考视频风格分析
  * POST /viral-video/{id}/generate-copy  阶段2:意图解析→文案融合→分镜→合规审核
  * POST /viral-video/{id}/confirm-copy  阶段3:TTS→渲染→上传
- GET /{id} 返回 image_analysis / storyboard / generated_copy_text / copy_result(对齐前端 CopyResult 结构)
- 幂等 ALTER TABLE 补列:storyboard(JSON) / generated_copy_text(TEXT),无 Alembic 环境下安全执行
- WS 新增 viral_video:image_analyzed / viral_video:copy_generated 事件类型
- 前端类型/API 客户端:新增 AnalyzeImagesRequest/GenerateCopyRequest/ConfirmCopyRequest/StoryboardSegment
  及 analyzeViralImages/generateViralCopy/confirmViralCopy 三个真实 API 函数,保留既有 mock 函数
- 状态值对齐前端 PR #2116 命名 copy_generated(非 copy_ready)
- 93 个 viral 相关单测全通过;musetalk 2 个超时用例失败为历史遗留与本 PR 无关
2026-10-01 11:38:31 +08:00

135 lines
4.2 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from __future__ import annotations
from sqlalchemy import create_engine, text
from sqlalchemy.engine import URL, make_url
from sqlalchemy.orm import sessionmaker
from packages.adapters.sqlalchemy_impl.models import Base
SCHEMA_INIT_LOCK_ID = 2026061501
SessionLocal = None
def build_engine(
database_url: str,
*,
pool_size: int = 20,
max_overflow: int = 40,
pool_timeout: int = 30,
pool_recycle: int = 3600,
):
# SQLite 不支持 QueuePool 的 pool_size/max_overflow/pool_timeout,
# 传了会在 create_engine 阶段直接 TypeError,这里只对非 SQLite 传连接池参数。
if _is_sqlite(database_url):
return create_engine(database_url, pool_recycle=pool_recycle)
return create_engine(
database_url,
pool_size=pool_size,
max_overflow=max_overflow,
pool_timeout=pool_timeout,
pool_recycle=pool_recycle,
)
def build_session_factory(
database_url: str,
*,
pool_size: int = 20,
max_overflow: int = 40,
pool_timeout: int = 30,
pool_recycle: int = 3600,
):
engine = build_engine(
database_url,
pool_size=pool_size,
max_overflow=max_overflow,
pool_timeout=pool_timeout,
pool_recycle=pool_recycle,
)
session_factory = sessionmaker(autocommit=False, autoflush=False, bind=engine)
global SessionLocal
SessionLocal = session_factory
return engine, session_factory
def _is_sqlite(database_url: str) -> bool:
"""检测是否为 SQLite 数据库 URL."""
return database_url.startswith("sqlite")
def _build_admin_url(database_url: str) -> URL:
url = make_url(database_url)
return url.set(database="postgres")
def ensure_database_exists(database_url: str) -> None:
"""确保数据库存在(仅 PostgreSQL 需要,SQLite 自动创建)."""
if _is_sqlite(database_url):
return
target_url = make_url(database_url)
admin_engine = create_engine(_build_admin_url(database_url), isolation_level="AUTOCOMMIT")
try:
with admin_engine.connect() as connection:
exists = connection.execute(
text("SELECT 1 FROM pg_database WHERE datname = :database_name"),
{"database_name": target_url.database},
).scalar()
if exists:
return
connection.execute(text(f'CREATE DATABASE "{target_url.database}"'))
finally:
admin_engine.dispose()
_VIRAL_VIDEO_BACKFILL_COLS = [
("storyboard", "JSON"),
("generated_copy_text", "TEXT NOT NULL DEFAULT ''"),
]
def _ensure_viral_video_columns(connection) -> None:
"""Idempotently add new columns to viral_video_jobs; create_all will not ALTER existing tables."""
from sqlalchemy import inspect as _inspect
try:
insp = _inspect(connection)
if not insp.has_table("viral_video_jobs"):
return
existing = {c["name"] for c in insp.get_columns("viral_video_jobs")}
except Exception:
return
import logging as _logging
_log = _logging.getLogger(__name__)
for col, ddl in _VIRAL_VIDEO_BACKFILL_COLS:
if col in existing:
continue
try:
connection.execute(text(f"ALTER TABLE viral_video_jobs ADD COLUMN {col} {ddl}"))
_log.info("added column viral_video_jobs.%s", col)
except Exception as e:
_log.warning("add column %s failed: %s", col, e)
def initialize_database(engine) -> None:
"""初始化数据库 schema。
PostgreSQL 使用 advisory lock 防止并发初始化冲突;
SQLite 直接 create_all(单文件,无并发风险)。
"""
if _is_sqlite(str(engine.url)):
Base.metadata.create_all(bind=engine)
return
with engine.connect() as connection:
connection.execute(text("SELECT pg_advisory_lock(:lock_id)"), {"lock_id": SCHEMA_INIT_LOCK_ID})
try:
Base.metadata.create_all(bind=connection)
connection.commit()
finally:
connection.execute(
text("SELECT pg_advisory_unlock(:lock_id)"),
{"lock_id": SCHEMA_INIT_LOCK_ID},
)
_ensure_viral_video_columns(connection)
connection.commit()