Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,10 @@
- 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`。
- README 平台表与 `data.platform` 枚举补充 Facebook,与运行时 `/api/platforms` 返回保持一致。

## [1.0.0]
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`)。

---
Expand Down
23 changes: 23 additions & 0 deletions app/services/platforms/wechat_channels.py
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,29 @@ 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:
# 只带 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]:
"""
解析视频号 TikHub 响应(raw=false 的扁平结构)。
Expand Down
193 changes: 131 additions & 62 deletions app/services/providers/tikhub.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 兜底,因此分类器不出终态。
Expand Down Expand Up @@ -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,命中即返回;终态立即短路;全失败抛错。
Expand All @@ -424,39 +438,78 @@ 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, per_timeout, retry_backoff,
):
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:
# 重试链有 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,
Expand All @@ -470,6 +523,36 @@ 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,
per_timeout: float, retry_backoff: float,
) -> bool:
"""重试前检查:须为下一次尝试的完整单端超时 + 退避留出余量。"""
if attempt >= max_attempts:
return False
elapsed = asyncio.get_event_loop().time() - start
return (total_budget - elapsed) > (per_timeout + retry_backoff)

@staticmethod
def _error_body_reason(http_status: object, body: Dict) -> Dict:
"""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]
return reason

async def _call_endpoint(
self, name: str, path: str, params: Dict, per_timeout: float
) -> Dict:
Expand Down Expand Up @@ -504,7 +587,9 @@ async def _call_endpoint(
except Exception:
body = None
if isinstance(body, dict):
return body # 交给分类器判定 terminal/retryable
# 状态码只在此处可得:带 body 抛给 _run_chain 跑 classify +
# has_playable,并记 http_status 证据供归因(与 verb 无关)。
raise EndpointHttpError(name, status, body)
raise ProviderError(f"{name} HTTP {status}")
except ProviderError:
raise
Expand Down Expand Up @@ -716,16 +801,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:
Expand Down Expand Up @@ -762,6 +852,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:
Expand Down Expand Up @@ -801,36 +893,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):
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
# 视频号瞬态重试 + 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 注入),然后按三段提交。

## 段落 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 最小注入红验后提交。
Loading
Loading