docs(chain): 多章工作流链正式设计文档(实施契约)
This commit is contained in:
211
docs/design/chain-workflow.md
Normal file
211
docs/design/chain-workflow.md
Normal file
@@ -0,0 +1,211 @@
|
||||
# 设计文档 · 多章工作流链(Chain Workflow)
|
||||
|
||||
> 状态:设计稿(实施前契约)· 分支 `feat/chain-workflow` · 锚点 ARCHITECTURE §5.2/§5.5、CLAUDE.md 架构不变量 #1/#3/#5/#9、DEV_PLAN §4
|
||||
>
|
||||
> 本文是实施的**唯一契约源**。所有实施 agent 先读本文,按其中签名/事件/DoD 落地;与现有 spec 冲突时以更具体者为准并回写。
|
||||
|
||||
---
|
||||
|
||||
## 1. 背景与目标
|
||||
|
||||
### 1.1 为什么做
|
||||
对标竞品「星月写作」的核心能力是**一键多章工作流**(导语→大纲→细纲→一键十章正文循环),但其纯走量、无一致性闸,靠提示词补"关联前文/不跑偏"。本产品的差异化:**同样的链式量产,但每一章都过 `assemble → write → 四审 → digest` 闭环 → 批量生成而长篇不崩**。这是星月架构上做不到的杀手锏。
|
||||
|
||||
### 1.2 为什么用 LangGraph(而非端点+DB 循环)
|
||||
多章循环是本项目**第一个真正需要 LangGraph 的场景**:单章流程是线性的(写一发、审一次 `ainvoke`、验收一个事务),无需图;而"写→审→判定→循环/中断、跨小时可恢复、遇冲突暂停交人"是**有环、可持久、可中断**的编排 = LangGraph 本职。本设计借此**激活当前休眠的资产**:`checkpointer`(现恒 `None`)、`interrupt()`(现未实现)、`build_write_graph`(现死代码——将被本链图取代/吸收)。
|
||||
|
||||
### 1.3 目标(DoD)
|
||||
1. 一条内置链 `draft_volume`:从指定起始章,循环写 K 章,每章 `write → 四审 → 判定 → 验收`,digest 入库。
|
||||
2. **杀手锏 interrupt-on-conflict**:链自动跑;某章四审报出冲突 → `interrupt()` 暂停 → 前端取该章冲突交人裁决 → resume 续跑。无冲突章自动 accept。
|
||||
3. 跑在 `run_job` 长任务壳里:`POST .../chains/{key}/run` → 202 `{job_id}`;`GET /jobs/{id}` 看进度("3/10 章已过审");`POST .../chains/runs/{job_id}/resume` 带裁决续跑。
|
||||
4. **后端门禁绿**(ruff/format/mypy/alembic 无漂移/pytest)+ **E2E**(mock 网关零 token)走通:无冲突全自动两章 + 冲突→interrupt→resume 两路径。
|
||||
|
||||
### 1.4 非目标(本期不做)
|
||||
- **前端 UI**(链发起页/进度/裁决续跑面板)→ 契约稳定后单独一刀(@frontend follow-up)。
|
||||
- **可编辑/可分享的 UGC 链市场**(星月式)→ 需多租户/Auth,已 deferred,单独立项。
|
||||
- **续写/扩写 节点** → 可作为链的可选步后补,本期链步固定为 write。
|
||||
- 多种 HITL 档位 → 本期只做 `auto`(无冲突自动过、有冲突才停)。
|
||||
|
||||
---
|
||||
|
||||
## 2. 架构总览
|
||||
|
||||
```
|
||||
POST .../chains/draft_volume/run {start_chapter_no, count}
|
||||
│ (写 jobs 行 status=queued, kind="chain") → 202 {job_id}
|
||||
▼ BackgroundTasks → run_chain_job(job_id, resume=None)
|
||||
┌──────────────────────── LangGraph ChainGraph (thread_id = job_id) ─────────────────────┐
|
||||
│ START → write_chapter → review_chapter → decide ──(无冲突)──→ accept_chapter ─┐ │
|
||||
│ ▲ │ │ │
|
||||
│ │ (有冲突) │ │
|
||||
│ │ ▼ │ │
|
||||
│ │ interrupt() ⏸ (job=awaiting_input) │ │
|
||||
│ │ │ │
|
||||
│ └────────────────(还有章: current<=last)──── advance ◄──────────────┘ │
|
||||
│ │ │
|
||||
│ (写满 K 章) → END │
|
||||
└────────────────────────────────────────────────────────────────────────────────────────┘
|
||||
│ checkpointer (Postgres runtime / MemorySaver tests) 持久控制流位置 + pending handle
|
||||
▼
|
||||
POST .../chains/runs/{job_id}/resume {decisions:[...]}
|
||||
│ (Command(resume=decisions)) → 202 → run_chain_job(job_id, resume=decisions)
|
||||
▼ 从 interrupt 续跑:accept_chapter 用 decisions → advance → 下一章…
|
||||
```
|
||||
|
||||
**关键纪律**:节点做 DB IO 时各自用**短事务、独立 session**(每章 accept 章原子提交,仿 `run_overdue_scan` 自建 session)。checkpointer **只存控制流**(`current_chapter_no` / `written` / `pending`),章正文落 `chapters`、审稿落 `chapter_reviews`——**resume 时重读领域表**为真相(不变量 #5)。
|
||||
|
||||
---
|
||||
|
||||
## 3. 状态与节点
|
||||
|
||||
### 3.1 State 类型 `ChainState`(新增 `orchestrator/chain/state.py`)
|
||||
TypedDict(langgraph 状态 = 检查点内容,**只放控制流,不放正文/上下文**):
|
||||
|
||||
| 字段 | 类型 | 说明 |
|
||||
|---|---|---|
|
||||
| `project_id` | `uuid.UUID` | |
|
||||
| `user_id` | `uuid.UUID` | STUB_OWNER_ID |
|
||||
| `chain_key` | `str` | 固定 `"draft_volume"` |
|
||||
| `start_chapter_no` | `int` | 起始章 |
|
||||
| `last_chapter_no` | `int` | = start + count - 1(含) |
|
||||
| `current_chapter_no` | `int` | 当前处理章;初值 = start |
|
||||
| `written` | `list[int]` | 本 run 已验收章号(进度来源) |
|
||||
| `has_conflicts` | `bool` | `decide` 节点写入;控制流标志,不存冲突明文 |
|
||||
|
||||
> 草稿正文/上下文**不进** state(落 DB)。冲突明文**不进** state(resume 时从 `chapter_reviews` 重读,不变量 #5)。
|
||||
|
||||
### 3.2 节点(`orchestrator/chain/nodes.py`,gateway/session_factory 经闭包绑入工厂)
|
||||
|
||||
1. **`write_chapter`**:新建读写 session → `assemble(project_id, current_chapter_no)` 组装记忆(确定性,不变量 #6/#9:硬规则+指纹进缓存前缀)→ `gateway.run`(收集版写章,复用 `write_node` 同一 prompt 组装路径,非 SSE)→ `chapter_repo.save_draft` 落 `chapters` draft → commit。State 不变(正文在 DB)。
|
||||
2. **`review_chapter`**:复用四审。两种实现取其一(实施者选更省的):(a) 内联调用 `build_review_graph(gateway, review_repo).ainvoke(review_state)`;(b) 直接 `run_review × REVIEW_SPECS` + `collect_reviews`。落 `chapter_reviews` 留痕(不变量 #3 只读,不写正文)→ commit。
|
||||
3. **`decide`**(确定性纯逻辑,无 LLM):读该章最近审稿留痕 → `has_conflicts = len(conflicts) > 0`。写入 state。
|
||||
4. **条件边 `decide → ?`**:
|
||||
- `has_conflicts == True` → `interrupt({chapter_no})`:`interrupt()` 暂停,恢复值 = 该章 `decisions` 列表。恢复后进 `accept_chapter`。
|
||||
- `has_conflicts == False` → 直接进 `accept_chapter`(`decisions = []`)。
|
||||
5. **`accept_chapter`**:新建 session → 从 `chapter_reviews` **重读**该章最近审稿(不信 state)→ `assert_conflicts_resolved(latest_review, decisions)`(缺判 → `CONFLICT_UNRESOLVED`,链失败)→ 事务外 `extract_digest_facts`(从终稿)→ `run_accept_transaction(...)` 章原子提交 → 提交后 `run_overdue_scan`(自建 session,伏笔到期)。`written.append(current_chapter_no)`。
|
||||
6. **条件边 `accept_chapter → ?`**:`current_chapter_no < last_chapter_no` → `advance`(`current_chapter_no += 1`)→ 回 `write_chapter`;否则 → `END`。
|
||||
|
||||
> `advance` 可并入 `accept_chapter` 的返回(`current_chapter_no+1`)+ 一个 `should_continue` 条件边判 START/END,避免独立节点。实施者择简。
|
||||
|
||||
### 3.3 图工厂 `build_chain_graph`(`orchestrator/chain/graph.py`)— 契约
|
||||
|
||||
```python
|
||||
def build_chain_graph(
|
||||
gateway: GatewayRun, # 复用 ._protocols.GatewayRun
|
||||
*,
|
||||
session_factory: SessionFactory, # 节点自建短事务用(services 注入)
|
||||
review_repo_factory, # session -> ReviewRecorder
|
||||
checkpointer: BaseCheckpointSaver | None = None, # 运行时 Postgres;测试 MemorySaver
|
||||
) -> CompiledStateGraph[ChainState, None, ChainState, ChainState]: ...
|
||||
```
|
||||
|
||||
- 节点 DB 依赖经 `session_factory` + `*_repo_factory` 闭包注入(与现有"gateway 闭包绑定"同模式;langgraph 不收 `functools.partial`,用默认参绑定,见 gotcha)。
|
||||
- **删除/吸收 `build_write_graph`**:本链图取代它;`write_node` 函数本体保留(被 `write_chapter` 节点复用),删掉无人调用的 `build_write_graph` 工厂(YAGNI,解决审计发现的死代码)。
|
||||
|
||||
---
|
||||
|
||||
## 4. checkpointer 接线(激活休眠资产)
|
||||
|
||||
- **运行时**:`AsyncPostgresSaver`(`langgraph.checkpoint.postgres.aio`)。新增注入缝 `get_checkpointer()`(`services/chain_deps.py`)→ 运行时返 Postgres saver(连 `DATABASE_URL`),测试注 `MemorySaver`。
|
||||
- **建表**:`setup_checkpointer(conn_string)`(已存在 `orchestrator/graph.py:106`,迁移本链图后移到 `orchestrator/chain/graph.py` 或保留)。**在 alembic 迁移/CI 跑一次**建 langgraph 检查点表(不变量:DDL 不在 app-runtime)。新增迁移 `xxxx_langgraph_checkpoint_setup`:`op` 内 `await setup_checkpointer(sync_or_async_conn)`(实施者按本仓 alembic async 模式处理;若 alembic 环境不便跑 async saver.setup,则在 CI bootstrap 步骤调用并在迁移注释说明)。
|
||||
- `thread_id = str(job_id)`:`config={"configurable": {"thread_id": str(job_id)}}`。初始 run `graph.ainvoke(initial, config)`;resume `graph.ainvoke(Command(resume=decisions), config)`。
|
||||
- **测试**:E2E 用**单个进程内 MemorySaver 实例**横跨"初始 run"与"resume run"两次模拟调用(同进程同 saver,故 interrupt 可恢复);绝不连真 Postgres、绝不真 LLM。
|
||||
|
||||
---
|
||||
|
||||
## 5. 服务层(`apps/api/ww_api/services/chain_runner.py`)
|
||||
|
||||
```python
|
||||
async def run_chain_job(
|
||||
session_factory, job_id, *, project_id, chain_key, start_chapter_no, count,
|
||||
gateway, checkpointer, resume_decisions=None, request_id=None,
|
||||
) -> None:
|
||||
# 复用 run_job 的生命周期壳:set_running → work → complete/fail。
|
||||
# work: 构 build_chain_graph(...) → config{thread_id=job_id}
|
||||
# - resume_decisions is None → graph.ainvoke(initial_state, config)
|
||||
# - else → graph.ainvoke(Command(resume=resume_decisions), config)
|
||||
# 返回值/进度: 从图最终 state 取 {written, last_chapter_no} → job.result
|
||||
# {"chain_key","written":[...],"completed":bool,"awaiting_chapter":int|None}
|
||||
# interrupt 命中 → graph 返回带 __interrupt__ → job 置 status="awaiting_input"
|
||||
# (jobs 复用现有状态机 + 新增/复用一个 awaiting 态,见 §7 数据)
|
||||
```
|
||||
|
||||
- **进度**:每验收一章可经 `job_repo.set_progress`(若有)更新 `progress`;否则 result 里 `written` 长度即进度。
|
||||
- **错误脱敏**:复用 `run_job` 的 `_classify_job_error`(AppError 透传安全文案,其余通用文案,不泄 `str(exc)`,P0-3)。
|
||||
- **token 不泄漏**:job result/status/日志绝不含 prompt/正文/token(只记章号、长度、计数)。
|
||||
|
||||
---
|
||||
|
||||
## 6. API 契约(@backend;改 OpenAPI → 记 contracts + 前端 follow-up `pnpm gen:api`)
|
||||
|
||||
| 端点 | 方法 | 请求 | 响应 | 说明 |
|
||||
|---|---|---|---|---|
|
||||
| `/projects/{pid}/chains/{chain_key}/run` | POST | `ChainRunRequest{start_chapter_no:int>=1, count:int 1..50}` | 202 `ChainRunAccepted{job_id, chain_key, start_chapter_no, count}` | 写 jobs(queued,kind=chain) + BackgroundTask `run_chain_job` |
|
||||
| `/jobs/{job_id}` | GET | — | 现有 `JobView` + `result{chain_key,written[],completed,awaiting_chapter}` | 复用现有 jobs 端点;轮询进度/awaiting |
|
||||
| `/projects/{pid}/chains/runs/{job_id}/resume` | POST | `ChainResumeRequest{decisions: list[ConflictDecision]}` | 202 `ChainRunAccepted` | 仅当 job=awaiting_input;BackgroundTask resume |
|
||||
|
||||
- 校验(系统边界,fail fast):`count` 1..50(>50 → 422 VALIDATION);`start_chapter_no >= 1`;`chain_key` 未知 → 404;resume 非 awaiting 态 → 409。
|
||||
- `ConflictDecision` 复用 `schemas/projects.py` 既有(`conflict_index/verdict/note`)。
|
||||
- schemas 新增 `schemas/chain.py`:`ChainRunRequest/ChainRunAccepted/ChainResumeRequest`。
|
||||
|
||||
---
|
||||
|
||||
## 7. 数据模型
|
||||
|
||||
- **无新业务表**。章正文 → `chapters`;审稿 → `chapter_reviews`;digest → `chapter_digests`;运行 → `jobs`。
|
||||
- `jobs`:复用。新增 `kind` 值 `"chain"`(若 `kind` 列已存在则直接用;若无则本设计**允许**一个最小迁移加 `jobs.kind TEXT` + `jobs.status` 增态 `awaiting_input`——实施者先查现表,能复用就零迁移)。
|
||||
- **langgraph 检查点表**:由 `setup_checkpointer` 在迁移/CI 建(§4)。这是唯一新增 DDL。
|
||||
|
||||
> 实施者 MUST 先查 `packages/db/ww_db/models.py` 的 `jobs` 现有列(是否已有 `kind`/`progress`/`result`/能否表达 awaiting)再决定是否需要迁移;优先零迁移方案(如用 `result.awaiting_chapter` 表达待裁决,`status` 复用现有枚举)。
|
||||
|
||||
---
|
||||
|
||||
## 8. 不变量符合性(逐条核对)
|
||||
|
||||
- **#1 DB 即单一真源**:链节点只经领域表通信;节点间不互调。✅
|
||||
- **#3 四审只读 + 验收 HITL 闸**:`review_chapter` 只写 `chapter_reviews`;唯一正文写经 `accept_chapter`→`run_accept_transaction`;有冲突必 `interrupt()` 交人,`assert_conflicts_resolved` 缺判即拦。✅
|
||||
- **#4 digest 从终稿提炼**:`accept_chapter` 事务外 `extract_digest_facts(final_text)`。✅
|
||||
- **#5 checkpoint 只存控制流**:`ChainState` 只含章号/计数/标志;冲突明文 + 正文重读领域表。✅
|
||||
- **#9 缓存前缀稳定**:`write_chapter` 复用 `assemble` 的 stable_core/volatile 切分,硬规则+指纹进前缀、无时间戳/UUID。✅
|
||||
- **节点确定性可测**:LLM 非确定性锁在 gateway 后;`decide` 是纯逻辑。✅
|
||||
|
||||
---
|
||||
|
||||
## 9. 任务拆分(多 agent,契约先行)
|
||||
|
||||
| 任务 | skill | 依赖 | 产出 | DoD |
|
||||
|---|---|---|---|---|
|
||||
| **C1 链图 + 节点 + 单测** | @llm | — | `orchestrator/chain/{state,nodes,graph}.py` + 删 `build_write_graph` + 单测(mock gateway + fake session_factory + MemorySaver;测无冲突全跑 + 有冲突 interrupt + resume→accept;`decide` 纯逻辑测) | `packages/core`+`packages/agents` ruff/format/mypy 干净 + pytest 绿;图可 ainvoke、可 interrupt、可 resume |
|
||||
| **C2 服务 + 端点 + schema + checkpointer 接线 + 单测** | @backend | C1 契约 | `services/{chain_runner,chain_deps}.py` + `routers/chain.py` + `schemas/chain.py` + main 注册 + jobs 复用/最小迁移 | `apps/api` 门禁绿;3 端点入 OpenAPI;run→202、resume→202、未知 key 404、count>50 422、resume 非 awaiting 409;token 不入 result/日志 |
|
||||
| **C3 检查点建表迁移** | @db/@devops | C1 | alembic 迁移调 `setup_checkpointer`(或 CI bootstrap)+ jobs 最小迁移(如需) | `alembic upgrade head` 应用、`alembic check` 无漂移 |
|
||||
| **C4 E2E** | @qa | C2 | `tests/test_chain_workflow_e2e.py`:真 pg + mock 网关零 token:① 两章无冲突全自动→`chapters` 两 accepted 版次 + `chapter_digests` 两行 + job done `written=[n,n+1]`;② 注入冲突→第一章 interrupt→job awaiting→resume 带裁决→续跑→accept;负向:token 不入 job result/状态 | E2E 绿;DB 真源逐项断言 |
|
||||
|
||||
**并行性**:C1 是契约源,先行(§3.3 签名 + §6 端点已在本文锁定,故 C2 可对照契约并行起,但建议 C1 先落地图签名)。C3 依赖 C1(用其 `setup_checkpointer`)。C4 依赖 C2。
|
||||
|
||||
---
|
||||
|
||||
## 10. 测试计划(守 TDD + mock LLM)
|
||||
|
||||
- **单元**(C1):节点函数直接调用(不起图运行时):`write_chapter`(fake gateway 返固定文 + fake session 断言 save_draft)、`decide`(构造有/无冲突的 review → 断 `has_conflicts`)、`accept_chapter`(fake repos 断 promote/digest/decisions + 冲突缺判抛 CONFLICT_UNRESOLVED)。**流程**:MemorySaver 跑全图,断无冲突→END 且 written 满、有冲突→`__interrupt__`、resume→END。
|
||||
- **集成/E2E**(C4):FastAPI 端点 + 真 pg + mock 网关;interrupt+resume 在**单进程单 MemorySaver**(经依赖覆盖注入)内横跨两次端点调用。
|
||||
- **绝不**:真 LLM、真 Postgres 检查点联网、把正文/token 写进断言外的日志。
|
||||
|
||||
---
|
||||
|
||||
## 11. 风险与回退
|
||||
|
||||
| 风险 | 缓解 |
|
||||
|---|---|
|
||||
| alembic 跑 async `saver.setup()` 不便 | 退路:CI/bootstrap 脚本调 `setup_checkpointer`,迁移注释指明;E2E 用 MemorySaver 不依赖该表 |
|
||||
| 一条 run 内多章长事务/连接占用 | 节点**各自短事务**,不跨章共享 session;图运行本身无 DB 事务,只 checkpointer 写 |
|
||||
| MemorySaver 跨请求不持久 | 仅测试用(同进程);运行时强制 Postgres saver,文档标注 |
|
||||
| 链中途崩溃 | checkpointer 持久 → 可经 resume/重跑同 thread_id 续;job 置 failed 带脱敏文案 |
|
||||
| 与既有单章 `draft/review/accept` 端点重复逻辑 | 复用同一 `run_accept_transaction`/`assemble`/审稿 spec,不复制;链是**编排**新增、非业务重写 |
|
||||
|
||||
---
|
||||
|
||||
## 12. 文档回写(实施完成后)
|
||||
- `ARCHITECTURE.md §5.2`:补"多章链图(cyclic + interrupt + Postgres checkpointer)",并更正现状描述(写章直连流式、单章 accept 为事务、链图首次启用 checkpointer)。
|
||||
- `memory/contracts.md`:登记 C-Chain 契约(图签名 + 3 端点)。
|
||||
- `memory/decisions.md`:记"多章编排选 LangGraph 链图而非端点循环"的理由。
|
||||
- `PROGRESS.md`:拉入"Chain Workflow"阶段台账。
|
||||
Reference in New Issue
Block a user