"""Per-run manifest.json — a queryable status snapshot for a campaign run. Turns what today required grepping structlog output and counting output files into a single file any tool (or a human) can read: overall status, per-shard timing/status, and video counts at each stage boundary. An in-process asyncio.Lock guards concurrent shard updates (same pattern as base.py's claim_video_id — safe because all shards run in one event loop). Deliberately NOT part of the StorageBackend ABC — a standalone, run-scoped status/coordination artifact, not analytics data. Still needs to work on S3 for the same reason resume needs to survive a full ECS task replacement, not just an in-process crash: read/write here branch on settings.storage_backend exactly like storage/s3.py, using the same key layout ("comments/{campaign_id}/{run_id}/manifest.json" / ".../request.json") so a fresh container can find and resume a "running" run regardless of which container originally owned it. """ from __future__ import annotations import asyncio import json from datetime import UTC, datetime from pathlib import Path from typing import Any, cast from marketing_intelligence.core.config import settings _lock = asyncio.Lock() _s3_client = None def _s3() -> Any: global _s3_client if _s3_client is None: import boto3 _s3_client = boto3.client("s3", region_name=settings.aws_region) return _s3_client def _manifest_path(campaign_id: str, run_id: str) -> Path: return ( Path(settings.output_dir) / "comments" / campaign_id / run_id / "manifest.json" ) def _request_path(campaign_id: str, run_id: str) -> Path: return ( Path(settings.output_dir) / "comments" / campaign_id / run_id / "request.json" ) def _manifest_key(campaign_id: str, run_id: str) -> str: return f"comments/{campaign_id}/{run_id}/manifest.json" def _request_key(campaign_id: str, run_id: str) -> str: return f"comments/{campaign_id}/{run_id}/request.json" def now_iso() -> str: return datetime.now(UTC).isoformat() def _get_json_local(path: Path) -> dict[str, Any] | None: if not path.exists(): return None return cast(dict[str, Any], json.loads(path.read_text())) def _put_json_local(path: Path, data: dict[str, Any]) -> None: path.parent.mkdir(parents=True, exist_ok=True) path.write_text(json.dumps(data, indent=2, ensure_ascii=False, default=str)) async def _get_json_s3(key: str) -> dict[str, Any] | None: from botocore.exceptions import ClientError def _get() -> dict[str, Any] | None: try: obj = _s3().get_object(Bucket=settings.s3_bucket, Key=key) return cast(dict[str, Any], json.loads(obj["Body"].read())) except ClientError as e: if e.response.get("Error", {}).get("Code") in ("NoSuchKey", "404"): return None raise return await asyncio.to_thread(_get) async def _put_json_s3(key: str, data: dict[str, Any]) -> None: body = json.dumps(data, indent=2, ensure_ascii=False, default=str).encode("utf-8") await asyncio.to_thread( _s3().put_object, Bucket=settings.s3_bucket, Key=key, Body=body, ContentType="application/json", ) async def _read(campaign_id: str, run_id: str) -> dict[str, Any]: if settings.storage_backend == "s3": data = await _get_json_s3(_manifest_key(campaign_id, run_id)) else: data = _get_json_local(_manifest_path(campaign_id, run_id)) return data or {} async def _write(campaign_id: str, run_id: str, data: dict[str, Any]) -> None: if settings.storage_backend == "s3": await _put_json_s3(_manifest_key(campaign_id, run_id), data) else: _put_json_local(_manifest_path(campaign_id, run_id), data) async def init_manifest( campaign_id: str, run_id: str, requested_sample_size: int, shard_labels: list[str], ) -> None: async with _lock: data = { "run_id": run_id, "campaign_id": campaign_id, "status": "running", "requested_sample_size": requested_sample_size, "started_at": now_iso(), "finished_at": None, "stage1": { "status": "running", "collected_videos": None, "shards": { label: { "status": "pending", "started_at": None, "finished_at": None, } for label in shard_labels }, }, "stage2": {"status": "pending", "started_at": None, "finished_at": None}, "stage3": {"status": "pending", "started_at": None, "finished_at": None}, } await _write(campaign_id, run_id, data) async def update_shard( campaign_id: str, run_id: str, shard_label: str, **fields: Any ) -> None: async with _lock: data = await _read(campaign_id, run_id) if not data: return data.setdefault("stage1", {}).setdefault("shards", {}).setdefault( shard_label, {} ).update(fields) await _write(campaign_id, run_id, data) async def update_stage( campaign_id: str, run_id: str, stage: str, **fields: Any ) -> None: async with _lock: data = await _read(campaign_id, run_id) if not data: return data.setdefault(stage, {}).update(fields) await _write(campaign_id, run_id, data) async def finalize(campaign_id: str, run_id: str, status: str) -> None: async with _lock: data = await _read(campaign_id, run_id) if not data: return data["status"] = status data["finished_at"] = now_iso() await _write(campaign_id, run_id, data) async def save_request( campaign_id: str, run_id: str, request_dict: dict[str, Any] ) -> None: """Persist the original SentimentRequest so a crashed run can be reconstructed and resumed.""" async with _lock: if settings.storage_backend == "s3": await _put_json_s3(_request_key(campaign_id, run_id), request_dict) else: _put_json_local(_request_path(campaign_id, run_id), request_dict) async def load_request(campaign_id: str, run_id: str) -> dict[str, Any] | None: if settings.storage_backend == "s3": return await _get_json_s3(_request_key(campaign_id, run_id)) return _get_json_local(_request_path(campaign_id, run_id)) async def mark_interrupted_shards( campaign_id: str, run_id: str, track_name: str ) -> None: """Mark any shard for this track still stuck at "running" as "interrupted".""" async with _lock: data = await _read(campaign_id, run_id) if not data: return shards = data.setdefault("stage1", {}).setdefault("shards", {}) for label, shard in shards.items(): if ( label == track_name or label.startswith(f"{track_name} (") ) and shard.get("status") == "running": shard["status"] = "interrupted" shard["finished_at"] = now_iso() await _write(campaign_id, run_id, data) async def count_runs(campaign_id: str) -> int: """Count existing run folders/prefixes for a campaign.""" if settings.storage_backend == "s3": def _count() -> int: prefixes: set[str] = set() paginator = _s3().get_paginator("list_objects_v2") for page in paginator.paginate( Bucket=settings.s3_bucket, Prefix=f"comments/{campaign_id}/", Delimiter="/", ): prefixes.update(p["Prefix"] for p in page.get("CommonPrefixes", [])) return len(prefixes) return await asyncio.to_thread(_count) base = Path(settings.output_dir) / "comments" / campaign_id if not base.is_dir(): return 0 return sum(1 for e in base.iterdir() if e.is_dir()) async def get_manifest(campaign_id: str, run_id: str) -> dict[str, Any] | None: """Read a run's manifest.json as-is. Returns None if the run doesn't exist.""" if settings.storage_backend == "s3": return await _get_json_s3(_manifest_key(campaign_id, run_id)) return _get_json_local(_manifest_path(campaign_id, run_id)) async def find_running_runs() -> list[tuple[str, str]]: """Scan every manifest.json for runs still marked "running".""" if settings.storage_backend == "s3": def _list_manifest_keys() -> list[str]: keys: list[str] = [] paginator = _s3().get_paginator("list_objects_v2") for page in paginator.paginate( Bucket=settings.s3_bucket, Prefix="comments/" ): keys.extend( obj["Key"] for obj in page.get("Contents", []) if obj["Key"].endswith("/manifest.json") ) return keys keys = await asyncio.to_thread(_list_manifest_keys) found: list[tuple[str, str]] = [] for key in keys: data = await _get_json_s3(key) if data and data.get("status") == "running": _, campaign_id, run_id, _ = key.split("/") found.append((campaign_id, run_id)) return found base = Path(settings.output_dir) / "comments" if not base.is_dir(): return [] found = [] for manifest_path in base.glob("*/*/manifest.json"): try: data = json.loads(manifest_path.read_text()) except json.JSONDecodeError, OSError: continue if data.get("status") == "running": campaign_id, run_id = ( manifest_path.parent.parent.name, manifest_path.parent.name, ) found.append((campaign_id, run_id)) return found