diff --git a/apps/worker/worker_app/tasks/generation.py b/apps/worker/worker_app/tasks/generation.py index a2f211170..a089726af 100644 --- a/apps/worker/worker_app/tasks/generation.py +++ b/apps/worker/worker_app/tasks/generation.py @@ -1493,9 +1493,9 @@ def generate_video(self, task_id: str) -> dict: ) if _cover_model: _cover_model.cover_url = cover_frame_url - meta = dict(_cover_model.metadata or {}) + meta = dict(_cover_model.extra_meta or {}) meta["cover_candidates"] = cover_candidates - _cover_model.metadata = meta + _cover_model.extra_meta = meta _cover_session.commit() finally: _cover_session.close() @@ -1680,10 +1680,27 @@ def generate_video(self, task_id: str) -> dict: logger.info("[task_id=%s] MediaKit 视频理解完成: %d 个素材", task_id, len(asset_urls)) # 保存分析结果到 extra_meta - if asset_analyses and gen_task: - gen_task.extra_meta = {**(gen_task.extra_meta or {}), "asset_analyses": asset_analyses} - _repo.update(gen_task) - _flush_logs(task_id, gen_task) + if asset_analyses: + _meta_session = SessionLocal() + try: + from packages.adapters.sqlalchemy_impl.models import ( + GenerationTaskModel, + ) + + _m = ( + _meta_session.query(GenerationTaskModel) + .filter(GenerationTaskModel.id == task_id) + .first() + ) + if _m: + existing = dict(_m.extra_meta or {}) + existing["asset_analyses"] = asset_analyses + _m.extra_meta = existing + _meta_session.commit() + finally: + _meta_session.close() + if gen_task: + _flush_logs(task_id, gen_task) except Exception: logger.warning("[task_id=%s] MediaKit 视频理解失败,继续渲染", task_id, exc_info=True) @@ -1777,10 +1794,10 @@ def generate_video(self, task_id: str) -> dict: ) if _cover_model: _cover_model.cover_url = cover_frame_url - # 持久化完整候选列表到 metadata - meta = dict(_cover_model.metadata or {}) + # 持久化完整候选列表到 extra_meta + meta = dict(_cover_model.extra_meta or {}) meta["cover_candidates"] = cover_candidates - _cover_model.metadata = meta + _cover_model.extra_meta = meta _cover_session.commit() logger.info( "[task_id=%s] 封面帧已持久化(ffmpeg本地抽帧): cover_url=%s candidates=%d", diff --git a/tests/unit/test_worker_cover_meta_and_status.py b/tests/unit/test_worker_cover_meta_and_status.py new file mode 100644 index 000000000..9d6fa1ed3 --- /dev/null +++ b/tests/unit/test_worker_cover_meta_and_status.py @@ -0,0 +1,70 @@ +"""防回归测试:P1 修复 +- Bug 1: gen_task 过期内存对象 _repo.update() 覆盖 DB status 为 pending +- Bug 2: 封面模型误用 .metadata(SQLAlchemy 保留属性),应为 .extra_meta +""" + +import ast +from pathlib import Path + +import pytest + +GENERATION_FILE = Path(__file__).resolve().parents[2] / "apps" / "worker" / "worker_app" / "tasks" / "generation.py" + + +def _read_source() -> str: + return GENERATION_FILE.read_text(encoding="utf-8") + + +class TestCoverModelUsesExtraMeta: + """封面持久化必须使用 ORM 属性 extra_meta,而不是 SQLAlchemy 保留的 .metadata。""" + + def test_no_metadata_attribute_access_on_cover_model(self): + source = _read_source() + # 禁止对 _cover_model.metadata 进行读或写 + assert "_cover_model.metadata" not in source, ( + "_cover_model.metadata is the SQLAlchemy reserved MetaData object, " + "not the JSON column. Use _cover_model.extra_meta instead." + ) + + def test_extra_meta_used_for_cover_candidates(self): + source = _read_source() + assert "_cover_model.extra_meta" in source + assert 'meta["cover_candidates"]' in source + + +class TestAssetAnalysesDoesNotOverwriteStatus: + """保存 asset_analyses 时不能用过期的 gen_task 内存对象整体 _repo.update, + 否则会把已被 _update_task_status 改为 running 的 status 覆盖回 pending。""" + + def test_no_stale_repo_update_with_gen_task(self): + source = _read_source() + # 旧代码:gen_task.extra_meta = {...}; _repo.update(gen_task) + # 这行会把内存中的 pending status 写回 DB + assert "_repo.update(gen_task)" not in source, ( + "_repo.update(gen_task) writes a stale in-memory object back to DB, " + "overwriting status set by _update_task_status. " + "Use an independent session to update only extra_meta." + ) + + def test_asset_analyses_uses_independent_session(self): + """asset_analyses 持久化必须用独立 session 查询最新模型再提交。""" + source = _read_source() + assert "_meta_session" in source + assert "GenerationTaskModel" in source + # 必须只更新 extra_meta 字段 + assert 'existing["asset_analyses"]' in source + + +class TestGenerationTaskModelOrmAttribute: + """确认 ORM 属性映射:Python 属性 extra_meta -> DB 列 metadata。""" + + def test_orm_attribute_is_extra_meta(self): + from packages.adapters.sqlalchemy_impl.models import GenerationTaskModel + + # ORM 属性必须存在 + assert hasattr(GenerationTaskModel, "extra_meta") + # .metadata 是 SQLAlchemy 声明基类保留的 MetaData,不是列描述符 + # 它不应该是我们的 JSON 字段 + from sqlalchemy import MetaData + + assert isinstance(GenerationTaskModel.metadata, MetaData)