Files
life-echo/api/app/features/memoir/service.py
Kevin 7f57f96c25 重构回忆录为 story-first / markdown-first 架构并整合图片意图与前端 UI 修复
本次 squash merge 将 codex-story-first-image-intent 的整体改动合入 development,核心内容包括:

1. 后端数据与迁移:新增 stories、story_versions、story_image_intents、chapter_cover_intents、assets 等模型与 Alembic 迁移,建立 story-first、markdown-first、asset-first 的主数据链路。

2. 生成与任务链:引入 StoryBuilderOrchestrator、ChapterComposerOrchestrator、story_image_tasks、chapter_cover_tasks,图片生成从正文占位符改为结构化 intent -> asset -> markdown 回填。

3. 并发与一致性:为 story/chapter intent 增加 claim_token、claimed_at、attempt_count,采用数据库原子 claim 为主、Redis 锁为辅,避免重复生成、锁误删和 processing 卡死。

4. Memoir 读写路径:章节 canonical_markdown 成为正文真源,列表/详情接口补齐 markdown、cover_asset、word_count 等字段,PDF 与 asset 解析链路同步升级。

5. Memory / Retrieval:扩展 transcript ingest、chunking、evidence 检索与 story 聚合基础设施,为后续 story-first RAG 与多 agent 编排提供底座。

6. App 端体验:章节页继续走 MarkdownRenderer 阅读链,同时吸收 fix3-19 的跨平台 UI glitch 修复;更新对话页、首页、文案资源与章节列表映射逻辑。

7. 测试与文档:补充 asset resolver、story image task、章节封面派发、markdown 映射等回归测试,并加入图片占位符退役设计文档。
2026-03-20 10:31:51 +08:00

