from __future__ import annotations import threading import time from datetime import datetime, timedelta, timezone from content_agent.integrations.oss_archive import ( AsyncArchiveDispatcher, archive_due_records, archive_pending_for_run, mark_archive_pending, ) from content_agent.integrations.runtime_files import LocalRuntimeFileStore NOW = datetime(2026, 6, 16, 10, 0, tzinfo=timezone.utc) def _record(**overrides): base = { "record_schema_version": "runtime_record.v1", "run_id": "run_001", "policy_run_id": "policy_run_001", "platform": "douyin", "platform_content_id": "content_001", "content_media_status": "metadata_only", "content_metadata_source": "douyin_keyword_search", "play_url": "https://source.example/video.mp4", "local_path": None, "oss_url": None, "raw_payload": {"run_id": "run_001", "platform_content_id": "content_001"}, "created_at": NOW.isoformat(), } base.update(overrides) return base def test_mark_archive_pending_for_judged_video_with_play_url(): [row] = mark_archive_pending([_record()], now_fn=lambda: NOW) assert row["content_media_status"] == "oss_upload_pending" assert row["raw_payload"]["oss_archive_status"] == "pending" assert row["raw_payload"]["oss_archive_attempt_count"] == 0 assert row["raw_payload"]["oss_archive_next_retry_at"] == NOW.isoformat() assert row["raw_payload"]["oss_archive_deadline_at"] == (NOW + timedelta(hours=24)).isoformat() def test_archive_due_records_updates_success_to_oss_uploaded(): [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW) def upload(src_url, **kwargs): return { "status": "ok", "oss_url": "https://res.example/video.mp4", "oss_object_key": "crawler/video/content_001.mp4", "save_oss_timestamp": 123, "oss_payload_mode": "no_referer", } [row] = archive_due_records([pending], upload_fn=upload, now_fn=lambda: NOW) assert row["content_media_status"] == "oss_uploaded" assert row["oss_url"] == "https://res.example/video.mp4" assert row["raw_payload"]["oss_archive_status"] == "uploaded" assert row["raw_payload"]["oss_archive_attempt_count"] == 1 assert row["raw_payload"]["oss_object_key"] == "crawler/video/content_001.mp4" assert row["raw_payload"]["oss_payload_mode"] == "no_referer" assert row["raw_payload"]["oss_timing_metrics"]["oss_upload_duration_ms"] >= 0 def test_archive_due_records_uploads_due_records_with_bounded_workers(): due_records = [ mark_archive_pending( [ _record( platform_content_id=f"content_{index:03d}", play_url=f"https://source.example/video_{index}.mp4", ) ], now_fn=lambda: NOW, )[0] for index in range(6) ] future_retry = mark_archive_pending( [_record(platform_content_id="content_future", play_url="https://source.example/future.mp4")], now_fn=lambda: NOW, )[0] future_retry["raw_payload"]["oss_archive_next_retry_at"] = (NOW + timedelta(minutes=5)).isoformat() no_play = _record(platform_content_id="content_no_play", play_url=None) uploaded = _record( platform_content_id="content_uploaded", content_media_status="oss_uploaded", play_url="https://source.example/uploaded.mp4", oss_url="https://res.example/uploaded.mp4", ) records = [due_records[0], future_retry, due_records[1], no_play, *due_records[2:], uploaded] active = 0 max_active = 0 calls: list[str] = [] lock = threading.Lock() def upload(src_url, **kwargs): nonlocal active, max_active with lock: active += 1 max_active = max(max_active, active) calls.append(src_url) time.sleep(0.03) with lock: active -= 1 return { "status": "ok", "oss_url": src_url.replace("source.example", "res.example"), "oss_object_key": src_url.rsplit("/", 1)[-1], } archived = archive_due_records( records, upload_fn=upload, now_fn=lambda: NOW, max_workers=3, ) assert sorted(calls) == sorted(record["play_url"] for record in due_records) assert 1 < max_active <= 3 assert [row["platform_content_id"] for row in archived] == [ row["platform_content_id"] for row in records ] assert archived[1]["content_media_status"] == "oss_upload_pending" assert archived[3]["content_media_status"] == "metadata_only" assert archived[-1]["content_media_status"] == "oss_uploaded" def test_archive_pending_for_run_uses_env_worker_default(monkeypatch, tmp_path): monkeypatch.setenv("CONTENT_AGENT_OSS_ARCHIVE_MAX_WORKERS", "bad") runtime = LocalRuntimeFileStore(tmp_path) runtime.prepare_run("run_001") [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW) runtime.append_jsonl("run_001", "content_media_records.jsonl", [pending]) archived = archive_pending_for_run( runtime, "run_001", upload_fn=lambda src_url, **kwargs: { "status": "ok", "oss_url": "https://res.example/video.mp4", }, now_fn=lambda: NOW, ) assert archived[0]["content_media_status"] == "oss_uploaded" def test_async_archive_dispatcher_uploads_without_waiting_and_writes_after_enabled(tmp_path): runtime = LocalRuntimeFileStore(tmp_path) runtime.prepare_run("run_001") pending_records = [ mark_archive_pending( [ _record( platform_content_id=f"content_{index:03d}", play_url=f"https://source.example/video_{index}.mp4", ) ], now_fn=lambda: NOW, )[0] for index in range(6) ] active = 0 max_active = 0 calls: list[str] = [] lock = threading.Lock() def upload(src_url, **kwargs): nonlocal active, max_active with lock: active += 1 max_active = max(max_active, active) calls.append(src_url) time.sleep(0.03) with lock: active -= 1 return { "status": "ok", "oss_url": src_url.replace("source.example", "res.example"), } dispatcher = AsyncArchiveDispatcher( runtime, "run_001", upload_fn=upload, now_fn=lambda: NOW, max_workers=3, ) dispatcher.submit_records(pending_records) assert runtime.read_jsonl("run_001", "content_media_records.jsonl") == [] dispatcher.enable_runtime_writes() dispatcher.shutdown(wait=True) rows = runtime.read_jsonl("run_001", "content_media_records.jsonl") assert sorted(calls) == sorted(row["play_url"] for row in pending_records) assert 1 < max_active <= 3 assert {row["content_media_status"] for row in rows} == {"oss_uploaded"} def test_archive_due_records_keeps_failed_attempt_pending_before_deadline(): [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW) def upload(src_url, **kwargs): assert "referer" not in kwargs assert kwargs["timeout_seconds"] == 300.0 # 修永久卡死:OSS 单次尝试 3600→300 return { "status": "failed", "failure_type": "oss_upload_http_error", "exception_type": "ReadTimeout", "error_message": "proxy timed out", "oss_payload_mode": "no_referer", "endpoint_host": "crawler-upload-v2.aiddit.com", "read_timeout_seconds": 300.0, "attempt_timeout_seconds": 300.0, "response_absent": True, } [row] = archive_due_records([pending], upload_fn=upload, now_fn=lambda: NOW) assert row["content_media_status"] == "oss_upload_pending" assert row["raw_payload"]["oss_payload_mode"] == "no_referer" assert row["failure_reason"] == "oss_upload_http_error" assert row["raw_payload"]["oss_archive_attempt_count"] == 1 assert row["raw_payload"]["oss_archive_last_error"] == "oss_upload_http_error" assert row["raw_payload"]["oss_archive_last_error_message"] == "proxy timed out" assert row["raw_payload"]["oss_endpoint_host"] == "crawler-upload-v2.aiddit.com" assert row["raw_payload"]["oss_read_timeout_seconds"] == 300.0 assert row["raw_payload"]["oss_attempt_timeout_seconds"] == 300.0 assert row["raw_payload"]["oss_response_absent"] is True assert row["raw_payload"]["oss_archive_next_retry_at"] == (NOW + timedelta(minutes=15)).isoformat() assert row["raw_payload"]["oss_timing_metrics"]["oss_upload_duration_ms"] >= 0 def test_archive_due_records_tries_fallback_candidate_after_first_url_fails(): [pending] = mark_archive_pending( [ _record( raw_payload={ "run_id": "run_001", "platform_content_id": "content_001", "video_url_candidates": [ { "url": "https://source.example/video.mp4", "host": "source.example", "path": "$.search.video_url_list[0].video_url", "source": "search", "candidate_index": 0, }, { "url": "https://backup.example/video.mp4", "host": "backup.example", "path": "$.detail.video_url_list[0].video_url", "source": "detail", "candidate_index": 1, }, ], }, ) ], now_fn=lambda: NOW, ) calls: list[str] = [] def upload(src_url, **kwargs): calls.append(src_url) if src_url == "https://source.example/video.mp4": return { "status": "failed", "failure_type": "oss_upload_http_error", "exception_type": "ReadTimeout", "read_timeout_seconds": 300.0, "attempt_timeout_seconds": kwargs["timeout_seconds"], "response_absent": True, } return { "status": "ok", "oss_url": "https://res.example/video.mp4", "oss_object_key": "crawler/video/content_001.mp4", } [row] = archive_due_records([pending], upload_fn=upload, now_fn=lambda: NOW) assert calls == ["https://source.example/video.mp4", "https://backup.example/video.mp4"] assert row["content_media_status"] == "oss_uploaded" assert row["oss_url"] == "https://res.example/video.mp4" assert row["raw_payload"]["oss_archive_attempt_count"] == 2 assert row["raw_payload"]["oss_archive_attempted_candidate_count"] == 2 assert row["raw_payload"]["oss_archive_selected_video_url_host"] == "backup.example" assert row["raw_payload"]["oss_archive_selected_video_url_path"] == "$.detail.video_url_list[0].video_url" assert row["raw_payload"]["oss_url_attempts"][0]["failure_type"] == "oss_upload_http_error" assert row["raw_payload"]["oss_url_attempts"][1]["status"] == "ok" def test_archive_due_records_keeps_invalid_response_summary(): [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW) def upload(src_url, **kwargs): assert "referer" not in kwargs return { "status": "failed", "failure_type": "oss_upload_response_invalid", "oss_payload_mode": "no_referer", "oss_response_summary": { "status": 10000, "msg": "bad", "oss_object_present": False, "oss_object_has_cdn_url": False, }, } [row] = archive_due_records([pending], upload_fn=upload, now_fn=lambda: NOW) assert row["content_media_status"] == "oss_upload_pending" assert row["raw_payload"]["oss_archive_last_error"] == "oss_upload_response_invalid" assert row["raw_payload"]["oss_payload_mode"] == "no_referer" assert row["raw_payload"]["oss_response_summary"]["status"] == 10000 def test_archive_due_records_marks_failed_after_deadline(): [pending] = mark_archive_pending([_record()], now_fn=lambda: NOW) pending["raw_payload"]["oss_archive_deadline_at"] = (NOW - timedelta(seconds=1)).isoformat() [row] = archive_due_records( [pending], upload_fn=lambda src_url, **kwargs: {"status": "failed", "failure_type": "oss_upload_http_error"}, now_fn=lambda: NOW, ) assert row["content_media_status"] == "oss_upload_failed" assert row["failure_reason"] == "oss_upload_failed" assert row["raw_payload"]["oss_archive_status"] == "failed" assert row["raw_payload"]["oss_archive_next_retry_at"] is None