Files
planet/backend/app/services/collectors/news_live_streams.py
2026-04-12 04:36:38 +08:00

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