Files
writer-work-flow/apps/api/ww_api/services/chain_deps.py
Yaojia Wang b5002a9864 fix(chain): checkpointer 连接串剥离 SQLAlchemy 驱动后缀
多章链 job 一发起即 failed(exc_type=ProgrammingError)。根因:runtime
checkpointer 工厂把 settings 的 sync URL `postgresql+psycopg://…` 原样喂给
`AsyncPostgresSaver.from_conn_string`,psycopg 不识别 `+psycopg` 方言段 →
`ProgrammingError: missing "=" …`,链在打开 saver 时即崩。

E2E 用 MemorySaver 覆盖 get_checkpointer_factory,从不走真 conn string,故此
路径无回归保护、bug 未被测出(迁移 d3e4f5a6b7c8 里 `_psycopg_conn_string` 已
正确剥离,runtime 缝漏了同款处理)。

修复:chain_deps 加 `_libpq_conn_string()`,剥离 `+psycopg`/`+asyncpg` 还原为
标准 libpq 串后再交 from_conn_string。+ 单测 test_chain_deps.py(3 参数化用例)。
已直连真 pg 复验:saver.aget OK,ProgrammingError 消除。
2026-06-24 12:40:40 +02:00

58 lines
2.9 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 _libpq_conn_string() -> str:
"""从 settings 的 sync URL 派生 psycopg 的 libpq conn string。
`database_url_sync` 带 SQLAlchemy 驱动后缀(``postgresql+psycopg://…``),但
`AsyncPostgresSaver.from_conn_string` 直接交给 psycopg 解析——psycopg 不识别
``+psycopg``/``+asyncpg`` 方言段,会抛 ``ProgrammingError: missing "=" …``。
故剥离后缀还原为标准 libpq 串(``postgresql://…``)。与迁移 `d3e4f5a6b7c8`
的 `_psycopg_conn_string()` 同义。
"""
return get_settings().database_url_sync.replace("+psycopg", "").replace("+asyncpg", "")
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(同步驱动串);须剥离 SQLAlchemy 驱动
# 后缀(+psycopg否则 psycopg 解析连接串报 ProgrammingError。
return AsyncPostgresSaver.from_conn_string(_libpq_conn_string())
return _factory