Files
KiraStream b969b7e5af Mejoras de estabilidad de stream IPTV y nuevo registro de fallos
- restream.py: eliminar -reconnect_streamed de ffmpeg para evitar rebobinados
  al reconectar (el proveedor re-enviaba desde keyframe anterior causando PTS
  backward). El loop de Python gestiona las reconexiones con conexión fresca.
- restream.py: eliminar _drain_queue del loop de reconexión para que el buffer
  de la cola cubra el tiempo de reconexión sin pantalla negra en el player.
- restream.py: añadir buffer_server_bytes/secs/capacity a stats() y método
  _log_stream_event() para persistir eventos de caída y recuperación en BD.
- models/stream_event.py: nuevo modelo StreamEvent para registro de fallos.
- database.py: registrar StreamEvent en init_db.
- api/admin/logs.py: nuevos endpoints GET/DELETE /logs/events.
- Dashboard.tsx: mostrar reserva de buffer del servidor (capacidad + estado).
- Logs.tsx: añadir pestaña "Fallos de stream" con tabla de eventos persistidos.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-19 14:43:31 +00:00

929 lines
36 KiB
Python

"""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"},
}