80 lines
3.0 KiB
Python
80 lines
3.0 KiB
Python
from __future__ import annotations
|
|
|
|
from datetime import UTC, datetime
|
|
from typing import Any
|
|
|
|
import httpx
|
|
|
|
from app.services.collectors.base import BaseCollector
|
|
|
|
|
|
class NewsLiveStreamsCollector(BaseCollector):
|
|
"""Collect normalized news live-stream sources from a JSON endpoint."""
|
|
|
|
name = "news_live_streams"
|
|
priority = "P2"
|
|
module = "L4"
|
|
frequency_hours = 12
|
|
data_type = "news_live_stream"
|
|
fail_on_empty = False
|
|
|
|
async def fetch(self) -> list[dict[str, Any]]:
|
|
request_url = (self._resolved_url or "").strip()
|
|
if not request_url:
|
|
return []
|
|
|
|
async with httpx.AsyncClient(timeout=45.0, follow_redirects=True) as client:
|
|
response = await client.get(
|
|
request_url,
|
|
headers={
|
|
"User-Agent": "Planet-Intelligence-System/1.0 (Python/collector)",
|
|
"Accept": "application/json",
|
|
},
|
|
)
|
|
response.raise_for_status()
|
|
return self.parse_response(response.json())
|
|
|
|
def parse_response(self, response: Any) -> list[dict[str, Any]]:
|
|
if isinstance(response, dict):
|
|
candidates = response.get("sources") or response.get("streams") or response.get("data") or []
|
|
elif isinstance(response, list):
|
|
candidates = response
|
|
else:
|
|
candidates = []
|
|
|
|
normalized: list[dict[str, Any]] = []
|
|
for index, item in enumerate(candidates):
|
|
if not isinstance(item, dict):
|
|
continue
|
|
|
|
stream_id = item.get("id") or item.get("source_id") or item.get("slug") or f"news-live-{index + 1}"
|
|
name = str(item.get("name") or item.get("title") or f"News Live {index + 1}").strip()
|
|
if not name:
|
|
continue
|
|
|
|
metadata = {
|
|
"provider": item.get("provider") or item.get("publisher") or "Collector",
|
|
"region": item.get("region") or item.get("country") or "Global",
|
|
"language": item.get("language") or "und",
|
|
"source_type": item.get("source_type") or "iframe",
|
|
"embed_url": item.get("embed_url") or item.get("url") or "",
|
|
"stream_url": item.get("stream_url") or "",
|
|
"homepage_url": item.get("homepage_url") or item.get("source_url") or "",
|
|
"poster_url": item.get("poster_url") or "",
|
|
"sort_order": item.get("sort_order", 200 + index),
|
|
"notes": item.get("notes") or item.get("description") or "",
|
|
"is_enabled": item.get("is_enabled", True),
|
|
}
|
|
|
|
normalized.append(
|
|
{
|
|
"source_id": str(stream_id),
|
|
"name": name,
|
|
"description": metadata["notes"],
|
|
"metadata": metadata,
|
|
"reference_date": item.get("reference_date", datetime.now(UTC).isoformat()),
|
|
}
|
|
)
|
|
|
|
return normalized
|