Files
writer-work-flow/apps/api/ww_api/services/project_deps.py
Yaojia Wang 1652ad9d20 feat(backend): AI 反问澄清预检端点——refine 侧结构化 clarify(WFW-9 M1,路线A 两阶段)
润色「再沟通」意见含糊时先反问给选项(路线A:问题走独立非流式 JSON 预检端点,正文仍走
既有 refine 一字不改)。新增:
- ClarifyDecision/ClarifyQuestion/ClarifyOption 结构化 schema(既有 output-schema 处,供
  producer 与端点共用);clarify_refine.md 教条(含糊→need_clarification+≤1问+2–4锚定选项+
  自由输入;明确→verification 放行;防循环);注册 clarify_refine_spec(analyst 档,#24)+
  SCHEMA_CATALOG + 重生成金标准。
- clarify_node:build_clarify_request 纯函数(缓存前缀不含易变) + run_clarify(gateway.run
  结构化,判别/校验失败确定性回退 need_clarification=false,只读不写库)。
- POST /projects/{id}/chapters/{no}/refine/clarify → ClarifyDecision(analyst 网关,404/503,
  末尾 commit 记账)。**既有 refine 端点/RefineRequest/Response 完全未改**。
门禁绿:ruff/mypy 227/pytest 900(+test_clarify_spec/_node/style clarify)/alembic 无漂移。
2026-07-08 08:47:13 +02:00

588 lines
25 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.

"""项目/章节端点的依赖装配(运行时实现)。
- `get_project_repo` / `get_chapter_repo`:把请求 session 装配成 SQLAlchemy repo。
- `get_writer_gateway`:据 writer 档位路由从已存凭据解密 → 建 OpenAI 兼容适配器
→ `Gateway`(注入 `SqlAlchemyLedgerSink` + `resolve_route`)。这是 **draft SSE 的可注入缝**——
测试经 `app.dependency_overrides[get_writer_gateway]` 注入 mock 网关(产 `Delta`,绝不联网)。
- `seed_stub_user`:幂等 seed 单用户 stubowner_id FK 依赖它,见 memory/gotchas
"""
from __future__ import annotations
from collections.abc import Awaitable, Callable
from typing import Annotated, Any
import httpx
from fastapi import Depends
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from ww_config import get_settings
from ww_core.domain import ForeshadowLedgerRepo, SqlForeshadowLedgerRepo
from ww_core.domain.chapter_repo import ChapterRepo, SqlChapterRepo
from ww_core.domain.character_repo import CharacterWriteRepo, SqlCharacterWriteRepo
from ww_core.domain.digest_repo import DigestAppendRepo, SqlDigestAppendRepo
from ww_core.domain.injection_repo import InjectionOverrideRepo, SqlInjectionOverrideRepo
from ww_core.domain.job_repo import JobRepo, SqlJobRepo
from ww_core.domain.outline_write_repo import OutlineWriteRepo, SqlOutlineWriteRepo
from ww_core.domain.project_repo import ProjectRepo, SqlProjectRepo
from ww_core.domain.repositories import MemoryRepos, OutlineRepo, RulesRepo
from ww_core.domain.review_repo import ReviewRepo, SqlReviewRepo
from ww_core.domain.rule_repo import RuleWriteRepo, SqlRuleWriteRepo
from ww_core.domain.style_repo import SqlStyleFingerprintWriteRepo, StyleFingerprintWriteRepo
from ww_core.domain.template_repo import SqlTemplateRepo, TemplateRepo
from ww_core.domain.world_entity_repo import SqlWorldEntityWriteRepo, WorldEntityWriteRepo
from ww_core.memory.sql_repositories import SqlOutlineRepo, SqlRulesRepo, sql_memory_repos
from ww_db import get_session, get_sessionmaker
from ww_db.models import User
from ww_llm_gateway import (
Gateway,
ProviderAdapter,
Route,
SqlAlchemyLedgerSink,
build_adapter,
chain_from_routing,
resolve_route,
)
from ww_llm_gateway.adapters.kimi_code import KIMI_CODE_PROVIDER
from ww_llm_gateway.types import Tier
from ww_shared import AppError, ErrorCode
from ww_skills import SkillRegistry, SqlSkillRepo
from ww_api.security.credentials import (
CredentialKeyError,
decrypt_api_key,
)
from ww_api.services.credentials import (
AUTH_TYPE_OAUTH,
STUB_OWNER_ID,
CredentialStore,
SqlCredentialStore,
StoredCredential,
)
from ww_api.services.foreshadow_scan import SessionFactory
from ww_api.services.kimi_oauth import (
decrypt_oauth_bundle,
encrypt_oauth_bundle,
needs_refresh,
)
from ww_api.services.kimi_oauth import refresh as kimi_refresh
from ww_api.services.provider_deps import _PROVIDER_BASE_URLS
# 单用户 stub 的占位邮箱(多租户化时由真实主体替换)。
_STUB_USER_EMAIL = "stub@local"
async def seed_stub_user(session: AsyncSession) -> None:
"""幂等 seed 单用户 stub 行——所有 owner_id FKprojects/usage_ledger/...)依赖它。
无该行时插入;已存在则跳过。在 app lifespan 启动时调用一次。
"""
existing = (
await session.execute(select(User).where(User.id == STUB_OWNER_ID))
).scalar_one_or_none()
if existing is not None:
return
session.add(User(id=STUB_OWNER_ID, email=_STUB_USER_EMAIL, display_name="stub"))
await session.commit()
def get_project_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> ProjectRepo:
return SqlProjectRepo(session)
def get_chapter_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> ChapterRepo:
return SqlChapterRepo(session)
def get_memory_repos(
session: Annotated[AsyncSession, Depends(get_session)],
) -> MemoryRepos:
"""记忆组装的 7-repo 捆绑draft SSE 用)。测试覆盖此依赖注入内存 fake。"""
return sql_memory_repos(session)
def get_injection_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> InjectionOverrideRepo:
"""本章注入覆盖读/写 repoB0 可控版 GET/PUT injection + draft 同读)。
upsert 只 flush端点提交与其它写侧一致。测试经 `app.dependency_overrides` 注 fake。
"""
return SqlInjectionOverrideRepo(session)
def get_review_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> ReviewRepo:
"""审稿留痕 reporeview SSE collect / 历史 / accept 裁决)。"""
return SqlReviewRepo(session)
def get_digest_append_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> DigestAppendRepo:
"""章节摘要写侧 repo验收事务追加终稿 digest"""
return SqlDigestAppendRepo(session)
def get_foreshadow_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> ForeshadowLedgerRepo:
"""伏笔账本写侧 repo登记/状态变更端点;只 flush端点提交。测试经 override 注 fake。"""
return SqlForeshadowLedgerRepo(session)
def get_job_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> JobRepo:
"""长任务写侧 repo创建/进度/完成/失败;状态写只 flush提交归 run_job/端点)。
测试经 `app.dependency_overrides[get_job_repo]` 注入 fake避免真连 DB
"""
return SqlJobRepo(session)
def get_outline_write_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> OutlineWriteRepo:
"""大纲写侧 repo大纲生成端点逐章 upsert只 flush端点提交。测试经 override 注 fake。"""
return SqlOutlineWriteRepo(session)
def get_rule_write_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> RuleWriteRepo:
"""规则写侧 repoPOST /rules作者显式加规则只 flush端点提交。测试经 override 注。"""
return SqlRuleWriteRepo(session)
async def get_template_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> TemplateRepo:
"""模板库 repoGET/POST/DELETE /templatescreate/delete 只 flush端点提交
测试经 `app.dependency_overrides` 注入内存 fake。
"""
return SqlTemplateRepo(session)
async def get_skill_registry(
session: Annotated[AsyncSession, Depends(get_session)],
) -> SkillRegistry:
"""从 `skills` 表加载声明式 skill registryARCH §5.6;越权声明 → VALIDATION
每请求按 session 加载registry 不可变快照)。测试经 `app.dependency_overrides` 注 fake repo
或直接注 `SkillRegistry`。技能库 UIT5.6)经此读 builtin/custom/community。
"""
return await SkillRegistry.load(SqlSkillRepo(session))
def get_style_write_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> StyleFingerprintWriteRepo:
"""文风指纹写侧 repo`GET /style` 读最新 + 学文风后台任务 append 版本化指纹)。
注:学文风的 `work` 在 `run_job` 自建的独立 session 上自造 repo请求 session 已关闭),
故本依赖只服务于 `GET /style` 读侧。测试经 `app.dependency_overrides` 注 fake。
"""
return SqlStyleFingerprintWriteRepo(session)
def get_character_write_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> CharacterWriteRepo:
"""角色写侧 repoPOST /characters 入库schema→DB 形变;只 flush端点提交
测试经 `app.dependency_overrides` 注 fake。
"""
return SqlCharacterWriteRepo(session)
def get_world_entity_write_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> WorldEntityWriteRepo:
"""世界观实体写侧 repo预留对称入库当前生成端点只用其形变能力"""
return SqlWorldEntityWriteRepo(session)
def get_rules_read_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> RulesRepo:
"""规则读侧 repoGET /rules 列表,复用 C5 assemble 读侧;测试经 override 注 fake"""
return SqlRulesRepo(session)
def get_outline_read_repo(
session: Annotated[AsyncSession, Depends(get_session)],
) -> OutlineRepo:
"""大纲读侧 repoGET /outline 列表,复用 C5 assemble 读侧;测试经 override 注 fake"""
return SqlOutlineRepo(session)
# 创作工具箱通用端点的网关来源:按 tier 动态建网关(工具的 tier 在运行时才知)。
# 测试经 `app.dependency_overrides[get_tier_gateway_builder]` 注入返回 mock 网关的 builder。
TierGatewayBuilder = Callable[[Tier], Awaitable[Gateway]]
def get_tier_gateway_builder(
session: Annotated[AsyncSession, Depends(get_session)],
) -> TierGatewayBuilder:
"""返回「按 tier 建网关」的可注入缝(创作工具箱按工具 spec.tier 动态路由)。
无凭据时 builder 调用内 `build_gateway_for_tier` 抛 `LLM_UNAVAILABLE`503——
与固定档位网关缝语义一致。测试 override 返回固定 mock 网关的 builder绝不联网
"""
store = SqlCredentialStore(session)
async def _build(tier: Tier) -> Gateway:
return await build_gateway_for_tier(session, store, tier)
return _build
async def get_worldbuilder_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""世界观生成writer 档位)的可注入网关缝。测试经 override 注 mock产 WorldGenResult"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "writer")
async def get_character_gen_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""角色生成writer 档位)的可注入网关缝。测试经 override 注 mock产 CharacterGenResult"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "writer")
async def get_precheck_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""入库前 continuity 预检analyst 档位)的可注入网关缝。测试注 mock产 ContinuityReview"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "analyst")
async def get_project_plan_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""AI 立项方案生成analyst 档位)的可注入网关缝。测试注 mock产 ProjectPlanResult"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "analyst")
def get_session_factory() -> SessionFactory:
"""验收后到期扫描的**独立 session 工厂**缝。
BackgroundTask 在请求 session 关闭后才跑,必须自建 session不复用 `get_session`)。
返回的工厂 `()` → `async with` 得一个新 `AsyncSession`。测试经 `app.dependency_overrides`
注入 fake 工厂(避免真起后台线程/真连 DB
"""
return get_sessionmaker()
async def _build_provider_adapter(store: CredentialStore, provider: str) -> ProviderAdapter | None:
"""据 provider 解密凭据 → 经 `build_adapter` 工厂建对应 provider 适配器T5.4 follow-up
`build_adapter(provider, *, api_key, base_url=None)` 按 provider 选适配器类:
OpenAI 兼容deepseek/kimi/qwen/glm/openai经 `base_url` 走 OpenAI 兼容适配器;
Anthropic/Gemini 走各自原生适配器(无需 `base_url``kimi-code` 走 OAuth bearer +
coding base + 伪造头(**access token 经 `_resolve_kimi_code_token` 按需刷新**K1.3)。
返回 `None` 表示该 provider 未配置凭据——回退链里缺位时网关会跳到下一个,故宽容返回
None不直接抛
"""
cred = await store.get_credential(STUB_OWNER_ID, provider)
if cred is None:
return None
settings = get_settings()
if cred.auth_type == AUTH_TYPE_OAUTH or provider == KIMI_CODE_PROVIDER:
# OAuth 凭据Kimi Code解密 token 包 → 临近过期则刷新并持久化 → access token
# 当 api_key 喂工厂(工厂为 kimi-code 构建带伪造头 + coding base 的客户端)。
access_token = await _resolve_kimi_code_token(
store, cred, settings.credential_enc_key.get_secret_value()
)
return build_adapter(
provider, api_key=access_token, base_url=_PROVIDER_BASE_URLS.get(provider)
)
if cred.api_key_enc is None:
# api_key 凭据但密文缺失(数据不一致)——视作未配置,回退链跳过。
return None
try:
api_key = decrypt_api_key(
cred.api_key_enc, key=settings.credential_enc_key.get_secret_value()
)
except CredentialKeyError as exc:
raise AppError(ErrorCode.INTERNAL, str(exc)) from exc
# OpenAI 兼容 provider 需 base_urlAnthropic/Gemini 走原生 SDKbase_url=None
base_url = _PROVIDER_BASE_URLS.get(provider)
return build_adapter(provider, api_key=api_key, base_url=base_url)
async def _resolve_kimi_code_token(
store: CredentialStore, cred: StoredCredential, enc_key: str
) -> str:
"""解密 Kimi Code OAuth token 包 → 临近过期时刷新并持久化 → 返回当前 access token。
刷新经一个**临时 httpx 客户端**(与 token 端点交互);新 token 包经
`store.upsert_oauth_credential` 持久化(下次建网关复用刷新结果)。明文 token 绝不进
日志/响应。无 `oauth_enc` → `LLM_UNAVAILABLE`(未连接 Kimi Code
"""
if cred.oauth_enc is None:
raise AppError(
ErrorCode.LLM_UNAVAILABLE,
f"{KIMI_CODE_PROVIDER} 未连接(无 OAuth 凭据),请先在设置中连接 Kimi Code",
{"provider": KIMI_CODE_PROVIDER},
)
try:
token = decrypt_oauth_bundle(cred.oauth_enc, key=enc_key)
except CredentialKeyError as exc:
raise AppError(ErrorCode.INTERNAL, str(exc)) from exc
if not needs_refresh(token):
return token.access_token
# 临近过期 → 刷新并持久化新包。
async with httpx.AsyncClient(timeout=30.0) as http:
refreshed = await kimi_refresh(http, token.refresh_token)
new_blob = encrypt_oauth_bundle(refreshed, key=enc_key)
await store.upsert_oauth_credential(STUB_OWNER_ID, KIMI_CODE_PROVIDER, new_blob)
# token 刷新是独立可持久的副作用(下次建网关复用),须立即提交,不依赖请求后续是否提交。
await store.commit()
return refreshed.access_token
async def build_gateway_for_tier(
session: AsyncSession, store: CredentialStore, tier: Tier
) -> Gateway:
"""据指定档位路由 + DB `tier_routing.fallback` 装配**多 provider 回退链**网关T5.4 接线)。
流程§4.3 三级解析 / §4.5 回退链):
1. 读 DB `tier_routing` 取该 tier 的 primary `provider:model` + fallback 列表(缺则退回
全局 `resolve_route`,单 provider向后兼容
2. 为 primary + fallback 里**每个能建出适配器**的 provider 预备 OpenAI 兼容适配器
(未知 base_url / 未配凭据的 provider 跳过——回退链自然绕过它)。
3. 至少要有一个可用适配器,否则 `LLM_UNAVAILABLE`(无任何凭据可用)。
4. 注入 `chain_resolver=chain_from_routing(...)`(多元素链,启用回退);无 DB 路由时
用 `resolver=resolve_route`单路由M1 行为不变)。
单 provider 配置仍走单元素链(网关把它当无回退处理),**不破既有行为**。
"""
ledger = SqlAlchemyLedgerSink(session)
stored = next((r for r in await store.list_routing() if r.tier == tier), None)
if stored is None:
# 无 DB 路由行:退回全局默认(单 providerM1 兼容)。
route = resolve_route(tier)
adapter = await _build_provider_adapter(store, route.provider)
if adapter is None:
raise AppError(
ErrorCode.LLM_UNAVAILABLE,
f"{tier} 档位 provider {route.provider} 未配置凭据,请先在设置中配置",
{"provider": route.provider, "tier": tier},
)
return Gateway(adapters={route.provider: adapter}, ledger=ledger, resolver=resolve_route)
# DB 路由primary + fallback 构链;为每个可建的 provider 预备适配器。
primary_spec = f"{stored.provider}:{stored.model}"
chain: list[Route] = chain_from_routing(tier, primary_spec, list(stored.fallback))
adapters: dict[str, ProviderAdapter] = {}
for route in chain:
if route.provider in adapters:
continue
built = await _build_provider_adapter(store, route.provider)
if built is not None:
adapters[route.provider] = built
if not adapters:
raise AppError(
ErrorCode.LLM_UNAVAILABLE,
f"{tier} 档位无任何已配置凭据的 provider请先在设置中配置",
{"providers": [r.provider for r in chain], "tier": tier},
)
def _resolver(_tier: Tier) -> list[Route]:
return chain_from_routing(_tier, primary_spec, list(stored.fallback))
return Gateway(adapters=adapters, ledger=ledger, chain_resolver=_resolver)
async def build_writer_gateway(session: AsyncSession, store: CredentialStore) -> Gateway:
"""据 writer 档位路由解密对应 provider 凭据 → 建网关。"""
return await build_gateway_for_tier(session, store, "writer")
# 多章链一条 run 内跨 writer(write)/analyst(review)/light(digest) 三档——单档网关的
# chain_resolver 恒返该档链、忽略请求 tier会把 review/digest 错路由到 writer。故链需一个
# **按请求 tier 分派**的网关:并 union 三档所有 provider 的适配器 + 逐 tier 解析其链。
_CHAIN_TIERS: tuple[Tier, ...] = ("writer", "analyst", "light")
def _chain_for_tier(stored_by_tier: dict[Tier, Any], tier: Tier) -> list[Route]:
"""据该 tier 的 DB 路由行构链;无行则退回全局默认(单 providerM1 兼容)。"""
stored = stored_by_tier.get(tier)
if stored is None:
return [resolve_route(tier)]
primary_spec = f"{stored.provider}:{stored.model}"
return chain_from_routing(tier, primary_spec, list(stored.fallback))
async def build_chain_gateway(session: AsyncSession, store: CredentialStore) -> Gateway:
"""装配多章链网关:按请求 tierwriter/analyst/light分派union 三档适配器 + 回退链。
与 `build_gateway_for_tier` 同凭据/适配器装配,但 `chain_resolver` 据 `req.tier` 返回
对应档的链(不锁死单档)。至少一个可用适配器,否则 `LLM_UNAVAILABLE`(无任何凭据)。
"""
ledger = SqlAlchemyLedgerSink(session)
routing_rows = await store.list_routing()
stored_by_tier: dict[Tier, Any] = {r.tier: r for r in routing_rows if r.tier in _CHAIN_TIERS}
# union 三档链里所有 provider逐个建可用适配器未配凭据/未知 base_url 的跳过)。
adapters: dict[str, ProviderAdapter] = {}
for tier in _CHAIN_TIERS:
for route in _chain_for_tier(stored_by_tier, tier):
if route.provider in adapters:
continue
built = await _build_provider_adapter(store, route.provider)
if built is not None:
adapters[route.provider] = built
if not adapters:
raise AppError(
ErrorCode.LLM_UNAVAILABLE,
"多章链无任何已配置凭据的 provider请先在设置中配置",
{"tiers": list(_CHAIN_TIERS)},
)
def _resolver(tier: Tier) -> list[Route]:
return _chain_for_tier(stored_by_tier, tier)
return Gateway(adapters=adapters, ledger=ledger, chain_resolver=_resolver)
async def get_writer_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""draft SSE 的可注入网关缝。测试覆盖此依赖注入 mock产 `Delta`,绝不联网)。"""
store = SqlCredentialStore(session)
return await build_writer_gateway(session, store)
async def get_review_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""续审analyst 档位)的可注入网关缝。测试经 override 注 mock产 `parsed`,绝不联网)。"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "analyst")
async def get_digest_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""验收终稿 digest 提炼light 档位)的可注入网关缝。测试经 override 注 mock。"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "light")
async def get_outline_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""大纲生成analyst 档位)的可注入网关缝。测试经 override 注 mock产 OutlineResult"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "analyst")
async def get_style_extract_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""学文风提取analyst 档位)的可注入网关缝。
`POST /style` 在 dep 解析阶段构建网关 → 无凭据时这里抛 `LLM_UNAVAILABLE`503
调度 job 之前拦下,避免凭空写一行注定失败的 job。提取本体在 BackgroundTask 里
用 `run_job` 自建的独立 session 重新构网关跑(请求 session 已关闭);本依赖确保
凭据探测在请求阶段发生。测试经 override 注 mock产 StyleFingerprintResult
"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "analyst")
async def get_refine_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""回炉writer 档位)的可注入网关缝。测试经 override 注 mock产纯文本重写段"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "writer")
async def get_clarify_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""润色预检澄清analyst 档位)的可注入网关缝。
`POST .../refine/clarify` 在 dep 解析阶段构建网关 → 无凭据时这里抛 `LLM_UNAVAILABLE`
503进入端点前拦下。测试经 override 注 mock产 `ClarifyDecision`)。
"""
store = SqlCredentialStore(session)
return await build_gateway_for_tier(session, store, "analyst")
async def get_chain_gateway(
session: Annotated[AsyncSession, Depends(get_session)],
) -> Gateway:
"""多章链write writer + review analyst + digest light按请求 tier 分派)的可注入网关缝。
`POST .../chains/{key}/run` 在 dep 解析阶段构建网关 → 无凭据时这里抛 `LLM_UNAVAILABLE`
503调度 job 之前拦下,避免凭空写注定失败的 job。链本体在 BackgroundTask 里用
`run_chain_job` 自建独立 session 重新构网关跑(请求 session 已关闭);本依赖只做请求阶段
凭据探测。测试经 `app.dependency_overrides[get_chain_gateway]` 注 mock绝不联网
"""
store = SqlCredentialStore(session)
return await build_chain_gateway(session, store)
# 「按 session 建 digest 网关」缝类型chain accept_op 终稿提炼按节点自建短事务的 session 建)。
GatewayDigestBuilder = Callable[[AsyncSession], Awaitable[Gateway]]
def _digest_gateway_builder(session: AsyncSession) -> Awaitable[Gateway]:
"""链 accept 节点在自建短事务里建 light 档 digest 网关(终稿提炼)。"""
return build_gateway_for_tier(session, SqlCredentialStore(session), "light")
def get_digest_gateway_builder() -> GatewayDigestBuilder:
"""返回「按 session 建 digest 网关」的缝chain accept_op 终稿提炼用)。
BackgroundTask 自建短事务 → 节点内据该 session 现建 light 档网关。测试经
`app.dependency_overrides[get_digest_gateway_builder]` 注返回 mock 的 builder绝不联网
"""
return _digest_gateway_builder
# 「按 session 建链网关」缝类型(链 write/review 节点在 BackgroundTask 自建短事务的 session 建)。
GatewayChainBuilder = Callable[[AsyncSession], Awaitable[Gateway]]
def _chain_gateway_builder(session: AsyncSession) -> Awaitable[Gateway]:
"""链 write/review 节点在自建短事务里建按请求 tier 分派的链网关。
关键(守不变量 #1网关的 `SqlAlchemyLedgerSink` 必须绑**节点当前的** session——
BackgroundTask 跑时请求 session 已关闭,故不能复用请求阶段的网关,必须按节点新 session 现建,
否则 `gateway.run()` 的 `usage_ledger` 写会落到已死 session 上被静默丢弃(成本记账断裂)。
"""
return build_chain_gateway(session, SqlCredentialStore(session))
def get_chain_gateway_builder() -> GatewayChainBuilder:
"""返回「按 session 建链网关」的缝(链 write/review 节点在 BackgroundTask 内现建网关用)。
与 `get_chain_gateway`(请求阶段凭据探测,返单实例,仅做 503 拦截)不同:本缝供
`run_chain_job` 在节点自建的**新鲜** session 上重建网关,使 usage_ledger 落到活 session。
测试经 `app.dependency_overrides[get_chain_gateway_builder]` 注返 mock 的 builder绝不联网
"""
return _chain_gateway_builder