"""Xtream Codes compatible API — used by IPTV player apps.""" import asyncio import logging import time from datetime import datetime, timezone from typing import AsyncIterator import aiohttp from fastapi import APIRouter, Depends, HTTPException, Request, Response from fastapi.responses import StreamingResponse, PlainTextResponse from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import selectinload from ..config import settings from ..core.catalog import get_user_categories, get_user_channels, get_channel_provider_maps, can_user_access_channel, get_user_preferred_provider_ids from ..core.epg_manager import filter_epg_for_channels from ..core.pool import pool, NoSlotsAvailableError, MaxConnectionsExceededError from ..core.xtream_client import XtreamClient from ..core.vod_tracker import vod_tracker, VodSession from ..database import get_db, AsyncSessionLocal from ..models.channel import Channel, ContentType from ..models.custom_category import CustomCategory, CustomCategoryItem from ..models.jellyfin import JellyfinConfig, JellyfinItem from ..models.log import ConnectionLog from ..models.provider import ProviderAccount from ..models.user import User, UserCustomCatalogEntry CC_PREFIX = "cc_" # prefix for custom category IDs in Xtream protocol BLOCKED_IMAGE_PATH = "/opt/kirastream/backend/static/blocked_image.jpg" # episode_id → provider_account_id; populated by get_series_info calls so # series_stream can route episode requests without storing episodes in the DB _episode_provider_cache: dict[str, int] = {} # jellyfin fake_ep_id → direct-play URL; populated in _get_jellyfin_series_info _jellyfin_episode_cache: dict[str, str] = {} logger = logging.getLogger(__name__) router = APIRouter(tags=["xtream"]) SERVER_URL_PLACEHOLDER = "{server_url}" def get_client_ip(request: Request) -> str: xff = request.headers.get("X-Forwarded-For", "") if xff: return xff.split(",")[0].strip() xri = request.headers.get("X-Real-IP", "") if xri: return xri return request.client.host if request.client else "unknown" async def _create_log( user_id: int, username: str, channel_id: int | None, channel_name: str | None, content_type: str, stream_id: str | None, client_ip: str, ) -> int: async with AsyncSessionLocal() as db: log = ConnectionLog( user_id=user_id, username=username, channel_id=channel_id, channel_name=channel_name, content_type=content_type, stream_id=stream_id, client_ip=client_ip, ) db.add(log) await db.commit() await db.refresh(log) return log.id async def _finish_log(log_id: int, bytes_transferred: int, started: float) -> None: ended = datetime.now(timezone.utc) duration = int(time.time() - started) async with AsyncSessionLocal() as db: result = await db.execute(select(ConnectionLog).where(ConnectionLog.id == log_id)) log = result.scalar_one_or_none() if log: log.ended_at = ended log.duration_seconds = duration log.bytes_transferred = bytes_transferred await db.commit() async def _stream_blocked_image() -> StreamingResponse: """Return a continuous MPEG-TS video stream of the blocked image via FFmpeg.""" import os if not os.path.exists(BLOCKED_IMAGE_PATH): raise HTTPException(503, "No provider slots available — try again later") async def generate(): proc = await asyncio.create_subprocess_exec( "ffmpeg", "-re", "-loop", "1", "-i", BLOCKED_IMAGE_PATH, "-c:v", "libx264", "-preset", "ultrafast", "-tune", "stillimage", "-b:v", "400k", "-pix_fmt", "yuv420p", "-vf", "scale=1280:720", "-f", "mpegts", "pipe:1", "-loglevel", "quiet", stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.DEVNULL, ) try: while True: chunk = await proc.stdout.read(65536) if not chunk: break yield chunk finally: try: proc.terminate() await asyncio.wait_for(proc.wait(), timeout=2.0) except (asyncio.TimeoutError, ProcessLookupError): try: proc.kill() except ProcessLookupError: pass return StreamingResponse( generate(), media_type="video/mp2t", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"}, ) async def _get_user_custom_cat_ids(db: AsyncSession, user_id: int) -> set[int] | None: """Returns set of custom category IDs for the user, or None if unrestricted (show all).""" result = await db.execute( select(UserCustomCatalogEntry.custom_category_id) .where(UserCustomCatalogEntry.user_id == user_id) ) rows = result.scalars().all() return set(rows) if rows else None # None means "no restriction — show all" async def _can_access_via_custom_category(db: AsyncSession, channel_id: int, user_id: int) -> bool: allowed_ids = await _get_user_custom_cat_ids(db, user_id) q = select(CustomCategoryItem).where(CustomCategoryItem.channel_id == channel_id) if allowed_ids is not None: q = q.where(CustomCategoryItem.custom_category_id.in_(allowed_ids)) result = await db.execute(q.limit(1)) return result.scalar_one_or_none() is not None async def _get_custom_category_channels(db: AsyncSession, cc_id: int) -> list[Channel]: result = await db.execute( select(Channel) .join(CustomCategoryItem, CustomCategoryItem.channel_id == Channel.id) .options(selectinload(Channel.provider_maps)) .where(CustomCategoryItem.custom_category_id == cc_id, Channel.is_active == True) # noqa: E712 .order_by(CustomCategoryItem.position, CustomCategoryItem.id) ) return list(result.scalars().all()) # --------------------------------------------------------------------------- # Auth helper # --------------------------------------------------------------------------- async def authenticate_user(username: str, password: str, db: AsyncSession) -> User: result = await db.execute(select(User).where(User.username == username, User.is_active == True)) # noqa: E712 user = result.scalar_one_or_none() if not user: raise HTTPException(403, "Invalid credentials") if not user.verify_password(password): raise HTTPException(403, "Invalid credentials") if user.expiry_date and user.expiry_date < datetime.now(timezone.utc): raise HTTPException(403, "Account expired") return user def server_url(request: Request) -> str: return f"{request.url.scheme}://{request.url.netloc}" # --------------------------------------------------------------------------- # player_api.php (main Xtream Codes endpoint) # --------------------------------------------------------------------------- @router.get("/player_api.php") async def player_api( request: Request, username: str = "", password: str = "", action: str = "", category_id: str | None = None, stream_id: str | None = None, series_id: str | None = None, vod_id: str | None = None, db: AsyncSession = Depends(get_db), ): if not username or not password: return {"user_info": {"status": "Disabled"}, "server_info": {}} user = await authenticate_user(username, password, db) base = server_url(request) if not action: return await _user_info(user, password, base, db) if action in ("user_info", "get_account_info"): return await _user_info(user, password, base, db) if action == "get_live_categories": return await _get_live_categories(user, db) if action == "get_live_streams": return await _get_live_streams(user, db, category_id, base) if action == "get_vod_categories": return await _get_categories(user, db, ContentType.movie) if action == "get_vod_streams": return await _get_channels(user, db, ContentType.movie, category_id, base) if action == "get_series_categories": return await _get_categories(user, db, ContentType.series) if action == "get_series": return await _get_channels(user, db, ContentType.series, category_id, base) if action == "get_series_info" and series_id: return await _get_series_info(series_id, user, db) if action == "get_vod_info" and vod_id: return await _get_vod_info(vod_id, user, db) if action == "get_short_epg" and stream_id: return {"epg_listings": []} return {} # --------------------------------------------------------------------------- # Playlist endpoint # --------------------------------------------------------------------------- @router.get("/get.php") async def get_playlist( request: Request, username: str = "", password: str = "", type: str = "m3u_plus", output: str = "ts", db: AsyncSession = Depends(get_db), ): user = await authenticate_user(username, password, db) base = server_url(request) channels = await get_user_channels(db, user) lines = ["#EXTM3U"] for ch in channels: cat_name = ch.category.name if ch.category else "" tvg_id = ch.tvg_id or "" logo = ch.tvg_logo or "" if ch.type == ContentType.live: url = f"{base}/{username}/{password}/{ch.stream_id_at_provider}" elif ch.type == ContentType.movie: ext = "mkv" url = f"{base}/movie/{username}/{password}/{ch.stream_id_at_provider}.{ext}" else: url = f"{base}/series/{username}/{password}/{ch.stream_id_at_provider}.mkv" lines.append( f'#EXTINF:-1 tvg-id="{tvg_id}" tvg-logo="{logo}" group-title="{cat_name}",{ch.name}' ) lines.append(url) return PlainTextResponse("\n".join(lines), media_type="application/x-mpegurl") # --------------------------------------------------------------------------- # EPG endpoint # --------------------------------------------------------------------------- @router.get("/xmltv.php") async def xmltv( username: str = "", password: str = "", db: AsyncSession = Depends(get_db), ): user = await authenticate_user(username, password, db) channels = await get_user_channels(db, user) tvg_ids = {ch.tvg_id for ch in channels if ch.tvg_id} epg_data = filter_epg_for_channels(tvg_ids) return Response(content=epg_data, media_type="application/xml") # --------------------------------------------------------------------------- # Live stream proxy /{username}/{password}/{stream_id} # Also handles /live/{username}/{password}/{stream_id} used by Smarters/IPTVnator # --------------------------------------------------------------------------- @router.get("/live/{username}/{password}/{stream_id}") @router.get("/{username}/{password}/{stream_id}") async def live_stream( username: str, password: str, stream_id: str, request: Request, db: AsyncSession = Depends(get_db), ): # Strip extension if present (e.g. "123456.ts" -> "123456") stream_id_clean = stream_id.rsplit(".", 1)[0] if "." in stream_id else stream_id user = await authenticate_user(username, password, db) # Resolve channel by stream_id_at_provider result = await db.execute( select(Channel) .options(selectinload(Channel.provider_maps)) .where(Channel.stream_id_at_provider == stream_id_clean, Channel.type == ContentType.live, Channel.is_active == True) # noqa: E712 ) channel = result.scalar_one_or_none() if not channel: raise HTTPException(404, "Channel not found") # Check user can access this channel (regular catalog or custom category) if not await can_user_access_channel(db, user, channel.id): if not await _can_access_via_custom_category(db, channel.id, user.id): raise HTTPException(403, "Channel not in your catalog") provider_maps = await get_channel_provider_maps(db, channel.id) if not provider_maps: raise HTTPException(503, "No provider available for this channel") preferred_ids = await get_user_preferred_provider_ids(db, user) try: handle = await pool.acquire( channel_id=channel.id, stream_url_resolver=None, provider_maps=provider_maps, user_id=str(user.id), user_max_conns=user.max_connections, is_priority=user.is_priority, preferred_provider_ids=preferred_ids, username=user.username, channel_name=channel.name, ) except MaxConnectionsExceededError: return await _stream_blocked_image() except NoSlotsAvailableError: raise HTTPException(503, "No provider slots available — try again later") client_ip = get_client_ip(request) log_id = await _create_log(user.id, user.username, channel.id, channel.name, "live", stream_id_clean, client_ip) stream_started = time.time() async def stream_generator() -> AsyncIterator[bytes]: bytes_sent = 0 try: async for chunk in handle.read(): bytes_sent += len(chunk) yield chunk finally: try: await asyncio.shield(pool.release(channel.id, handle.client_id, str(user.id))) except asyncio.CancelledError: pass asyncio.create_task(_finish_log(log_id, bytes_sent, stream_started)) return StreamingResponse( stream_generator(), media_type="video/mp2t", headers={ "Cache-Control": "no-cache", "X-Accel-Buffering": "no", }, ) # --------------------------------------------------------------------------- # VOD proxy /movie/{username}/{password}/{stream_id}.ext # --------------------------------------------------------------------------- @router.get("/movie/{username}/{password}/{stream_id}") async def vod_stream( username: str, password: str, stream_id: str, request: Request, db: AsyncSession = Depends(get_db), ): # Parse extension parts = stream_id.rsplit(".", 1) sid = parts[0] ext = parts[1] if len(parts) > 1 else "mkv" user = await authenticate_user(username, password, db) result = await db.execute( select(Channel) .options(selectinload(Channel.provider_maps)) .where(Channel.stream_id_at_provider == sid, Channel.type == ContentType.movie) ) channel = result.scalar_one_or_none() # Fallback: Jellyfin channels have non-numeric stream_id_at_provider so the app # receives ch.id (DB int) as the stream_id — look it up by primary key. if not channel: try: ch_id = int(sid) fb = await db.execute( select(Channel) .options(selectinload(Channel.provider_maps)) .where(Channel.id == ch_id, Channel.type == ContentType.movie) ) channel = fb.scalar_one_or_none() except (ValueError, TypeError): pass if not channel: raise HTTPException(404, "VOD not found") if not await can_user_access_channel(db, user, channel.id): if not await _can_access_via_custom_category(db, channel.id, user.id): raise HTTPException(403, "Not in your catalog") if not channel.provider_maps: raise HTTPException(503, "No provider available") pmap = channel.provider_maps[0] prov_res = await db.execute(select(ProviderAccount).where(ProviderAccount.id == pmap.provider_account_id)) provider = prov_res.scalar_one_or_none() vod_session = vod_tracker.start( user_id=str(user.id), username=user.username, channel_id=channel.id, channel_name=channel.name, provider_account_id=pmap.provider_account_id, provider_name=provider.name if provider else f"#{pmap.provider_account_id}", stream_type="movie", stream_url=pmap.stream_url, user_max_connections=user.max_connections, ) client_ip = get_client_ip(request) log_id = await _create_log(user.id, user.username, channel.id, channel.name, "movie", sid, client_ip) return await _proxy_vod(pmap.stream_url, request, vod_session=vod_session, log_id=log_id) # --------------------------------------------------------------------------- # Series proxy /series/{username}/{password}/{stream_id}.ext # --------------------------------------------------------------------------- @router.get("/series/{username}/{password}/{stream_id}") async def series_stream( username: str, password: str, stream_id: str, request: Request, db: AsyncSession = Depends(get_db), ): parts = stream_id.rsplit(".", 1) sid = parts[0] ext = parts[1] if len(parts) > 1 else "mkv" user = await authenticate_user(username, password, db) # First: try the episode as a stored channel (future: when episodes are synced) result = await db.execute( select(Channel) .options(selectinload(Channel.provider_maps)) .where(Channel.stream_id_at_provider == sid, Channel.type == ContentType.series) ) channel = result.scalar_one_or_none() client_ip = get_client_ip(request) if channel and channel.provider_maps: if not await can_user_access_channel(db, user, channel.id): if not await _can_access_via_custom_category(db, channel.id, user.id): raise HTTPException(403, "Not in your catalog") pmap = channel.provider_maps[0] prov_res = await db.execute(select(ProviderAccount).where(ProviderAccount.id == pmap.provider_account_id)) provider = prov_res.scalar_one_or_none() vod_session = vod_tracker.start( user_id=str(user.id), username=user.username, channel_id=channel.id, channel_name=channel.name, provider_account_id=pmap.provider_account_id, provider_name=provider.name if provider else f"#{pmap.provider_account_id}", stream_type="series", stream_url=pmap.stream_url, user_max_connections=user.max_connections, ) log_id = await _create_log(user.id, user.username, channel.id, channel.name, "series", sid, client_ip) return await _proxy_vod(pmap.stream_url, request, vod_session=vod_session, log_id=log_id) # Second: episode_id was cached when the app called get_series_info provider_id = _episode_provider_cache.get(sid) if provider_id: prov_result = await db.execute(select(ProviderAccount).where(ProviderAccount.id == provider_id)) provider = prov_result.scalar_one_or_none() if provider: client = XtreamClient(provider.base_url, provider.username, provider.password) url = client.build_series_stream_url(sid, ext) vod_session = vod_tracker.start( user_id=str(user.id), username=user.username, channel_id=0, channel_name=f"Episode {sid}", provider_account_id=provider_id, provider_name=provider.name, stream_type="series", stream_url=url, user_max_connections=user.max_connections, ) log_id = await _create_log(user.id, user.username, 0, f"Episode {sid}", "series", sid, client_ip) return await _proxy_vod(url, request, vod_session=vod_session, log_id=log_id) # Third: check if it's a cached Jellyfin episode (fake_id → direct-play URL) jelly_url = _jellyfin_episode_cache.get(sid) if jelly_url: vod_session = vod_tracker.start( user_id=str(user.id), username=user.username, channel_id=0, channel_name=f"Episode {sid}", provider_account_id=0, provider_name="Jellyfin", stream_type="series", stream_url=jelly_url, user_max_connections=user.max_connections, ) log_id = await _create_log(user.id, user.username, 0, f"Episode {sid}", "series", sid, client_ip) return await _proxy_vod(jelly_url, request, vod_session=vod_session, log_id=log_id) raise HTTPException(404, "Episode not found — open the series detail first to load episode info") async def _proxy_vod( upstream_url: str, request: Request, vod_session: VodSession | None = None, log_id: int | None = None, ) -> StreamingResponse: """Transparent proxy for VOD/series — forwards status, Content-Type, and range headers.""" req_headers = {} if "Range" in request.headers: req_headers["Range"] = request.headers["Range"] timeout = aiohttp.ClientTimeout(connect=10, sock_read=300) http_session = aiohttp.ClientSession(timeout=timeout) if vod_session: vod_session._http_session = http_session try: resp = await http_session.get(upstream_url, headers=req_headers, ssl=False) except Exception as e: await http_session.close() if vod_session: vod_tracker.end(vod_session.session_id) raise HTTPException(502, f"Provider unreachable: {e}") content_type = resp.headers.get("Content-Type", "application/octet-stream") forward_headers: dict[str, str] = {"Cache-Control": "no-cache"} for h in ("Content-Length", "Content-Range", "Accept-Ranges"): if h in resp.headers: forward_headers[h] = resp.headers[h] log_started = time.time() async def generate(): prev_bytes = 0 prev_time = time.time() bytes_sent = 0 try: async for chunk in resp.content.iter_chunked(65536): if vod_session and vod_session.should_stop: break bytes_sent += len(chunk) if vod_session: now = time.time() vod_session.bytes_pumped += len(chunk) vod_session.last_active = now if now - prev_time >= 2.0: bps = (vod_session.bytes_pumped - prev_bytes) / (now - prev_time) vod_session.bps_down = bps vod_session.bps_up = bps prev_bytes = vod_session.bytes_pumped prev_time = now yield chunk finally: resp.release() await http_session.close() if vod_session: vod_session.running = False vod_tracker.end(vod_session.session_id) if log_id: asyncio.create_task(_finish_log(log_id, bytes_sent, log_started)) return StreamingResponse( generate(), status_code=resp.status, media_type=content_type, headers=forward_headers, ) # --------------------------------------------------------------------------- # Internal helpers # --------------------------------------------------------------------------- async def _user_info(user: User, password: str, base: str, db: AsyncSession) -> dict: exp = int(user.expiry_date.timestamp()) if user.expiry_date else 4102444800 # 2100-01-01 return { "user_info": { "username": user.username, "password": password, "message": "KiraStream", "auth": 1, "status": "Active" if user.is_active else "Disabled", "exp_date": str(exp), "is_trial": "0", "active_cons": "0", "created_at": str(int(user.created_at.timestamp())), "max_connections": str(user.max_connections), "allowed_output_formats": ["ts", "m3u8", "rtmp"], }, "server_info": { "url": base, "port": "80", "https_port": "443", "server_protocol": "http", "rtmp_port": "1935", "timezone": "Europe/Madrid", "timestamp_now": int(time.time()), "time_now": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), }, } async def _get_custom_cats(db: AsyncSession, ctype: ContentType, user_id: int) -> list[dict]: allowed_ids = await _get_user_custom_cat_ids(db, user_id) q = ( select(CustomCategory) .where(CustomCategory.type == ctype, CustomCategory.is_visible == True) # noqa: E712 .order_by(CustomCategory.position, CustomCategory.id) ) if allowed_ids is not None: q = q.where(CustomCategory.id.in_(allowed_ids)) result = await db.execute(q) return [ {"category_id": f"{CC_PREFIX}{c.id}", "category_name": c.name, "parent_id": 0} for c in result.scalars().all() ] async def _get_live_categories(user: User, db: AsyncSession) -> list[dict]: cats = await get_user_categories(db, user, ContentType.live) result = [{"category_id": str(c.id), "category_name": c.name, "parent_id": 0} for c in cats] result += await _get_custom_cats(db, ContentType.live, user.id) return result async def _get_categories(user: User, db: AsyncSession, ctype: ContentType) -> list[dict]: cats = await get_user_categories(db, user, ctype) result = [{"category_id": str(c.id), "category_name": c.name, "parent_id": 0} for c in cats] result += await _get_custom_cats(db, ctype, user.id) return result async def _get_live_streams(user: User, db: AsyncSession, category_id: str | None, base: str) -> list[dict]: if category_id and category_id.startswith(CC_PREFIX): cc_id = int(category_id[len(CC_PREFIX):]) channels = await _get_custom_category_channels(db, cc_id) else: cat_id = int(category_id) if category_id else None channels = await get_user_channels(db, user, ContentType.live, cat_id) return [ { "num": i + 1, "name": ch.name, "stream_type": "live", "stream_id": int(ch.stream_id_at_provider) if ch.stream_id_at_provider.isdigit() else ch.id, "stream_icon": ch.tvg_logo or "", "epg_channel_id": ch.tvg_id or "", "added": "0", "category_id": str(ch.category_id) if ch.category_id else "0", "custom_sid": "", "tv_archive": 0, "direct_source": "", "tv_archive_duration": 0, } for i, ch in enumerate(channels) ] async def _get_channels(user: User, db: AsyncSession, ctype: ContentType, category_id: str | None, base: str) -> list[dict]: if category_id and category_id.startswith(CC_PREFIX): cc_id = int(category_id[len(CC_PREFIX):]) channels = await _get_custom_category_channels(db, cc_id) else: cat_id = int(category_id) if category_id else None channels = await get_user_channels(db, user, ctype, cat_id) if ctype == ContentType.movie: return [ { "num": i + 1, "name": ch.name, "stream_type": "movie", "stream_id": int(ch.stream_id_at_provider) if ch.stream_id_at_provider.isdigit() else ch.id, "stream_icon": ch.tvg_logo or "", "added": "0", "category_id": str(ch.category_id) if ch.category_id else "0", "category_ids": [str(ch.category_id) if ch.category_id else "0"], "container_extension": "mkv", } for i, ch in enumerate(channels) ] # Series: must return series_id (not stream_id) so apps can call get_series_info return [ { "num": i + 1, "name": ch.name, "series_id": int(ch.stream_id_at_provider) if ch.stream_id_at_provider.isdigit() else ch.id, "stream_id": int(ch.stream_id_at_provider) if ch.stream_id_at_provider.isdigit() else ch.id, "cover": ch.tvg_logo or "", "stream_icon": ch.tvg_logo or "", "plot": "", "cast": "", "director": "", "genre": "", "release_date": "", "last_modified": "0", "rating": "", "rating_5based": 0, "backdrop_path": [], "youtube_trailer": "", "episode_run_time": "", "added": "0", "category_id": str(ch.category_id) if ch.category_id else "0", "category_ids": [str(ch.category_id) if ch.category_id else "0"], "container_extension": "mkv", } for i, ch in enumerate(channels) ] async def _get_series_info(series_id: str, user: User, db: AsyncSession) -> dict: result = await db.execute( select(Channel) .options(selectinload(Channel.provider_maps)) .where(Channel.stream_id_at_provider == series_id, Channel.type == ContentType.series) ) ch = result.scalar_one_or_none() # Fallback: Jellyfin series have non-numeric stream_id_at_provider so the app # receives ch.id (DB int) as series_id — look it up by primary key. if not ch: try: ch_id = int(series_id) fb = await db.execute( select(Channel) .options(selectinload(Channel.provider_maps)) .where(Channel.id == ch_id, Channel.type == ContentType.series) ) ch = fb.scalar_one_or_none() except (ValueError, TypeError): pass if not ch: return {"info": {}, "episodes": {}} # Check if this channel belongs to a Jellyfin config ji_result = await db.execute( select(JellyfinItem).where(JellyfinItem.channel_id == ch.id) ) ji = ji_result.scalar_one_or_none() if ji: cfg_result = await db.execute( select(JellyfinConfig).where(JellyfinConfig.id == ji.jellyfin_config_id) ) cfg = cfg_result.scalar_one_or_none() if cfg: return await _get_jellyfin_series_info(ch, ji, cfg) if not ch.provider_maps: return {"info": {}, "episodes": {}} pmap = ch.provider_maps[0] prov_result = await db.execute(select(ProviderAccount).where(ProviderAccount.id == pmap.provider_account_id)) provider = prov_result.scalar_one_or_none() if not provider: return {"info": {}, "episodes": {}} client = XtreamClient(provider.base_url, provider.username, provider.password) try: info = await client.get_series_info(series_id) # Cache episode_id → provider_id so series_stream can route the playback request for season_eps in (info.get("episodes") or {}).values(): if isinstance(season_eps, list): for ep in season_eps: ep_id = str(ep.get("id", "")) if ep_id: _episode_provider_cache[ep_id] = provider.id # Ensure cover is populated (some providers return it here, some in the list) if info.get("info") is not None and not info["info"].get("cover"): info["info"]["cover"] = ch.tvg_logo or "" return info except Exception as e: logger.warning(f"Series info fetch failed for series_id={series_id}: {e}") return { "info": {"name": ch.name, "cover": ch.tvg_logo or "", "plot": "", "cast": "", "director": "", "genre": "", "rating": ""}, "episodes": {}, } async def _get_jellyfin_series_info(ch: Channel, ji: JellyfinItem, cfg: JellyfinConfig) -> dict: """Build get_series_info response from Jellyfin seasons/episodes.""" from ..core.jellyfin_client import JellyfinClient as _JClient client = _JClient(cfg.url, cfg.api_key) episodes_by_season: dict[str, list[dict]] = {} try: seasons = await client.get_seasons(ji.jellyfin_item_id) for season in seasons: season_id = season.get("Id", "") season_num = str(season.get("IndexNumber", 1)) eps = await client.get_episodes(ji.jellyfin_item_id, season_id) ep_list = [] for ep_num, ep in enumerate(eps, 1): ep_jf_id = ep.get("Id", "") if not ep_jf_id: continue # Stable numeric ID derived from the Jellyfin episode ID fake_id = abs(hash(ep_jf_id)) % (10 ** 9) stream_url = client.build_stream_url(ep_jf_id) _jellyfin_episode_cache[str(fake_id)] = stream_url ep_list.append({ "id": fake_id, "episode_num": ep.get("IndexNumber", ep_num), "title": ep.get("Name", f"Episodio {ep_num}"), "container_extension": "mkv", "info": { "plot": ep.get("Overview", ""), "duration_secs": 0, "rating": "", "name": ep.get("Name", ""), "air_date": "", }, }) if ep_list: episodes_by_season[season_num] = ep_list except Exception as e: logger.warning(f"Jellyfin seasons fetch failed for {ji.jellyfin_item_id}: {e}") return { "info": { "name": ch.name, "cover": ch.tvg_logo or "", "plot": ji.plot or "", "cast": "", "director": "", "genre": ji.genres or "", "release_date": str(ji.year) if ji.year else "", "rating": ji.rating or "", "backdrop_path": ch.tvg_logo or "", "youtube_trailer": "", "episode_run_time": "", }, "episodes": episodes_by_season, } async def _get_vod_info(vod_id: str, user: User, db: AsyncSession) -> dict: result = await db.execute( select(Channel) .options(selectinload(Channel.provider_maps)) .where(Channel.stream_id_at_provider == vod_id, Channel.type == ContentType.movie) ) ch = result.scalar_one_or_none() # Fallback: Jellyfin movies return ch.id as vod_id (non-numeric stream_id_at_provider) if not ch: try: ch_id = int(vod_id) fb = await db.execute( select(Channel) .options(selectinload(Channel.provider_maps)) .where(Channel.id == ch_id, Channel.type == ContentType.movie) ) ch = fb.scalar_one_or_none() except (ValueError, TypeError): pass if not ch: return {"info": {}, "movie_data": {}} # If this is a Jellyfin movie, return its metadata from JellyfinItem ji_result = await db.execute(select(JellyfinItem).where(JellyfinItem.channel_id == ch.id)) ji = ji_result.scalar_one_or_none() if ji: return { "info": { "name": ch.name, "cover": ch.tvg_logo or "", "plot": ji.plot or "", "cast": "", "director": "", "genre": ji.genres or "", "release_date": str(ji.year) if ji.year else "", "rating": ji.rating or "", "duration_secs": 0, "backdrop_path": ch.tvg_logo or "", "youtube_trailer": "", "tmdb_id": "", }, "movie_data": { "stream_id": str(ch.id), "name": ch.name, "added": "0", "category_id": str(ch.category_id) if ch.category_id else "0", "container_extension": "mkv", "custom_sid": "", "direct_source": "", }, } if ch.provider_maps: pmap = ch.provider_maps[0] prov_result = await db.execute(select(ProviderAccount).where(ProviderAccount.id == pmap.provider_account_id)) provider = prov_result.scalar_one_or_none() if provider: client = XtreamClient(provider.base_url, provider.username, provider.password) try: return await client.get_vod_info(vod_id) except Exception as e: logger.warning(f"VOD info fetch failed for vod_id={vod_id}: {e}") return { "info": {"name": ch.name, "cover": ch.tvg_logo or "", "plot": "", "cast": "", "director": "", "genre": "", "release_date": "", "rating": ""}, "movie_data": {"stream_id": ch.stream_id_at_provider, "container_extension": "mkv"}, }