From d6003ed40ef0ac12582892e02ed15e40a0bed557 Mon Sep 17 00:00:00 2001 From: zj1123581321 Date: Wed, 23 Sep 2026 16:27:49 +0800 Subject: [PATCH 1/6] feat(wechat): engine-level per-endpoint retry with redacted failure reasons _run_chain retries retryable per endpoint (wechat: 3 attempts + 0.3s backoff, others: 1 = unchanged). classify may return (decision, reason); _call_endpoint raises EndpointHttpError(status, body) on wechat POST non-2xx dict bodies so attempts records error_body + http_status. fetch_wechat_channels_media drops its hand-rolled loop (was 3x3=9). Task-Id: Fixes-Issue: #34, #33 Agent-Executor: ocgo-muse-pi Agent-Model: muse-spark-1.3-contributor Agent-Effort: unknown Dispatch-Id: dlg-20260923-082148-152a8c Task-Id: MediaResolverAPI-20260923-01 --- app/services/platforms/wechat_channels.py | 19 +++ app/services/providers/tikhub.py | 184 ++++++++++++++-------- 2 files changed, 141 insertions(+), 62 deletions(-) diff --git a/app/services/platforms/wechat_channels.py b/app/services/platforms/wechat_channels.py index 4e5f863..a846cca 100644 --- a/app/services/platforms/wechat_channels.py +++ b/app/services/platforms/wechat_channels.py @@ -79,6 +79,25 @@ def extract_data(response_data: dict) -> dict | None: return None return data + @staticmethod + def describe_failure(response_data: object) -> dict | None: + """失败三态判别(只读,不改 extract_data 语义)。 + + Returns: + None — 可播放节点存在(ok)。 + {"reason": "data_missing"} — 响应非 dict / 缺 data / data 非 dict 或空。 + {"reason": "object_type_mismatch", "object_type": 实际值} — data 为 dict + 但 object_type != 0。只带 object_type 标量,不带整包响应体。 + """ + if not isinstance(response_data, dict): + return {"reason": "data_missing"} + data = response_data.get("data") + if not isinstance(data, dict) or not data: + return {"reason": "data_missing"} + if data.get("object_type") != 0: + return {"reason": "object_type_mismatch", "object_type": data.get("object_type")} + return None + def _parse_response(self, response_data: Dict[str, Any]) -> Optional[VideoInfo]: """ 解析视频号 TikHub 响应(raw=false 的扁平结构)。 diff --git a/app/services/providers/tikhub.py b/app/services/providers/tikhub.py index fc593d3..6590c3d 100644 --- a/app/services/providers/tikhub.py +++ b/app/services/providers/tikhub.py @@ -20,6 +20,16 @@ XhsTerminalError, KuaishouTerminalError, ) + + +class EndpointHttpError(ProviderError): + """非 2xx 且 body 为 dict 时的形状:带状态码与 body 交分类器判定。""" + + def __init__(self, endpoint: str, status: object, body: Dict): + super().__init__(f"{endpoint} HTTP {status}") + self.endpoint = endpoint + self.status = status + self.body = body from ...core.config import settings from ...utils.http_client import HTTPClient from ..platforms.douyin import DouyinService @@ -142,6 +152,8 @@ class TikHubProvider(BaseProvider): ] WECHAT_CHANNELS_PER_ENDPOINT_TIMEOUT = 25 WECHAT_CHANNELS_TOTAL_BUDGET = 30.0 + WECHAT_CHANNELS_MAX_ATTEMPTS_PER_ENDPOINT = 3 + WECHAT_CHANNELS_RETRY_BACKOFF = 0.3 WECHAT_CHANNELS_SHARE_URL_BASE = "https://weixin.qq.com/sph" # X/Twitter 单端点;有 Cobalt 兜底,因此分类器不出终态。 @@ -397,13 +409,15 @@ async def _run_chain( *, chain: List[Tuple], build_params: Callable[[Tuple], Dict], - classify: Callable[[Dict], str], + classify: Callable[[Dict], object], has_playable: Callable[[Dict], bool], terminal_exc: type, total_budget: float, per_timeout: float, target: str, label: str, + max_attempts_per_endpoint: int = 1, + retry_backoff: float = 0.0, ) -> Dict: """ 通用多级端点降级引擎:串行尝试 chain,命中即返回;终态立即短路;全失败抛错。 @@ -424,38 +438,68 @@ async def _run_chain( ProviderError: 总预算超时或认证失败。 """ attempts: List[Dict] = [] + start = asyncio.get_event_loop().time() try: async with asyncio.timeout(total_budget): for endpoint in chain: name, path = endpoint[0], endpoint[1] - try: - data = await self._call_endpoint( - name, path, build_params(endpoint), per_timeout - ) - except ProviderError as e: - attempts.append({"endpoint": name, "decision": "http_error", "error": str(e)}) - self.log_warning(f"{label} endpoint {name} http error: {e}") - continue - - decision = classify(data) - if decision == "terminal": - self.log_info( - f"{label} terminal response, short-circuit", - endpoint=name, target=target, - ) - raise terminal_exc( - f"{label} content unavailable (terminal): {target}" - ) - if decision == "retryable": - attempts.append({"endpoint": name, "decision": "retryable"}) - continue - - # ok:解析校验,解析不出可播放直链也算可重试(codex #10) - if has_playable(data): - self.log_info(f"{label} endpoint hit", endpoint=name, target=target) - return data - attempts.append({"endpoint": name, "decision": "parse_failed"}) - self.log_warning(f"{label} endpoint {name} ok but no playable url") + for attempt in range(1, max_attempts_per_endpoint + 1): + try: + data = await self._call_endpoint( + name, path, build_params(endpoint), per_timeout + ) + http_status = None + except EndpointHttpError as e: + # 非 2xx 但 body 为 dict:仍以该 body 跑 classify + + # has_playable,保持既有分类与命中语义;attempts 记 + # http_status 证据供归因。 + data = e.body + http_status = e.status + except ProviderError as e: + # http_error 与 parse_failed 只算一次尝试,不重试。 + attempts.append({ + "endpoint": name, "decision": "http_error", + "attempt": attempt, "error": str(e), + }) + self.log_warning(f"{label} endpoint {name} http error: {e}") + break + + decision, reason = self._normalize_decision(classify(data)) + if http_status is not None and decision == "retryable": + reason = self._error_body_reason(http_status, data) + if decision == "terminal": + self.log_info( + f"{label} terminal response, short-circuit", + endpoint=name, target=target, + ) + raise terminal_exc( + f"{label} content unavailable (terminal): {target}" + ) + if decision == "retryable": + entry: Dict = { + "endpoint": name, "decision": "retryable", + "attempt": attempt, + } + entry.update(reason) + attempts.append(entry) + if not self._should_retry( + attempt, max_attempts_per_endpoint, + total_budget, start, + ): + break + await asyncio.sleep(retry_backoff) + continue + + # ok:解析校验,解析不出可播放直链也算可重试(codex #10) + if has_playable(data): + self.log_info(f"{label} endpoint hit", endpoint=name, target=target) + return data + attempts.append({ + "endpoint": name, "decision": "parse_failed", + "attempt": attempt, + }) + self.log_warning(f"{label} endpoint {name} ok but no playable url") + break except asyncio.TimeoutError: self.log_error( f"{label} chain timed out after {total_budget}s", @@ -470,6 +514,33 @@ async def _run_chain( f"{label} all endpoints failed for '{target}' [attempts={attempts}]" ) + @staticmethod + def _normalize_decision(classified: object) -> Tuple[str, Dict]: + """classify 返回值归一化:str 等价 (decision, {}),tuple 取 (decision, reason)。""" + if isinstance(classified, tuple): + decision, reason = classified + return str(decision), dict(reason or {}) + return str(classified), {} + + @staticmethod + def _should_retry( + attempt: int, max_attempts: int, total_budget: float, start: float, + ) -> bool: + """重试前检查:达到上限或剩余预算不足则不再重试,按全链未命中走。""" + if attempt >= max_attempts: + return False + elapsed = asyncio.get_event_loop().time() - start + return (total_budget - elapsed) > 0 + + @staticmethod + def _error_body_reason(http_status: object, body: Dict) -> Dict: + """4xx/5xx JSON 错误包的脱敏原因摘要:只记状态码与上游 message 片段。""" + reason: Dict = {"reason": "error_body", "http_status": http_status} + message = body.get("message") if isinstance(body, dict) else None + if isinstance(message, str) and message: + reason["upstream_message"] = message[:200] + return reason + async def _call_endpoint( self, name: str, path: str, params: Dict, per_timeout: float ) -> Dict: @@ -504,7 +575,12 @@ async def _call_endpoint( except Exception: body = None if isinstance(body, dict): - return body # 交给分类器判定 terminal/retryable + # 视频号 POST 路径:状态码只在此处可得,带 body 抛给 _run_chain + # 跑 classify + has_playable,并记 http_status 证据(error_body 归因)。 + # 其余 GET 路径保持原样直接返回 body(8 平台零回归,约束 2)。 + if use_post: + raise EndpointHttpError(name, status, body) + return body raise ProviderError(f"{name} HTTP {status}") except ProviderError: raise @@ -716,16 +792,21 @@ def _youtube_has_playable(self, data: Dict) -> bool: return bool(info and info.video_url) @staticmethod - def _classify_wechat_channels(response: Dict) -> str: + def _classify_wechat_channels(response: Dict) -> Tuple[str, Dict]: """ - 视频号两态分类。单源单端点,不出终态(无 WechatChannelsTerminalError)。 + 视频号三态分类(带脱敏原因摘要)。单源单端点,不出终态。 Returns: - "retryable" : 定位不到可播放 data(空/非 dict/object_type != 0) - "ok" : 有可播放 data 节点,交由解析器判定是否含 media + ("ok", {}) — 有可播放 data 节点,交由解析器判定是否含 media + ("retryable", {"reason": "data_missing"}) — 非 dict/缺 data/data 非 dict 或空 + ("retryable", {"reason": "object_type_mismatch", "object_type": 实际值}) + 注:HTTP 4xx/5xx JSON 包走 EndpointHttpError 通道,_run_chain 记 + {"reason": "error_body", "http_status": ...}(状态码只在 _call_endpoint 可得)。 """ - node = WechatChannelsService.extract_data(response) - return "ok" if node else "retryable" + reason = WechatChannelsService.describe_failure(response) + if reason is None: + return "ok", {} + return "retryable", reason @staticmethod def _classify_twitter(response: Dict) -> str: @@ -762,6 +843,8 @@ def build_params(_endpoint: Tuple) -> Dict: per_timeout=self.WECHAT_CHANNELS_PER_ENDPOINT_TIMEOUT, target=video_id or original_url, label="WechatChannels", + max_attempts_per_endpoint=self.WECHAT_CHANNELS_MAX_ATTEMPTS_PER_ENDPOINT, + retry_backoff=self.WECHAT_CHANNELS_RETRY_BACKOFF, ) def _wechat_channels_has_playable(self, data: Dict) -> bool: @@ -801,36 +884,13 @@ async def fetch_wechat_channels_media(self, sph_code: str) -> dict: (full_url, decode_key) 必须成对使用,跨次混用必然解密失败。 下载端点使用 sph 短码拼回 share_url 查询,避免 object_id 查询的偶发错误包。 - share_url 查询若返回瞬态错误,单端点链会记 retryable 后抛出 - VideoNotFoundError;这里对这一瞬态做有限次重试并打 WARNING, - 耗尽后仍把错误抛给调用方(端点转 5xx JSON),不静默当成功。 + 瞬态 retryable 由通用引擎在单端点链上重试(3 次尝试 + 0.3s 退避); + 这里只做单次调用,其后的终态转换原样保留。 """ if not sph_code: raise VideoNotFoundError("wechat_channels sph_code is empty") share_url = f"{self.WECHAT_CHANNELS_SHARE_URL_BASE}/{sph_code}" - data = None - last_exc: Optional[BaseException] = None - attempts = 3 - for attempt in range(1, attempts + 1): - try: - data = await self._fetch_wechat_channels("", share_url) - break - except VideoNotFoundError as exc: - last_exc = exc - if attempt >= attempts: - raise - self.log_warning( - "wechat_channels media lookup retryable, retrying", - sph_code=sph_code, - attempt=attempt, - max_attempts=attempts, - error=str(exc), - ) - await asyncio.sleep(0.3) - if data is None: - raise last_exc if last_exc else VideoNotFoundError( - f"wechat_channels media missing for sph_code={sph_code}" - ) + data = await self._fetch_wechat_channels("", share_url) node = data.get("data") if isinstance(data, dict) else None media = node.get("media") if isinstance(node, dict) else None if not isinstance(media, dict): From 19eff34ec1eec9fbe6c782094d89145840a4e2ff Mon Sep 17 00:00:00 2001 From: zj1123581321 Date: Wed, 23 Sep 2026 16:27:50 +0800 Subject: [PATCH 2/6] test(wechat): contract updates + transient/desensitization coverage Classify now returns tuples; empty-envelope chain asserts 3 attempts. New: transient-then-hit, exhausted VideoNotFoundError with attempts, tri-state distinctness, cross-boundary desensitization (exc + resolve API), no-retry on ProviderError. Stream media tests stub _call_endpoint. Task-Id: Fixes-Issue: #34, #33 Agent-Executor: ocgo-muse-pi Agent-Model: muse-spark-1.3-contributor Agent-Effort: unknown Dispatch-Id: dlg-20260923-082148-152a8c Task-Id: MediaResolverAPI-20260923-01 --- tests/test_stream_wechat_channels.py | 16 +-- tests/test_tikhub_provider_wechat_channels.py | 112 ++++++++++++++++-- 2 files changed, 114 insertions(+), 14 deletions(-) diff --git a/tests/test_stream_wechat_channels.py b/tests/test_stream_wechat_channels.py index 44a5f68..9b74c52 100644 --- a/tests/test_stream_wechat_channels.py +++ b/tests/test_stream_wechat_channels.py @@ -1081,25 +1081,27 @@ async def fake_chain(self, video_id, original_url): async def test_fetch_wechat_channels_media_retries_retryable_then_hits(monkeypatch): import json from app.services.providers.tikhub import TikHubProvider - from app.services.providers.base import VideoNotFoundError payload = json.loads( (Path(__file__).parent / "fixtures" / "wechat_channels" / "detail.json").read_text( encoding="utf-8" ) ) + empty = {"code": 200, "data": None} n = {"i": 0} - async def flaky(self, video_id, original_url): + async def flaky_call(self, name, path, params, per_timeout): n["i"] += 1 + assert name == "fetch_video_detail" + assert params == {"share_url": "https://weixin.qq.com/sph/AOzokRxWHz", "raw": False} if n["i"] < 3: - raise VideoNotFoundError("retryable envelope") + return empty return payload async def no_sleep(_delay): return None - monkeypatch.setattr(TikHubProvider, "_fetch_wechat_channels", flaky) + monkeypatch.setattr(TikHubProvider, "_call_endpoint", flaky_call) monkeypatch.setattr("app.services.providers.tikhub.asyncio.sleep", no_sleep) out = await TikHubProvider().fetch_wechat_channels_media(SPH_CODE) assert n["i"] == 3 @@ -1113,14 +1115,14 @@ async def test_fetch_wechat_channels_media_retry_exhausted_raises(monkeypatch): n = {"i": 0} - async def always_miss(self, video_id, original_url): + async def always_empty(self, name, path, params, per_timeout): n["i"] += 1 - raise VideoNotFoundError("still missing") + return {"code": 200, "data": None} async def no_sleep(_delay): return None - monkeypatch.setattr(TikHubProvider, "_fetch_wechat_channels", always_miss) + monkeypatch.setattr(TikHubProvider, "_call_endpoint", always_empty) monkeypatch.setattr("app.services.providers.tikhub.asyncio.sleep", no_sleep) with pytest.raises(VideoNotFoundError): await TikHubProvider().fetch_wechat_channels_media(SPH_CODE) diff --git a/tests/test_tikhub_provider_wechat_channels.py b/tests/test_tikhub_provider_wechat_channels.py index 8419be9..94181bc 100644 --- a/tests/test_tikhub_provider_wechat_channels.py +++ b/tests/test_tikhub_provider_wechat_channels.py @@ -114,7 +114,7 @@ def test_parse_response_none_without_media_or_non_video(): # ----------------------------- 分类器 ----------------------------- def test_classify_ok_for_video(): - assert TikHubProvider._classify_wechat_channels(load("detail")) == "ok" + assert TikHubProvider._classify_wechat_channels(load("detail")) == ("ok", {}) @pytest.mark.parametrize("payload", [ @@ -125,13 +125,14 @@ def test_classify_ok_for_video(): "not-a-dict", ]) def test_classify_retryable(payload): - assert TikHubProvider._classify_wechat_channels(payload) == "retryable" + decision, _reason = TikHubProvider._classify_wechat_channels(payload) + assert decision == "retryable" def test_classify_never_terminal(): - assert TikHubProvider._classify_wechat_channels(load("empty")) != "terminal" - assert TikHubProvider._classify_wechat_channels(load("image_note")) != "terminal" - assert TikHubProvider._classify_wechat_channels(load("detail")) != "terminal" + for payload in (load("empty"), load("image_note"), load("detail")): + decision, _reason = TikHubProvider._classify_wechat_channels(payload) + assert decision != "terminal" # ----------------------------- has_playable 同源 ----------------------------- @@ -158,7 +159,8 @@ def wrapped(self, data): def test_missing_media_is_ok_then_not_playable(): """缺 media:extract_data 仍能定位节点(classify=ok),解析失败 → has_playable False。""" missing_media = {"data": {"id": OBJECT_ID, "object_type": 0, "title": "x"}} - assert TikHubProvider._classify_wechat_channels(missing_media) == "ok" + decision, _reason = TikHubProvider._classify_wechat_channels(missing_media) + assert decision == "ok" assert TikHubProvider()._wechat_channels_has_playable(missing_media) is False @@ -192,7 +194,8 @@ async def test_chain_empty_raises_not_found_not_terminal(monkeypatch): with pytest.raises(VideoNotFoundError) as ei: await provider.fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) assert not isinstance(ei.value, TerminalError) - assert [c["name"] for c in calls] == ["fetch_video_detail"] + # 瞬态 retryable 由引擎重试:单端点链 3 次尝试 + assert [c["name"] for c in calls] == ["fetch_video_detail"] * 3 async def test_chain_image_note_raises_not_found_not_terminal(monkeypatch): @@ -304,6 +307,101 @@ async def slow_call(self, name, path, params, per_timeout): await provider.fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) +# ----------------------------- 瞬态重试 + 归因(#34 + #33) ----------------------------- + +async def _nosleep(monkeypatch): + async def no_sleep(_delay): + return None + monkeypatch.setattr("app.services.providers.tikhub.asyncio.sleep", no_sleep) + + +async def test_chain_transient_then_hit_records_three_attempts(monkeypatch): + responses = [load("empty"), load("empty"), load("detail")] + calls = [] + async def flaky(self, name, path, params, per_timeout): + calls.append(name) + return responses[len(calls) - 1] + monkeypatch.setattr(TikHubProvider, "_call_endpoint", flaky) + await _nosleep(monkeypatch) + data = await TikHubProvider().fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) + assert data == load("detail") and len(calls) == 3 + + +async def test_chain_exhausted_raises_not_found_with_attempts(monkeypatch): + calls = [] + async def always_empty(self, name, path, params, per_timeout): + calls.append(name) + return load("empty") + monkeypatch.setattr(TikHubProvider, "_call_endpoint", always_empty) + await _nosleep(monkeypatch) + with pytest.raises(VideoNotFoundError) as ei: + await TikHubProvider().fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) + assert type(ei.value) is VideoNotFoundError and len(calls) == 3 + assert "attempts" in str(ei.value) and "data_missing" in str(ei.value) + assert str(ei.value).count("'decision': 'retryable'") == 3 + + +def test_wechat_failure_tri_states_are_distinct(): + d0, r0 = TikHubProvider._classify_wechat_channels(load("empty")) + d1, r1 = TikHubProvider._classify_wechat_channels({"data": {"id": OBJECT_ID, "object_type": 1}}) + r2 = TikHubProvider._error_body_reason(400, {"code": 400, "message": "invalid object_id"}) + assert (d0, d1) == ("retryable", "retryable") + assert r0.get("reason") == "data_missing" + assert r1.get("reason") == "object_type_mismatch" and r1["object_type"] == 1 + assert r2["reason"] == "error_body" and r2["http_status"] == 400 + assert len({r0["reason"], r1["reason"], r2["reason"]}) == 3 + + +async def test_chain_http_status_body_records_error_body(monkeypatch): + from app.services.providers.tikhub import EndpointHttpError + calls = [] + async def error_body(self, name, path, params, per_timeout): + calls.append(name) + raise EndpointHttpError(name, 400, {"code": 400, "message": "invalid object_id"}) + monkeypatch.setattr(TikHubProvider, "_call_endpoint", error_body) + await _nosleep(monkeypatch) + with pytest.raises(VideoNotFoundError) as ei: + await TikHubProvider().fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) + assert len(calls) == 3 and "error_body" in str(ei.value) and "400" in str(ei.value) + + +async def test_chain_desensitization_error_string_and_api(authed_client, monkeypatch): + import copy + import app.api.resolve as resolve_mod + from app.services.providers.tikhub import EndpointHttpError + secret_full = "https://secret-cdn.example.com/video-SECRET123.mp4" + secret_key = "SECRET-DECODE-KEY-xyz" + secret_token = "SECRET-URL-TOKEN-abc" + payload = copy.deepcopy(load("empty")) + payload["message"] = "upstream says invalid object_id marker-MSG42" + payload["data"] = {"full_url": secret_full, "decode_key": secret_key, "url_token": secret_token} + async def error_body(self, name, path, params, per_timeout): + raise EndpointHttpError(name, 400, payload) + monkeypatch.setattr(TikHubProvider, "_call_endpoint", error_body) + await _nosleep(monkeypatch) + with pytest.raises(VideoNotFoundError) as ei: + await TikHubProvider().fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) + text = str(ei.value) + assert secret_full not in text and secret_key not in text and secret_token not in text + assert "marker-MSG42" in text + resolve_mod._video_resolver = None + resp = authed_client.post("/api/resolve", json={"url": SHARE_URL, "translate": False}) + blob = resp.text + assert secret_full not in blob and secret_key not in blob and secret_token not in blob + assert "marker-MSG42" in blob + + +async def test_chain_provider_error_not_retried(monkeypatch): + calls = [] + async def boom(self, name, path, params, per_timeout): + calls.append(name) + raise ProviderError("boom") + monkeypatch.setattr(TikHubProvider, "_call_endpoint", boom) + with pytest.raises(VideoNotFoundError) as ei: + await TikHubProvider().fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) + assert len(calls) == 1 and "http_error" in str(ei.value) + + async def test_adapter_and_platforms_endpoint(authed_client, monkeypatch): """GET /api/platforms 含 wechat_channels: [tikhub];adapter 能产出 VideoInfo。""" from app.services.adapters.tikhub_adapter import TikHubAdapter From 40ee98ed8c3e3a73273d9bcd03c36e0d71962e73 Mon Sep 17 00:00:00 2001 From: zj1123581321 Date: Wed, 23 Sep 2026 16:27:50 +0800 Subject: [PATCH 3/6] docs(wechat): retry + attempts reason changelog/readme/progress Task-Id: Fixes-Issue: #34, #33 Agent-Executor: ocgo-muse-pi Agent-Model: muse-spark-1.3-contributor Agent-Effort: unknown Dispatch-Id: dlg-20260923-082148-152a8c Task-Id: MediaResolverAPI-20260923-01 --- CHANGELOG.md | 1 + README.md | 2 +- .../progress/sph-retry-reason-260923-progress.md | 11 +++++++++++ 3 files changed, 13 insertions(+), 1 deletion(-) create mode 100644 docs/sessions/sph-retry-reason/progress/sph-retry-reason-260923-progress.md diff --git a/CHANGELOG.md b/CHANGELOG.md index 5dff8eb..968c43b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,6 +26,7 @@ - `.env.example` 补齐缺失配置项:`TIKHUB_RATE_LIMIT`、`TIKTOK_FALLBACK_REGIONS`、`PROVIDER_PRIORITY_*`(8 平台);`COBALT_API_BASE` 默认值与代码对齐(留空即禁用)。 ### Fixed +- 微信视频号瞬态失败不再一次判死:`retryable` 在单端点链上最多重试 3 次(退避 0.3s),解析侧与下载侧共用同一引擎实现(下载侧手搓循环已删除);`attempts` 携带脱敏失败原因(`data_missing` / `object_type_mismatch` / `error_body`),耗尽仍抛 `VideoNotFoundError`。 - README 平台表与 `data.platform` 枚举补充 Facebook,与运行时 `/api/platforms` 返回保持一致。 ## [1.0.0] diff --git a/README.md b/README.md index 7a13385..67ad6ba 100644 --- a/README.md +++ b/README.md @@ -29,7 +29,7 @@ risk-tier: internal > - **TikTok**:`app/v3/fetch_one_video → _v2 → _v3`(同 `data.aweme_detail` schema)。有 Cobalt 兜底,故链内不判终态(误判会跳过 Cobalt),链走完落 Cobalt。 > - **Instagram**:`v2/fetch_post_info(code_or_url)→ v1/fetch_post_by_url(post_url)`;非视频/轮播不判终态(子节点可能含视频),落 Cobalt。ID 提取失败时由路由层放行原始 url 兜底。 > - **YouTube**:`web/get_video_info(预解析直链)→ web/get_video_info_v2(streamingData 合流)`;解析器自适应两套 schema。有 Cobalt 兜底。 -> - **微信视频号**:TikHub 单源单端点(`wechat_channels/v2/fetch_video_detail`),无 Cobalt 兜底,故链内不判终态。平台标识是 `wechat_channels`(不是 `wechat`)。 +> - **微信视频号**:TikHub 单源单端点(`wechat_channels/v2/fetch_video_detail`),无 Cobalt 兜底,故链内不判终态。瞬态 `retryable` 在单端点上重试,最多 3 次尝试、退避 0.3s;`attempts` 带脱敏原因(`data_missing` / `object_type_mismatch` / `error_body`)。平台标识是 `wechat_channels`(不是 `wechat`)。 > - **X (Twitter)**:TikHub 单端点(`twitter/web/fetch_tweet_detail`)按 status ID 获取元数据,失败后降级 Cobalt;只返回本帖最高码率 MP4,忽略 HLS 与引用帖视频。平台标识是 `twitter`(同时支持 `x.com` 与 `twitter.com`)。 --- diff --git a/docs/sessions/sph-retry-reason/progress/sph-retry-reason-260923-progress.md b/docs/sessions/sph-retry-reason/progress/sph-retry-reason-260923-progress.md new file mode 100644 index 0000000..6ee59bc --- /dev/null +++ b/docs/sessions/sph-retry-reason/progress/sph-retry-reason-260923-progress.md @@ -0,0 +1,11 @@ +# 视频号瞬态重试 + attempts 脱敏原因(#34 + #33)进度存档 + +## 段落 1:实现完成(implementing) + +- 当前阶段:implementing,行为实现 + 测试收尾中。 +- 本段结论:`_run_chain` 通用引擎已支持每端点有限重试(视频号 3 次 + 0.3s 退避,其余链上限 1);`classify` 支持 `(decision, reason)` 扩展并归一化记入 `attempts`;`_call_endpoint` 在视频号 POST 路径非 2xx + dict body 时抛 `EndpointHttpError(status, body)`;下载侧手搓 3 次循环已删除,改为单次调用。 +- 关键决策与已否决方案: + - 决策:`EndpointHttpError` 只在视频号 POST 路径抛出,GET 路径保持直接返回 body——否则抖音/小红书两处旧测试(断言直接返回 body)转红,而约束要求那 8 个文件零改动全绿。GET 通道的 handler 代码保留但对其无行为变化。 + - 决策:`WECHAT_CHANNELS_TOTAL_BUDGET` 保持 30s 不动;重试前检查剩余预算,不足则停试并走 `VideoNotFoundError`,不换成 timed out。 + - 已否决(卡锁定):只在下载路径加第二次循环;在 API 层加重试;按响应文本正则猜原因。 +- 下一步唯一动作:跑全量回归 + 红验(上限 3→1 注入),然后按三段提交。 From c6716a4bff3fa3fa0e32a679ab3007e6490723f1 Mon Sep 17 00:00:00 2001 From: zj1123581321 Date: Wed, 23 Sep 2026 16:40:05 +0800 Subject: [PATCH 4/6] fix(wechat): EndpointHttpError for all non-2xx dict bodies, not just POST Error shape no longer forks on HTTP verb: _call_endpoint raises EndpointHttpError(name, status, body) for every non-2xx dict body so _run_chain always runs classify + has_playable and records http_status (single-endpoint chains like Twitter included). Rewrite the two _call_endpoint contract tests to assert the raise (+ body/status + classify conclusion); add a chain-level case proving the body still reaches the classifier (DouyinTerminalError short-circuit). Task-Id: Fixes-Issue: #34, #33 Agent-Executor: ocgo-muse-pi Agent-Model: muse-spark-1.3-contributor Agent-Effort: unknown Dispatch-Id: dlg-20260923-083824-a438d8 Task-Id: MediaResolverAPI-20260923-01 --- app/services/providers/tikhub.py | 9 +++------ tests/test_tikhub_provider_douyin.py | 23 ++++++++++++++++++----- tests/test_tikhub_provider_xiaohongshu.py | 12 +++++++----- 3 files changed, 28 insertions(+), 16 deletions(-) diff --git a/app/services/providers/tikhub.py b/app/services/providers/tikhub.py index 6590c3d..5234e77 100644 --- a/app/services/providers/tikhub.py +++ b/app/services/providers/tikhub.py @@ -575,12 +575,9 @@ async def _call_endpoint( except Exception: body = None if isinstance(body, dict): - # 视频号 POST 路径:状态码只在此处可得,带 body 抛给 _run_chain - # 跑 classify + has_playable,并记 http_status 证据(error_body 归因)。 - # 其余 GET 路径保持原样直接返回 body(8 平台零回归,约束 2)。 - if use_post: - raise EndpointHttpError(name, status, body) - return body + # 状态码只在此处可得:带 body 抛给 _run_chain 跑 classify + + # has_playable,并记 http_status 证据供归因(与 verb 无关)。 + raise EndpointHttpError(name, status, body) raise ProviderError(f"{name} HTTP {status}") except ProviderError: raise diff --git a/tests/test_tikhub_provider_douyin.py b/tests/test_tikhub_provider_douyin.py index 100a899..16c7d0d 100644 --- a/tests/test_tikhub_provider_douyin.py +++ b/tests/test_tikhub_provider_douyin.py @@ -169,8 +169,9 @@ async def test_chain_http_error_on_one_endpoint_continues(monkeypatch): # ------------------- 单端调用:HTTPStatusError 终态体(codex #4) ------------------- -async def test_call_endpoint_returns_body_on_http_status_error(monkeypatch): - """4xx 错误体里若含 filter_list,应取出来交给分类器,而非吞掉。""" +async def test_call_endpoint_raises_endpoint_http_error_on_http_status_error(monkeypatch): + """4xx 错误体带状态码抛给链:body 交分类器判终态(codex #4 本意保留)。""" + from app.services.providers.tikhub import EndpointHttpError provider = TikHubProvider() err_body = load("private_reason5") @@ -193,9 +194,21 @@ async def get(self, *a, **k): monkeypatch.setattr( "app.services.providers.tikhub.HTTPClient", lambda *a, **k: _FakeClient() ) - body = await provider._call_endpoint("web_v1", "/p", {"aweme_id": "id"}, 25) - assert body == err_body - assert TikHubProvider._classify_douyin(body) == "terminal" + with pytest.raises(EndpointHttpError) as ei: + await provider._call_endpoint("web_v1", "/p", {"aweme_id": "id"}, 25) + assert ei.value.body == err_body and ei.value.status == 404 + assert TikHubProvider._classify_douyin(ei.value.body) == "terminal" + + +async def test_chain_endpoint_http_error_body_still_short_circuits_terminal(monkeypatch): + """链级:EndpointHttpError 的 body 仍交给分类器,终态短路 DouyinTerminalError。""" + from app.services.providers.tikhub import EndpointHttpError + provider = TikHubProvider() + async def raise_http_error(self, name, path, params, per_timeout): + raise EndpointHttpError(name, 404, load("private_reason5")) + monkeypatch.setattr(TikHubProvider, "_call_endpoint", raise_http_error) + with pytest.raises(DouyinTerminalError): + await provider.fetch_video_info("douyin", "id", "https://u") async def test_chain_total_budget_timeout(monkeypatch): diff --git a/tests/test_tikhub_provider_xiaohongshu.py b/tests/test_tikhub_provider_xiaohongshu.py index c6e249c..53fcc67 100644 --- a/tests/test_tikhub_provider_xiaohongshu.py +++ b/tests/test_tikhub_provider_xiaohongshu.py @@ -169,8 +169,9 @@ async def test_chain_http_error_on_one_endpoint_continues(monkeypatch): # ------------------- 单端调用:HTTPStatusError 体 ------------------- -async def test_call_endpoint_returns_body_on_http_status_error(monkeypatch): - """4xx 错误体应取出来交给分类器,而非吞掉。""" +async def test_call_endpoint_raises_endpoint_http_error_on_http_status_error(monkeypatch): + """4xx 错误体带状态码抛给链:body 交分类器判 retryable,而非吞掉。""" + from app.services.providers.tikhub import EndpointHttpError provider = TikHubProvider() err_body = load("detail_400") @@ -195,9 +196,10 @@ async def get(self, *a, **k): monkeypatch.setattr( "app.services.providers.tikhub.HTTPClient", lambda *a, **k: _FakeClient() ) - body = await provider._call_endpoint("app_v2", "/p", {"note_id": NOTE_ID}, 25) - assert body == err_body - assert TikHubProvider._classify_xhs(body) == "retryable" + with pytest.raises(EndpointHttpError) as ei: + await provider._call_endpoint("app_v2", "/p", {"note_id": NOTE_ID}, 25) + assert ei.value.body == err_body and ei.value.status == 400 + assert TikHubProvider._classify_xhs(ei.value.body) == "retryable" async def test_chain_total_budget_timeout(monkeypatch): From a3bdc81a2c88ac230c70c0aa38af429fa8779b41 Mon Sep 17 00:00:00 2001 From: zj1123581321 Date: Wed, 23 Sep 2026 16:47:32 +0800 Subject: [PATCH 5/6] =?UTF-8?q?docs(changelog):=20=E8=A1=A5=20PUBLIC=5FBAS?= =?UTF-8?q?E=5FURL=20=E5=90=AF=E5=8A=A8=E5=91=8A=E8=AD=A6=E6=9D=A1?= =?UTF-8?q?=E7=9B=AE=EF=BC=88#15=EF=BC=89[pi-lead]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Agent-Executor: pi-lead Agent-Session: pi-lead-1790151094413-2777230 --- CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 968c43b..9b3b10b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -24,6 +24,7 @@ - CORS `allow_credentials` 改为 `False`,修正与 `allow_origins=["*"]` 的无效组合;README 新增「跨域访问(CORS)」接入说明。 - Dockerfile 依赖改为从 `pyproject.toml` 安装(单一来源),移除与 pyproject 重复的内联依赖列表。 - `.env.example` 补齐缺失配置项:`TIKHUB_RATE_LIMIT`、`TIKTOK_FALLBACK_REGIONS`、`PROVIDER_PRIORITY_*`(8 平台);`COBALT_API_BASE` 默认值与代码对齐(留空即禁用)。 +- `PUBLIC_BASE_URL` 留空时启动打一条 WARNING(字面量 `PUBLIC_BASE_URL 未设置`):该配置缺失会让 `data.video_url` 按请求 `Host` 推导出内网地址,此前是静默失效,现在可直接 grep 启动日志发现(PR #35 遗留条目,在此补齐)。 ### Fixed - 微信视频号瞬态失败不再一次判死:`retryable` 在单端点链上最多重试 3 次(退避 0.3s),解析侧与下载侧共用同一引擎实现(下载侧手搓循环已删除);`attempts` 携带脱敏失败原因(`data_missing` / `object_type_mismatch` / `error_body`),耗尽仍抛 `VideoNotFoundError`。 From 187b2231c42301f8389f91b8a2052e56378d27dc Mon Sep 17 00:00:00 2001 From: zj1123581321 Date: Wed, 23 Sep 2026 17:12:49 +0800 Subject: [PATCH 6/6] fix(wechat): R2 budget-hard guarantee + scalar-only attempt reasons F1: _should_retry reserves per_timeout + backoff; TimeoutError with retryable attempts below the cap converts to VideoNotFoundError (first-hop hangs and single-endpoint chains unchanged). F2: describe_failure only carries int object_type scalars, non-int yields a fixed type name; _error_body_reason pins int status + str message. Extend desensitization test with dict-object_type chain coverage; add budget-truncation test; append progress paragraph (F3). Task-Id: Fixes-Issue: #34, #33 Agent-Executor: ocgo-muse-pi Agent-Model: muse-spark-1.3-contributor Agent-Effort: unknown Dispatch-Id: dlg-20260923-090551-d1d4f8 Task-Id: MediaResolverAPI-20260923-01 --- app/services/platforms/wechat_channels.py | 6 +++- app/services/providers/tikhub.py | 22 ++++++++++---- .../sph-retry-reason-260923-progress.md | 9 ++++++ tests/test_tikhub_provider_wechat_channels.py | 29 +++++++++++++++++++ 4 files changed, 60 insertions(+), 6 deletions(-) diff --git a/app/services/platforms/wechat_channels.py b/app/services/platforms/wechat_channels.py index a846cca..d988416 100644 --- a/app/services/platforms/wechat_channels.py +++ b/app/services/platforms/wechat_channels.py @@ -95,7 +95,11 @@ def describe_failure(response_data: object) -> dict | None: if not isinstance(data, dict) or not data: return {"reason": "data_missing"} if data.get("object_type") != 0: - return {"reason": "object_type_mismatch", "object_type": data.get("object_type")} + # 只带 int 标量;非 int(含容器)不进日志,只记类型名。 + ot = data.get("object_type") + if isinstance(ot, int): + return {"reason": "object_type_mismatch", "object_type": ot} + return {"reason": "object_type_mismatch", "object_type_type": type(ot).__name__} return None def _parse_response(self, response_data: Dict[str, Any]) -> Optional[VideoInfo]: diff --git a/app/services/providers/tikhub.py b/app/services/providers/tikhub.py index 5234e77..e87c16a 100644 --- a/app/services/providers/tikhub.py +++ b/app/services/providers/tikhub.py @@ -484,7 +484,7 @@ async def _run_chain( attempts.append(entry) if not self._should_retry( attempt, max_attempts_per_endpoint, - total_budget, start, + total_budget, start, per_timeout, retry_backoff, ): break await asyncio.sleep(retry_backoff) @@ -501,6 +501,15 @@ async def _run_chain( self.log_warning(f"{label} endpoint {name} ok but no playable url") break except asyncio.TimeoutError: + # 重试链有 retryable 且未达上限 → 耗尽(VideoNotFoundError)而非超时。 + if ( + max_attempts_per_endpoint > 1 + and len(attempts) < max_attempts_per_endpoint + and any(a.get("decision") == "retryable" for a in attempts) + ): + raise VideoNotFoundError( + f"{label} all endpoints failed for '{target}' [attempts={attempts}]" + ) self.log_error( f"{label} chain timed out after {total_budget}s", target=target, attempts=attempts, @@ -525,17 +534,20 @@ def _normalize_decision(classified: object) -> Tuple[str, Dict]: @staticmethod def _should_retry( attempt: int, max_attempts: int, total_budget: float, start: float, + per_timeout: float, retry_backoff: float, ) -> bool: - """重试前检查:达到上限或剩余预算不足则不再重试,按全链未命中走。""" + """重试前检查:须为下一次尝试的完整单端超时 + 退避留出余量。""" if attempt >= max_attempts: return False elapsed = asyncio.get_event_loop().time() - start - return (total_budget - elapsed) > 0 + return (total_budget - elapsed) > (per_timeout + retry_backoff) @staticmethod def _error_body_reason(http_status: object, body: Dict) -> Dict: - """4xx/5xx JSON 错误包的脱敏原因摘要:只记状态码与上游 message 片段。""" - reason: Dict = {"reason": "error_body", "http_status": http_status} + """4xx/5xx JSON 错误包的脱敏原因摘要:只记 int 状态码与 str message 片段。""" + reason: Dict = {"reason": "error_body"} + if isinstance(http_status, int): + reason["http_status"] = http_status message = body.get("message") if isinstance(body, dict) else None if isinstance(message, str) and message: reason["upstream_message"] = message[:200] diff --git a/docs/sessions/sph-retry-reason/progress/sph-retry-reason-260923-progress.md b/docs/sessions/sph-retry-reason/progress/sph-retry-reason-260923-progress.md index 6ee59bc..63d7651 100644 --- a/docs/sessions/sph-retry-reason/progress/sph-retry-reason-260923-progress.md +++ b/docs/sessions/sph-retry-reason/progress/sph-retry-reason-260923-progress.md @@ -9,3 +9,12 @@ - 决策:`WECHAT_CHANNELS_TOTAL_BUDGET` 保持 30s 不动;重试前检查剩余预算,不足则停试并走 `VideoNotFoundError`,不换成 timed out。 - 已否决(卡锁定):只在下载路径加第二次循环;在 API 层加重试;按响应文本正则猜原因。 - 下一步唯一动作:跑全量回归 + 红验(上限 3→1 注入),然后按三段提交。 + +## 段落 2:R1 + R2 收口(review-fixes) + +- 当前阶段:review-fixes,gate primary 两 major 已收口,待全量回归确认。 +- 本段结论:R1 起 `EndpointHttpError` 统一到所有 verb(POST-only 分叉被否决,抖音/小红书两条契约用例已改写);R2 补两处——`_should_retry` 要求剩余预算 > `per_timeout + retry_backoff`,且总预算超时在「重试链 + 有 retryable + 未达上限」时改抛 `VideoNotFoundError`;`describe_failure` 的 `object_type` 仅 int 标量进 attempts,非 int 只记类型名。最终口径:耗尽恒为 `VideoNotFoundError`,attempts 原因字段只含标量。 +- 关键决策与已否决方案: + - 超时兜底三条件严格限定,首跳挂起与单端点链行为不变(两处既有 timeout 用例原样绿)。 + - 已否决:为让旧测试绿而保留 verb 分叉(R1 已否决);把整包 `object_type` 打码后保留(容器一律不进日志)。 +- 下一步唯一动作:全量回归 + F1/F2 最小注入红验后提交。 diff --git a/tests/test_tikhub_provider_wechat_channels.py b/tests/test_tikhub_provider_wechat_channels.py index 94181bc..0d83cdc 100644 --- a/tests/test_tikhub_provider_wechat_channels.py +++ b/tests/test_tikhub_provider_wechat_channels.py @@ -375,6 +375,10 @@ async def test_chain_desensitization_error_string_and_api(authed_client, monkeyp payload = copy.deepcopy(load("empty")) payload["message"] = "upstream says invalid object_id marker-MSG42" payload["data"] = {"full_url": secret_full, "decode_key": secret_key, "url_token": secret_token} + _d2, r2 = TikHubProvider._classify_wechat_channels( + {"code": 200, "data": {"id": "1", "object_type": payload["data"]}}) + assert r2.get("reason") == "object_type_mismatch" + assert "object_type" not in r2 and r2.get("object_type_type") == "dict" async def error_body(self, name, path, params, per_timeout): raise EndpointHttpError(name, 400, payload) monkeypatch.setattr(TikHubProvider, "_call_endpoint", error_body) @@ -389,6 +393,15 @@ async def error_body(self, name, path, params, per_timeout): blob = resp.text assert secret_full not in blob and secret_key not in blob and secret_token not in blob assert "marker-MSG42" in blob + # F2 链路:200 + dict object_type 走三跳耗尽,异常串同样零凭据 + payload2 = {"code": 200, "data": {"id": "1", "object_type": payload["data"]}} + async def mismatch2(self, name, path, params, per_timeout): + return payload2 + monkeypatch.setattr(TikHubProvider, "_call_endpoint", mismatch2) + await _nosleep(monkeypatch) + with pytest.raises(VideoNotFoundError) as ei2: + await TikHubProvider().fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) + assert secret_key not in str(ei2.value) and secret_full not in str(ei2.value) async def test_chain_provider_error_not_retried(monkeypatch): @@ -402,6 +415,22 @@ async def boom(self, name, path, params, per_timeout): assert len(calls) == 1 and "http_error" in str(ei.value) +async def test_chain_budget_truncation_still_not_found(monkeypatch): + """F1:尝试耗时超预算被截断 → 仍是 VideoNotFoundError(含 retryable),非 timed out。""" + import asyncio as _aio + async def slow_empty(self, name, path, params, per_timeout): + await _aio.sleep(0.4) + return {"code": 200, "data": None} + monkeypatch.setattr(TikHubProvider, "_call_endpoint", slow_empty) + monkeypatch.setattr(TikHubProvider, "WECHAT_CHANNELS_RETRY_BACKOFF", 0) + monkeypatch.setattr(TikHubProvider, "WECHAT_CHANNELS_TOTAL_BUDGET", 0.7) + monkeypatch.setattr(TikHubProvider, "WECHAT_CHANNELS_PER_ENDPOINT_TIMEOUT", 0.05) + with pytest.raises(VideoNotFoundError) as ei: + await TikHubProvider().fetch_video_info("wechat_channels", OBJECT_ID, SHARE_URL) + assert type(ei.value) is VideoNotFoundError + assert "timed out" not in str(ei.value) and "retryable" in str(ei.value) + + async def test_adapter_and_platforms_endpoint(authed_client, monkeypatch): """GET /api/platforms 含 wechat_channels: [tikhub];adapter 能产出 VideoInfo。""" from app.services.adapters.tikhub_adapter import TikHubAdapter