refactor(wechat_woc): 优化WocBridgeClient HTTP请求逻辑,减少凭证查询次数并调整重试策略
- 新增_resolve_endpoint方法一次性解析bridge URL与鉴权token,避免重复调用_get_credentials - 重构_execute_http,替换原有的URL构造与token获取逻辑,统一通过新方法处理 - 调整重试策略:仅对瞬时网络错误做1次快速重试,移除超时与5xx错误的重试,避免双重重试放大延迟 - 更新相关注释与常量命名,对齐现有逻辑说明
This commit is contained in:
parent
76981cd9bc
commit
17abcd4502
@ -433,6 +433,51 @@ class WocBridgeClient:
|
|||||||
# 通用 HTTP 执行
|
# 通用 HTTP 执行
|
||||||
# ------------------------------------------------------------------
|
# ------------------------------------------------------------------
|
||||||
|
|
||||||
|
async def _resolve_endpoint(
|
||||||
|
self,
|
||||||
|
account_id: str,
|
||||||
|
path: str,
|
||||||
|
override_bridge_url: str | None,
|
||||||
|
bridge_token: str | None,
|
||||||
|
) -> tuple[str, str | None]:
|
||||||
|
"""一次性解析 bridge URL 与鉴权 token,避免双重凭证缓存查询。
|
||||||
|
|
||||||
|
非 wizard 模式(``override_bridge_url`` 为空)时从 CachePort 凭证缓存
|
||||||
|
读取一次凭证字典,同时提取 ``bridge_url`` 与 ``bridge_token``,避免
|
||||||
|
``_build_bridge_url`` 与 ``_get_bridge_token`` 各调一次 ``_get_credentials``
|
||||||
|
导致同一份凭证查 2 次。wizard 模式下直接使用调用方传入的
|
||||||
|
``override_bridge_url``,``bridge_token`` 透传。
|
||||||
|
|
||||||
|
@pre
|
||||||
|
- ``account_id`` 非空
|
||||||
|
- ``path`` 以 ``/`` 开头
|
||||||
|
- ``override_bridge_url`` 非空时必须以 ``http://`` 或 ``https://`` 开头
|
||||||
|
|
||||||
|
@post
|
||||||
|
- 返回 ``(url, bridge_token_or_none)`` 元组
|
||||||
|
|
||||||
|
@failure
|
||||||
|
- 透传 ``_get_credentials`` 的 ``DependencyError``
|
||||||
|
- ``DependencyError``:bridge_url 为空字符串
|
||||||
|
"""
|
||||||
|
if override_bridge_url is not None:
|
||||||
|
base = override_bridge_url.rstrip("/")
|
||||||
|
url = f"{base}{path}" if path.startswith("/") else f"{base}/{path}"
|
||||||
|
return url, bridge_token
|
||||||
|
# 非 wizard 模式:一次读取凭证字典,同时提取 bridge_url 与 bridge_token
|
||||||
|
credentials = await self._get_credentials(account_id)
|
||||||
|
bridge_url = credentials.get("bridge_url", "").rstrip("/")
|
||||||
|
if not bridge_url:
|
||||||
|
raise DependencyError(
|
||||||
|
dep=WOC_BRIDGE_DEP,
|
||||||
|
cause=Exception(f"bridge_url_empty: account_id={account_id}"),
|
||||||
|
)
|
||||||
|
if not path.startswith("/"):
|
||||||
|
path = f"/{path}"
|
||||||
|
# 调用方未显式传入 bridge_token 时从凭证读取;传入时优先使用调用方值
|
||||||
|
resolved_token = bridge_token if bridge_token is not None else credentials.get("bridge_token", "")
|
||||||
|
return f"{bridge_url}{path}", resolved_token
|
||||||
|
|
||||||
async def _execute_http(
|
async def _execute_http(
|
||||||
self,
|
self,
|
||||||
method: str,
|
method: str,
|
||||||
@ -454,19 +499,15 @@ class WocBridgeClient:
|
|||||||
在途请求计数:进入时 +1,退出时 -1,归零时设置 ``_drain_event`` 供
|
在途请求计数:进入时 +1,退出时 -1,归零时设置 ``_drain_event`` 供
|
||||||
lifecycle 的 ``drain`` 等待。
|
lifecycle 的 ``drain`` 等待。
|
||||||
|
|
||||||
URL 构造:默认通过 ``_build_bridge_url(account_id, path)`` 从 CachePort
|
URL 构造与鉴权注入:通过 ``_resolve_endpoint`` 一次性从凭证缓存读取
|
||||||
凭证缓存读取 ``bridge_url`` 拼接;当 ``override_bridge_url`` 非空时直接与
|
``bridge_url`` 与 ``bridge_token``(避免双重缓存查询),非空 token 注入
|
||||||
``path`` 拼接,绕过缓存读取(用于 wizard 阶段账户未创建时探活,对齐
|
``Authorization: Bearer <token>`` 头。``override_bridge_url`` 非空时
|
||||||
wechat_ilink ``getconfig_with_token`` 模式)。
|
绕过缓存读取(wizard 阶段账户未创建时探活)。
|
||||||
|
|
||||||
鉴权注入:``bridge_token`` 非空时注入 ``Authorization: Bearer <token>``
|
重试策略:仅对 ``httpx.NetworkError``(连接失败/DNS 失败等瞬时网络抖动)
|
||||||
头,用于通过 WechatOnCloud 面板 ``/api/bridge/:id/*`` 反代路由的 M2M
|
做 1 次快速重试(不等待),超时与 5xx 不重试,由调用方或 outbox 机制
|
||||||
鉴权(与面板管理员会话二选一)。``bridge_token`` 为 ``None`` 或空字符串
|
决策(P2-3 超时重试)。避免 bridge 客户端重试与 outbox 退避重试双重
|
||||||
时不注入该头,回退到面板会话鉴权(兼容未启用 M2M 鉴权的旧部署)。
|
放大延迟。
|
||||||
|
|
||||||
重试策略:``httpx.TimeoutException`` / ``httpx.NetworkError`` / 5xx
|
|
||||||
``HTTPStatusError`` 最多重试 ``_HTTP_RETRY_ATTEMPTS`` 次(指数退避),
|
|
||||||
4xx(含 429)不重试,由调用方或 outbox 决策(P2-3 超时重试)。
|
|
||||||
|
|
||||||
@pre
|
@pre
|
||||||
- ``method`` 为合法 HTTP 方法(GET / POST)
|
- ``method`` 为合法 HTTP 方法(GET / POST)
|
||||||
@ -493,16 +534,10 @@ class WocBridgeClient:
|
|||||||
)
|
)
|
||||||
# 生成请求级 trace_id,注入 X-Trace-Id 头与 error_translator 供异常追踪
|
# 生成请求级 trace_id,注入 X-Trace-Id 头与 error_translator 供异常追踪
|
||||||
trace_id = uuid4().hex
|
trace_id = uuid4().hex
|
||||||
if override_bridge_url is not None:
|
# 一次性解析 URL 与鉴权 token,避免双重凭证缓存查询
|
||||||
base = override_bridge_url.rstrip("/")
|
url, resolved_token = await self._resolve_endpoint(
|
||||||
url = f"{base}{path}" if path.startswith("/") else f"{base}/{path}"
|
account_id, path, override_bridge_url, bridge_token
|
||||||
else:
|
)
|
||||||
url = await self._build_bridge_url(account_id, path)
|
|
||||||
# 鉴权 token 统一获取:非 wizard 模式(override_bridge_url 为空)且
|
|
||||||
# 调用方未显式传入 bridge_token 时,从凭证缓存读取。wizard 模式下
|
|
||||||
# bridge_token 由调用方显式传入(可为 None,表示未配置 M2M 鉴权)。
|
|
||||||
if bridge_token is None and override_bridge_url is None:
|
|
||||||
bridge_token = await self._get_bridge_token(account_id)
|
|
||||||
# 合并请求头:调用方显式传入的 headers 优先(覆盖 User-Agent / X-Trace-Id)
|
# 合并请求头:调用方显式传入的 headers 优先(覆盖 User-Agent / X-Trace-Id)
|
||||||
request_headers: dict[str, str] = {
|
request_headers: dict[str, str] = {
|
||||||
"User-Agent": self._user_agent,
|
"User-Agent": self._user_agent,
|
||||||
@ -510,11 +545,11 @@ class WocBridgeClient:
|
|||||||
}
|
}
|
||||||
if json_body is not None:
|
if json_body is not None:
|
||||||
request_headers["Content-Type"] = "application/json"
|
request_headers["Content-Type"] = "application/json"
|
||||||
if bridge_token:
|
if resolved_token:
|
||||||
# 注入 M2M 鉴权头,通过 WechatOnCloud 面板反代路由的双鉴权校验。
|
# 注入 M2M 鉴权头,通过 WechatOnCloud 面板反代路由的双鉴权校验。
|
||||||
# bridge_token 为 None 或空字符串时不注入,回退到面板会话鉴权
|
# resolved_token 为 None 或空字符串时不注入,回退到面板会话鉴权
|
||||||
# (兼容未启用 M2M 鉴权的旧部署)。
|
# (兼容未启用 M2M 鉴权的旧部署)。
|
||||||
request_headers["Authorization"] = f"Bearer {bridge_token}"
|
request_headers["Authorization"] = f"Bearer {resolved_token}"
|
||||||
if headers is not None:
|
if headers is not None:
|
||||||
request_headers.update(headers)
|
request_headers.update(headers)
|
||||||
# 使用 path 作为 operation 上下文,供 translate_call_error 丰富超时错误信息
|
# 使用 path 作为 operation 上下文,供 translate_call_error 丰富超时错误信息
|
||||||
@ -523,7 +558,7 @@ class WocBridgeClient:
|
|||||||
self._drain_event.clear()
|
self._drain_event.clear()
|
||||||
last_exc: httpx.HTTPError | None = None
|
last_exc: httpx.HTTPError | None = None
|
||||||
try:
|
try:
|
||||||
for attempt in range(1, _HTTP_RETRY_ATTEMPTS + 1):
|
for attempt in range(1, _HTTP_NETWORK_ERROR_MAX_ATTEMPTS + 1):
|
||||||
# 每次重试前重新检查活跃状态,stop/pause 后立即拒绝新请求
|
# 每次重试前重新检查活跃状态,stop/pause 后立即拒绝新请求
|
||||||
if not self._active:
|
if not self._active:
|
||||||
raise DependencyError(
|
raise DependencyError(
|
||||||
@ -554,26 +589,12 @@ class WocBridgeClient:
|
|||||||
resp.raise_for_status()
|
resp.raise_for_status()
|
||||||
last_exc = None
|
last_exc = None
|
||||||
break
|
break
|
||||||
except httpx.TimeoutException as exc:
|
|
||||||
last_exc = exc
|
|
||||||
if attempt == _HTTP_RETRY_ATTEMPTS:
|
|
||||||
break
|
|
||||||
await self._logger.warn(
|
|
||||||
"woc bridge HTTP request timeout, retrying",
|
|
||||||
method=method,
|
|
||||||
path=path,
|
|
||||||
attempt=attempt,
|
|
||||||
trace_id=trace_id,
|
|
||||||
)
|
|
||||||
await asyncio.sleep(
|
|
||||||
min(
|
|
||||||
_HTTP_RETRY_WAIT_MAX_SECONDS,
|
|
||||||
_HTTP_RETRY_WAIT_MIN_SECONDS * (2 ** (attempt - 1)),
|
|
||||||
)
|
|
||||||
)
|
|
||||||
except httpx.NetworkError as exc:
|
except httpx.NetworkError as exc:
|
||||||
|
# 仅对网络错误(连接失败/DNS 失败等瞬时抖动)做 1 次快速重试,
|
||||||
|
# 不等待。超时与 5xx 不重试,交给 outbox 退避重试机制处理,
|
||||||
|
# 避免双重重试放大延迟。
|
||||||
last_exc = exc
|
last_exc = exc
|
||||||
if attempt == _HTTP_RETRY_ATTEMPTS:
|
if attempt == _HTTP_NETWORK_ERROR_MAX_ATTEMPTS:
|
||||||
break
|
break
|
||||||
await self._logger.warn(
|
await self._logger.warn(
|
||||||
"woc bridge HTTP request network error, retrying",
|
"woc bridge HTTP request network error, retrying",
|
||||||
@ -582,33 +603,14 @@ class WocBridgeClient:
|
|||||||
attempt=attempt,
|
attempt=attempt,
|
||||||
trace_id=trace_id,
|
trace_id=trace_id,
|
||||||
)
|
)
|
||||||
await asyncio.sleep(
|
|
||||||
min(
|
|
||||||
_HTTP_RETRY_WAIT_MAX_SECONDS,
|
|
||||||
_HTTP_RETRY_WAIT_MIN_SECONDS * (2 ** (attempt - 1)),
|
|
||||||
)
|
|
||||||
)
|
|
||||||
except httpx.HTTPStatusError as exc:
|
except httpx.HTTPStatusError as exc:
|
||||||
if exc.response.status_code < 500:
|
if exc.response.status_code < 500:
|
||||||
# 4xx(含 429)为客户端/限流错误,不重试
|
# 4xx(含 429)为客户端/限流错误,不重试
|
||||||
raise
|
raise
|
||||||
|
# 5xx 服务端错误不重试,交给 outbox 退避重试机制处理
|
||||||
last_exc = exc
|
last_exc = exc
|
||||||
if attempt == _HTTP_RETRY_ATTEMPTS:
|
break
|
||||||
break
|
# httpx.TimeoutException 不重试,交给 outbox 退避重试机制处理
|
||||||
await self._logger.warn(
|
|
||||||
"woc bridge HTTP request server error, retrying",
|
|
||||||
method=method,
|
|
||||||
path=path,
|
|
||||||
attempt=attempt,
|
|
||||||
status=exc.response.status_code,
|
|
||||||
trace_id=trace_id,
|
|
||||||
)
|
|
||||||
await asyncio.sleep(
|
|
||||||
min(
|
|
||||||
_HTTP_RETRY_WAIT_MAX_SECONDS,
|
|
||||||
_HTTP_RETRY_WAIT_MIN_SECONDS * (2 ** (attempt - 1)),
|
|
||||||
)
|
|
||||||
)
|
|
||||||
if last_exc is not None:
|
if last_exc is not None:
|
||||||
raise last_exc
|
raise last_exc
|
||||||
except httpx.HTTPError as exc:
|
except httpx.HTTPError as exc:
|
||||||
@ -1030,8 +1032,10 @@ class WocBridgeClient:
|
|||||||
cause=Exception("client_inactive_after_stop"),
|
cause=Exception("client_inactive_after_stop"),
|
||||||
)
|
)
|
||||||
trace_id = uuid4().hex
|
trace_id = uuid4().hex
|
||||||
bridge_url = await self._build_bridge_url(account_id, "/api/messages/stream")
|
# 一次性解析 URL 与鉴权 token,避免双重凭证缓存查询
|
||||||
bridge_token = await self._get_bridge_token(account_id)
|
bridge_url, bridge_token = await self._resolve_endpoint(
|
||||||
|
account_id, "/api/messages/stream", None, None
|
||||||
|
)
|
||||||
request_headers: dict[str, str] = {
|
request_headers: dict[str, str] = {
|
||||||
"User-Agent": self._user_agent,
|
"User-Agent": self._user_agent,
|
||||||
"X-Trace-Id": trace_id,
|
"X-Trace-Id": trace_id,
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user