Files
writer-work-flow/apps/api/ww_api/services/chain_deps.py
Yaojia Wang 29349dc7ee feat(api): C2 多章链 服务+端点+schema+checkpointer 接线
承 C1 链图(build_chain_graph),落地多章工作流链的 apps/api 壳:
- 3 端点 routers/chain.py:POST .../chains/{key}/run→202 ChainRunAccepted;
  POST .../chains/runs/{job_id}/resume→202;GET /jobs/{id} 复用。校验:
  count 1..50→422、未知 chain_key→404、resume 非 awaiting→409、无凭据→503。
- schemas/chain.py:ChainRunRequest/ChainRunAccepted/ChainResumeRequest
  (ConflictDecision 复用 schemas/projects)。
- services/chain_runner.py:run_chain_job 仿 run_job 壳自建独立 session 驱动链图
  (set_running→ainvoke→据 __interrupt__ 置 awaiting_input/done/failed);
  build_accept_op 在 apps/api 装配验收事务闭包注入图节点(守 #3/#4);
  token 不入 result/日志。
- services/chain_deps.py:get_checkpointer_factory(运行时 AsyncPostgresSaver
  上下文 / 测试 MemorySaver)。
- 零迁移(设计 §7):复用 jobs,新增 status="awaiting_input" + JobRepo.set_awaiting,
  awaiting 章经 result.awaiting_chapter;新错误码 ErrorCode.CONFLICT(409)。
- project_deps:build_chain_gateway/get_chain_gateway(按请求 tier writer/analyst/light
  分派——单档网关恒返该档会错路由 review/digest)+ get_digest_gateway_builder。

单测 apps/api/tests/test_chain.py 12 用例(mock 网关 + MemorySaver + fake session/
accept_op,无 DB/无网络/无真 LLM):run/resume→202、未知 key 404、count 越界 422、
resume 非 awaiting 409、run_chain_job 无冲突→done、冲突→awaiting→resume→done、
错误脱敏、accept_op 冲突缺判→CONFLICT_UNRESOLVED。

门禁绿:ruff/format 干净 · mypy 193 Success · alembic 无漂移 · pytest 600 passed。
守不变量 #1/#3/#4/#5/#9。唯一新增 DDL(langgraph 检查点表)= C3 迁移。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 17:12:53 +02:00

45 lines
2.2 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.

"""多章链的 checkpointer 注入缝chain-workflow 设计 §4——激活休眠资产。
`get_checkpointer()` 是可注入缝:
- **运行时**:返 `AsyncPostgresSaver``langgraph.checkpoint.postgres.aio`),连 `DATABASE_URL`
持久控制流位置 + pending handle跨小时可恢复。检查点表由迁移/CI 调 `setup_checkpointer`
C3本缝**绝不**触碰 DDLCLAUDE.md「LangGraph」`setup()` 只在 migrations/CI
- **测试**:经 `app.dependency_overrides[get_checkpointer]` 注 `MemorySaver`(同进程跨两次
invoke 可恢复 interrupt绝不连真 Postgres、绝不真 LLM。
注:`AsyncPostgresSaver.from_conn_string` 是 async 上下文管理器——其连接生命周期须横跨整条
链 run含 interrupt 暂停后的 resume。故 `run_chain_job` 在自身 task 内打开/关闭 saver
上下文(见 `chain_runner`),本缝只提供「按需建一个 saver 上下文管理器」的工厂。
"""
from __future__ import annotations
from collections.abc import Callable
from contextlib import AbstractAsyncContextManager
from typing import TYPE_CHECKING, Any
from ww_config import get_settings
if TYPE_CHECKING:
from langgraph.checkpoint.base import BaseCheckpointSaver
# checkpointer 工厂:`()` → async 上下文管理器,进入得一个 `BaseCheckpointSaver`。
# 运行时 = Postgres saver 上下文(连接随上下文开/关);测试注一个产 MemorySaver 的工厂。
CheckpointerFactory = Callable[[], AbstractAsyncContextManager[Any]]
def get_checkpointer_factory() -> CheckpointerFactory:
"""运行时 checkpointer 工厂:建连 `DATABASE_URL` 的 `AsyncPostgresSaver` 上下文。
`from_conn_string` 返回的 async-CM 在 `run_chain_job` 内 `async with` 打开——其连接
存活整条链(含 interrupt→resumerun 结束关闭。测试经 dep 覆盖注 MemorySaver 工厂。
"""
def _factory() -> AbstractAsyncContextManager[BaseCheckpointSaver[Any]]:
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
# langgraph 的 Postgres saver 走 psycopg(同步驱动串);用 sync 串建连接。
return AsyncPostgresSaver.from_conn_string(get_settings().database_url_sync)
return _factory