Files
life-echo/api/app/features/story/sync_write.py
Kevin 309a051038 feat: 回忆录证据血缘与内部评测可追溯,顺带对齐本地评测台与 CI
数据库与模型:新增多版迁移(章节证据快照、对话血缘、记忆事实/时间线 lineage 等),把「成稿 ↔ 对话/记忆」的溯源信息落到表结构里。
业务链路:会话与 WS、回忆录/故事流水线、记忆写入与 enrichment 等跟着接上线索与快照;新增章节证据快照与评测侧 EvalTraceService 等模块,方便组评审用的证据包。
内部评测:自动化 run 与手工 memoir 评审共用可追溯证据;rubric/ judge 相关脚本与文档有配套调整。
app-eval-web:Memoir/实验详情里能展开看证据摘要与 evidence_trace(含对话轮次 id);Vite 代理与 development.sh 注入的 API 端口与当前默认内部评测端口一致,避免改端口后页面连错服务。
工程杂项:GitHub Actions / 仓库说明有更新;各适配器与支付/配额/plan 等多处为小改动或跟随主改动的收尾;新增/扩充了?
2026-04-08 15:37:09 +08:00

291 lines
8.5 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.
"""
Story 同步写入Celery / sync Session
与 StoryService 行为对齐:版本链、主图 intent、章节 dirty不 commit由调用方提交。
"""
from __future__ import annotations
import uuid
from datetime import datetime, timezone
from sqlalchemy import delete, func, select
from sqlalchemy.orm import Session, joinedload
from app.core.logging import get_logger
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.models import ChapterStoryLink
from app.features.story.image_intent_extractor import extract_primary_image_intent
from app.features.story.models import (
Story,
StoryEvidenceLink,
StoryImageIntent,
StoryVersion,
)
from app.features.story.time_hints import apply_infer_story_time_start_to_model
logger = get_logger(__name__)
def count_story_versions_sync(session: Session, story_id: str) -> int:
stmt = select(func.count(StoryVersion.id)).where(StoryVersion.story_id == story_id)
return int(session.execute(stmt).scalar() or 0)
def _delete_pending_failed_intents_sync(session: Session, story_id: str) -> None:
session.execute(
delete(StoryImageIntent).where(
StoryImageIntent.story_id == story_id,
StoryImageIntent.intent_role == "primary",
StoryImageIntent.status.in_(["pending", "failed"]),
)
)
def _get_primary_intent_sync(
session: Session, story_id: str
) -> StoryImageIntent | None:
stmt = select(StoryImageIntent).where(
StoryImageIntent.story_id == story_id,
StoryImageIntent.intent_role == "primary",
)
return session.execute(stmt).scalar_one_or_none()
def _extract_and_store_image_intent_sync(
session: Session,
*,
story: Story,
version: StoryVersion,
markdown: str,
) -> None:
_delete_pending_failed_intents_sync(session, story.id)
result = extract_primary_image_intent(
markdown,
title=story.title or "",
stage=story.stage,
summary=story.summary,
people_refs=story.people_refs or [],
place_refs=story.place_refs or [],
time_start=story.time_start,
time_end=story.time_end,
)
existing = _get_primary_intent_sync(session, story.id)
now = datetime.now(timezone.utc)
if existing and existing.story_version_id == version.id:
st = (existing.status or "").strip()
if st in ("processing", "completed"):
return
existing.caption = result.caption
existing.prompt_brief = result.prompt_brief
existing.style_profile = result.style_profile
existing.status = "pending"
existing.error = None
existing.asset_id = None
existing.updated_at = now
return
if existing and existing.story_version_id != version.id:
existing.story_version_id = version.id
existing.caption = result.caption
existing.prompt_brief = result.prompt_brief
existing.style_profile = result.style_profile
existing.status = "pending"
existing.error = None
existing.asset_id = None
existing.updated_at = now
return
session.add(
StoryImageIntent(
id=str(uuid.uuid4()),
story_id=story.id,
story_version_id=version.id,
intent_role="primary",
caption=result.caption,
prompt_brief=result.prompt_brief,
style_profile=result.style_profile,
status="pending",
)
)
def replace_story_evidence_links_sync(
session: Session,
*,
story_id: str,
chunk_ids: list[str],
fact_ids: list[str],
timeline_event_ids: list[str],
summary_ids: list[str],
) -> None:
"""以当前生成所用的检索闭包覆盖 story 证据关联artifact 当前态绑定)。"""
session.execute(
delete(StoryEvidenceLink).where(StoryEvidenceLink.story_id == story_id)
)
for cid in chunk_ids:
session.add(
StoryEvidenceLink(
id=str(uuid.uuid4()),
story_id=story_id,
evidence_type="chunk",
evidence_id=cid,
role="primary",
)
)
for fid in fact_ids:
session.add(
StoryEvidenceLink(
id=str(uuid.uuid4()),
story_id=story_id,
evidence_type="fact",
evidence_id=fid,
role="supporting",
)
)
for tid in timeline_event_ids:
session.add(
StoryEvidenceLink(
id=str(uuid.uuid4()),
story_id=story_id,
evidence_type="timeline_event",
evidence_id=tid,
role="supporting",
)
)
for sid in summary_ids:
session.add(
StoryEvidenceLink(
id=str(uuid.uuid4()),
story_id=story_id,
evidence_type="summary",
evidence_id=sid,
role="background",
)
)
def create_story_with_version_sync(
session: Session,
*,
user_id: str,
title: str,
canonical_markdown: str,
stage: str | None = None,
prompt_meta: dict | None = None,
) -> Story:
md = strip_asset_image_refs_from_markdown(canonical_markdown or "")
story = Story(
id=str(uuid.uuid4()),
user_id=user_id,
title=title,
stage=stage,
canonical_markdown=md,
)
session.add(story)
session.flush()
vid = str(uuid.uuid4())
version = StoryVersion(
id=vid,
story_id=story.id,
version_no=1,
markdown_snapshot=md,
actor_type="ai",
source_type="generate",
prompt_meta=prompt_meta,
)
session.add(version)
session.flush()
story.current_version_id = vid
apply_infer_story_time_start_to_model(story)
if md.strip():
_extract_and_store_image_intent_sync(
session, story=story, version=version, markdown=md
)
memoir_repo.mark_chapters_dirty_for_story_sync(session, story.id)
return story
def append_story_version_sync(
session: Session,
story_id: str,
markdown_snapshot: str,
*,
actor_type: str = "ai",
source_type: str = "generate",
prompt_meta: dict | None = None,
) -> StoryVersion:
story = session.get(Story, story_id)
if not story:
raise ValueError(f"Story {story_id} not found")
md = strip_asset_image_refs_from_markdown(markdown_snapshot or "")
parent_id = story.current_version_id
version_no = count_story_versions_sync(session, story_id) + 1
vid = str(uuid.uuid4())
version = StoryVersion(
id=vid,
story_id=story_id,
version_no=version_no,
markdown_snapshot=md,
actor_type=actor_type,
source_type=source_type,
parent_version_id=parent_id,
prompt_meta=prompt_meta,
)
session.add(version)
session.flush()
story.current_version_id = vid
story.canonical_markdown = md
apply_infer_story_time_start_to_model(story)
_extract_and_store_image_intent_sync(
session, story=story, version=version, markdown=md
)
memoir_repo.mark_chapters_dirty_for_story_sync(session, story_id)
return version
def ensure_chapter_story_link_sync(
session: Session,
*,
chapter_id: str,
story_id: str,
) -> None:
"""若章节尚未关联该 story则在末尾追加一条 chapter_story_link。"""
exists = session.scalars(
select(ChapterStoryLink)
.where(
ChapterStoryLink.chapter_id == chapter_id,
ChapterStoryLink.story_id == story_id,
)
.limit(1)
).first()
if exists is not None:
return
max_stmt = select(func.coalesce(func.max(ChapterStoryLink.order_index), -1)).where(
ChapterStoryLink.chapter_id == chapter_id
)
max_idx = int(session.execute(max_stmt).scalar() or -1)
session.add(
ChapterStoryLink(
id=str(uuid.uuid4()),
chapter_id=chapter_id,
story_id=story_id,
order_index=max_idx + 1,
)
)
session.flush()
def list_active_stories_for_user_sync(session: Session, user_id: str) -> list[Story]:
stmt = (
select(Story)
.where(Story.user_id == user_id, Story.status == "active")
.options(
joinedload(Story.chapter_links).joinedload(ChapterStoryLink.chapter),
)
.order_by(Story.updated_at.desc())
)
return list(session.execute(stmt).unique().scalars().all())