douyin.py 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571
  1. from __future__ import annotations
  2. from pathlib import Path
  3. from typing import Any
  4. import httpx
  5. # 共享 crawapi 基座(V3-M1A):HTTP/限流/限流错误识别/env helper 集中于 crawapi_http,
  6. # 下方 re-export 保持既有外部 import(测试、smoke 脚本)零改。
  7. from content_agent.integrations.crawapi_http import (
  8. CrawapiBusinessError,
  9. CrawapiTransientError,
  10. RATE_LIMIT_MESSAGE_TOKENS,
  11. RateLimiter,
  12. _env,
  13. _load_env_file,
  14. _optional_positive_int,
  15. content_format as _content_format,
  16. is_rate_limit_business_error,
  17. post_crawapi_json,
  18. score_from_statistics as _score_from_statistics,
  19. search_previous_discovery_step as _search_previous_discovery_step,
  20. )
  21. from content_agent.errors import ContentAgentError
  22. from content_agent.integrations import platform_video_url
  23. RAW_CONTENT_ID_KEY = "_".join(["aweme", "id"])
  24. RAW_AUTHOR_ID_KEY = "_".join(["sec", "uid"])
  25. RAW_AUTHOR_ACCOUNT_KEY = "_".join(["account", "id"])
  26. # 已证实的限流 business code 白名单。当前没有任何已证实的限流 code,
  27. # 识别先依靠 HTTP 429 与 message token;live smoke / 真实运行发现新 code 后补入并加用例。
  28. RATE_LIMIT_BUSINESS_CODES: set[str] = set()
  29. SEARCH_RATE_LIMIT_BUCKET = "douyin_search"
  30. BLOGGER_RATE_LIMIT_BUCKET = "douyin_blogger"
  31. DOUYIN_SEARCH_MIN_INTERVAL_SECONDS = 10.0
  32. DOUYIN_SEARCH_MAX_INTERVAL_SECONDS = 12.0
  33. NON_FALLBACK_BUSINESS_CODES = {"10001"}
  34. NON_FALLBACK_BUSINESS_MESSAGE_TOKENS = ("参数异常", "参数校验", "invalid parameter")
  35. class CrawapiDouyinClient:
  36. requires_progressive_search_rate_limit = True
  37. def __init__(
  38. self,
  39. base_url: str,
  40. keyword_path: str,
  41. fallback_base_url: str = "",
  42. blogger_path: str = "",
  43. detail_path: str = "",
  44. portrait_path: str = "",
  45. timeout_seconds: float = 60.0,
  46. default_crawapi_account_ref: str = "",
  47. default_content_type: str = "视频",
  48. default_sort_type: str = "综合排序",
  49. default_publish_time: str = "不限",
  50. default_cursor: str = "0",
  51. piaoquantv_account_id: str = "",
  52. piaoquantv_duration: str = "不限",
  53. piaoquantv_cookie_batch: str = "default",
  54. crawler_sort_type: str = "最多点赞",
  55. crawler_cursor_default: str = "",
  56. default_account_works_sort_type: str = "最新",
  57. max_results_per_query: int | None = 5,
  58. http_client: Any | None = None,
  59. rate_limiter: RateLimiter | None = None,
  60. video_url_probe_fn: platform_video_url.ProbeFn | None = None,
  61. ) -> None:
  62. self.base_url = base_url.rstrip("/") + "/"
  63. self.fallback_base_url = fallback_base_url.rstrip("/") + "/" if fallback_base_url else ""
  64. self._crawapi_base_urls = _dedupe_crawapi_base_urls(
  65. [("primary", self.base_url), ("candidate", self.fallback_base_url)]
  66. )
  67. self._legacy_crawapi_base_urls = _legacy_crawapi_base_urls(self._crawapi_base_urls)
  68. self.keyword_path = keyword_path.lstrip("/")
  69. self.blogger_path = blogger_path.lstrip("/")
  70. self.detail_path = detail_path.lstrip("/")
  71. self.portrait_path = portrait_path.lstrip("/")
  72. self.timeout_seconds = timeout_seconds
  73. self.default_crawapi_account_ref = default_crawapi_account_ref
  74. self.default_content_type = default_content_type
  75. self.default_sort_type = default_sort_type
  76. self.default_publish_time = default_publish_time
  77. self.default_cursor = default_cursor
  78. self.piaoquantv_account_id = piaoquantv_account_id or default_crawapi_account_ref
  79. self.piaoquantv_duration = piaoquantv_duration
  80. self.piaoquantv_cookie_batch = piaoquantv_cookie_batch
  81. self.crawler_sort_type = crawler_sort_type
  82. self.crawler_cursor_default = crawler_cursor_default
  83. self.default_account_works_sort_type = default_account_works_sort_type
  84. self.max_results_per_query = max_results_per_query
  85. self.http_client = http_client or httpx.Client(timeout=timeout_seconds)
  86. self.rate_limiter = rate_limiter
  87. self.video_url_probe_fn = video_url_probe_fn
  88. @classmethod
  89. def from_env(cls, env_path: str | Path = ".env") -> "CrawapiDouyinClient":
  90. env = _load_env_file(env_path)
  91. return cls(
  92. base_url=_env("CONTENTFIND_API_CRAWAPI_BASE_URL", env, required=True),
  93. fallback_base_url=_env("CONTENTFIND_API_CRAWAPI_FALLBACK_BASE_URL", env, default=""),
  94. keyword_path=_env("CONTENTFIND_DOUYIN_KEYWORD_PATH", env, required=True),
  95. blogger_path=_env("CONTENTFIND_DOUYIN_BLOGGER_PATH", env, required=True),
  96. detail_path=_env(
  97. "CONTENTFIND_DOUYIN_DETAIL_PATH", env, default="/crawler/dou_yin/detail"
  98. ),
  99. portrait_path=_env(
  100. "CONTENTFIND_DOUYIN_PORTRAIT_PATH",
  101. env,
  102. default="/crawler/dou_yin/re_dian_bao/account_fans_portrait",
  103. ),
  104. timeout_seconds=float(
  105. _env("CONTENTFIND_API_CRAWAPI_TIMEOUT_SECONDS", env, default="180")
  106. ),
  107. default_crawapi_account_ref=_env("CONTENTFIND_DOUYIN_DEFAULT_ACCOUNT_ID", env, default=""),
  108. default_content_type=_env("CONTENTFIND_DOUYIN_DEFAULT_CONTENT_TYPE", env, default="视频"),
  109. default_sort_type=_env("CONTENTFIND_DOUYIN_DEFAULT_SORT_TYPE", env, default="综合排序"),
  110. default_publish_time=_env("CONTENTFIND_DOUYIN_DEFAULT_PUBLISH_TIME", env, default="不限"),
  111. default_cursor=_env("CONTENTFIND_DOUYIN_DEFAULT_CURSOR", env, default="0"),
  112. piaoquantv_account_id=_env(
  113. "CONTENTFIND_DOUYIN_PIAOQUANTV_ACCOUNT_ID",
  114. env,
  115. default=_env("CONTENTFIND_DOUYIN_DEFAULT_ACCOUNT_ID", env, default=""),
  116. ),
  117. piaoquantv_duration=_env("CONTENTFIND_DOUYIN_PIAOQUANTV_DURATION", env, default="不限"),
  118. piaoquantv_cookie_batch=_env(
  119. "CONTENTFIND_DOUYIN_PIAOQUANTV_COOKIE_BATCH", env, default="default"
  120. ),
  121. crawler_sort_type=_env("CONTENTFIND_DOUYIN_CRAWLER_SORT_TYPE", env, default="最多点赞"),
  122. crawler_cursor_default=_env("CONTENTFIND_DOUYIN_CRAWLER_CURSOR_DEFAULT", env, default=""),
  123. default_account_works_sort_type=_env(
  124. "CONTENTFIND_DOUYIN_ACCOUNT_WORKS_DEFAULT_SORT_TYPE", env, default="最新"
  125. ),
  126. max_results_per_query=_optional_positive_int(
  127. _env("CONTENTFIND_DOUYIN_MAX_RESULTS_PER_QUERY", env, default="5")
  128. ),
  129. rate_limiter=RateLimiter(
  130. min_interval_seconds=DOUYIN_SEARCH_MIN_INTERVAL_SECONDS,
  131. max_interval_seconds=DOUYIN_SEARCH_MAX_INTERVAL_SECONDS,
  132. ),
  133. )
  134. def search(self, query: dict[str, Any]) -> list[dict[str, Any]]:
  135. return self._search(query, self.max_results_per_query)
  136. def search_full_page(self, query: dict[str, Any]) -> list[dict[str, Any]]:
  137. return self._search(query, None)
  138. def _search(
  139. self,
  140. query: dict[str, Any],
  141. max_results_per_query: int | None,
  142. ) -> list[dict[str, Any]]:
  143. data = self._post_keyword_json(query)
  144. data_block = data.get("data", {}) if isinstance(data.get("data"), dict) else {}
  145. items = data_block.get("data", []) if isinstance(data_block.get("data"), list) else []
  146. has_more = bool(data_block.get("has_more", False))
  147. next_cursor = str(data_block.get("next_cursor") or "")
  148. results: list[dict[str, Any]] = []
  149. selected_items = items[:max_results_per_query] if max_results_per_query else items
  150. for index, item in enumerate(selected_items, start=1):
  151. selection = self._select_video_url_with_detail_fallback(item)
  152. results.append(
  153. self._normalize_content_item(query, item, index, has_more, next_cursor, selection)
  154. )
  155. return results
  156. def fetch_author_works(self, query: dict[str, Any]) -> list[dict[str, Any]]:
  157. payload = {
  158. RAW_AUTHOR_ACCOUNT_KEY: str(query.get("platform_author_id") or ""),
  159. "sort_type": self.default_account_works_sort_type,
  160. "cursor": str(query.get("page_cursor") or ""),
  161. }
  162. data = self._post_json(
  163. self.blogger_path, payload, operation="author_works",
  164. rate_limit_bucket=BLOGGER_RATE_LIMIT_BUCKET,
  165. )
  166. data_block = data.get("data", {}) if isinstance(data.get("data"), dict) else {}
  167. items = data_block.get("data", []) if isinstance(data_block.get("data"), list) else []
  168. has_more = bool(data_block.get("has_more", False))
  169. next_cursor = str(data_block.get("next_cursor") or "")
  170. selected_items = items[: self.max_results_per_query] if self.max_results_per_query else items
  171. results: list[dict[str, Any]] = []
  172. for index, item in enumerate(selected_items, start=1):
  173. selection = self._select_video_url_with_detail_fallback(item)
  174. normalized = self._normalize_content_item(query, item, index, has_more, next_cursor, selection)
  175. normalized["previous_discovery_step"] = "author_works"
  176. normalized["content_metadata_source"] = "douyin_blogger"
  177. results.append(normalized)
  178. return results
  179. def fetch_account_fans_portrait(self, account_id: str) -> dict[str, Any]:
  180. """M9A:热点宝作者粉丝画像(account_id=sec_uid)。返回 data.data(含 account/fans/posts)。
  181. 错误自然上抛(post_crawapi_json 已分类 rate_limit/transient/business),不在此吞。
  182. """
  183. payload = {
  184. RAW_AUTHOR_ACCOUNT_KEY: str(account_id or ""),
  185. "need_age": True,
  186. "need_gender": True,
  187. "need_province": True,
  188. }
  189. data = self._post_json(
  190. self.portrait_path, payload, operation="account_fans_portrait",
  191. )
  192. outer = data.get("data", {})
  193. inner = outer.get("data", {}) if isinstance(outer, dict) else {}
  194. return inner if isinstance(inner, dict) else {}
  195. def _normalize_content_item(
  196. self,
  197. query: dict[str, Any],
  198. item: dict[str, Any],
  199. index: int,
  200. has_more: bool,
  201. next_cursor: str,
  202. video_url_selection: dict[str, Any] | None = None,
  203. ) -> dict[str, Any]:
  204. video_url_selection = video_url_selection or self._select_video_url([("search", item)])
  205. author = item.get("author", {}) if isinstance(item.get("author"), dict) else {}
  206. statistics = item.get("statistics", {}) if isinstance(item.get("statistics"), dict) else {}
  207. platform_content_id = str(item.get(RAW_CONTENT_ID_KEY) or "")
  208. platform_author_id = str(author.get(RAW_AUTHOR_ID_KEY) or "")
  209. result = {
  210. "content_discovery_id": f"{query['search_query_id']}_content_{index:03d}",
  211. "search_query_id": query["search_query_id"],
  212. "platform": "douyin",
  213. "platform_content_id": platform_content_id,
  214. "platform_content_format": _content_format(self.default_content_type),
  215. "play_url": video_url_selection.get("play_url"),
  216. "description": item.get("desc") or item.get("item_title") or "",
  217. "platform_author_id": platform_author_id,
  218. "author_display_name": author.get("nickname") or "",
  219. "statistics": {
  220. "digg_count": int(statistics.get("digg_count") or 0),
  221. "comment_count": int(statistics.get("comment_count") or 0),
  222. "share_count": int(statistics.get("share_count") or 0),
  223. "collect_count": int(statistics.get("collect_count") or 0),
  224. "play_count": int(statistics.get("play_count") or 0),
  225. },
  226. "tags": _extract_tags(item),
  227. "text_extra": item.get("text_extra") or [],
  228. "create_time": item.get("create_time"),
  229. "has_more": has_more,
  230. "next_cursor": next_cursor,
  231. "score": _score_from_statistics(statistics),
  232. "risk_level": "unknown",
  233. "discovery_relation": "derived_from_pattern_demand",
  234. "discovery_start_source": query["discovery_start_source"],
  235. "previous_discovery_step": _search_previous_discovery_step(query),
  236. "content_metadata_source": "douyin_keyword_search",
  237. "platform_auth_mode": "no_bearer",
  238. "platform_raw_payload": {
  239. RAW_CONTENT_ID_KEY: platform_content_id,
  240. "author": {RAW_AUTHOR_ID_KEY: platform_author_id},
  241. **dict(video_url_selection.get("platform_raw_payload") or {}),
  242. },
  243. }
  244. if video_url_selection.get("media_failure_reason"):
  245. result["media_failure_reason"] = video_url_selection["media_failure_reason"]
  246. if video_url_selection.get("video_url_candidates"):
  247. result["video_url_candidates"] = video_url_selection["video_url_candidates"]
  248. return result
  249. def fetch_detail(self, content_id: str) -> dict[str, Any]:
  250. detail = self._fetch_detail_item(content_id)
  251. statistics = {
  252. "digg_count": int(detail.get("like_count") or 0),
  253. "comment_count": int(detail.get("comment_count") or 0),
  254. "share_count": int(detail.get("share_count") or 0),
  255. "collect_count": int(detail.get("collect_count") or 0),
  256. "play_count": int(detail.get("play_count") or 0),
  257. }
  258. topic_list = detail.get("topic_list") or []
  259. tags = [t if str(t).startswith("#") else f"#{t}" for t in topic_list if t]
  260. selection = self._select_video_url([("detail", detail)])
  261. publish_ms = detail.get("publish_timestamp")
  262. result = {
  263. "platform": "douyin",
  264. "platform_content_id": str(detail.get("channel_content_id") or content_id),
  265. "platform_content_url": detail.get("content_link"),
  266. "description": detail.get("body_text") or detail.get("title") or "",
  267. "platform_author_id": str(detail.get("channel_account_id") or ""),
  268. "author_display_name": detail.get("channel_account_name") or "",
  269. "statistics": statistics,
  270. "tags": tags,
  271. "play_url": selection.get("play_url"),
  272. "create_time": int(publish_ms) // 1000 if publish_ms else None,
  273. "content_metadata_source": "douyin_detail",
  274. "platform_raw_payload": dict(selection.get("platform_raw_payload") or {}),
  275. }
  276. if selection.get("media_failure_reason"):
  277. result["media_failure_reason"] = selection["media_failure_reason"]
  278. if selection.get("video_url_candidates"):
  279. result["video_url_candidates"] = selection["video_url_candidates"]
  280. return result
  281. def _select_video_url(self, sources: list[tuple[str, dict[str, Any]]]) -> dict[str, Any]:
  282. return platform_video_url.select_video_url(
  283. "douyin",
  284. sources,
  285. probe_fn=self.video_url_probe_fn or self._probe_video_url,
  286. )
  287. def _select_video_url_with_detail_fallback(self, item: dict[str, Any]) -> dict[str, Any]:
  288. search_selection = self._select_video_url([("search", item)])
  289. content_id = str(item.get(RAW_CONTENT_ID_KEY) or "")
  290. if not self._should_fetch_detail_for_video_url(search_selection) or not self.detail_path or not content_id:
  291. return search_selection
  292. try:
  293. detail = self._fetch_detail_item(content_id)
  294. except Exception as exc: # noqa: BLE001 - keep search diagnostics if detail is unavailable.
  295. payload = dict(search_selection.get("platform_raw_payload") or {})
  296. payload.update(
  297. {
  298. "douyin_detail_fallback_attempted": True,
  299. "douyin_detail_fallback_status": "failed",
  300. "douyin_detail_fallback_exception_type": type(exc).__name__,
  301. "douyin_detail_fallback_error": str(exc)[:300],
  302. }
  303. )
  304. return {**search_selection, "platform_raw_payload": payload}
  305. detail_selection = self._select_video_url([("search", item), ("detail", detail)])
  306. payload = dict(detail_selection.get("platform_raw_payload") or {})
  307. payload.update(
  308. {
  309. "douyin_detail_fallback_attempted": True,
  310. "douyin_detail_fallback_status": (
  311. "used" if detail_selection.get("play_url") else "no_valid_play_url"
  312. ),
  313. }
  314. )
  315. return {**detail_selection, "platform_raw_payload": payload}
  316. def _should_fetch_detail_for_video_url(self, selection: dict[str, Any]) -> bool:
  317. payload = selection.get("platform_raw_payload") if isinstance(selection.get("platform_raw_payload"), dict) else {}
  318. if not selection.get("play_url"):
  319. return True
  320. host = str(payload.get("selected_video_url_host") or "").lower()
  321. if host in {"v11-weba.douyinvod.com", "v96-hcc.douyinvod.com"}:
  322. return True
  323. return str(payload.get("selected_video_url_probe_status") or "") in {"failed", "failed_fallback"}
  324. def _fetch_detail_item(self, content_id: str) -> dict[str, Any]:
  325. data = self._post_json(
  326. self.detail_path,
  327. {"content_id": str(content_id)},
  328. operation="detail",
  329. rate_limit_bucket=SEARCH_RATE_LIMIT_BUCKET,
  330. )
  331. block = data.get("data", {}) if isinstance(data.get("data"), dict) else {}
  332. return block.get("data", {}) if isinstance(block.get("data"), dict) else {}
  333. def _probe_video_url(self, url: str, platform: str) -> dict[str, Any]:
  334. return platform_video_url.probe_url_with_httpx(
  335. url,
  336. platform,
  337. http_client=self.http_client,
  338. )
  339. def _post_keyword_json(self, query: dict[str, Any]) -> dict[str, Any]:
  340. attempts: list[dict[str, Any]] = []
  341. for index, (role, base_url) in enumerate(self._crawapi_base_urls):
  342. provider = _keyword_provider_for_base_url(role, base_url)
  343. payload = self._keyword_payload_for_provider(provider, query)
  344. try:
  345. return post_crawapi_json(
  346. http_client=self.http_client,
  347. base_url=base_url,
  348. path=self.keyword_path,
  349. payload=payload,
  350. operation="keyword_search",
  351. timeout_seconds=self.timeout_seconds,
  352. rate_limiter=self.rate_limiter,
  353. rate_limit_bucket=SEARCH_RATE_LIMIT_BUCKET,
  354. business_codes=RATE_LIMIT_BUSINESS_CODES,
  355. )
  356. except ContentAgentError as exc:
  357. attempt = _crawapi_host_attempt_summary(role, base_url, exc, provider=provider)
  358. self._raise_crawapi_error(exc, attempts + [attempt])
  359. except (CrawapiBusinessError, CrawapiTransientError, RuntimeError) as exc:
  360. attempt = _crawapi_host_attempt_summary(role, base_url, exc, provider=provider)
  361. if index == len(self._crawapi_base_urls) - 1 or not _should_try_crawapi_candidate(exc):
  362. self._raise_crawapi_error(exc, attempts + [attempt])
  363. attempts.append(attempt)
  364. raise RuntimeError("crawapi keyword_search failed: no base_url configured")
  365. def _keyword_payload_for_provider(self, provider: str, query: dict[str, Any]) -> dict[str, Any]:
  366. if provider == "crawler":
  367. return {
  368. "content_type": self.default_content_type,
  369. "keyword": query["search_query"],
  370. "cursor": _query_cursor(query, self.crawler_cursor_default),
  371. "sort_type": self.crawler_sort_type,
  372. }
  373. return {
  374. RAW_AUTHOR_ACCOUNT_KEY: self.piaoquantv_account_id,
  375. "keyword": query["search_query"],
  376. "content_type": self.default_content_type,
  377. "sort_type": self.default_sort_type,
  378. "publish_time": self.default_publish_time,
  379. "duration": self.piaoquantv_duration,
  380. "cursor": _query_cursor(query, self.default_cursor),
  381. "cookie_batch": self.piaoquantv_cookie_batch,
  382. }
  383. def _post_json(
  384. self,
  385. path: str,
  386. payload: dict[str, Any],
  387. operation: str,
  388. rate_limit_bucket: str | None = None,
  389. ) -> dict[str, Any]:
  390. attempts: list[dict[str, Any]] = []
  391. base_urls = self._legacy_crawapi_base_urls
  392. for index, (role, base_url) in enumerate(base_urls):
  393. try:
  394. return post_crawapi_json(
  395. http_client=self.http_client,
  396. base_url=base_url,
  397. path=path,
  398. payload=payload,
  399. operation=operation,
  400. timeout_seconds=self.timeout_seconds,
  401. rate_limiter=self.rate_limiter,
  402. rate_limit_bucket=rate_limit_bucket,
  403. business_codes=RATE_LIMIT_BUSINESS_CODES,
  404. )
  405. except ContentAgentError as exc:
  406. attempt = _crawapi_host_attempt_summary(role, base_url, exc)
  407. self._raise_crawapi_error(exc, attempts + [attempt])
  408. except (CrawapiBusinessError, CrawapiTransientError, RuntimeError) as exc:
  409. attempt = _crawapi_host_attempt_summary(role, base_url, exc)
  410. if index == len(base_urls) - 1 or not _should_try_crawapi_candidate(exc):
  411. self._raise_crawapi_error(exc, attempts + [attempt])
  412. attempts.append(attempt)
  413. raise RuntimeError(f"crawapi {operation} failed: no base_url configured")
  414. def _raise_crawapi_error(self, exc: Exception, attempts: list[dict[str, Any]]) -> None:
  415. if len(self._crawapi_base_urls) > 1:
  416. detail = getattr(exc, "detail", None)
  417. if not isinstance(detail, dict):
  418. detail = {}
  419. setattr(exc, "detail", detail)
  420. detail["crawapi_host_attempts"] = attempts
  421. raise exc
  422. def _dedupe_crawapi_base_urls(candidates: list[tuple[str, str]]) -> list[tuple[str, str]]:
  423. result: list[tuple[str, str]] = []
  424. seen: set[str] = set()
  425. for role, base_url in candidates:
  426. if not base_url:
  427. continue
  428. normalized = base_url.rstrip("/") + "/"
  429. if normalized in seen:
  430. continue
  431. seen.add(normalized)
  432. result.append((role, normalized))
  433. return result
  434. def _legacy_crawapi_base_urls(candidates: list[tuple[str, str]]) -> list[tuple[str, str]]:
  435. aiddit_candidates = [
  436. (role, base_url) for role, base_url in candidates if _is_aiddit_crawapi_base_url(base_url)
  437. ]
  438. return aiddit_candidates or candidates
  439. def _is_aiddit_crawapi_base_url(base_url: str) -> bool:
  440. return "crawler.aiddit.com" in base_url.lower()
  441. def _keyword_provider_for_base_url(role: str, base_url: str) -> str:
  442. if _is_aiddit_crawapi_base_url(base_url):
  443. return "crawler"
  444. normalized = base_url.lower()
  445. if "piaoquantv" in normalized:
  446. return "piaoquantv"
  447. return "crawler" if role == "candidate" else "piaoquantv"
  448. def _query_cursor(query: dict[str, Any], default: str) -> str:
  449. if "page_cursor" in query:
  450. return str(query.get("page_cursor") or "")
  451. return str(default)
  452. def _should_try_crawapi_candidate(exc: Exception) -> bool:
  453. if not isinstance(exc, CrawapiBusinessError):
  454. return True
  455. detail = getattr(exc, "detail", None)
  456. if not isinstance(detail, dict):
  457. return True
  458. business_code = str(detail.get("business_code") or "")
  459. business_message = str(detail.get("business_message") or "").lower()
  460. if business_code in NON_FALLBACK_BUSINESS_CODES:
  461. return False
  462. return not any(
  463. token.lower() in business_message
  464. for token in NON_FALLBACK_BUSINESS_MESSAGE_TOKENS
  465. )
  466. def _crawapi_host_attempt_summary(
  467. role: str,
  468. base_url: str,
  469. exc: Exception,
  470. provider: str | None = None,
  471. ) -> dict[str, Any]:
  472. summary: dict[str, Any] = {
  473. "role": role,
  474. "provider": provider or _keyword_provider_for_base_url(role, base_url),
  475. "base_url": base_url.rstrip("/"),
  476. "exception_type": type(exc).__name__,
  477. }
  478. message = str(exc)
  479. if message:
  480. summary["exception_message"] = message[:300]
  481. http_status = _http_status_from_exception_message(message)
  482. if http_status is not None:
  483. summary["status_code"] = http_status
  484. if isinstance(exc, ContentAgentError):
  485. summary["error_code"] = exc.error_code.value
  486. detail = getattr(exc, "detail", None)
  487. if isinstance(detail, dict):
  488. for key in ("operation", "business_code", "business_message", "status_code"):
  489. if key in detail:
  490. summary[key] = detail.get(key)
  491. response_summary = detail.get("response_summary")
  492. if isinstance(response_summary, dict):
  493. for key in ("status_code", "content_type"):
  494. if key in response_summary:
  495. summary[key] = response_summary.get(key)
  496. request_payload_summary = detail.get("request_payload_summary")
  497. if isinstance(request_payload_summary, dict):
  498. summary["request_payload_summary"] = request_payload_summary
  499. return summary
  500. def _http_status_from_exception_message(message: str) -> int | None:
  501. if "HTTP " not in message:
  502. return None
  503. status_text = message.split("HTTP ", 1)[1].split(";", 1)[0].strip()
  504. return int(status_text) if status_text.isdigit() else None
  505. def _extract_play_url(item: dict[str, Any]) -> str | None:
  506. video = item.get("video") if isinstance(item.get("video"), dict) else {}
  507. play_addr = video.get("play_addr") if isinstance(video.get("play_addr"), dict) else {}
  508. url_list = play_addr.get("url_list") or []
  509. return str(url_list[0]) if url_list else None
  510. def _extract_tags(item: dict[str, Any]) -> list[str]:
  511. tags: list[str] = []
  512. for tag in item.get("cha_list") or []:
  513. if isinstance(tag, str):
  514. tags.append(tag if tag.startswith("#") else f"#{tag}")
  515. elif isinstance(tag, dict):
  516. name = tag.get("cha_name") or tag.get("hashtag_name") or tag.get("name")
  517. if name:
  518. tags.append(str(name) if str(name).startswith("#") else f"#{name}")
  519. for text in item.get("text_extra") or []:
  520. if isinstance(text, dict) and text.get("hashtag_name"):
  521. tags.append(f"#{text['hashtag_name']}")
  522. return list(dict.fromkeys(tags))