feat(api): 拆分章节物化与 Story 后处理,并加固 Redis 锁与腾讯 ASR

回忆录 Story 流水线(同步)
- 同步路径仅写入 Story 与章节关联,改为 mark_chapter_dirty_sync,不再内联 compose
- 物化由 Celery recompose_chapter 异步完成;compose 不变量与异常时保留 dirty 的语义在 repo 中补充说明
- Evidence:大批次时降低 top_k;路由候选 story 携带 char_count/version_count;append 超长/版本过多时强制新开 story
- 叙事 prompt:relevant_chunks 去重,减少重复证据噪声
- 叙事回退与忠实度 gate:返回 fallback 类型并记录结构化日志(含耗时、JSON 有效性等)

Post-commit 与任务编排
- 新增 post_commit.enqueue_story_post_commit_effects:统一派发 generate_story_image(Redis 去重)、延迟 recompose_chapter、可选 memory compaction
- memoir_tasks / story_service / story_image_tasks 改为调用 post-commit 入口;主图回填后按关联章节重算并调度物化与 compacs(锁委托、Redis 单例、ASR to_thread)
- 更新 test_narrative_pipeline 以适配 _apply_narrative_fallbacks 返回值
This commit is contained in:
Kevin
2026-03-30 11:53:04 +08:00
parent e884409410
commit aac484463d
15 changed files with 775 additions and 144 deletions

View File

@@ -11,11 +11,10 @@ from datetime import datetime, timezone
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.logging import get_logger
from app.features.memoir.asset_resolver import strip_asset_image_refs_from_markdown
from app.features.memoir import repo as memoir_repo
from app.features.memoir.asset_resolver import strip_asset_image_refs_from_markdown
from app.features.memoir.memoir_images.settings import MemoirImageSettings
from app.features.story.image_intent_extractor import extract_primary_image_intent
from app.features.story.time_hints import apply_infer_story_time_start_to_model
from app.features.story.repo import (
count_story_versions,
create_story,
@@ -27,6 +26,7 @@ from app.features.story.repo import (
get_story_by_id,
get_story_image_intent_by_story,
)
from app.features.story.time_hints import apply_infer_story_time_start_to_model
logger = get_logger(__name__)
@@ -156,17 +156,26 @@ class StoryService:
await memoir_repo.mark_chapters_dirty_for_story(self._db, story.id)
await self._db.commit()
if md.strip():
from app.tasks.chapter_compose_tasks import recompose_chapters_for_story
from app.tasks.story_image_tasks import generate_story_image
from app.features.memoir.repo import get_chapter_ids_linked_to_story
from app.features.story.post_commit import enqueue_story_post_commit_effects
try:
generate_story_image.delay(story.id)
except Exception as exc:
logger.warning("派发 generate_story_image 失败: {}", exc)
try:
recompose_chapters_for_story.delay(story.id)
except Exception as exc:
logger.warning("派发 recompose_chapters_for_story 失败: {}", exc)
chapter_ids = set(await get_chapter_ids_linked_to_story(self._db, story.id))
pc = enqueue_story_post_commit_effects(
user_id=user_id,
story_ids={story.id},
chapter_ids=chapter_ids,
trigger_source="manual_api",
need_compaction=False,
)
logger.info(
"event=story_post_commit user_id={} trigger=manual_api "
"enqueued_story_image_count={} enqueued_chapter_recompose_count={} "
"errors={}",
user_id,
pc.enqueued_story_image_count,
pc.enqueued_chapter_recompose_count,
pc.errors,
)
return story.id
async def append_version(
@@ -208,17 +217,25 @@ class StoryService:
)
await memoir_repo.mark_chapters_dirty_for_story(self._db, story_id)
await self._db.commit()
from app.tasks.chapter_compose_tasks import recompose_chapters_for_story
from app.tasks.story_image_tasks import generate_story_image
from app.features.memoir.repo import get_chapter_ids_linked_to_story
from app.features.story.post_commit import enqueue_story_post_commit_effects
try:
generate_story_image.delay(story_id)
except Exception as exc:
logger.warning("派发 generate_story_image 失败: {}", exc)
try:
recompose_chapters_for_story.delay(story_id)
except Exception as exc:
logger.warning("派发 recompose_chapters_for_story 失败: {}", exc)
chapter_ids = set(await get_chapter_ids_linked_to_story(self._db, story_id))
pc = enqueue_story_post_commit_effects(
user_id=story.user_id,
story_ids={story_id},
chapter_ids=chapter_ids,
trigger_source="manual_api",
need_compaction=False,
)
logger.info(
"event=story_post_commit user_id={} trigger=manual_api_append "
"enqueued_story_image_count={} enqueued_chapter_recompose_count={} errors={}",
story.user_id,
pc.enqueued_story_image_count,
pc.enqueued_chapter_recompose_count,
pc.errors,
)
return version.id
async def link_evidence(