# AI 写作 — 后端开发文档 > **本文档目的**:帮助后端开发者快速理解 AI 写作子系统的代码结构、Graph 拓扑、数据流、持久化、容错策略和扩展点。 > > **配套文档**: > - 前端:`frontend-web/docs/ai-writing-frontend-dev.md` > - 故障排查:`frontend-web/docs/ai-writing-说明.md` > - 前端 API 字段索引:`frontend-web/docs/ai-writing-assistant-backend-api.md` --- ## 一、产品定位与技术架构 ### 产品定位 **杂志社协作写作子系统**:由 4 个 AI Agent(**素材收集专家 / 作家 / 编辑 / 系统**)协作完成一篇文章,通过 4 个用户干预点参与决策。 ### 技术栈 - **编排框架**:LangGraph(独立 StateGraph,不复用主 lead agent 的 middleware 链) - **节点驱动**:异步 LLM 调用(`langchain_openai` 等),支持流式 token 推送到 SSE - **状态持久化**:LangGraph checkpointer(SQLite 文件,跨重启可恢复) - **业务数据**:MySQL `ai_writing_sessions` 表 - **协议**:SSE(Server-Sent Events)推进度 + REST 控制流 ### 与主 Agent 的关系 **完全独立**。AI 写作 graph 不通过 `make_lead_agent` 也不走主 agent 的 18 层 middleware。它有自己的 StateGraph、自己的 checkpointer、自己的路由前缀(`/api/ai-writing/*`),跟主 agent 解耦。 --- ## 二、目录结构 ``` offline-backend-20260512/backend/ ├── app/gateway/ │ ├── routers/ │ │ └── ai_writing.py # FastAPI 路由 + 模块级状态字典 + SSE 生成器(1000+ 行,核心) │ ├── ai_writing_cleanup.py # 定期清理 + cron 调度(独立模块) │ └── app.py # lifespan 启动/停止清理任务、检测多 worker、关闭 checkpointer │ └── packages/harness/deerflow/ ├── agents/ai_writing/ │ ├── __init__.py # 暴露 AI_WRITING_GRAPH + 工厂函数 │ ├── graph.py # ⭐ StateGraph 定义 + 4 个 pause 节点 + checkpointer 工厂(538 行) │ ├── state.py # AIWritingState TypedDict(121 行) │ │ │ ├── nodes/ # ── Graph 节点 ── │ │ ├── _utils.py # LLM 工具:parse_json, ThinkTagParser, llm_json │ │ ├── intent_parser.py # 节点 1:意图解析 │ │ ├── researcher.py # 节点 2:素材检索(支持流式 progress) │ │ ├── writer_outline.py # 节点 3:大纲规划 │ │ ├── writer_draft.py # 节点 4:草稿生成(按章节流式) │ │ └── editor.py # 节点 5:编辑审核(3 维度并行流式) │ │ │ ├── prompts/ # ── Prompt 模板 ── │ │ ├── __init__.py │ │ ├── researcher_prompts.py │ │ ├── writer_prompts.py │ │ └── editor_prompts.py │ │ │ └── search/ # ── 素材检索 provider ── │ ├── base.py # SearchResult / SearchProvider 接口 │ ├── duckduckgo_search.py │ ├── internet_search.py # 互联网检索(默认) │ ├── intranet_stub.py # 内网检索占位(按部署接入) │ └── demo_search.py # 离线 demo 数据 │ └── persistence/ai_writing_sessions/ ├── __init__.py # AIWritingSessionRepository 工厂 ├── model.py # ORM:AIWritingSessionRow └── sql.py # CRUD 实现(含清理用的 list_older_than / delete_by_ids) ``` --- ## 三、Graph 拓扑(最重要) ### StateGraph 节点 + 边 ``` START │ ↓ ┌─────────────────┐ │ intent_parser │ 解析用户意图(确定 article_type / target_audience) └────────┬────────┘ ↓ ┌─────────────────┐ │ researcher │ ← (re_search 回到这里) └────────┬────────┘ ↓ (after_researcher: ⼟材 ready) ┌──────────────────────┐ │ pause_material │ ⚙️ PAUSE-1:素材确认 └─┬──────────┬─────────┘ │ confirm │ re_search ↓ └→ researcher ┌─────────────────┐ │ writer_outline │ ← (re_outline 回到这里) └────────┬────────┘ ↓ ┌──────────────────────┐ │ pause_outline │ ⚙️ PAUSE-2:大纲确认 └─┬──────────┬─────────┘ │ confirm │ re_outline ↓ └→ writer_outline ┌─────────────────┐ │ writer_draft │ ← (素材补充后回到这里) └─┬───────────────┘ ↓ (after_writer_draft) ┌──────┴────────────────┐ ↓ ↓ ┌──────────────────┐ ┌────────────────────┐ │ pause_section_ │ │ pause_draft │ ⚙️ PAUSE-3:草稿确认 │ help(素材不足)│ │ │ └─┬─────────────┬──┘ └─┬────────┬───┬─────┘ │ supplement │ del/ │ │ │ to_editor │↓ │ loose │ │ user_revise research writer_ │ final ↓ draft │ ize writer_draft │↓ done │ (to_editor)│ ↓ ┌─────────────────┐ │ editor │ 3 维度并行流式审核 └────────┬────────┘ ↓ (after_review) ┌──────────────────────┐ │ pause_review │ ⚙️ PAUSE-4:审核确认 └─┬──────────┬─────────┘ │ confirm │ revise │↓ └→ writer_draft done ``` ### 节点定义位置一览 | 节点名 | 实现文件 | 类型 | 流式? | |---|---|---|---| | `intent_parser` | `nodes/intent_parser.py` | async function | 否 | | `researcher` | `nodes/researcher.py` | async function | 是(research_progress) | | `pause_material` | `graph.py` 顶部 | async, 调 `interrupt()` | — | | `writer_outline` | `nodes/writer_outline.py` | async function | 是(outline_chunk) | | `pause_outline` | `graph.py` 顶部 | async, 调 `interrupt()` | — | | `writer_draft` | `nodes/writer_draft.py` | async function | 是(draft_chunk,按章节) | | `pause_section_help` | `graph.py` 顶部 | async, 调 `interrupt()` | — | | `pause_draft` | `graph.py` 顶部 | async, 调 `interrupt()` | — | | `editor` | `nodes/editor.py` | async function | 是(review_chunk,3 路并行) | | `pause_review` | `graph.py` 顶部 | async, 调 `interrupt()` | — | ### 条件路由函数(graph.py) | 函数 | 在哪个节点之后 | 路由依据 | |---|---|---| | `after_researcher` | researcher | 是否找到素材 | | `after_material_pause` | pause_material | `state["status"]` 是 researching 还是 writing_outline | | `after_outline_pause` | pause_outline | 同上,writing_outline 还是 writing_draft | | `after_writer_draft` | writer_draft | 是否有 blocked_sections | | `after_section_help_pause` | pause_section_help | status | | `after_draft_pause` | pause_draft | done / revising / reviewing | | `after_review` | editor | review_result.verdict 是 pass 还是 reject + revision_count 上限 | | `after_review_pause` | pause_review | done / writing_draft | --- ## 四、State 定义(state.py) ```python class AIWritingState(TypedDict): # —— 写作配置(启动时填充)—— user_intent: str article_type: str target_audience: str word_count_target: int keyword_count: int # 1-10,默认 6 writing_mode: str # "strict" | "loose" # —— Agent 产出 —— material_package: Optional[MaterialPackage] # 素材收集专家 current_outline: Optional[Outline] # 作家 drafts: Annotated[List[DraftArticle], operator.add] # 累计所有版本 current_draft: Optional[DraftArticle] # 最新草稿 review_result: Optional[ReviewResult] # 编辑 # —— 用户干预 —— last_user_intervention: Optional[UserInterventionPayload] blocked_sections: Optional[List[BlockedSection]] # 素材不足的章节 # —— 流程控制 —— revision_count: int max_revisions: int # 默认 3 status: WritingStatus # —— SSE 推送队列 —— progress_events: Annotated[List[dict], operator.add] # 累加,不覆盖 ``` ### State 字段的特殊语义 - **`Annotated[..., operator.add]`** — LangGraph 的累加 reducer,节点 return 时 `progress_events: [evt]` 会**追加**到现有列表,而不是替换 - **TypedDict** — 运行时是普通 dict,类型注解仅供 IDE 提示 - **节点输出** — 节点 return 的 dict 只需要包含**要更新的字段**,未指定的字段保持不变 --- ## 五、路由层(app/gateway/routers/ai_writing.py) ### 14 个接口一览 | HTTP | 路径 | 用途 | 鉴权 | |---|---|---|---| | POST | `/api/ai-writing/start` | 启动写作,返回 session_id | 用户 | | GET | `/api/ai-writing/{id}/stream` | SSE 进度流 | 用户 | | POST | `/api/ai-writing/{id}/resume` | 提交用户干预 | 用户 | | GET | `/api/ai-writing/{id}/state` | 查 graph 状态快照(调试用) | 用户 | | GET | `/api/ai-writing/sessions` | 列出当前用户的历史会话 | 用户 | | GET | `/api/ai-writing/sessions/{id}` | 获取单个历史会话 | 用户 | | PATCH | `/api/ai-writing/sessions/{id}` | 重命名会话 | 用户 | | DELETE | `/api/ai-writing/sessions/{id}` | 删除会话 | 用户 | | PUT | `/api/ai-writing/sessions/{id}/transcript` | 保存对话时间线(**fire-and-forget**) | 用户 | | GET | `/api/ai-writing/article-types` | 列出文章类型 | 用户 | | POST | `/api/ai-writing/article-types` | 新增文章类型 | 用户(建议改成 admin) | | PATCH | `/api/ai-writing/article-types/{id}` | 更新文章类型 | 用户 | | DELETE | `/api/ai-writing/article-types/{id}` | 删除文章类型 | 用户 | | POST | `/api/ai-writing/admin/cleanup` | **手动触发清理** | admin only | ### 模块级状态字典(关键设计 — 也是 worker 限制的根源) ```python # 每个 session_id 对应一份内存状态。多 worker 部署时每个 worker 一份,互不共享。 _event_queues: dict[str, asyncio.Queue] # SSE 事件队列 _session_pending_pause: dict[str, tuple] # 挂起的 await_user 事件(供重连重投) _session_model_names: dict[str, str] # 用户选的模型 _session_interventions: dict[str, list] # 累计干预历史 _session_repos: dict[str, repo] # DB 仓库引用 _session_user_ids: dict[str, str] # 用户 id _session_article_types: dict[str, list] # 动态文章类型列表 _session_error_flags: dict[str, bool] _inflight_graph_tasks: set[asyncio.Task] # in-flight task 注册表(lifespan drain 用) _recovery_lock: asyncio.Lock # 跨 worker 会话恢复的并发锁 ``` ⚠️ **关键限制**:这些都是 per-process 内存。多 worker 部署下,worker B 看不到 worker A 的状态 → SSE 事件分发失败。详见 `frontend-web/docs/ai-writing-说明.md` 第五、六章。 ### `_stream_graph_segment`(graph 推进函数) `/start` 和 `/resume` 都调用这个函数(用 `_spawn_graph_task` 包装),负责: 1. 调 `graph.astream(input_or_command, config=config, stream_mode="updates")` 2. 把每个 chunk 解析成 SSE 事件 put 进 `_event_queues[session_id]` 3. 遇到 `__interrupt__` → 推 `await_user` 事件 → 返回(不关闭 queue) 4. 遇到 END → 推 `done` 事件 → 关闭 queue(put None 哨兵) 5. 任何异常 → 推 `error` 事件 → 关闭 queue 6. finally 段:更新 DB 最终状态、清理所有模块字典 ### SSE 生成器 `_sse_generator` `GET /stream` 的实现,从 `_event_queues[session_id]` 拿事件 → `_sse(event, data)` 序列化 → yield 给 fastapi `StreamingResponse`。 **关键逻辑**: - 队列不存在时**等待 5 秒**,期间触发 `_recover_session`(跨 worker / 进程重启恢复) - `_sse()` 三层兜底:`default=str` → 手写 notice 帧 → 永不抛 TypeError - 收到 None 哨兵 → 关闭流 ### `_recover_session`(跨 worker 恢复) 当 worker 内存里没有 session_id 但 sqlite checkpointer 有快照时: 1. 从 checkpointer 读出 state snapshot 2. 重建 `_event_queues[sid]` 队列 3. 从 snapshot 提取当前 interrupt 负载,put 进队列 4. 同时尽量恢复 `_session_repos`、`_session_user_ids` 等 返回三种状态: - `"interrupted"` — 图停在干预点,已重建队列 - `"finished"` — 图已完成或中断在节点中间,推 done/error 事件后关闭 - `"unrecoverable"` — sqlite 也没有这个 thread_id → 真的丢了 ### `_spawn_graph_task` / `wait_for_inflight_graph_tasks`(优雅关闭) - `_spawn_graph_task(coro)` — 启 task 并自动登记到 `_inflight_graph_tasks`,完成时反登记 - `wait_for_inflight_graph_tasks(timeout)` — lifespan shutdown 等所有 task 跑到下一个 interrupt / END(默认 30s),超时则 cancel 强制结束 --- ## 六、Graph 与持久化(graph.py) ### Lazy 单例 + 自动降级 ```python AI_WRITING_GRAPH = None # 模块级槽位(测试可 monkeypatch) _async_checkpointer_singleton = None # checkpointer 单例 async def get_ai_writing_graph(): """懒加载 graph。第一次调用时按 config.yaml 决定 checkpointer。""" global AI_WRITING_GRAPH if AI_WRITING_GRAPH is not None: return AI_WRITING_GRAPH checkpointer = await _get_async_checkpointer() AI_WRITING_GRAPH = build_ai_writing_graph(checkpointer=checkpointer) return AI_WRITING_GRAPH ``` ### Checkpointer 优先级(_get_async_checkpointer) ``` 1. config.yaml 的 checkpointer 段 ├── type: sqlite → AsyncSqliteSaver(aiosqlite.connect(...)) ├── type: postgres → AsyncPostgresSaver(AsyncConnectionPool(...)) └── type: memory → InMemorySaver 2. 任何初始化失败 → InMemorySaver(降级,打 warning) 3. 未配置 checkpointer 段 → InMemorySaver(降级,打 warning) 4. 未知 type(如 "mysql") → InMemorySaver(降级,打 warning) ``` **核心契约**:`_get_async_checkpointer()` **绝不抛异常**,永远返回一个可用的 saver。这是「写作流程能完整继续」契约的关键。 ### Shutdown `shutdown_ai_writing_graph()` 在 lifespan shutdown 时调用: - 关闭 `aiosqlite.Connection` / `AsyncConnectionPool` - 清空 `AI_WRITING_GRAPH` / `_async_checkpointer_singleton` 调用顺序(详见 `app/gateway/app.py` lifespan): 1. 停清理任务 2. `wait_for_inflight_graph_tasks(30s)` 3. `shutdown_ai_writing_graph()` --- ## 七、节点详解(nodes/) ### `nodes/_utils.py`(公共工具) | 函数 / 类 | 用途 | |---|---| | `get_model(config, model_name)` | 从 config.configurable 拿 model_name 或回退默认 | | `parse_json(text, fallback)` | 鲁棒 JSON 解析(支持代码块、不完整 JSON、转义错误) | | `_extract_balanced(text)` | 找出第一对平衡的 `{...}` | | `_repair_json(text)` | 修复常见的 JSON 错误(缺尾引号等) | | `ThinkTagParser` | 流式解析 `...` 标签(思考 / 正文分流) | | `llm_json(system, user, fallback, config)` | 调 LLM 并解析 JSON 输出 | | `count_words(text)` | 中文 + 英文混合计字数 | ### `nodes/intent_parser.py` — 意图解析 - 调 LLM 把用户意图(`user_intent` 字符串)解析成结构化字段:`article_type`、`target_audience`、`word_count_target` - 用户已经在 `AIWritingRequest` 里指定的字段会**覆盖** LLM 推断 - 输出:更新 state 的对应字段,推 `step_done` 事件 ### `nodes/researcher.py` — 素材检索 - 根据 `user_intent` 用 LLM 生成 N 个关键词(N = `keyword_count`,默认 6) - 用 `search/*` 的 provider 并行检索 - 流式推 `research_progress` 事件("正在搜索关键词 X 中...") - 输出:`material_package = {keywords, materials, summary}`,推 `materials_ready` ### `nodes/writer_outline.py` — 大纲规划 - 输入:`user_intent` + `material_package` - 调 LLM 流式生成大纲 JSON,期间推 `outline_chunk`(正文 + thinking 分流) - 解析最终 JSON → `current_outline`,推 `outline_ready` ### `nodes/writer_draft.py` — 草稿生成(最复杂) - 按大纲**逐章节**调 LLM 生成正文,每章节独立流式推 `draft_chunk` - 检测 LLM 是否「拒答」(`_looks_like_refusal`),如果是则换提示重试 - **严格模式下**:如果章节关联的素材不足,把该章节加入 `blocked_sections`,触发 `pause_section_help` - 完成所有章节后拼成 `current_draft.full_markdown`,推 `draft_ready` ### `nodes/editor.py` — 编辑审核(3 维度并行) - 3 个维度(**fact / logic / language**)并行调 LLM - 每个维度独立流式推 `review_chunk`,完成后推 `review_item_done` - 三个都完成后聚合成 `review_result`,推 `review_ready` - 评分加权(fact 40% + logic 30% + language 30%),超过 `pass_threshold`(默认 80)verdict = `pass` --- ## 八、持久化(persistence/ai_writing_sessions/) ### `model.py` — ORM ```python class AIWritingSessionRow(Base): __tablename__ = "ai_writing_sessions" id: Mapped[str] # = session_id(uuid) user_id: Mapped[str | None] # NULL = no-auth 模式 title: Mapped[str] # 从 user_intent 截取前 50 字 user_intent: Mapped[str | None] status: Mapped[str] # in_progress / draft_ready / review_ready / done / error draft_title: Mapped[str | None] draft_markdown: Mapped[str | None] completed_interventions: Mapped[str | None] # JSON 字符串 review_result: Mapped[str | None] # JSON 字符串 transcript: Mapped[str | None] # PortableLongText(MySQL→LONGTEXT,4GB) created_at / updated_at: BeijingDateTime ``` ### `sql.py` — Repository | 方法 | 用途 | |---|---| | `create(id, title, user_intent, user_id)` | 启动写作时建一行 | | `get(session_id, user_id)` | 取单个(带 user_id 权限校验) | | `list(user_id, limit)` | 列出当前用户的会话 | | `update(session_id, user_id, ...)` | 更新指定字段(其他保持不变) | | `delete(session_id, user_id)` | 删除(带权限校验) | | `list_older_than(cutoff, only_finished, limit)` | 清理用:列出过期 session_id | | `delete_by_ids(session_ids)` | 清理用:批量删除 | ### Alembic 迁移 - `20260517_03_ai_writing_transcript.py` — 加 transcript 列(TEXT) - `20260519_02_ai_writing_transcript_longtext.py` — MySQL 上把 transcript 改成 LONGTEXT(修 64KB 限制导致的 500) --- ## 九、定期清理(ai_writing_cleanup.py) ### 设计目标 - 删超过 `retention_days`(默认 7)的旧会话 - 联动清理:MySQL 业务行 + SQLite checkpointer 里的 thread 状态 - **绝对不删正在写的会话**(活跃保护) ### 关键函数 ```python async def run_cleanup_once(*, retention_days, only_finished=False, batch_size=500) -> dict ``` 返回 `{deleted_db, deleted_checkpoints, skipped_active}`。 执行流程: 1. 调 `repo.list_older_than(cutoff, only_finished, batch_size)` 拿过期 id 2. 从 `_event_queues.keys()` 拿活跃 session_id 集合 → **过滤掉所有仍在跑的** 3. 对每个待删 id:调 `checkpointer.adelete_thread(id)` 删 checkpoint 4. 调 `repo.delete_by_ids(ids)` 批量删 MySQL 5. 循环直到本批 < batch_size ### Cron 调度 `cleanup_scheduler_loop()` 常驻后台任务: - 用 `croniter` 解析 `ai_writing.cleanup.cron`(默认 `"0 3 * * *"`,每天凌晨 3 点) - `asyncio.sleep` 到下次触发时间 - 多 worker 用文件锁(`runtime_home() / ".ai_writing_cleanup.lock"`)做 leader election ### Admin 手动触发 `POST /api/ai-writing/admin/cleanup`: - 鉴权:`system_role == "admin"` - 可选覆盖 `retention_days` / `only_finished` / `batch_size` - 即使传 `retention_days: 0` 也不会秒杀活跃会话(活跃保护硬保证) --- ## 十、配置(config.yaml) ```yaml # Checkpointer:graph 状态持久化 checkpointer: type: sqlite # sqlite | postgres | memory connection_string: .deer-flow/data/checkpoints.db # 清理调度 ai_writing: cleanup: enabled: true retention_days: 7 cron: "0 3 * * *" # 每天凌晨 3 点 delete_only_finished: false batch_size: 500 # (可选)素材检索 provider 配置 search: provider: deerflow.agents.ai_writing.search.internet_search:WebSearchProvider ... ``` `_ai_writing_config()` 函数从 `get_app_config().ai_writing` 里读 dict,节点通过 `config.configurable.ai_writing_config` 拿到。 --- ## 十一、SSE 事件协议(后端 → 前端契约) 后端通过 `_sse(event_type, data)` 推 `event: \ndata: \n\n` 格式。 ### 推出事件的位置 | 事件 | 推出位置 | |---|---| | `connected` | `_sse_generator` 入口 | | `step_start` / `step_done` | 节点内部主动 emit | | `research_progress` | `researcher_node` 内的 `_emit` | | `materials_ready` | `_stream_graph_segment` 检测到 `material_package` 字段 | | `outline_chunk` | `writer_outline_node` 流式推(通过 `sse_queue`) | | `outline_ready` | `_stream_graph_segment` 检测到 `current_outline` | | `draft_chunk` | `writer_draft_node` 按章节流式推 | | `draft_ready` | `_stream_graph_segment` 检测到 `current_draft` | | `review_chunk` / `review_item_done` | `editor_node` 3 维度并行流式 | | `review_ready` | `_stream_graph_segment` 检测到 `review_result` | | `await_user` | `_stream_graph_segment` 捕获 `__interrupt__` | | `resumed` | pause 节点 return 时塞进 `progress_events` | | `notice` | 节点内手动推(非致命告警) | | `done` | `_stream_graph_segment` 正常结束 | | `error` | `_stream_graph_segment` except 分支 | ### 节点如何流式推 chunk 通过 `config.configurable.sse_queue`: ```python async def writer_outline_node(state, config): sse_queue = config.get("configurable", {}).get("sse_queue") sse_session_id = config.get("configurable", {}).get("sse_session_id") async for chunk in llm.astream(prompt): if sse_queue: await sse_queue.put(("outline_chunk", { "session_id": sse_session_id, "text": chunk.content, "is_thinking": False, })) return {"current_outline": parsed, ...} ``` `sse_queue` 是路由层在 `_make_config` 里注入的 `_event_queues[session_id]`。 --- ## 十二、扩展指南:常见开发任务 ### 任务 1:新增一个 Graph 节点 **例**:在 `editor` 之后加一个「术语统一」节点。 1. **写节点**:在 `nodes/` 新建 `terminology.py`: ```python async def terminology_node(state, config): # ... 读 current_draft 改术语 return { "current_draft": new_draft, "progress_events": [{"type": "step_done", "agent_name": "术语统一", "message": "..."}], } ``` 2. **注册到 graph**:在 `graph.py` 的 `build_ai_writing_graph` 里: ```python from .nodes.terminology import terminology_node g.add_node("terminology", terminology_node) # 改 editor 后的边 g.add_edge("editor", "terminology") g.add_edge("terminology", "pause_review") ``` 3. **State 加字段**(如果需要):`state.py` 加新字段到 `AIWritingState` 4. **前端事件支持**(如果要新事件类型):见前端文档「扩展指南」 5. **测试**:在 `tests/` 加单元测试,至少测节点能正常 return 和 graph 拓扑能编译 ### 任务 2:新增一个用户干预点 1. **写 pause 节点**:在 `graph.py` 顶部加: ```python async def pause_xxx(state): user_input = interrupt({ "pause_point": "xxx_confirm", "message": "...", "context_field": state["..."], }) action = user_input.get("action", "confirm") if action == "confirm": return {"status": "next_status", "progress_events": [{"type": "resumed", ...}]} # ... 其他分支 ``` 2. **路由层 ResumeRequest 加字段**(如果干预需要新字段): ```python class ResumeRequest(BaseModel): # ... new_field: Optional[str] = None ``` 3. **`resume_writing` 路由把字段塞进 Command(resume=...)**: ```python intervention = { # ... "newField": body.new_field, } ``` 4. **graph 拓扑加 pause 节点的边**:在 build_ai_writing_graph 里 5. **前端配合**:见前端文档 ### 任务 3:新增一种素材检索 provider 1. 在 `agents/ai_writing/search/` 加新文件,实现 `SearchProvider` 接口(参考 `internet_search.py`): ```python class MyProvider(SearchProvider): async def search(self, keyword: str, limit: int) -> list[SearchResult]: ... ``` 2. 配置 `config.yaml`: ```yaml ai_writing: search: provider: my_pkg.search:MyProvider # provider 特定的配置... ``` 3. `researcher_node` 通过 `config.configurable.ai_writing_config.search.provider` 用 `resolve_class` 动态加载 ### 任务 4:修改 LLM Prompt 所有 prompt 集中在 `agents/ai_writing/prompts/`: - `researcher_prompts.py` — 检索专家 - `writer_prompts.py` — 大纲 + 草稿 - `editor_prompts.py` — 审核 改 prompt 后**强烈建议**: 1. 在 `tests/` 加针对该 prompt 输出的回归测试(用 LLM mock 或固定输入) 2. 跑现有测试确认拓扑兼容 ### 任务 5:扩展持久化字段 例:给 `ai_writing_sessions` 加一个 `total_tokens` 列。 1. `model.py` 加列: ```python total_tokens: Mapped[int | None] = mapped_column(Integer, nullable=True) ``` 2. **必须**新建 alembic 迁移: ```bash PYTHONPATH=. uv run alembic revision -m "add total_tokens column" ``` 然后填 upgrade/downgrade 逻辑(参考 `20260519_02_ai_writing_transcript_longtext.py`) 3. `sql.py` 的 `update()` / `to_dict()` 加字段处理 4. `routers/ai_writing.py` 的 `SessionResponse` 加字段 5. 测试:在 `tests/test_ai_writing_session_repo.py` 加新字段的 round-trip 测试 --- ## 十三、容错策略(关键设计原则) ### 「写作流程能完整继续」契约 后端承诺:**除非 Pydantic 校验失败(4xx),AI 写作的 4 个核心接口任何后端组件失败都不阻断当次写作的完整执行**。 ### 落实位置 | 失败点 | 兜底位置 | 行为 | |---|---|---| | sqlite checkpointer 初始化失败 | `_get_async_checkpointer` | 自动降级 InMemorySaver | | SSE JSON 序列化失败 | `_sse()` | 三层兜底(default=str → notice 帧) | | DB `create()` session 行失败 | `start_writing` | try/except 吞掉,session_id 仍返回 | | DB `update(draft)` 失败 | `_stream_graph_segment` | try/except 吞掉,graph 继续推 | | DB `update(review)` 失败 | 同上 | 同上 | | DB `update(transcript)` 失败 | `save_transcript` 路由 | 返回 `success=false`,不抛 5xx | | 文章类型加载失败 | `start_writing` | 退回 `[]` | | Graph 编译失败 | `start_writing` | 抛 503(极端情况,节点 import 错) | | Graph 运行时节点抛错 | `_stream_graph_segment` | except 推 error 事件,session 结束 | | Cleanup 误删活跃会话 | `run_cleanup_once` | 活跃保护硬过滤 | ### 容器优雅关闭 lifespan shutdown 顺序: 1. 停 cleanup scheduler 2. `wait_for_inflight_graph_tasks(30s)` — 等所有 in-flight 写作跑到下一个 checkpoint 3. 超时未结束的 task `cancel()` + 2s ack 等待 4. `shutdown_ai_writing_graph()` — 关闭 aiosqlite Connection / pg pool --- ## 十四、测试覆盖 | 测试文件 | 数量 | 覆盖范围 | |---|---|---| | `tests/test_ai_writing_session_repo.py` | 7 | Repo CRUD + 清理用的 list_older_than / delete_by_ids | | `tests/test_ai_writing_session_recovery.py` | 8 | `_recover_session` / `_extract_interrupt` / `resume_writing` 跨 worker 恢复 | | `tests/test_ai_writing_transcript_endpoint.py` | 5 | transcript 保存兜底 | | `tests/test_ai_writing_cleanup.py` | 26 | 清理 + drain + checkpointer 降级 + SSE 序列化 + admin 接口 + 多 worker 检测 | | `tests/test_portable_long_text.py` | 5 | LONGTEXT 列类型 + 迁移链 | 跑全部: ```bash PYTHONPATH=. uv run pytest tests/test_ai_writing*.py tests/test_portable_long_text.py -v ``` 期望 **51 个测试全过**。 --- ## 十五、关键调试 / 排查命令 ### 看 graph 当前状态 ```python # 进 python shell from deerflow.agents.ai_writing.graph import get_ai_writing_graph graph = await get_ai_writing_graph() snap = await graph.aget_state({"configurable": {"thread_id": ""}}) print(snap.values, snap.next, snap.tasks) ``` ### 直接读 sqlite checkpointer ```bash sqlite3 /app/.deer-flow/data/checkpoints.db .tables # 看表(checkpoints / checkpoint_writes 等) SELECT thread_id, checkpoint_ns, checkpoint_id FROM checkpoints WHERE thread_id = '' LIMIT 5; ``` ### 查 MySQL 业务行 ```sql SELECT id, status, draft_title, updated_at, JSON_LENGTH(transcript) AS transcript_size FROM ai_writing_sessions WHERE id = ''; ``` ### 查内存中的活跃会话(调试时手动加) ```python # 临时在路由加一个 admin 接口 @router.get("/admin/inflight") async def admin_inflight(): return { "event_queues": list(_event_queues.keys()), "inflight_tasks": len(_inflight_graph_tasks), "pending_pause": list(_session_pending_pause.keys()), } ``` ### 日志关键字 ```bash docker logs 2>&1 | grep -E 'ai_writing|AsyncSqliteSaver|checkpointer|recovered' ``` 期望看到(启动时): - `ai_writing: 使用 AsyncSqliteSaver (xxx.db)` — 持久化生效 - `ai_writing_cleanup: 已启动 cron=...` — 清理调度起来 警告: - `ai_writing: 持久化 checkpointer 不可用,退回 InMemorySaver` — sqlite 故障 - `⚠ 检测到 WEB_CONCURRENCY=N` — 多 worker 警告(详见 `ai-writing-说明.md`) --- ## 十六、性能与限制 ### 单 session 资源占用 - 内存:~50KB(state + progress_events 累计) - sqlite checkpoint:每次 state update 写一行(一篇完整文章约 20-50 行) - LLM tokens:约 30K-100K(取决于素材量 + 字数 + 修订次数) ### 并发上限 - 单 worker 内:理论上百级 session 并发(IO bound),实际 LLM API rate limit 才是瓶颈 - 多 worker:**当前 in-memory dict 隔离,详见 ai-writing-说明.md 第六章** ### 长会话风险 - progress_events 累加无界 → 长会话 state 可能膨胀到几 MB - transcript JSON 同理 → MySQL 行可能上 MB(但 LONGTEXT 列上限 4GB,影响读取性能) ### 优化方向 - progress_events 节流(聚合 chunk 事件) - transcript 增量保存(diff 而非全量覆盖) - 节点级 LLM 调用加 retry(tenacity) --- ## 十七、与主 Agent 的关系 | 维度 | 主 Agent | AI 写作 | |---|---|---| | Graph 入口 | `make_lead_agent` | `AI_WRITING_GRAPH` | | Middleware | 18 层中间件 | 无(直接节点链) | | State | `ThreadState` | `AIWritingState` | | Checkpointer | `make_checkpointer()`(共享) | `_get_async_checkpointer()`(同一文件,不同 thread_id 命名空间) | | 路由前缀 | `/api/langgraph/*` `/api/threads/*` | `/api/ai-writing/*` | | 流式协议 | LangGraph stream_mode(messages-tuple/values) | 自定义 SSE 事件 | | 用户干预 | 通过 `ask_clarification` 工具 | 通过 LangGraph `interrupt()` | | 持久化业务表 | `threads_meta`、`runs` | `ai_writing_sessions` | 二者**完全独立**,仅在「共享同一个 sqlite checkpointer 文件」这一点上有间接关联(也可以拆开,性能层面没区别)。 --- ## 十八、未来 RFC 方向 | 主题 | 描述 | |---|---| | 多 worker 部署支持 | 把 `_event_queues` 等搬 Redis Pub/Sub,参考 `ai-writing-说明.md` 方案 B | | Graph 拓扑可配置 | 现在节点链路写死,未来按 article_type 走不同拓扑(如「研报」vs「散文」) | | 节点级重试 | LLM 抽风时自动重试 N 次,参考 `tenacity` | | 素材跨 session 复用 | 同一用户多次写作可能用同一批素材 → 缓存层 | | 多模态素材 | 支持图片素材输入和图文混合输出 | | 中文标点 / 排版规则注入 | 在节点 prompt 里强化中文出版规范 | --- ## 附录 A:相关文件速查 | 文件 | 行数 | 用途 | |---|---|---| | `routers/ai_writing.py` | 1006 | 路由 + 状态字典 + SSE + drain | | `agents/ai_writing/graph.py` | 538 | StateGraph + checkpointer 工厂 | | `agents/ai_writing/state.py` | 121 | AIWritingState 定义 | | `agents/ai_writing/nodes/_utils.py` | 242 | LLM 工具 | | `agents/ai_writing/nodes/intent_parser.py` | 131 | 节点:意图解析 | | `agents/ai_writing/nodes/researcher.py` | 271 | 节点:素材检索 | | `agents/ai_writing/nodes/writer_outline.py` | 147 | 节点:大纲规划 | | `agents/ai_writing/nodes/writer_draft.py` | 318 | 节点:草稿生成 | | `agents/ai_writing/nodes/editor.py` | 180 | 节点:编辑审核 | | `app/gateway/ai_writing_cleanup.py` | ~280 | 清理 + 调度 | | `persistence/ai_writing_sessions/sql.py` | ~200 | Repository | | `persistence/ai_writing_sessions/model.py` | ~35 | ORM | --- ## 附录 B:常见踩坑 | 现象 | 根因 | 解决 | |---|---|---| | `/start` 报 500 | graph 编译失败(多半是 import 错) | 看启动日志 stack trace | | `/resume` 报 404 | session_id 不存在或 checkpointer 没快照 | 看 `ai_writing recover` 日志 | | transcript 500 | DB 列容量 / 网络抖 | ✅ 已修,前端 fire-and-forget | | SSE 卡住不动 | 多 worker SSE 隔离 / 反代超时 | 见 `ai-writing-说明.md` | | `draft_chunk` 丢失内容 | `ThinkTagParser` 在 `` 标签边界处理 | 看 `_utils.py:ThinkTagParser.feed` 单测 | | graph 卡在某节点不退出 | LLM API 卡住 / 无限重试 | 看节点日志,加 timeout | | cleanup 没跑 | cron 表达式配错 / 文件锁被其他 worker 持有 | 看 `ai_writing_cleanup` 日志 | | 草稿在前端看不到但 MySQL 有 | 多 worker 导致 SSE 事件丢失 | 同上 | | 配 mysql 做 checkpointer 报错 | LangGraph 不支持 MySQL checkpointer | 用 sqlite 或 postgres | --- 如需补充内容,请直接编辑本文件,并在 git commit message 里注明改动原因。