303 lines
12 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.
"""Memoir service — 回忆录编排(章节生成、状态流转);通过 MemoryService 获取 evidence。"""
from datetime import datetime, timezone
from typing import List, Optional
from app.core.logging import get_logger
from fastapi import HTTPException
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import joinedload
from app.agents.memoir.prompts import (
CHAPTER_CATEGORIES,
CHAPTER_ORDER,
STAGE_TO_ORDER,
)
from app.features.memoir import repo
from app.features.memoir.asset_resolver import (
collect_asset_ids_for_chapter,
collect_asset_ids_for_chapters,
strip_legacy_image_placeholders,
)
from app.features.memoir.asset_urls import signed_urls_for_asset_ids
from app.features.memoir.helpers import (
chapter_to_dict,
chapter_to_list_dict,
is_image_permanently_unavailable,
)
from app.features.memoir.models import Book, Chapter, ChapterSection
from app.features.memoir.memoir_images.settings import MemoirImageSettings
from app.features.memory.service import MemoryService
logger = get_logger(__name__)
async def get_or_create_book(user_id: str, db: AsyncSession):
"""Get the user's current book or return None."""
return await repo.get_current_book(user_id, db)
class MemoirService:
def __init__(
self,
db: AsyncSession,
memory_service: Optional[MemoryService] = None,
):
self._db = db
self._memory = memory_service
async def get_evidence(self, user_id: str, query: str, *, top_k: int = 10) -> dict:
"""通过 MemoryService 获取检索证据(章节生成时优先使用)。"""
if self._memory is None:
return {
"relevant_chunks": [],
"relevant_summaries": [],
"relevant_facts": [],
"timeline_hints": [],
}
return await self._memory.retrieve(user_id, query, top_k=top_k)
async def _cleanup_unavailable_images(self, ch: Chapter) -> None:
sections = getattr(ch, "sections", None) or []
cleaned = False
for s in sections:
rec = getattr(s, "image_record", None)
if rec and is_image_permanently_unavailable(rec):
logger.info("清理不可用配图: chapter=%s, section=%s", ch.id, s.id)
s.image_id = None
await self._db.delete(rec)
cleaned = True
if cleaned:
await self._db.commit()
await self._db.refresh(ch)
async def get_current_book(self, user_id: str) -> dict:
book = await repo.get_current_book(user_id, self._db)
if not book:
return {"message": "No book found"}
return {
"id": book.id,
"title": book.title,
"total_pages": book.total_pages,
"total_words": book.total_words,
"cover_image_url": book.cover_image_url,
"has_update": book.has_update,
"last_update_chapter_id": book.last_update_chapter_id,
}
async def clear_book_update(self, user_id: str) -> dict:
book = await repo.get_current_book(user_id, self._db)
if not book:
return {"status": "ok", "message": "No book found"}
book.has_update = False
await self._db.commit()
return {"status": "ok"}
async def update_book(self, book_id: str, user_id: str, title: str) -> dict:
book = await self._db.get(Book, book_id)
if not book:
raise HTTPException(status_code=404, detail="Book not found")
if book.user_id != user_id:
raise HTTPException(status_code=403, detail="无权更新此回忆录")
book.title = title
await self._db.commit()
await self._db.refresh(book)
return {
"id": book.id,
"title": book.title,
"total_pages": book.total_pages,
"total_words": book.total_words,
"cover_image_url": book.cover_image_url,
"has_update": book.has_update,
"last_update_chapter_id": book.last_update_chapter_id,
}
async def export_pdf(self, user_id: str, book_id: str) -> dict:
from app.features.memoir.pdf_service import pdf_service
book = await self._db.get(Book, book_id)
if not book:
raise HTTPException(status_code=404, detail="Book not found")
if book.user_id != user_id:
raise HTTPException(status_code=403, detail="无权导出此回忆录")
stmt = (
select(Chapter)
.where(Chapter.user_id == user_id, Chapter.is_active == True)
.options(
joinedload(Chapter.sections).joinedload(ChapterSection.image_record),
joinedload(Chapter.images),
)
.order_by(Chapter.order_index)
)
result = await self._db.execute(stmt)
chapters = list(result.unique().scalars().all())
asset_ids = collect_asset_ids_for_chapters(chapters)
asset_map = await signed_urls_for_asset_ids(self._db, asset_ids)
pdf_bytes = await pdf_service.generate_pdf(
book, chapters, asset_url_map=asset_map
)
return {
"pdf_base64": pdf_bytes.decode("latin1"),
"filename": f"{book.title}.pdf",
}
async def get_chapters(
self, user_id: str, is_new: bool | None = None
) -> List[dict]:
chapters = await repo.get_chapters_with_sections(
user_id, self._db, is_new_only=is_new
)
asset_ids: set[str] = set()
for ch in chapters:
asset_ids |= collect_asset_ids_for_chapter(ch)
asset_map = await signed_urls_for_asset_ids(self._db, asset_ids)
chapter_by_category: dict[str, Chapter] = {}
for ch in chapters:
if ch.category and ch.category not in chapter_by_category:
chapter_by_category[ch.category] = ch
all_chapters: List[dict] = []
for category in CHAPTER_ORDER:
ch = chapter_by_category.pop(category, None)
if ch:
await self._cleanup_unavailable_images(ch)
all_chapters.append(chapter_to_list_dict(ch, asset_url_map=asset_map))
else:
if is_new is True:
continue
all_chapters.append(
{
"id": f"placeholder_{category}",
"title": CHAPTER_CATEGORIES[category],
"category": category,
"order_index": STAGE_TO_ORDER.get(category, 999),
"status": "empty",
"summary": "",
"canonical_markdown": "",
"content": "",
"cover_asset": None,
"cover_image": None,
"images": [],
"sections": [],
"word_count": 0,
"updated_at": None,
"is_new": False,
"source_segments": [],
}
)
for ch in chapter_by_category.values():
await self._cleanup_unavailable_images(ch)
all_chapters.append(chapter_to_list_dict(ch, asset_url_map=asset_map))
return all_chapters
async def get_chapter(self, chapter_id: str, user_id: str) -> dict:
chapter = await repo.get_chapter_by_id(chapter_id, self._db)
if not chapter:
raise HTTPException(status_code=404, detail="Chapter not found")
if chapter.user_id != user_id:
raise HTTPException(status_code=403, detail="无权访问此章节")
if not chapter.is_active:
raise HTTPException(status_code=404, detail="Chapter not found")
await self._cleanup_unavailable_images(chapter)
asset_map = await signed_urls_for_asset_ids(
self._db, collect_asset_ids_for_chapter(chapter)
)
return chapter_to_dict(chapter, asset_url_map=asset_map)
async def disable_chapter(self, chapter_id: str, user_id: str) -> dict:
chapter = await self._db.get(Chapter, chapter_id)
if not chapter:
raise HTTPException(status_code=404, detail="Chapter not found")
if chapter.user_id != user_id:
raise HTTPException(status_code=403, detail="无权操作此章节")
chapter.is_active = False
await self._db.commit()
return {"status": "ok", "message": "章节已清除"}
async def regenerate_chapter(self, chapter_id: str, user_id: str) -> dict:
chapter = await self._db.get(Chapter, chapter_id)
if not chapter:
raise HTTPException(status_code=404, detail="Chapter not found")
if chapter.user_id != user_id:
raise HTTPException(status_code=403, detail="无权操作此章节")
# TODO: 实现重新整理逻辑
return {"status": "ok", "message": "Chapter regeneration triggered"}
async def get_memoir_state(self, user_id: str) -> dict:
from app.features.memoir.state_service import get_or_create_state
state = await get_or_create_state(user_id, self._db)
return state.model_dump()
async def get_next_question_context(self, user_id: str) -> dict:
from app.features.memoir.state_service import get_or_create_state
state = await get_or_create_state(user_id, self._db)
return {
"current_stage": state.current_stage,
"empty_slots": state.empty_slots_for_current_stage(),
"covered_stages": state.covered_stages,
}
async def check_and_trigger_cover_generation(self, user_id: str) -> dict:
"""
有正文、尚无 cover_asset、且 legacy 封面 MemoirImage 未 completed 时,
派发 generate_chapter_cover由 intent/asset 闭环完成)。
"""
from app.tasks.chapter_cover_tasks import generate_chapter_cover
img_settings = MemoirImageSettings.from_env()
if not img_settings.enabled:
return {"triggered": []}
chapters = await repo.get_chapters_with_sections(
user_id, self._db, active_only=True, is_new_only=None
)
triggered: List[str] = []
for ch in chapters:
if not ch.category or ch.status == "empty":
continue
if getattr(ch, "cover_asset_id", None):
continue
md = (ch.canonical_markdown or "").strip()
section_blob = "\n\n".join(
(s.content or "").strip()
for s in getattr(ch, "sections", None) or []
if (s.content or "").strip()
)
body = md or strip_legacy_image_placeholders(section_blob).strip()
if not body:
continue
images = getattr(ch, "images", None) or []
cover_rec = next(
(m for m in images if getattr(m, "section_id", None) is None),
None,
)
if (
cover_rec
and (getattr(cover_rec, "status") or "").strip() == "completed"
):
continue
try:
generate_chapter_cover.delay(ch.id)
triggered.append(ch.id)
logger.info("触发生成章节封面(asset): chapter=%s", ch.id)
except Exception as exc:
logger.warning("封面任务派发失败: chapter=%s, error=%s", ch.id, exc)
return {"triggered": triggered}
async def mark_memoir_read(self, user_id: str) -> dict:
stmt = select(Chapter).where(Chapter.user_id == user_id, Chapter.is_new == True)
result = await self._db.execute(stmt)
for chapter in result.scalars().all():
chapter.is_new = False
stmt_book = (
select(Book).where(Book.user_id == user_id).order_by(Book.updated_at.desc())
)
result_book = await self._db.execute(stmt_book)
book = result_book.scalar_one_or_none()
if book:
book.has_update = False
await self._db.commit()
return {"status": "ok"}