"""用量记账落库(ARCH §4.8)。 `LedgerSink` 为接口,便于测试注入内存替身;生产用 SQLAlchemy 实现写 usage_ledger。 """ from __future__ import annotations from typing import Protocol from sqlalchemy.ext.asyncio import AsyncSession from ww_db.models import UsageLedger from .types import Scope, Usage class LedgerSink(Protocol): async def record(self, scope: Scope, usage: Usage) -> None: ... class SqlAlchemyLedgerSink: """把每次调用写入 usage_ledger(owner_id 取 scope.user_id,单用户 stub)。""" def __init__(self, session: AsyncSession) -> None: self._session = session async def record(self, scope: Scope, usage: Usage) -> None: """把一次调用追加进 session(**只 add,不 flush/commit**)。 并发安全(T3.8 bugfix):三审作为同一 LangGraph superstep 的并行分支 共用一个请求 session 的本 sink。`AsyncSession` 非并发安全——若在此 `await self._session.flush()`,第二/三审的 flush 会撞上「Session is already flushing」(flush 重入),被 `run_review` 失败隔离吞成 `incomplete`, foreshadow/pace 结果静默丢失。`session.add()` 是同步、无 await、不让出 控制权,故并行协程不会在 add 中途交错。 持久化靠调用方:每个端点/验收事务在流耗尽/事务末尾 `commit()`(commit 自动 flush 全部待决行),与「网关 ledger 只 add/flush、端点或事务负责 commit」的既有边界一致(见 gotcha「网关 ledger 只 flush、调用方必须 commit」)。 """ row = UsageLedger( owner_id=scope.user_id, project_id=scope.project_id, provider=usage.provider, model=usage.model, input_tokens=usage.input_tokens, output_tokens=usage.output_tokens, cache_read=usage.cache_read_tokens, cost_minor=usage.cost_minor, currency=usage.currency, ) self._session.add(row)