runner.py 40 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034
  1. """Formal acquisition runner backed by the repository boundary."""
  2. from __future__ import annotations
  3. from dataclasses import dataclass
  4. from typing import Any, Callable, Iterable
  5. from uuid import UUID
  6. from acquisition.classification.coarse import ClassificationResult, coarse_classify_item
  7. from acquisition.content_mode import guard_content_for_processing, infer_content_mode
  8. from acquisition.crawler import RateLimiter
  9. from acquisition.domain import AcquisitionJob, Query
  10. from acquisition.media.service import StabilizedMedia, stabilize_media_urls
  11. from acquisition.platforms import PlatformAdapter, get_platform_adapter
  12. from acquisition.platforms.base import PlatformCandidate, PlatformSearchPage
  13. from acquisition.repositories.base import AcquisitionRepository
  14. from acquisition.unique_key import build_unique_key
  15. from core.config import Settings
  16. from core.text_limits import (
  17. BODY_TEXT_MAX_CHARS,
  18. ERROR_MESSAGE_MAX_CHARS,
  19. RAW_SUMMARY_MAX_CHARS,
  20. clip_text,
  21. )
  22. from pipeline.tracing import NoopTraceWriter, TraceContext, TraceWriter
  23. DEFAULT_PLATFORMS = ("xiaohongshu", "weixin", "douyin")
  24. PAGINATION_STRATEGY = "creation_ratio_two_page_v1"
  25. PAGINATION_CREATION_THRESHOLD = 0.4
  26. PAGINATION_MAX_PAGES = 2
  27. LOW_RESULT_RETRY_THRESHOLD = 5
  28. LOW_RESULT_MAX_RETRIES = 1
  29. AdapterFactory = Callable[[str], PlatformAdapter]
  30. Classifier = Callable[..., ClassificationResult]
  31. MediaStabilizer = Callable[..., list[StabilizedMedia]]
  32. RateLimiterFactory = Callable[[str], Any]
  33. @dataclass(frozen=True)
  34. class RunBatchResult:
  35. run_id: UUID
  36. total_jobs: int
  37. done: int
  38. partial: int
  39. failed: int
  40. skipped: int
  41. @dataclass(frozen=True)
  42. class RecordCandidateResult:
  43. displayed: bool
  44. skipped_processing: bool = False
  45. classified: bool = False
  46. is_creation_knowledge: bool | None = None
  47. class PlatformRateLimiter:
  48. """Share one rate bucket per platform across search and detail calls."""
  49. def __init__(self, platform: str, delegate: RateLimiter | None = None) -> None:
  50. self.platform = platform
  51. self.delegate = delegate or RateLimiter(
  52. min_interval_seconds=10.0,
  53. max_interval_seconds=12.0,
  54. )
  55. def wait(self, bucket: str) -> None:
  56. self.delegate.wait(self.platform)
  57. def _query_id(query: Query) -> UUID:
  58. if query.id is None:
  59. raise RuntimeError("formal acquisition queries must have id before running")
  60. return query.id
  61. def _job_id(job: AcquisitionJob) -> UUID:
  62. if job.id is None:
  63. raise RuntimeError("formal acquisition jobs must have id before running")
  64. return job.id
  65. def _best_video_url(media: list[StabilizedMedia]) -> str:
  66. for row in media:
  67. if row.media_type == "video" and row.status == "ready" and row.cdn_url:
  68. return row.cdn_url
  69. return ""
  70. def _image_urls(media: list[StabilizedMedia]) -> list[str]:
  71. return [row.cdn_url for row in media if row.media_type == "image" and row.cdn_url]
  72. def _source_payload(candidate: Any, detail: Any) -> dict[str, Any]:
  73. return {
  74. "candidate": candidate.model_dump() if hasattr(candidate, "model_dump") else {},
  75. "detail": detail.raw if isinstance(getattr(detail, "raw", None), dict) else {},
  76. }
  77. def _candidate_provider(candidate: Any) -> str:
  78. provider = getattr(candidate, "provider", "") or ""
  79. raw = getattr(candidate, "raw", None)
  80. if not provider and isinstance(raw, dict):
  81. provider = raw.get("search_provider") or ""
  82. return provider
  83. def _detail_provider(detail: Any, candidate: Any) -> str:
  84. provider = getattr(detail, "provider", "") or ""
  85. if not provider:
  86. raw = getattr(detail, "raw", None)
  87. if isinstance(raw, dict):
  88. provider = raw.get("detail_provider") or ""
  89. return provider or _candidate_provider(candidate)
  90. def _candidate_page_metadata(candidate: Any) -> dict[str, Any]:
  91. raw = getattr(candidate, "raw", None)
  92. if not isinstance(raw, dict):
  93. return {}
  94. out: dict[str, Any] = {}
  95. for key in ("page_index", "page_rank", "source_cursor"):
  96. value = raw.get(key)
  97. if value not in (None, ""):
  98. out[key] = value
  99. return out
  100. def _content_mode(platform: str, detail: Any) -> str:
  101. explicit = getattr(detail, "content_mode", "") or ""
  102. if explicit:
  103. return explicit
  104. return infer_content_mode(
  105. platform=platform,
  106. content_type=getattr(detail, "content_type", "") or "",
  107. body_text=getattr(detail, "body_text", "") or "",
  108. image_urls=getattr(detail, "image_urls", []) or [],
  109. video_urls=getattr(detail, "video_urls", []) or [],
  110. raw=getattr(detail, "raw", {}) if isinstance(getattr(detail, "raw", None), dict) else {},
  111. )
  112. def _skip_reason_message(reason: str) -> str:
  113. if reason == "video_url_missing":
  114. return "Video post has no usable video URL; skipped before classification."
  115. if reason == "video_oss_failed":
  116. return "Video OSS transfer failed after retries; skipped before classification."
  117. if reason == "unsupported_content_mode":
  118. return "Unsupported content mode; skipped before classification."
  119. return reason or "Skipped before classification."
  120. def _requires_strict_video_oss(platform: str, content_mode: str) -> bool:
  121. return content_mode == "video_post"
  122. def _video_media_errors(media_rows: list[StabilizedMedia]) -> list[dict[str, Any]]:
  123. errors: list[dict[str, Any]] = []
  124. for row in media_rows:
  125. if row.media_type != "video" or row.status == "ready":
  126. continue
  127. errors.append(
  128. {
  129. "source_url": row.source_url,
  130. "status": row.status,
  131. "error_message": row.error_message,
  132. }
  133. )
  134. return errors
  135. def _record_candidate(
  136. repo: AcquisitionRepository,
  137. *,
  138. job: AcquisitionJob,
  139. query: Query,
  140. platform: str,
  141. candidate: Any,
  142. detail: Any,
  143. settings: Settings,
  144. media_stabilizer: MediaStabilizer,
  145. classifier: Classifier,
  146. classify: bool,
  147. attempt_index: int = 1,
  148. trace_writer: TraceWriter | None = None,
  149. trace_context: TraceContext | None = None,
  150. ) -> RecordCandidateResult:
  151. source_payload = _source_payload(candidate, detail)
  152. platform_item_id = detail.source_id or candidate.source_id or None
  153. canonical_url = detail.url or candidate.url or None
  154. content_mode = _content_mode(platform, detail)
  155. guard = guard_content_for_processing(
  156. content_mode=content_mode,
  157. video_urls=detail.video_urls,
  158. )
  159. metadata = {
  160. "candidate_rank": candidate.rank,
  161. "acquisition_match_status": "created",
  162. "content_mode": content_mode,
  163. "retry_attempt_index": attempt_index,
  164. **_candidate_page_metadata(candidate),
  165. }
  166. search_provider = _candidate_provider(candidate)
  167. detail_provider = _detail_provider(detail, candidate)
  168. if search_provider:
  169. metadata["search_provider"] = search_provider
  170. if detail_provider:
  171. metadata["detail_provider"] = detail_provider
  172. if guard.metadata:
  173. metadata.update(guard.metadata)
  174. body_text = clip_text(detail.body_text or "", BODY_TEXT_MAX_CHARS)
  175. item = repo.upsert_candidate_item(
  176. platform=platform,
  177. job_id=_job_id(job),
  178. query_id=_query_id(query),
  179. platform_item_id=platform_item_id,
  180. unique_key=build_unique_key(
  181. platform,
  182. platform_item_id=platform_item_id,
  183. url=canonical_url,
  184. raw_payload=source_payload,
  185. ),
  186. canonical_url=canonical_url,
  187. content_type=detail.content_type or None,
  188. content_mode=content_mode,
  189. title=detail.title or None,
  190. author_name=detail.author or None,
  191. body_text=body_text or None,
  192. raw_summary=clip_text(body_text, RAW_SUMMARY_MAX_CHARS) or None,
  193. status="candidate" if guard.can_process else "skipped",
  194. source_payload=source_payload,
  195. metadata=metadata,
  196. error_message=None if guard.can_process else guard.reason,
  197. )
  198. if item.id is None:
  199. raise RuntimeError("repository returned candidate without id")
  200. _record_candidate_hit(
  201. trace_writer,
  202. trace_context,
  203. item_id=item.id,
  204. platform=platform,
  205. unique_key=item.unique_key,
  206. platform_item_id=item.platform_item_id,
  207. search_provider=search_provider,
  208. detail_provider=detail_provider,
  209. candidate=candidate,
  210. attempt_index=attempt_index,
  211. is_duplicate_hit=False,
  212. metadata={"content_mode": content_mode, "status": item.status},
  213. )
  214. if not guard.can_process:
  215. _trace_event(
  216. trace_writer,
  217. trace_context.child(item_id=item.id, platform=platform) if trace_context else None,
  218. stage="classify",
  219. event_type="item_skipped_before_classification",
  220. status="skipped",
  221. target_table="candidate_items",
  222. target_id=item.id,
  223. payload={"reason": guard.reason, "content_mode": content_mode},
  224. )
  225. repo.add_item_classification(
  226. item_id=item.id,
  227. is_creation_knowledge=None,
  228. label=guard.label,
  229. confidence=None,
  230. reason=_skip_reason_message(guard.reason),
  231. model_name=None,
  232. prompt_version="content_mode_guard",
  233. result_payload={
  234. "content_mode": content_mode,
  235. "skip_reason": guard.reason,
  236. **(guard.metadata or {}),
  237. },
  238. status="skipped",
  239. error_message=guard.reason,
  240. )
  241. return RecordCandidateResult(displayed=True, skipped_processing=True)
  242. strict_video_oss = _requires_strict_video_oss(platform, content_mode)
  243. media_rows = media_stabilizer(
  244. image_urls=detail.image_urls,
  245. video_urls=detail.video_urls,
  246. settings=settings,
  247. video_fallback_to_source=not strict_video_oss,
  248. video_retry_delays_seconds=(
  249. settings.video_oss_retry_delays_seconds if strict_video_oss else ()
  250. ),
  251. )
  252. for row in media_rows:
  253. repo.add_media_asset(
  254. item_id=item.id,
  255. media_type=row.media_type,
  256. source_url=row.source_url,
  257. oss_url=row.cdn_url,
  258. cdn_url=row.cdn_url,
  259. position=row.position,
  260. status=row.status,
  261. source_payload={},
  262. metadata={"error_message": row.error_message} if row.error_message else {},
  263. )
  264. if strict_video_oss and not _best_video_url(media_rows):
  265. repo.add_item_classification(
  266. item_id=item.id,
  267. is_creation_knowledge=None,
  268. label="video_oss_failed",
  269. confidence=None,
  270. reason=_skip_reason_message("video_oss_failed"),
  271. model_name=None,
  272. prompt_version="media_stabilization_guard",
  273. result_payload={
  274. "content_mode": content_mode,
  275. "skip_reason": "video_oss_failed",
  276. "video_oss_retry_delays_seconds": list(
  277. settings.video_oss_retry_delays_seconds
  278. ),
  279. "media_errors": _video_media_errors(media_rows),
  280. },
  281. status="skipped",
  282. error_message="video_oss_failed",
  283. )
  284. _trace_event(
  285. trace_writer,
  286. trace_context.child(item_id=item.id, platform=platform) if trace_context else None,
  287. stage="classify",
  288. event_type="item_skipped_before_classification",
  289. status="skipped",
  290. target_table="candidate_items",
  291. target_id=item.id,
  292. payload={"reason": "video_oss_failed", "content_mode": content_mode},
  293. )
  294. return RecordCandidateResult(displayed=True, skipped_processing=True)
  295. if classify:
  296. classifier_kwargs = {
  297. "platform": platform,
  298. "content_mode": content_mode,
  299. "title": detail.title,
  300. "body_text": detail.body_text,
  301. "image_urls": _image_urls(media_rows),
  302. "video_url": _best_video_url(media_rows),
  303. "settings": settings,
  304. }
  305. if classifier is coarse_classify_item:
  306. classifier_kwargs.update(
  307. {
  308. "trace_writer": trace_writer,
  309. "trace_context": trace_context.child(item_id=item.id, platform=platform) if trace_context else None,
  310. }
  311. )
  312. result = classifier(**classifier_kwargs)
  313. repo.add_item_classification(
  314. item_id=item.id,
  315. is_creation_knowledge=result.is_creation_knowledge,
  316. label=result.label,
  317. confidence=result.confidence,
  318. reason=result.reason,
  319. model_name=settings.video_model,
  320. prompt_version=result.prompt_version,
  321. result_payload=result.result_payload,
  322. status=result.status,
  323. error_message=result.error_message,
  324. )
  325. _trace_event(
  326. trace_writer,
  327. trace_context.child(item_id=item.id, platform=platform) if trace_context else None,
  328. stage="classify",
  329. event_type="classified",
  330. status=result.status,
  331. target_table="candidate_items",
  332. target_id=item.id,
  333. payload={
  334. "is_creation_knowledge": result.is_creation_knowledge,
  335. "label": result.label,
  336. "confidence": result.confidence,
  337. },
  338. error_message=result.error_message,
  339. )
  340. return RecordCandidateResult(
  341. displayed=True,
  342. classified=result.status == "classified" and result.is_creation_knowledge is not None,
  343. is_creation_knowledge=result.is_creation_knowledge,
  344. )
  345. return RecordCandidateResult(displayed=True)
  346. def _candidate_unique_key(platform: str, candidate: PlatformCandidate) -> str | None:
  347. return build_unique_key(
  348. platform,
  349. platform_item_id=candidate.source_id or None,
  350. url=candidate.url or None,
  351. raw_payload={"candidate": candidate.model_dump()},
  352. )
  353. def _candidate_dedupe_key(platform: str, candidate: PlatformCandidate) -> str | None:
  354. unique_key = _candidate_unique_key(platform, candidate)
  355. if unique_key:
  356. return f"unique:{unique_key}"
  357. if candidate.source_id:
  358. return f"source:{platform}:{candidate.source_id}"
  359. if candidate.url:
  360. return f"url:{platform}:{candidate.url}"
  361. return None
  362. def _attach_existing_candidate(
  363. repo: AcquisitionRepository,
  364. *,
  365. job: AcquisitionJob,
  366. query: Query,
  367. candidate: PlatformCandidate,
  368. unique_key: str,
  369. attempt_index: int = 1,
  370. trace_writer: TraceWriter | None = None,
  371. trace_context: TraceContext | None = None,
  372. ) -> bool:
  373. get_by_key = getattr(repo, "get_candidate_item_by_unique_key", None)
  374. attach = getattr(repo, "attach_existing_candidate_item", None)
  375. if get_by_key is None or attach is None:
  376. return False
  377. existing = get_by_key(unique_key)
  378. if existing is None or existing.id is None:
  379. return False
  380. attach(
  381. existing.id,
  382. job_id=_job_id(job),
  383. query_id=_query_id(query),
  384. metadata={
  385. "acquisition_match_status": "existing",
  386. "matched_unique_key": unique_key,
  387. "matched_candidate": candidate.model_dump(),
  388. "candidate_rank": candidate.rank,
  389. "retry_attempt_index": attempt_index,
  390. **_candidate_page_metadata(candidate),
  391. **({"search_provider": _candidate_provider(candidate)} if _candidate_provider(candidate) else {}),
  392. },
  393. )
  394. _record_candidate_hit(
  395. trace_writer,
  396. trace_context,
  397. item_id=existing.id,
  398. platform=candidate.platform,
  399. unique_key=unique_key,
  400. platform_item_id=candidate.source_id or None,
  401. search_provider=_candidate_provider(candidate),
  402. detail_provider=None,
  403. candidate=candidate,
  404. attempt_index=attempt_index,
  405. is_duplicate_hit=True,
  406. metadata={"acquisition_match_status": "existing"},
  407. )
  408. _trace_event(
  409. trace_writer,
  410. trace_context.child(item_id=existing.id, platform=candidate.platform) if trace_context else None,
  411. stage="search",
  412. event_type="duplicate_candidate_reused",
  413. status="skipped",
  414. target_table="candidate_items",
  415. target_id=existing.id,
  416. payload={"unique_key": unique_key},
  417. )
  418. return True
  419. def _trace_event(
  420. trace_writer: TraceWriter | None,
  421. context: TraceContext | None,
  422. *,
  423. stage: str,
  424. event_type: str,
  425. **kwargs: Any,
  426. ) -> None:
  427. if trace_writer is None or context is None:
  428. return
  429. trace_writer.event(context=context, stage=stage, event_type=event_type, **kwargs)
  430. def _record_candidate_hit(
  431. trace_writer: TraceWriter | None,
  432. context: TraceContext | None,
  433. *,
  434. item_id: UUID | None,
  435. platform: str,
  436. unique_key: str | None,
  437. platform_item_id: str | None,
  438. search_provider: str | None,
  439. detail_provider: str | None,
  440. candidate: PlatformCandidate,
  441. attempt_index: int,
  442. is_duplicate_hit: bool,
  443. metadata: dict[str, Any] | None = None,
  444. ) -> None:
  445. if trace_writer is None or context is None:
  446. return
  447. raw = candidate.model_dump() if hasattr(candidate, "model_dump") else {}
  448. trace_writer.candidate_hit(
  449. context=context.child(item_id=item_id, platform=platform),
  450. item_id=item_id,
  451. platform=platform,
  452. unique_key=unique_key,
  453. platform_item_id=platform_item_id,
  454. search_provider=search_provider,
  455. detail_provider=detail_provider,
  456. attempt_index=attempt_index,
  457. page_index=raw.get("raw", {}).get("page_index") if isinstance(raw.get("raw"), dict) else None,
  458. page_rank=raw.get("raw", {}).get("page_rank") if isinstance(raw.get("raw"), dict) else None,
  459. candidate_rank=candidate.rank,
  460. source_cursor=raw.get("raw", {}).get("source_cursor") if isinstance(raw.get("raw"), dict) else None,
  461. is_duplicate_hit=is_duplicate_hit,
  462. raw_candidate=raw,
  463. metadata=metadata or {},
  464. )
  465. def _pagination_config(threshold: float, max_pages: int) -> dict[str, Any]:
  466. return {
  467. "strategy": PAGINATION_STRATEGY,
  468. "creation_threshold": threshold,
  469. "max_pages": max_pages,
  470. }
  471. def _ratio(numerator: int, denominator: int) -> float | None:
  472. if denominator <= 0:
  473. return None
  474. return numerator / denominator
  475. def _search_pages(
  476. adapter: PlatformAdapter,
  477. query_text: str,
  478. *,
  479. settings: Settings,
  480. limit: int,
  481. rate_limiter: Any,
  482. max_pages: int,
  483. ) -> Iterable[PlatformSearchPage]:
  484. search_pages = getattr(adapter, "search_pages", None)
  485. if callable(search_pages):
  486. return search_pages(
  487. query_text,
  488. settings=settings,
  489. limit=limit,
  490. rate_limiter=rate_limiter,
  491. max_pages=max_pages,
  492. )
  493. candidates = adapter.search(
  494. query_text,
  495. settings=settings,
  496. limit=limit,
  497. rate_limiter=rate_limiter,
  498. )
  499. return [
  500. PlatformSearchPage(
  501. page_index=1,
  502. candidates=candidates,
  503. raw_count=len(candidates),
  504. has_more=False,
  505. )
  506. ]
  507. def run_batch(
  508. repo: AcquisitionRepository,
  509. *,
  510. batch_id: UUID,
  511. settings: Settings,
  512. platforms: tuple[str, ...] | list[str] = DEFAULT_PLATFORMS,
  513. search_limit: int = 10,
  514. display_limit: int = 5,
  515. classify: bool = True,
  516. resume: bool = True,
  517. skip_done: bool = True,
  518. run_key: str | None = None,
  519. adapter_factory: AdapterFactory = get_platform_adapter,
  520. media_stabilizer: MediaStabilizer = stabilize_media_urls,
  521. classifier: Classifier = coarse_classify_item,
  522. rate_limiter_factory: RateLimiterFactory | None = None,
  523. pagination_creation_threshold: float = PAGINATION_CREATION_THRESHOLD,
  524. pagination_max_pages: int = PAGINATION_MAX_PAGES,
  525. low_result_retry_threshold: int = LOW_RESULT_RETRY_THRESHOLD,
  526. low_result_max_retries: int = LOW_RESULT_MAX_RETRIES,
  527. trace_writer: TraceWriter | None = None,
  528. trace_context: TraceContext | None = None,
  529. ) -> RunBatchResult:
  530. """Run query x platform acquisition and write formal cloud-state rows."""
  531. queries = repo.list_queries_for_batch(batch_id, keep=True)
  532. run = repo.create_acquisition_run(
  533. batch_id=batch_id,
  534. run_key=run_key or f"acquisition:{batch_id}",
  535. status="running",
  536. metadata={
  537. "platforms": list(platforms),
  538. "search_limit": search_limit,
  539. "display_limit": display_limit,
  540. "classify": classify,
  541. "resume": resume,
  542. "pagination": _pagination_config(
  543. pagination_creation_threshold,
  544. pagination_max_pages,
  545. ),
  546. "low_result_retry": {
  547. "threshold": low_result_retry_threshold,
  548. "max_retries": low_result_max_retries,
  549. },
  550. },
  551. )
  552. if run.id is None:
  553. raise RuntimeError("repository returned acquisition run without id")
  554. trace_writer = trace_writer or NoopTraceWriter()
  555. run_trace_context = (trace_context or TraceContext()).child(
  556. acquisition_run_id=run.id,
  557. stage="search",
  558. )
  559. _trace_event(
  560. trace_writer,
  561. run_trace_context,
  562. stage="search",
  563. event_type="acquisition_started",
  564. status="running",
  565. target_table="acquisition_runs",
  566. target_id=run.id,
  567. payload={"batch_id": str(batch_id), "platforms": list(platforms)},
  568. )
  569. total_jobs = len(queries) * len(platforms)
  570. done = partial = failed = skipped = 0
  571. gates: dict[str, Any] = {}
  572. for query in queries:
  573. query_id = _query_id(query)
  574. for platform in platforms:
  575. job = repo.ensure_acquisition_job(
  576. run_id=run.id,
  577. query_id=query_id,
  578. platform=platform,
  579. search_limit=search_limit,
  580. display_limit=display_limit,
  581. status="pending",
  582. metadata={"query_text": query.query_text},
  583. )
  584. if skip_done and job.status == "done":
  585. _trace_event(
  586. trace_writer,
  587. run_trace_context.child(
  588. acquisition_job_id=_job_id(job),
  589. query_id=query_id,
  590. platform=platform,
  591. ),
  592. stage="search",
  593. event_type="job_skipped_done",
  594. status="skipped",
  595. target_table="acquisition_jobs",
  596. target_id=_job_id(job),
  597. )
  598. skipped += 1
  599. continue
  600. attempts = job.attempt_count + 1
  601. job = repo.update_acquisition_job(
  602. _job_id(job),
  603. status="running",
  604. attempt_count=attempts,
  605. error_message=None,
  606. )
  607. job_trace_context = run_trace_context.child(
  608. acquisition_job_id=_job_id(job),
  609. query_id=query_id,
  610. platform=platform,
  611. )
  612. _trace_event(
  613. trace_writer,
  614. job_trace_context,
  615. stage="search",
  616. event_type="job_started",
  617. status="running",
  618. target_table="acquisition_jobs",
  619. target_id=_job_id(job),
  620. payload={"query_text": query.query_text, "platform": platform},
  621. )
  622. errors: list[str] = []
  623. display_count = 0
  624. searched_count = 0
  625. raw_searched_count = 0
  626. skipped_item_count = 0
  627. search_page_count = 0
  628. page_summaries: list[dict[str, Any]] = []
  629. first_page_classified_count = 0
  630. first_page_creation_count = 0
  631. first_page_creation_ratio: float | None = None
  632. requested_second_page = False
  633. pagination_stop_reason = ""
  634. try:
  635. adapter = adapter_factory(platform)
  636. gate = gates.get(platform)
  637. if gate is None:
  638. gate = (
  639. rate_limiter_factory(platform)
  640. if rate_limiter_factory
  641. else PlatformRateLimiter(platform)
  642. )
  643. gates[platform] = gate
  644. retry_enabled = (
  645. low_result_max_retries > 0
  646. and low_result_retry_threshold >= 0
  647. and search_limit > low_result_retry_threshold
  648. )
  649. max_attempts = 1 + (low_result_max_retries if retry_enabled else 0)
  650. seen_candidate_keys: set[str] = set()
  651. retry_attempts: list[dict[str, Any]] = []
  652. retry_triggered = False
  653. for attempt_index in range(1, max_attempts + 1):
  654. attempt_display_start = display_count
  655. attempt_searched_start = searched_count
  656. attempt_raw_start = raw_searched_count
  657. attempt_pages_start = search_page_count
  658. attempt_page_summaries: list[dict[str, Any]] = []
  659. attempt_first_page_classified_count = 0
  660. attempt_first_page_creation_count = 0
  661. attempt_first_page_creation_ratio: float | None = None
  662. attempt_requested_second_page = False
  663. attempt_stop_reason = ""
  664. pages = _search_pages(
  665. adapter,
  666. query.query_text,
  667. settings=settings,
  668. limit=search_limit,
  669. rate_limiter=gate,
  670. max_pages=pagination_max_pages,
  671. )
  672. for page in pages:
  673. search_page_count += 1
  674. raw_searched_count += page.raw_count
  675. page_classified_count = 0
  676. page_creation_count = 0
  677. unique_candidate_count = 0
  678. for candidate in page.candidates[:search_limit]:
  679. dedupe_key = _candidate_dedupe_key(platform, candidate)
  680. if dedupe_key and dedupe_key in seen_candidate_keys:
  681. continue
  682. if dedupe_key:
  683. seen_candidate_keys.add(dedupe_key)
  684. unique_candidate_count += 1
  685. searched_count += 1
  686. try:
  687. unique_key = _candidate_unique_key(platform, candidate)
  688. if unique_key and _attach_existing_candidate(
  689. repo,
  690. job=job,
  691. query=query,
  692. candidate=candidate,
  693. unique_key=unique_key,
  694. attempt_index=attempt_index,
  695. trace_writer=trace_writer,
  696. trace_context=job_trace_context,
  697. ):
  698. display_count += 1
  699. continue
  700. detail = adapter.fetch_detail(
  701. candidate,
  702. settings=settings,
  703. rate_limiter=gate,
  704. )
  705. recorded = _record_candidate(
  706. repo,
  707. job=job,
  708. query=query,
  709. platform=platform,
  710. candidate=candidate,
  711. detail=detail,
  712. settings=settings,
  713. media_stabilizer=media_stabilizer,
  714. classifier=classifier,
  715. classify=classify,
  716. attempt_index=attempt_index,
  717. trace_writer=trace_writer,
  718. trace_context=job_trace_context,
  719. )
  720. if recorded.displayed:
  721. display_count += 1
  722. if recorded.skipped_processing:
  723. skipped_item_count += 1
  724. if recorded.classified:
  725. page_classified_count += 1
  726. if recorded.is_creation_knowledge is True:
  727. page_creation_count += 1
  728. except Exception as exc:
  729. message = clip_text(exc, ERROR_MESSAGE_MAX_CHARS)
  730. errors.append(message)
  731. _trace_event(
  732. trace_writer,
  733. job_trace_context,
  734. stage="search",
  735. event_type="candidate_failed",
  736. status="failed",
  737. severity="warning",
  738. payload={
  739. "candidate": candidate.model_dump() if hasattr(candidate, "model_dump") else {},
  740. "attempt_index": attempt_index,
  741. },
  742. error_message=message,
  743. )
  744. continue
  745. page_ratio = _ratio(page_creation_count, page_classified_count)
  746. page_summary = {
  747. "attempt_index": attempt_index,
  748. "page_index": page.page_index,
  749. "raw_count": page.raw_count,
  750. "candidate_count": len(page.candidates),
  751. "unique_candidate_count": unique_candidate_count,
  752. "classified_count": page_classified_count,
  753. "creation_count": page_creation_count,
  754. "creation_ratio": page_ratio,
  755. "has_more": page.has_more,
  756. "next_cursor": page.next_cursor,
  757. "provider": page.provider,
  758. }
  759. page_summaries.append(page_summary)
  760. attempt_page_summaries.append(page_summary)
  761. _trace_event(
  762. trace_writer,
  763. job_trace_context,
  764. stage="search",
  765. event_type="page_searched",
  766. status="done",
  767. payload=page_summary,
  768. attempt_index=attempt_index,
  769. )
  770. if page.page_index == 1:
  771. attempt_first_page_classified_count = page_classified_count
  772. attempt_first_page_creation_count = page_creation_count
  773. attempt_first_page_creation_ratio = page_ratio
  774. if pagination_max_pages <= 1:
  775. attempt_stop_reason = "max_pages_reached"
  776. break
  777. if not page.has_more or not page.next_cursor:
  778. attempt_stop_reason = "no_more_pages"
  779. break
  780. if page_classified_count <= 0:
  781. attempt_stop_reason = "no_classified_candidates"
  782. break
  783. if page_ratio is None or page_ratio < pagination_creation_threshold:
  784. attempt_stop_reason = "below_threshold"
  785. break
  786. attempt_requested_second_page = True
  787. continue
  788. attempt_stop_reason = "max_pages_reached"
  789. break
  790. if not attempt_stop_reason:
  791. attempt_stop_reason = "pages_exhausted"
  792. first_page_classified_count = attempt_first_page_classified_count
  793. first_page_creation_count = attempt_first_page_creation_count
  794. first_page_creation_ratio = attempt_first_page_creation_ratio
  795. requested_second_page = attempt_requested_second_page
  796. pagination_stop_reason = attempt_stop_reason
  797. attempt_display_count = display_count - attempt_display_start
  798. retry_attempts.append(
  799. {
  800. "attempt_index": attempt_index,
  801. "display_count": attempt_display_count,
  802. "searched_count": searched_count - attempt_searched_start,
  803. "raw_searched_count": raw_searched_count - attempt_raw_start,
  804. "search_page_count": search_page_count - attempt_pages_start,
  805. "stop_reason": attempt_stop_reason,
  806. "requested_second_page": attempt_requested_second_page,
  807. "first_page_classified_count": attempt_first_page_classified_count,
  808. "first_page_creation_count": attempt_first_page_creation_count,
  809. "first_page_creation_ratio": attempt_first_page_creation_ratio,
  810. "pages": attempt_page_summaries,
  811. }
  812. )
  813. if (
  814. attempt_index <= low_result_max_retries
  815. and display_count <= low_result_retry_threshold
  816. ):
  817. retry_triggered = True
  818. continue
  819. break
  820. status = "done" if display_count else "failed"
  821. if status == "done":
  822. done += 1
  823. else:
  824. failed += 1
  825. repo.update_acquisition_job(
  826. _job_id(job),
  827. status=status,
  828. attempt_count=attempts,
  829. error_message=None if status != "failed" else "; ".join(errors[-3:]),
  830. metadata={
  831. "query_text": query.query_text,
  832. "searched_count": searched_count,
  833. "raw_searched_count": raw_searched_count,
  834. "display_count": display_count,
  835. "skipped_item_count": skipped_item_count,
  836. "search_page_count": search_page_count,
  837. "pagination": {
  838. **_pagination_config(
  839. pagination_creation_threshold,
  840. pagination_max_pages,
  841. ),
  842. "first_page_classified_count": first_page_classified_count,
  843. "first_page_creation_count": first_page_creation_count,
  844. "first_page_creation_ratio": first_page_creation_ratio,
  845. "requested_second_page": requested_second_page,
  846. "stop_reason": pagination_stop_reason,
  847. "pages": page_summaries,
  848. },
  849. "low_result_retry": {
  850. "threshold": low_result_retry_threshold,
  851. "max_retries": low_result_max_retries if retry_enabled else 0,
  852. "triggered": retry_triggered,
  853. "attempt_count": len(retry_attempts),
  854. "attempts": retry_attempts,
  855. },
  856. "errors": errors[-3:],
  857. },
  858. )
  859. _trace_event(
  860. trace_writer,
  861. job_trace_context,
  862. stage="search",
  863. event_type="job_finished",
  864. status=status,
  865. target_table="acquisition_jobs",
  866. target_id=_job_id(job),
  867. payload={
  868. "searched_count": searched_count,
  869. "raw_searched_count": raw_searched_count,
  870. "display_count": display_count,
  871. "skipped_item_count": skipped_item_count,
  872. "search_page_count": search_page_count,
  873. },
  874. error_message=None if status != "failed" else "; ".join(errors[-3:]),
  875. )
  876. except Exception as exc:
  877. message = clip_text(exc, ERROR_MESSAGE_MAX_CHARS)
  878. failed += 1
  879. repo.update_acquisition_job(
  880. _job_id(job),
  881. status="failed",
  882. attempt_count=attempts,
  883. error_message=message,
  884. metadata={
  885. "query_text": query.query_text,
  886. "searched_count": searched_count,
  887. "raw_searched_count": raw_searched_count,
  888. "display_count": display_count,
  889. "skipped_item_count": skipped_item_count,
  890. "search_page_count": search_page_count,
  891. "pagination": {
  892. **_pagination_config(
  893. pagination_creation_threshold,
  894. pagination_max_pages,
  895. ),
  896. "first_page_classified_count": first_page_classified_count,
  897. "first_page_creation_count": first_page_creation_count,
  898. "first_page_creation_ratio": first_page_creation_ratio,
  899. "requested_second_page": requested_second_page,
  900. "stop_reason": pagination_stop_reason or "failed",
  901. "pages": page_summaries,
  902. },
  903. "errors": errors[-3:],
  904. },
  905. )
  906. _trace_event(
  907. trace_writer,
  908. job_trace_context,
  909. stage="search",
  910. event_type="job_failed",
  911. status="failed",
  912. severity="error",
  913. target_table="acquisition_jobs",
  914. target_id=_job_id(job),
  915. error_message=message,
  916. )
  917. run_status = "done" if done > 0 else "failed"
  918. update_run = getattr(repo, "update_acquisition_run", None)
  919. if update_run:
  920. update_run(
  921. run.id,
  922. status=run_status,
  923. metadata={
  924. "total_jobs": total_jobs,
  925. "done": done,
  926. "partial": partial,
  927. "failed": failed,
  928. "skipped": skipped,
  929. "pagination": _pagination_config(
  930. pagination_creation_threshold,
  931. pagination_max_pages,
  932. ),
  933. },
  934. )
  935. _trace_event(
  936. trace_writer,
  937. run_trace_context,
  938. stage="search",
  939. event_type="acquisition_finished",
  940. status=run_status,
  941. target_table="acquisition_runs",
  942. target_id=run.id,
  943. payload={
  944. "total_jobs": total_jobs,
  945. "done": done,
  946. "partial": partial,
  947. "failed": failed,
  948. "skipped": skipped,
  949. },
  950. )
  951. return RunBatchResult(
  952. run_id=run.id,
  953. total_jobs=total_jobs,
  954. done=done,
  955. partial=partial,
  956. failed=failed,
  957. skipped=skipped,
  958. )