fix(backend): Kimi OAuth 轮询移出 DB 会话——只在最终持久化开短会话(CR-H1)

This commit is contained in:
Yaojia Wang
2026-07-08 10:37:42 +02:00
parent ed0a679322
commit ca8fbc89a1
2 changed files with 176 additions and 37 deletions

View File

@@ -17,6 +17,7 @@ token 仅以 Fernet 密文存 `provider_credentials.oauth_enc`。
from __future__ import annotations
import asyncio
import uuid
from typing import Annotated, Any
import httpx
@@ -47,6 +48,7 @@ from ww_api.services.kimi_oauth import (
AuthorizationPending,
DeviceAuth,
SlowDown,
TokenSet,
decrypt_oauth_bundle,
encrypt_oauth_bundle,
poll_token,
@@ -74,49 +76,85 @@ def _default_http_client() -> AsyncHttpClient:
return httpx.AsyncClient(timeout=30.0)
def _make_poll_work(device: DeviceAuth) -> Any:
"""构造后台轮询工作闭包:循环 poll token 直到成功/过期/拒绝。
# SlowDown 退避每次增大的间隔量(秒)。
_SLOW_DOWN_STEP_SECONDS = 5
`work(session)` 自建 httpx 客户端 + 凭据 store用 `run_job` 传入的独立 session
按 `interval` 轮询 token → 成功则加密存 `oauth_enc` 并返回非密 job 结果 → 过期/拒绝则
抛 AppError`run_job` 置 job failed。**token 绝不进 job 结果/日志**
async def _poll_for_token(http: AsyncHttpClient, device: DeviceAuth) -> TokenSet:
"""设备授权轮询循环(**纯 HTTP无 DB 会话**CR-H1直到拿到 token 或超时
`authorization_pending` → 继续;`slow_down` → 增大 interval成功 → 返回 `TokenSet`
次数耗尽 → 抛 `AppError`(视作过期,由调用方经 `run_job` 置 job failed
"""
interval = max(1, device.interval)
for _ in range(_MAX_POLL_ATTEMPTS):
await asyncio.sleep(interval)
try:
return await poll_token(http, device.device_code)
except AuthorizationPending:
continue
except SlowDown:
interval += _SLOW_DOWN_STEP_SECONDS
continue
raise AppError(
ErrorCode.LLM_UNAVAILABLE,
"Kimi 设备授权轮询超时",
{"provider": KIMI_CODE_PROVIDER},
)
def _make_persist_work(token: TokenSet) -> Any:
"""构造**最终持久化**工作闭包:只做「加密 token → upsert oauth 凭据」的一次短事务。
`run_job` 用独立 session 调用此 work轮询已在会话外完成故此处会话极短。
**明文 token 绝不进 job 结果/日志。**
"""
async def work(session: AsyncSession) -> dict[str, Any]:
enc_key = get_settings().credential_enc_key.get_secret_value()
store = SqlCredentialStore(session)
interval = max(1, device.interval)
http = _default_http_client()
try:
for _ in range(_MAX_POLL_ATTEMPTS):
await asyncio.sleep(interval)
try:
token = await poll_token(http, device.device_code)
except AuthorizationPending:
continue
except SlowDown:
interval += 5
continue
# 成功:加密 token 包入库auth_type=oauth。明文 token 不进结果/日志。
blob = encrypt_oauth_bundle(token, key=enc_key)
await store.upsert_oauth_credential(STUB_OWNER_ID, KIMI_CODE_PROVIDER, blob)
return {"connected": True, "provider": KIMI_CODE_PROVIDER}
# 轮询次数耗尽(视作过期)。
raise AppError(
ErrorCode.LLM_UNAVAILABLE,
"Kimi 设备授权轮询超时",
{"provider": KIMI_CODE_PROVIDER},
)
finally:
# `httpx.AsyncClient` 有 `aclose`(最小 Protocol 无;测试 fake 也实现了它)。
aclose = getattr(http, "aclose", None)
if aclose is not None:
await aclose()
blob = encrypt_oauth_bundle(token, key=enc_key)
await store.upsert_oauth_credential(STUB_OWNER_ID, KIMI_CODE_PROVIDER, blob)
return {"connected": True, "provider": KIMI_CODE_PROVIDER}
return work
def _raise_work(exc: Exception) -> Any:
"""构造只把 `exc` 重抛的 work轮询失败路径交由 `run_job` 分类脱敏并置 job failed"""
async def work(_session: AsyncSession) -> dict[str, Any]:
raise exc
return work
async def _run_kimi_oauth_poll(
session_factory: SessionFactory,
job_id: uuid.UUID,
device: DeviceAuth,
*,
request_id: str | None = None,
) -> None:
"""后台任务:**先在会话外**轮询 token再经 `run_job` 开短会话持久化CR-H1
轮询期间不占用任何 DB 会话(防连接池饥饿);成功 → 交由 `run_job` 落库并置 done
失败 → 把异常带进 `run_job`,走其 `_classify_job_error`/`_mark_failed` 脱敏置 failed。
"""
http = _default_http_client()
try:
token = await _poll_for_token(http, device)
work = _make_persist_work(token)
except Exception as exc: # noqa: BLE001 — 失败经 run_job 统一脱敏置 job failed不在此处冒泡。
work = _raise_work(exc)
finally:
# `httpx.AsyncClient` 有 `aclose`(最小 Protocol 无;测试 fake 也实现了它)。
aclose = getattr(http, "aclose", None)
if aclose is not None:
await aclose()
await run_job(session_factory, job_id, work, request_id=request_id)
@router.post("/start", status_code=202)
async def start_oauth(
request: Request,
@@ -139,11 +177,12 @@ async def start_oauth(
job = await job_repo.create(None, _JOB_KIND_KIMI_OAUTH)
await session.commit() # job 行需在 202 前持久化(供前端立即轮询)。
# 轮询移出 DB 会话CR-H1后台先在会话外轮询 token仅最终持久化开一个短会话。
background_tasks.add_task(
run_job,
_run_kimi_oauth_poll,
session_factory,
job.id,
_make_poll_work(device),
device,
request_id=request_id,
)