多章链 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 消除。
58 lines
2.9 KiB
Python
58 lines
2.9 KiB
Python
"""多章链的 checkpointer 注入缝(chain-workflow 设计 §4)——激活休眠资产。
|
||
|
||
`get_checkpointer()` 是可注入缝:
|
||
- **运行时**:返 `AsyncPostgresSaver`(`langgraph.checkpoint.postgres.aio`),连 `DATABASE_URL`,
|
||
持久控制流位置 + pending handle(跨小时可恢复)。检查点表由迁移/CI 调 `setup_checkpointer`
|
||
建(C3),本缝**绝不**触碰 DDL(CLAUDE.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→resume);run 结束关闭。测试经 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
|