test_oss_archive.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340
  1. from __future__ import annotations
  2. import threading
  3. import time
  4. from datetime import datetime, timedelta, timezone
  5. from content_agent.integrations.oss_archive import (
  6. AsyncArchiveDispatcher,
  7. archive_due_records,
  8. archive_pending_for_run,
  9. mark_archive_pending,
  10. )
  11. from content_agent.integrations.runtime_files import LocalRuntimeFileStore
  12. NOW = datetime(2026, 6, 16, 10, 0, tzinfo=timezone.utc)
  13. def _record(**overrides):
  14. base = {
  15. "record_schema_version": "runtime_record.v1",
  16. "run_id": "run_001",
  17. "policy_run_id": "policy_run_001",
  18. "platform": "douyin",
  19. "platform_content_id": "content_001",
  20. "content_media_status": "metadata_only",
  21. "content_metadata_source": "douyin_keyword_search",
  22. "play_url": "https://source.example/video.mp4",
  23. "local_path": None,
  24. "oss_url": None,
  25. "raw_payload": {"run_id": "run_001", "platform_content_id": "content_001"},
  26. "created_at": NOW.isoformat(),
  27. }
  28. base.update(overrides)
  29. return base
  30. def test_mark_archive_pending_for_judged_video_with_play_url():
  31. [row] = mark_archive_pending([_record()], now_fn=lambda: NOW)
  32. assert row["content_media_status"] == "oss_upload_pending"
  33. assert row["raw_payload"]["oss_archive_status"] == "pending"
  34. assert row["raw_payload"]["oss_archive_attempt_count"] == 0
  35. assert row["raw_payload"]["oss_archive_next_retry_at"] == NOW.isoformat()
  36. assert row["raw_payload"]["oss_archive_deadline_at"] == (NOW + timedelta(hours=24)).isoformat()
  37. def test_archive_due_records_updates_success_to_oss_uploaded():
  38. [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW)
  39. def upload(src_url, **kwargs):
  40. return {
  41. "status": "ok",
  42. "oss_url": "https://res.example/video.mp4",
  43. "oss_object_key": "crawler/video/content_001.mp4",
  44. "save_oss_timestamp": 123,
  45. "oss_payload_mode": "no_referer",
  46. }
  47. [row] = archive_due_records([pending], upload_fn=upload, now_fn=lambda: NOW)
  48. assert row["content_media_status"] == "oss_uploaded"
  49. assert row["oss_url"] == "https://res.example/video.mp4"
  50. assert row["raw_payload"]["oss_archive_status"] == "uploaded"
  51. assert row["raw_payload"]["oss_archive_attempt_count"] == 1
  52. assert row["raw_payload"]["oss_object_key"] == "crawler/video/content_001.mp4"
  53. assert row["raw_payload"]["oss_payload_mode"] == "no_referer"
  54. assert row["raw_payload"]["oss_timing_metrics"]["oss_upload_duration_ms"] >= 0
  55. def test_archive_due_records_uploads_due_records_with_bounded_workers():
  56. due_records = [
  57. mark_archive_pending(
  58. [
  59. _record(
  60. platform_content_id=f"content_{index:03d}",
  61. play_url=f"https://source.example/video_{index}.mp4",
  62. )
  63. ],
  64. now_fn=lambda: NOW,
  65. )[0]
  66. for index in range(6)
  67. ]
  68. future_retry = mark_archive_pending(
  69. [_record(platform_content_id="content_future", play_url="https://source.example/future.mp4")],
  70. now_fn=lambda: NOW,
  71. )[0]
  72. future_retry["raw_payload"]["oss_archive_next_retry_at"] = (NOW + timedelta(minutes=5)).isoformat()
  73. no_play = _record(platform_content_id="content_no_play", play_url=None)
  74. uploaded = _record(
  75. platform_content_id="content_uploaded",
  76. content_media_status="oss_uploaded",
  77. play_url="https://source.example/uploaded.mp4",
  78. oss_url="https://res.example/uploaded.mp4",
  79. )
  80. records = [due_records[0], future_retry, due_records[1], no_play, *due_records[2:], uploaded]
  81. active = 0
  82. max_active = 0
  83. calls: list[str] = []
  84. lock = threading.Lock()
  85. def upload(src_url, **kwargs):
  86. nonlocal active, max_active
  87. with lock:
  88. active += 1
  89. max_active = max(max_active, active)
  90. calls.append(src_url)
  91. time.sleep(0.03)
  92. with lock:
  93. active -= 1
  94. return {
  95. "status": "ok",
  96. "oss_url": src_url.replace("source.example", "res.example"),
  97. "oss_object_key": src_url.rsplit("/", 1)[-1],
  98. }
  99. archived = archive_due_records(
  100. records,
  101. upload_fn=upload,
  102. now_fn=lambda: NOW,
  103. max_workers=3,
  104. )
  105. assert sorted(calls) == sorted(record["play_url"] for record in due_records)
  106. assert 1 < max_active <= 3
  107. assert [row["platform_content_id"] for row in archived] == [
  108. row["platform_content_id"] for row in records
  109. ]
  110. assert archived[1]["content_media_status"] == "oss_upload_pending"
  111. assert archived[3]["content_media_status"] == "metadata_only"
  112. assert archived[-1]["content_media_status"] == "oss_uploaded"
  113. def test_archive_pending_for_run_uses_env_worker_default(monkeypatch, tmp_path):
  114. monkeypatch.setenv("CONTENT_AGENT_OSS_ARCHIVE_MAX_WORKERS", "bad")
  115. runtime = LocalRuntimeFileStore(tmp_path)
  116. runtime.prepare_run("run_001")
  117. [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW)
  118. runtime.append_jsonl("run_001", "content_media_records.jsonl", [pending])
  119. archived = archive_pending_for_run(
  120. runtime,
  121. "run_001",
  122. upload_fn=lambda src_url, **kwargs: {
  123. "status": "ok",
  124. "oss_url": "https://res.example/video.mp4",
  125. },
  126. now_fn=lambda: NOW,
  127. )
  128. assert archived[0]["content_media_status"] == "oss_uploaded"
  129. def test_async_archive_dispatcher_uploads_without_waiting_and_writes_after_enabled(tmp_path):
  130. runtime = LocalRuntimeFileStore(tmp_path)
  131. runtime.prepare_run("run_001")
  132. pending_records = [
  133. mark_archive_pending(
  134. [
  135. _record(
  136. platform_content_id=f"content_{index:03d}",
  137. play_url=f"https://source.example/video_{index}.mp4",
  138. )
  139. ],
  140. now_fn=lambda: NOW,
  141. )[0]
  142. for index in range(6)
  143. ]
  144. active = 0
  145. max_active = 0
  146. calls: list[str] = []
  147. lock = threading.Lock()
  148. def upload(src_url, **kwargs):
  149. nonlocal active, max_active
  150. with lock:
  151. active += 1
  152. max_active = max(max_active, active)
  153. calls.append(src_url)
  154. time.sleep(0.03)
  155. with lock:
  156. active -= 1
  157. return {
  158. "status": "ok",
  159. "oss_url": src_url.replace("source.example", "res.example"),
  160. }
  161. dispatcher = AsyncArchiveDispatcher(
  162. runtime,
  163. "run_001",
  164. upload_fn=upload,
  165. now_fn=lambda: NOW,
  166. max_workers=3,
  167. )
  168. dispatcher.submit_records(pending_records)
  169. assert runtime.read_jsonl("run_001", "content_media_records.jsonl") == []
  170. dispatcher.enable_runtime_writes()
  171. dispatcher.shutdown(wait=True)
  172. rows = runtime.read_jsonl("run_001", "content_media_records.jsonl")
  173. assert sorted(calls) == sorted(row["play_url"] for row in pending_records)
  174. assert 1 < max_active <= 3
  175. assert {row["content_media_status"] for row in rows} == {"oss_uploaded"}
  176. def test_archive_due_records_keeps_failed_attempt_pending_before_deadline():
  177. [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW)
  178. def upload(src_url, **kwargs):
  179. assert "referer" not in kwargs
  180. assert kwargs["timeout_seconds"] == 300.0 # 修永久卡死:OSS 单次尝试 3600→300
  181. return {
  182. "status": "failed",
  183. "failure_type": "oss_upload_http_error",
  184. "exception_type": "ReadTimeout",
  185. "error_message": "proxy timed out",
  186. "oss_payload_mode": "no_referer",
  187. "endpoint_host": "crawler-upload-v2.aiddit.com",
  188. "read_timeout_seconds": 300.0,
  189. "attempt_timeout_seconds": 300.0,
  190. "response_absent": True,
  191. }
  192. [row] = archive_due_records([pending], upload_fn=upload, now_fn=lambda: NOW)
  193. assert row["content_media_status"] == "oss_upload_pending"
  194. assert row["raw_payload"]["oss_payload_mode"] == "no_referer"
  195. assert row["failure_reason"] == "oss_upload_http_error"
  196. assert row["raw_payload"]["oss_archive_attempt_count"] == 1
  197. assert row["raw_payload"]["oss_archive_last_error"] == "oss_upload_http_error"
  198. assert row["raw_payload"]["oss_archive_last_error_message"] == "proxy timed out"
  199. assert row["raw_payload"]["oss_endpoint_host"] == "crawler-upload-v2.aiddit.com"
  200. assert row["raw_payload"]["oss_read_timeout_seconds"] == 300.0
  201. assert row["raw_payload"]["oss_attempt_timeout_seconds"] == 300.0
  202. assert row["raw_payload"]["oss_response_absent"] is True
  203. assert row["raw_payload"]["oss_archive_next_retry_at"] == (NOW + timedelta(minutes=15)).isoformat()
  204. assert row["raw_payload"]["oss_timing_metrics"]["oss_upload_duration_ms"] >= 0
  205. def test_archive_due_records_tries_fallback_candidate_after_first_url_fails():
  206. [pending] = mark_archive_pending(
  207. [
  208. _record(
  209. raw_payload={
  210. "run_id": "run_001",
  211. "platform_content_id": "content_001",
  212. "video_url_candidates": [
  213. {
  214. "url": "https://source.example/video.mp4",
  215. "host": "source.example",
  216. "path": "$.search.video_url_list[0].video_url",
  217. "source": "search",
  218. "candidate_index": 0,
  219. },
  220. {
  221. "url": "https://backup.example/video.mp4",
  222. "host": "backup.example",
  223. "path": "$.detail.video_url_list[0].video_url",
  224. "source": "detail",
  225. "candidate_index": 1,
  226. },
  227. ],
  228. },
  229. )
  230. ],
  231. now_fn=lambda: NOW,
  232. )
  233. calls: list[str] = []
  234. def upload(src_url, **kwargs):
  235. calls.append(src_url)
  236. if src_url == "https://source.example/video.mp4":
  237. return {
  238. "status": "failed",
  239. "failure_type": "oss_upload_http_error",
  240. "exception_type": "ReadTimeout",
  241. "read_timeout_seconds": 300.0,
  242. "attempt_timeout_seconds": kwargs["timeout_seconds"],
  243. "response_absent": True,
  244. }
  245. return {
  246. "status": "ok",
  247. "oss_url": "https://res.example/video.mp4",
  248. "oss_object_key": "crawler/video/content_001.mp4",
  249. }
  250. [row] = archive_due_records([pending], upload_fn=upload, now_fn=lambda: NOW)
  251. assert calls == ["https://source.example/video.mp4", "https://backup.example/video.mp4"]
  252. assert row["content_media_status"] == "oss_uploaded"
  253. assert row["oss_url"] == "https://res.example/video.mp4"
  254. assert row["raw_payload"]["oss_archive_attempt_count"] == 2
  255. assert row["raw_payload"]["oss_archive_attempted_candidate_count"] == 2
  256. assert row["raw_payload"]["oss_archive_selected_video_url_host"] == "backup.example"
  257. assert row["raw_payload"]["oss_archive_selected_video_url_path"] == "$.detail.video_url_list[0].video_url"
  258. assert row["raw_payload"]["oss_url_attempts"][0]["failure_type"] == "oss_upload_http_error"
  259. assert row["raw_payload"]["oss_url_attempts"][1]["status"] == "ok"
  260. def test_archive_due_records_keeps_invalid_response_summary():
  261. [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW)
  262. def upload(src_url, **kwargs):
  263. assert "referer" not in kwargs
  264. return {
  265. "status": "failed",
  266. "failure_type": "oss_upload_response_invalid",
  267. "oss_payload_mode": "no_referer",
  268. "oss_response_summary": {
  269. "status": 10000,
  270. "msg": "bad",
  271. "oss_object_present": False,
  272. "oss_object_has_cdn_url": False,
  273. },
  274. }
  275. [row] = archive_due_records([pending], upload_fn=upload, now_fn=lambda: NOW)
  276. assert row["content_media_status"] == "oss_upload_pending"
  277. assert row["raw_payload"]["oss_archive_last_error"] == "oss_upload_response_invalid"
  278. assert row["raw_payload"]["oss_payload_mode"] == "no_referer"
  279. assert row["raw_payload"]["oss_response_summary"]["status"] == 10000
  280. def test_archive_due_records_marks_failed_after_deadline():
  281. [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW)
  282. pending["raw_payload"]["oss_archive_deadline_at"] = (NOW - timedelta(seconds=1)).isoformat()
  283. [row] = archive_due_records(
  284. [pending],
  285. upload_fn=lambda src_url, **kwargs: {"status": "failed", "failure_type": "oss_upload_http_error"},
  286. now_fn=lambda: NOW,
  287. )
  288. assert row["content_media_status"] == "oss_upload_failed"
  289. assert row["failure_reason"] == "oss_upload_failed"
  290. assert row["raw_payload"]["oss_archive_status"] == "failed"
  291. assert row["raw_payload"]["oss_archive_next_retry_at"] is None