Files
server_IPTV/server.py
T

1489 lines
54 KiB
Python

#!/usr/bin/env python
# -*- coding: utf-8 -*-
# pylint: disable=W0613, C0116
# type: ignore[union-attr]
from pymongo import MongoClient
from aiohttp import web
import socketio
import os, glob
import time
import json
import ssl
import sys
import uuid
import logging
from aiohttp_session import setup, get_session, new_session
from aiohttp_session.cookie_storage import EncryptedCookieStorage
from tornado.template import Loader
from bson.json_util import dumps
from urllib.parse import unquote
import aiohttp_cors
import aiofiles
import urllib.parse
import requests
from PIL import Image
import io, re
import xml.etree.ElementTree as ET
from bs4 import BeautifulSoup
import pathlib
from datetime import datetime
from bson import ObjectId
import urllib.request
from pymongo import MongoClient
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
handlers=[logging.StreamHandler(sys.stdout)],
)
logger = logging.getLogger(__name__)
# Connexion Database
MONGO_URI = os.getenv("MONGO_URI", "mongodb://root:example@mongodb:27017/?authSource=admin")
def connect_mongo():
for attempt in range(1, 6):
try:
client = MongoClient(MONGO_URI, serverSelectionTimeoutMS=5000, connectTimeoutMS=5000)
client.admin.command("ping")
logger.info("Connected to MongoDB")
return client
except Exception as exc: # retry if Mongo is not ready yet
logger.warning("Mongo connection failed (attempt %s/5): %s", attempt, exc)
time.sleep(5)
raise RuntimeError("MongoDB is unreachable after retries")
client = connect_mongo()
# Connect to the database and collection
db_iptv = client["db_iptv"]
channels = db_iptv['channels']
channels.create_index("id_playlist")
playlists = db_iptv['playlists']
logos = db_iptv['logos']
recents = db_iptv['recents']
settings_collection = db_iptv['settings']
devices_db = db_iptv['devices']
devices = []
active_xtream_imports = {}
# Configure Socket IO
sio = socketio.AsyncServer(logger=True, engineio_logger=True, cors_allowed_origins='*')
app = web.Application()
sio.attach(app)
# Variables
path = os.path.dirname(os.path.realpath(__file__))+'/'
ssl_cert_path = '/home/letsencrypt/live/iptv.mrk.ovh/cert.pem'
ssl_key_path = '/home/letsencrypt/live/iptv.mrk.ovh/privkey.pem'
UPLOAD_PLAYLIST_DIR = "playlists"
DOWNLOAD_EPG_DIR = "epg"
LOGO_DIR = "/home/logos"
M3U_LOGO_SOURCE = '/home/playlists_download/tv_channels_2a7e1ba958_plus.m3u'
URL_SERVER = 'https://iptv.mrk.ovh/'
DEVICE_SNAPSHOT_DIR = pathlib.Path(__file__).parent / 'static' / 'device_snapshots'
DEVICE_SNAPSHOT_DIR.mkdir(parents=True, exist_ok=True)
# Set up Tornado template loader
loader = Loader('')
# SSL context setup
ssl_context = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
ssl_context.load_cert_chain(certfile=ssl_cert_path, keyfile=ssl_key_path)
# Generic functions
def findInfo(tag, s):
end = '"'
if (s.find(tag)>-1):
t = s[s.find(tag)+len(tag)+2:s.rfind(end)]
return t[0:t.find(end)]
else:
return ""
# Helper function to convert MongoDB documents to JSON-serializable format
def serialize_doc(doc):
doc["_id"] = str(doc["_id"]) # Convert ObjectId to string
return doc
def clean_name_channel(name_channel):
new_name_channel = re.sub(r'\b\w{2}\b', '', name_channel)
new_name_channel = re.sub(r'\s+', ' ', new_name_channel).strip()
new_name_channel = new_name_channel.replace('- ','')
new_name_channel = new_name_channel.replace('| ','')
new_name_channel = new_name_channel.replace(' UHD','')
new_name_channel = new_name_channel.replace(' FHD','')
new_name_channel = new_name_channel.replace(' -','')
new_name_channel = new_name_channel.replace(' /','')
new_name_channel = new_name_channel.replace('','')
pattern = r".{2}\|"
new_name_channel = re.sub(pattern, "", new_name_channel)
new_name_channel = re.sub(r"[\(\{\[].*?[\)\}\]]", "", new_name_channel)
new_name_channel = new_name_channel.strip()
return new_name_channel
def get_playlist_collection(id_playlist):
return db_iptv['channels_' + id_playlist]
async def emit_analysis_progress(payload, sid_client=None):
if not sid_client:
return
data = {'typeCommand': 'progress_analysis', 'payload': payload}
await sio.emit('communication', data, room=sid_client, callback=messageReceived)
def get_existing_playlist(id_playlist):
return playlists.find_one({"id_playlist": id_playlist}) or {}
def find_existing_xtream_playlist(server_url, username, password):
return playlists.find_one({
"playlist_type": "xtream_codes",
"xtream_server": normalize_xtream_server(server_url),
"xtream_username": (username or '').strip(),
"xtream_password": (password or '').strip(),
}) or {}
def build_playlist_document(id_playlist, title_playlist, date_end_playlist, selected_color, extra_fields=None):
existing_playlist = get_existing_playlist(id_playlist)
playlist_doc = {
'id_playlist': id_playlist,
'title_playlist': title_playlist,
'date_end_playlist': date_end_playlist,
'selectedColor': selected_color,
'activation': existing_playlist.get('activation', 1),
}
if extra_fields:
playlist_doc.update(extra_fields)
return playlist_doc
def save_playlist_document(playlist_doc):
playlists.update_one(
{"id_playlist": playlist_doc["id_playlist"]},
{"$set": playlist_doc},
upsert=True
)
def delete_playlist_resources(id_playlist):
db_iptv.drop_collection("channels_" + id_playlist)
result = playlists.delete_one({"id_playlist": id_playlist})
file_path = os.path.join(UPLOAD_PLAYLIST_DIR, id_playlist + '.m3u')
if os.path.exists(file_path):
os.remove(file_path)
print(f"File {file_path} deleted.")
return result.deleted_count
def bulk_insert_channels(collection_channels, new_channels, num_channel):
if not new_channels:
return
collection_channels.insert_many(new_channels)
print("Save channels ", num_channel, flush=True)
def normalize_xtream_server(server_url):
server_url = (server_url or '').strip()
if not server_url:
return ''
if not re.match(r'^https?://', server_url, re.IGNORECASE):
server_url = 'http://' + server_url
return server_url.rstrip('/')
def build_xtream_playlist_url(server_url, username, password):
base = normalize_xtream_server(server_url)
if not base or not username or not password:
return ''
return f"{base}/get.php?username={username}&password={password}&type=m3u_plus&output=ts"
def build_xtream_api_url(server_url, username, password, action, **extra_params):
base = normalize_xtream_server(server_url)
query = {
'username': username,
'password': password,
'action': action,
}
query.update(extra_params)
return f"{base}/player_api.php?{urllib.parse.urlencode(query)}"
def fetch_xtream_json(server_url, username, password, action, **extra_params):
url = build_xtream_api_url(server_url, username, password, action, **extra_params)
response = requests.get(url, timeout=60)
response.raise_for_status()
return response.json()
def parse_xtream_category_map(items, id_key='category_id', name_key='category_name'):
category_map = {}
for item in items or []:
category_id = str(item.get(id_key, '')).strip()
if category_id:
category_map[category_id] = item.get(name_key, '') or ''
return category_map
def make_xtream_live_channel(id_playlist, item, category_map, username, password, server_url):
stream_id = item.get('stream_id')
if not stream_id:
return None
container_extension = item.get('container_extension') or 'ts'
return {
"id_playlist": id_playlist,
"id_channel": str(uuid.uuid4()),
"Name": item.get('name') or f"Live {stream_id}",
"Language": item.get('stream_type') or 'live',
"Topic": category_map.get(str(item.get('category_id', '')), ''),
"urlLogo": item.get('stream_icon', '') or '',
"url": f"{normalize_xtream_server(server_url)}/live/{username}/{password}/{stream_id}.{container_extension}",
"media_type": "live",
"stream_type": item.get('stream_type', 'live'),
"stream_id": stream_id,
"category_id": item.get('category_id'),
"epg_channel_id": item.get('epg_channel_id', ''),
}
def make_xtream_vod_channel(id_playlist, item, category_map, username, password, server_url):
stream_id = item.get('stream_id')
if not stream_id:
return None
container_extension = item.get('container_extension') or 'mp4'
return {
"id_playlist": id_playlist,
"id_channel": str(uuid.uuid4()),
"Name": item.get('name') or f"Movie {stream_id}",
"Language": item.get('stream_type') or 'movie',
"Topic": category_map.get(str(item.get('category_id', '')), ''),
"urlLogo": item.get('stream_icon', '') or item.get('cover', '') or '',
"url": f"{normalize_xtream_server(server_url)}/movie/{username}/{password}/{stream_id}.{container_extension}",
"media_type": "movie",
"stream_type": item.get('stream_type', 'movie'),
"stream_id": stream_id,
"category_id": item.get('category_id'),
"rating": item.get('rating', ''),
}
def flatten_xtream_episodes(episodes_payload):
episodes = []
if isinstance(episodes_payload, dict):
for season_episodes in episodes_payload.values():
if isinstance(season_episodes, list):
episodes.extend(season_episodes)
elif isinstance(episodes_payload, list):
episodes.extend(episodes_payload)
return episodes
def make_xtream_series_channels(id_playlist, series_item, series_info, category_map, username, password, server_url):
topic = category_map.get(str(series_item.get('category_id', '')), '')
series_name = series_item.get('name') or f"Series {series_item.get('series_id')}"
series_logo = series_item.get('cover') or series_item.get('cover_big') or series_item.get('stream_icon') or ''
channels_to_insert = []
for episode in flatten_xtream_episodes(series_info.get('episodes')):
episode_id = episode.get('id') or episode.get('episode_id')
if not episode_id:
continue
container_extension = episode.get('container_extension') or 'mp4'
episode_title = episode.get('title') or episode.get('name') or f"Episode {episode_id}"
episode_num = episode.get('episode_num')
season_num = episode.get('season')
episode_label = episode_title
try:
if season_num is not None and episode_num is not None:
episode_label = f"{series_name} - S{int(season_num):02d}E{int(episode_num):02d} - {episode_title}"
else:
episode_label = f"{series_name} - {episode_title}"
except (TypeError, ValueError):
episode_label = f"{series_name} - {episode_title}"
channels_to_insert.append({
"id_playlist": id_playlist,
"id_channel": str(uuid.uuid4()),
"Name": episode_label,
"Language": "series",
"Topic": topic,
"urlLogo": episode.get('info', {}).get('movie_image', '') or series_logo,
"url": f"{normalize_xtream_server(server_url)}/series/{username}/{password}/{episode_id}.{container_extension}",
"media_type": "series",
"stream_type": "series",
"series_id": series_item.get('series_id'),
"episode_id": episode_id,
"category_id": series_item.get('category_id'),
"series_name": series_name,
"season": season_num,
"episode_num": episode_num,
})
if channels_to_insert:
return channels_to_insert
return [{
"id_playlist": id_playlist,
"id_channel": str(uuid.uuid4()),
"Name": series_name,
"Language": "series",
"Topic": topic,
"urlLogo": series_logo,
"url": "",
"media_type": "series",
"stream_type": "series",
"series_id": series_item.get('series_id'),
"category_id": series_item.get('category_id'),
"series_name": series_name,
}]
def save_image_from_url(url, save_path):
try:
# Send a GET request to the URL
response = requests.get(url)
response.raise_for_status() # Raise an error for unsuccessful status codes
# Write the content of the response to a file
with open(save_path, 'wb') as file:
file.write(response.content)
print(f"Image saved to {save_path}")
except requests.exceptions.RequestException as e:
print(f"Error downloading the image: {e}")
def load_epg_root():
"""
Load EPG XML either from the URL stored in settings or from a local fallback file.
Returns an ElementTree root or None if nothing can be loaded.
"""
xml_content = None
# Try remote URL first (if configured)
doc = settings_collection.find_one({'key': 'epg_url'})
epg_url = doc.get('value') if doc else None
if epg_url:
try:
resp = requests.get(epg_url, timeout=10)
if resp.status_code == 200 and resp.content:
xml_content = resp.content
except requests.RequestException as exc:
logger.warning("Unable to fetch remote EPG (%s): %s", epg_url, exc)
# Fallback to local file
if xml_content is None:
local_path = os.path.join(DOWNLOAD_EPG_DIR, 'epg.xml')
if os.path.exists(local_path):
with open(local_path, 'rb') as f:
xml_content = f.read()
if not xml_content:
return None
try:
return ET.fromstring(xml_content)
except ET.ParseError as exc:
logger.error("Failed to parse EPG XML: %s", exc)
return None
# AIOHTTP functions
# INDEX
async def index(request):
template = loader.load('static/index.html')
return web.Response(text=template.generate().decode('utf-8'), content_type='text/html')
# Access to playlist file from the web
async def playlist_download(request):
data = dict(request.query)
file_path = '/home/playlists_download/'+data['filename']
if os.path.exists(file_path):
return web.FileResponse(file_path)
else:
raise web.HTTPNotFound(text='File not found')
async def init_db(request):
db_iptv.drop_collection("channels")
db_iptv.drop_collection("playlists")
return web.Response(text=str(res),content_type='text/html')
async def analyse_m3u_file(request):
NUM_MAX_CHANNEL = 300000
data = await request.post()
id_playlist = data['id_playlist']
title_playlist = data['title']
url_playlist = data['url']
date_end_playlist = data['date_end']
filename = data['filename']
sidClient = data['sidClient']
selectedColor = data['selectedColor']
playlist_comment = data.get('comment', '').strip()
new_path_filename = "./" + UPLOAD_PLAYLIST_DIR + '/' + id_playlist + '.m3u'
old_path_filename = "./" + UPLOAD_PLAYLIST_DIR + '/' +filename
os.rename(old_path_filename,new_path_filename)
collection_channels = get_playlist_collection(id_playlist)
collection_channels.drop()
new_playlist = build_playlist_document(
id_playlist,
title_playlist,
date_end_playlist,
selectedColor,
{
'playlist_type': 'm3u',
'url_playlist': url_playlist,
'comment': playlist_comment,
'xtream_server': '',
'xtream_username': '',
'xtream_password': '',
}
)
save_playlist_document(new_playlist)
# Process the m3u file
new_channels =[]
f = open(new_path_filename, "r", encoding="utf8")
Lines = f.readlines()
print('ouverture fichier', flush=True)
channel = {}
separator = '--------'
topic = ''
num_channel = 0
TypeLangChannel = ['EN','FR','US','CA', 'UK' ]
total_lines = len(Lines)
number_line = 0
for line in Lines:
number_line += 1
if separator in line:
topic = line[10:]
topic = topic.replace(separator,'')
# PROCESS URL LINE
if len(channel)>0 and line[0:4]=="http":
channel["url"] = line.strip('\n')
new_channels.append(channel)
#print(channel, flush=True)
if num_channel % 1000 == 0:
collection_channels.insert_many(new_channels)
new_channels = []
print("Save channels ", num_channel, flush=True)
print('%d / %d' % (number_line, total_lines), flush=True)
await emit_analysis_progress({'percentage': number_line / total_lines}, sidClient)
# PROCESS CHANNEL LINE
if line[0:4] == "#EXT":
channel = {}
# ID CHANNEL
id_channel = str(uuid.uuid4())
# DEFINE LANGUAGE
LanguageChannel = next((lang for lang in TypeLangChannel if lang in line), '')
# deternine url logo
url_logo = findInfo('tvg-logo', line)
channel_exist = False
# case separator is |
name_channel = line[line.rfind(',')+1:].strip('\n')
if len(name_channel)>0:
channel = {"id_playlist":id_playlist, "id_channel":id_channel, "Name":name_channel, "Language":LanguageChannel, "Topic":topic, "urlLogo" : url_logo}
num_channel += 1
channel_exist = True
# case separator is ,
if channel_exist == False:
name_channel = line[line.rfind('|')+2:].strip('\n')
if len(name_channel) > 0:
channel = {"id_playlist":id_playlist, "id_channel":id_channel, "Name":name_channel, "Language":LanguageChannel, "Topic":topic, "urlLogo" : url_logo}
channel_exist = True
num_channel += 1
bulk_insert_channels(collection_channels, new_channels, num_channel)
print("Save final channels ", num_channel, flush=True)
number_channels = collection_channels.count_documents({"id_playlist":id_playlist})
playlists.update_many({'id_playlist':id_playlist}, { "$set": { "num_channels": number_channels } } )
result = {'num_channel': number_channels}
return web.json_response(result)
async def analyse_m3u_file2(request):
NUM_MAX_CHANNEL = 300000
data = await request.post()
id_playlist = data['id_playlist']
title_playlist = data['title']
url_playlist = data['url']
date_end_playlist = data['date_end']
filename = data['filename']
sidClient = data['sidClient']
playlist_comment = data.get('comment', '').strip()
new_path_filename = "./" + UPLOAD_PLAYLIST_DIR + '/' + id_playlist + '.m3u'
old_path_filename = "./" + UPLOAD_PLAYLIST_DIR + '/' +filename
os.rename(old_path_filename,new_path_filename)
collection_channels = get_playlist_collection(id_playlist)
collection_channels.drop()
new_playlist = build_playlist_document(
id_playlist,
title_playlist,
date_end_playlist,
'',
{
'playlist_type': 'm3u',
'url_playlist': url_playlist,
'comment': playlist_comment,
'xtream_server': '',
'xtream_username': '',
'xtream_password': '',
}
)
save_playlist_document(new_playlist)
# Process the m3u file
TypeLangChannel = ['EN','FR','US','CA', 'UK' ]
new_channels =[]
f = open(new_path_filename, "r", encoding="utf8")
Lines = f.readlines()
print('ouverture fichier', flush=True)
channel = {}
separator = '--------'
topic = ''
num_channel = 0
total_lines = len(Lines)
number_line = 0
for line in Lines:
number_line += 1
if separator in line:
topic = line[10:]
topic = topic.replace(separator,'')
if len(channel)>0 and line[0:4]=="http":
channel["url"] = line.strip('\n')
if len(LanguageChannel) > 0:
new_channels.append(channel)
print(num_channel, flush=True)
if num_channel % 1000 == 0:
bulk_insert_channels(collection_channels, new_channels, num_channel)
new_channels = []
print('%d / %d' % (number_line, total_lines), flush=True)
if line[0:4] == "#EXT":
channel = {}
# DEFINE CHANNEL NAME
name_channel = line[line.rfind(',')+1:].strip('\n')
id_channel = str(uuid.uuid4())
LanguageChannel =''
for langChannel in TypeLangChannel:
if langChannel in line:
LanguageChannel = langChannel
# deternine url logo
url_logo = findInfo('tvg-logo', line)
#if len(url_logo) == 0:
# cleaned_channel_name = clean_name_channel(name_channel)
# logo_data = logos.find({ 'logo_name': { '$regex': cleaned_channel_name, '$options': "i" } })
# for l in logo_data:
# url_logo = l['logo_url']
# break
channel_exist = False
if len(name_channel)>0:
channel = {"id_playlist":id_playlist, "id_channel":id_channel, "Name":name_channel, "Language":LanguageChannel, "Topic":topic, "urlLogo" : url_logo}
num_channel += 1
channel_exist = True
if channel_exist == False:
name_channel = line[line.rfind('|')+2:].strip('\n')
if len(name_channel) > 0:
channel = {"id_playlist":id_playlist, "id_channel":id_channel, "Name":name_channel, "Language":LanguageChannel, "Topic":topic, "urlLogo" : url_logo}
bulk_insert_channels(collection_channels, new_channels, num_channel)
print("Save final channels ", num_channel, flush=True)
number_channels = collection_channels.count_documents({"id_playlist":id_playlist})
playlists.update_many({'id_playlist':id_playlist}, { "$set": { "num_channels": number_channels } } )
result = {'num_channel': number_channels}
return web.json_response(result)
async def import_xtream_codes(request):
data = await request.post()
id_playlist = data.get('id_playlist', '').strip()
title_playlist = data['title'].strip()
date_end_playlist = data['date_end']
selectedColor = data.get('selectedColor', '')
playlist_comment = data.get('comment', '').strip()
sidClient = data.get('sidClient', '')
xtream_server = normalize_xtream_server(data.get('xtream_server', ''))
xtream_username = data.get('xtream_username', '').strip()
xtream_password = data.get('xtream_password', '').strip()
if not xtream_server or not xtream_username or not xtream_password:
return web.json_response({'error': 'Xtream server, username and password are required.'}, status=400)
if not id_playlist:
existing_playlist = find_existing_xtream_playlist(xtream_server, xtream_username, xtream_password)
if existing_playlist:
id_playlist = existing_playlist['id_playlist']
logger.info("Reusing existing Xtream playlist id=%s for %s", id_playlist, xtream_server)
else:
id_playlist = str(uuid.uuid4())
import_key = sidClient.strip() or id_playlist
if active_xtream_imports.get(import_key):
logger.warning("Xtream import already running for key=%s", import_key)
running_import = active_xtream_imports.get(import_key) or {}
return web.json_response({
'status': 'already_running',
'message': 'An Xtream import is already running for this client. Please wait for it to finish.',
'id_playlist': running_import.get('id_playlist', id_playlist),
}, status=202)
active_xtream_imports[import_key] = {
'id_playlist': id_playlist,
'started_at': time.time(),
}
collection_channels = get_playlist_collection(id_playlist)
collection_channels.drop()
playlist_doc = build_playlist_document(
id_playlist,
title_playlist,
date_end_playlist,
selectedColor,
{
'playlist_type': 'xtream_codes',
'url_playlist': build_xtream_playlist_url(xtream_server, xtream_username, xtream_password),
'comment': playlist_comment,
'xtream_server': xtream_server,
'xtream_username': xtream_username,
'xtream_password': xtream_password,
}
)
save_playlist_document(playlist_doc)
try:
await emit_analysis_progress({'percentage': 0.02}, sidClient)
live_categories = parse_xtream_category_map(
fetch_xtream_json(xtream_server, xtream_username, xtream_password, 'get_live_categories')
)
vod_categories = parse_xtream_category_map(
fetch_xtream_json(xtream_server, xtream_username, xtream_password, 'get_vod_categories')
)
series_categories = parse_xtream_category_map(
fetch_xtream_json(xtream_server, xtream_username, xtream_password, 'get_series_categories')
)
live_streams = fetch_xtream_json(xtream_server, xtream_username, xtream_password, 'get_live_streams')
vod_streams = fetch_xtream_json(xtream_server, xtream_username, xtream_password, 'get_vod_streams')
series_list = fetch_xtream_json(xtream_server, xtream_username, xtream_password, 'get_series')
total_steps = max(1, len(live_streams) + len(vod_streams) + len(series_list))
processed_steps = 0
pending_channels = []
inserted_channels = 0
for item in live_streams:
channel = make_xtream_live_channel(id_playlist, item, live_categories, xtream_username, xtream_password, xtream_server)
if channel:
pending_channels.append(channel)
inserted_channels += 1
processed_steps += 1
if len(pending_channels) >= 1000:
bulk_insert_channels(collection_channels, pending_channels, inserted_channels)
pending_channels = []
if processed_steps % 250 == 0:
await emit_analysis_progress({'percentage': processed_steps / total_steps}, sidClient)
for item in vod_streams:
channel = make_xtream_vod_channel(id_playlist, item, vod_categories, xtream_username, xtream_password, xtream_server)
if channel:
pending_channels.append(channel)
inserted_channels += 1
processed_steps += 1
if len(pending_channels) >= 1000:
bulk_insert_channels(collection_channels, pending_channels, inserted_channels)
pending_channels = []
if processed_steps % 250 == 0:
await emit_analysis_progress({'percentage': processed_steps / total_steps}, sidClient)
for item in series_list:
series_channels = make_xtream_series_channels(
id_playlist,
item,
{},
series_categories,
xtream_username,
xtream_password,
xtream_server
)
pending_channels.extend(series_channels)
inserted_channels += len(series_channels)
processed_steps += 1
if len(pending_channels) >= 1000:
bulk_insert_channels(collection_channels, pending_channels, inserted_channels)
pending_channels = []
await emit_analysis_progress({'percentage': processed_steps / total_steps}, sidClient)
bulk_insert_channels(collection_channels, pending_channels, inserted_channels)
number_channels = collection_channels.count_documents({"id_playlist": id_playlist})
playlists.update_many({'id_playlist': id_playlist}, {"$set": {"num_channels": number_channels}})
await emit_analysis_progress({'percentage': 1}, sidClient)
return web.json_response({'num_channel': number_channels, 'id_playlist': id_playlist})
except requests.RequestException as exc:
logger.exception("Xtream import failed")
return web.json_response({'error': f'Xtream request failed: {exc}'}, status=502)
except Exception as exc:
logger.exception("Xtream import failed")
return web.json_response({'error': f'Xtream import failed: {exc}'}, status=500)
finally:
active_xtream_imports.pop(import_key, None)
async def handle_upload(request):
reader = await request.multipart()
id_playlist = None
filename = None
filepath = None
# Handle both form field and file in one pass
async for part in reader:
if part.name == "id_playlist_to_action":
id_playlist_data = await part.text()
if len(id_playlist_data) == 0:
print('No playlist to update', flush=True)
id_playlist = str(uuid.uuid4())
else:
id_playlist = id_playlist_data
print('playlist to update :', id_playlist, flush=True)
elif part.name == "file":
filename = part.filename
if not filename:
return web.Response(text="No file provided.", status=400)
filepath = os.path.join(UPLOAD_PLAYLIST_DIR, filename)
with open(filepath, "wb") as f:
while chunk := await part.read_chunk():
f.write(chunk)
if filepath:
return web.json_response({'status': 'ok', 'id_playlist': id_playlist})
return web.Response(text="No file field found in the request.", status=400)
async def upload_snapshot(request):
reader = await request.multipart()
id_device = None
file_bytes = bytearray()
async for part in reader:
if part.name == 'id_device':
id_device = (await part.text()).strip()
elif part.name == 'file':
while True:
chunk = await part.read_chunk()
if not chunk:
break
file_bytes.extend(chunk)
if not id_device or not file_bytes:
return web.json_response({'status': 'error', 'message': 'missing id_device or file'}, status=400)
destination = DEVICE_SNAPSHOT_DIR / f'{id_device}.jpg'
with open(destination, 'wb') as output_file:
output_file.write(file_bytes)
snapshot_relpath = f'/static/device_snapshots/{id_device}.jpg'
devices_db.update_one(
{'id_device': id_device},
{'$set': {'snapshot_url': snapshot_relpath, 'snapshot_updated_at': time.time()}},
upsert=True
)
return web.json_response({'status': 'ok', 'snapshot_url': snapshot_relpath})
# COUNT CHANNELS
async def count_channels(request):
data = dict(request.query)
id_playlist = data['id_playlist']
print('id_playlist : ', id_playlist, flush=True)
collection_channels = db_iptv['channels_'+id_playlist]
number_channels = collection_channels.count_documents({"id_playlist":id_playlist})
print('number_channels : ', number_channels, flush=True)
return web.json_response({'number_channels': number_channels})
# DELETE PLAYLIST
async def delete_playlist(request):
query_params = request.query
id_playlist = query_params.get("id_playlist")
nb_playlist_deleted = delete_playlist_resources(id_playlist)
return web.json_response({
'nb_playlist_deleted': nb_playlist_deleted
})
async def cleanup_xtream_duplicates(request):
duplicate_groups = {}
for playlist in playlists.find({"playlist_type": "xtream_codes"}):
xtream_server = normalize_xtream_server(playlist.get("xtream_server", ""))
xtream_username = (playlist.get("xtream_username", "") or "").strip()
xtream_password = (playlist.get("xtream_password", "") or "").strip()
if not xtream_server or not xtream_username or not xtream_password:
continue
duplicate_key = (xtream_server, xtream_username, xtream_password)
duplicate_groups.setdefault(duplicate_key, []).append(playlist)
deleted_ids = []
kept_ids = []
for playlists_group in duplicate_groups.values():
if len(playlists_group) <= 1:
continue
playlists_group.sort(
key=lambda playlist: (
int(playlist.get("num_channels", 0) or 0),
int(playlist.get("activation", 0) or 0),
str(playlist.get("_id", "")),
),
reverse=True,
)
kept_playlist = playlists_group[0]
kept_ids.append(kept_playlist["id_playlist"])
for duplicate_playlist in playlists_group[1:]:
duplicate_id = duplicate_playlist["id_playlist"]
if duplicate_id == kept_playlist["id_playlist"]:
continue
if delete_playlist_resources(duplicate_id):
deleted_ids.append(duplicate_id)
return web.json_response({
"status": "ok",
"deleted_count": len(deleted_ids),
"deleted_ids": deleted_ids,
"kept_ids": kept_ids,
})
async def update_playlist_comment(request):
data = await request.post()
id_playlist = (data.get('id_playlist') or '').strip()
comment = (data.get('comment') or '').strip()
if not id_playlist:
return web.json_response({'error': 'id_playlist is required'}, status=400)
result = playlists.update_one(
{'id_playlist': id_playlist},
{'$set': {'comment': comment}}
)
if result.matched_count == 0:
return web.json_response({'error': 'playlist not found'}, status=404)
return web.json_response({'status': 'ok', 'id_playlist': id_playlist, 'comment': comment})
# Get List Playlists
async def get_list_playlists(request):
active_count = playlists.count_documents({})
#for p in playlists.find():
#id_playlist = p['id_playlist']
#number_channels = channels.count_documents({"id_playlist":id_playlist})
#playlists.update_many({'id_playlist':id_playlist}, { "$set": { "num_channels": number_channels } } )
#print("Document found:", p, flush=True)
#print("Number of channels:", number_channels, flush=True)
print(str(active_count), flush=True)
if active_count == 0:
return web.json_response({'status':'no data'})
else:
return web.json_response(dumps(playlists.find()))
async def get_source(request):
with open('static/IPTV.py') as f:
return web.Response(text=f.read(), content_type='text/plain')
# List Channels
async def list_channel(request):
MAX_RECORDS = 30
data = await request.post()
channel_to_search = unquote(data['keyword'])
current_page = int(data['num_page'])
logger.info("Search keyword='%s' page=%s", channel_to_search, current_page)
regex_pattern = ".*" + re.escape(channel_to_search) + ".*"
# First pass: compute totals without pulling all docs into memory
total_records = 0
active_playlists = list(playlists.find({"activation": 1}))
per_playlist_counts = []
for playlist in active_playlists:
collection_channels = db_iptv['channels_' + playlist['id_playlist']]
count = collection_channels.count_documents({
'id_playlist': playlist['id_playlist'],
'Name': {'$regex': regex_pattern, '$options': "i"}
})
per_playlist_counts.append((playlist, count))
total_records += count
if total_records == 0:
return web.json_response({'num_total_page': 0, 'current_page': 1, 'data': []})
num_total_page = max(1, (total_records + MAX_RECORDS - 1) // MAX_RECORDS)
current_page = max(1, min(current_page, num_total_page))
begin = (current_page - 1) * MAX_RECORDS
remaining = MAX_RECORDS
# Second pass: fetch only the page slice
serialized_data = []
offset = 0
for playlist, count in per_playlist_counts:
if remaining == 0:
break
if offset + count <= begin:
offset += count
continue
local_skip = max(0, begin - offset)
fetch_limit = min(remaining, count - local_skip)
collection_channels = db_iptv['channels_' + playlist['id_playlist']]
cursor = collection_channels.find({
'id_playlist': playlist['id_playlist'],
'Name': {'$regex': regex_pattern, '$options': "i"}
}).skip(local_skip).limit(fetch_limit)
serialized_data.extend(serialize_doc(doc) for doc in cursor)
remaining = MAX_RECORDS - len(serialized_data)
offset += count
return web.json_response({'num_total_page': num_total_page, 'current_page': current_page, 'data': serialized_data})
def get_image_urls(keyword, max_results=10):
search_url = f"https://duckduckgo.com/?q={keyword}&iax=images&ia=images"
headers = {
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36"
}
# Send the search request
response = requests.get(search_url, headers=headers)
soup = BeautifulSoup(response.text, 'html.parser')
# Find image containers
image_urls = []
for img in soup.find_all('img', {'class': 'tile--img__img'}):
src = img.get('src')
if src and src.startswith('http'):
image_urls.append(src)
if len(image_urls) >= max_results:
break
return image_urls
async def get_logo(request):
data = dict(request.query)
id_channel = data['id_channel']
id_playlist = data['id_playlist']
logo_channel_file_path = LOGO_DIR+'/'+id_channel+'.jpg'
collection_channels = db_iptv['channels_'+id_playlist]
d = collection_channels.find_one({ 'id_channel':id_channel})
url_logo = d['urlLogo']
name_channel = d['Name']
cleaned_channel_name = clean_name_channel(name_channel)
if not os.path.exists(logo_channel_file_path):
print(d['Name']+' :::: ' + id_channel + ' :::: ' + url_logo)
if len(url_logo) == 0:
print('Recherche url pour : ' + cleaned_channel_name, flush=True)
#words = cleaned_channel_name.split()
first_two_words = " ".join(cleaned_channel_name.split()[:3])
regex_pattern = ".*" + ".*".join(first_two_words.split()) + ".*"
#regex_pattern = ".*" + ".*".join(words) + ".*"
l = logos.find_one({ 'logo_name': { '$regex': regex_pattern, '$options': "i" } })
if l :
url_logo = l['logo_url']
print('url trouvée : ' + url_logo, flush=True)
save_image_from_url(url_logo, logo_channel_file_path)
query = { 'id_channel':id_channel}
update = { "$set": { "urlLogo": url_logo } }
collection_channels.update_one(query, update)
else:
print('url identifiee et sauvegarde du fichier : ' + url_logo, flush=True)
if len(url_logo) > 0:
save_image_from_url(url_logo, logo_channel_file_path)
image_data = None
if not os.path.exists(logo_channel_file_path):
with open(LOGO_DIR+'/none.jpg', "rb") as file:
image_data = file.read()
else:
with open(logo_channel_file_path, "rb") as file:
image_data = file.read()
print('image found '+logo_channel_file_path, flush=True)
# Return the image as a response
return web.Response(body=image_data, content_type="image/jpeg")
async def proxy_image(request):
url = request.query.get("url")
if not url:
return web.Response(status=400, text="Missing URL")
async with aiohttp.ClientSession() as session:
async with session.get(url) as resp:
if resp.status == 200:
headers = {"Content-Type": resp.headers["Content-Type"]}
return web.Response(body=await resp.read(), headers=headers)
else:
return web.Response(status=resp.status)
async def list_logo(request):
result = logos.find({}).limit(1000)
search_string = "harry potter"
words = search_string.split()
regex_pattern = ".*" + ".*".join(words) + ".*"
logo_data = logos.find({ 'logo_name': { '$regex': regex_pattern, '$options': "i" } })
html = ''
for r in logo_data:
html += '<p>' + r['logo_name'] + ' :::: ' + r['original_name'] + ' :::: ' + r['logo_url']
return web.Response(text=html,content_type='text/html')
async def list_devices(request):
json_devices = json.dumps(devices)
#return web.json_response(json_devices)
return web.json_response(dumps(devices_db.find({}).limit(50)))
async def clean_recents(request):
recents.delete_many({})
return web.json_response({'status':'ok'})
async def get_recents(request):
# Step 1: Get the IDs of the 50 most recent documents
top_50 = list(recents.find({}, {"_id": 1}).sort("datetime", -1).limit(50))
top_50_ids = [doc["_id"] for doc in top_50]
# Step 2: Delete all documents that are NOT in the top 50
result = recents.delete_many({
"_id": {"$nin": top_50_ids}
})
return web.json_response(dumps(recents.find({}).sort("datetime", -1)))
async def most_recents(request):
pipeline = [
{
"$sort": {"datetime": -1}
},
{
"$limit": 50
},
{
"$group": {
"_id": "$id_channel",
"id_channel": {"$first": "$id_channel"},
"id_playlist": {"$first": "$id_playlist"},
"name": {"$first": "$name_channel"},
"url": {"$first": "$url_channel"},
"datetime": {"$first": "$datetime"},
"count": {"$sum": 1}
}
},
{
"$sort": {"count": -1, "datetime": -1, "name": 1}
}
]
# Run the aggregation
results = recents.aggregate(pipeline)
return web.json_response(dumps(results))
async def activation_playlist(request):
data = dict(request.query)
id_playlist = data['id_playlist']
activation = int(data['activation'])
print('activation : ', activation, flush=True)
print(id_playlist,activation, flush=True )
playlists.update_many({'id_playlist':id_playlist}, { "$set": { "activation": activation } } )
return web.json_response({'id_playlist':id_playlist,'status':'ok'})
async def delete_device(request):
data = dict(request.query)
id_device = data.get('id_device') # Use .get() to avoid KeyError
print(" ID device to delete : ",id_device,flush=True)
if not id_device:
return web.json_response({'error': 'id_device is required'}, status=400)
result = devices_db.delete_many({"id_device": id_device})
return web.json_response({'deleted_count': result.deleted_count})
async def delete_recent(request):
data = await request.post()
oid = data.get('oid')
if not oid:
return web.json_response({'error': 'Missing oid'}, status=400)
result = recents.delete_one({'_id': ObjectId(oid)})
if result.deleted_count == 1:
return web.json_response({'status': 'success'})
async def get_epg_url(request):
doc = settings_collection.find_one({'key': 'epg_url'})
epg_url = doc.get('value', '') if doc else ''
return web.json_response({'epg_url': epg_url})
async def set_epg_url(request):
data = await request.post()
epg_url = data.get('epg_url', '').strip()
if not epg_url:
return web.json_response({'error': 'epg_url is required'}, status=400)
settings_collection.update_one(
{'key': 'epg_url'},
{'$set': {'value': epg_url, 'updated_at': datetime.utcnow()}},
upsert=True,
)
return web.json_response({'status': 'ok', 'epg_url': epg_url})
async def get_gemini_key(request):
doc = settings_collection.find_one({'key': 'gemini_key'})
key = doc.get('value', '') if doc else ''
return web.json_response({'gemini_key': key})
async def set_gemini_key(request):
data = await request.post()
key = data.get('gemini_key', '').strip()
if not key:
return web.json_response({'error': 'gemini_key is required'}, status=400)
settings_collection.update_one(
{'key': 'gemini_key'},
{'$set': {'value': key, 'updated_at': datetime.utcnow()}},
upsert=True,
)
return web.json_response({'status': 'ok'})
async def gemini_search(request):
data = await request.post()
query = data.get('query', '').strip()
if not query:
return web.json_response({'error': 'query is required'}, status=400)
doc = settings_collection.find_one({'key': 'gemini_key'})
api_key = doc.get('value', '') if doc else ''
if not api_key:
return web.json_response({'error': 'Gemini API key not configured. Go to Settings > Gemini API Key.'}, status=400)
url = f'https://generativelanguage.googleapis.com/v1beta/models/gemini-2.0-flash:generateContent?key={api_key}'
payload = {
'contents': [{'parts': [{'text': query}]}],
'tools': [{'google_search': {}}]
}
try:
resp = requests.post(url, json=payload, timeout=20)
resp.raise_for_status()
result = resp.json()
answer = ''
sources = []
candidates = result.get('candidates', [])
if candidates:
for part in candidates[0].get('content', {}).get('parts', []):
if 'text' in part:
answer += part['text']
grounding = candidates[0].get('groundingMetadata', {})
for chunk in grounding.get('groundingChunks', []):
web_info = chunk.get('web', {})
if web_info.get('uri'):
sources.append({'uri': web_info['uri'], 'title': web_info.get('title', '')})
return web.json_response({'answer': answer, 'sources': sources})
except requests.exceptions.HTTPError as e:
try:
err_msg = e.response.json().get('error', {}).get('message', str(e))
except Exception:
err_msg = str(e)
return web.json_response({'error': err_msg}, status=502)
except Exception as e:
return web.json_response({'error': str(e)}, status=502)
async def channels_epg(request):
root = load_epg_root()
if root is None:
return web.json_response({'error': 'No EPG data available'}, status=503)
channels_list = []
for channel in root.findall('channel'):
channel_id = channel.attrib.get('id', '')
display_name = ''
name_node = channel.find('display-name')
if name_node is not None and name_node.text:
display_name = name_node.text
if display_name:
channels_list.append({'id': channel_id, 'name': display_name})
return web.json_response({'channels': channels_list})
def parse_epg_datetime(raw_value: str) -> datetime:
ts = raw_value[:14]
return datetime.strptime(ts, "%Y%m%d%H%M%S")
async def epg_programs(request):
channel = request.query.get('channel')
if not channel:
return web.json_response({'error': 'channel is required'}, status=400)
root = load_epg_root()
if root is None:
return web.json_response({'error': 'No EPG data available'}, status=503)
programs = []
for programme in root.findall('programme'):
if programme.attrib.get('channel') != channel:
continue
start_raw = programme.attrib.get('start', '')
stop_raw = programme.attrib.get('stop', '')
title_node = programme.find('title')
desc_node = programme.find('desc')
try:
start_dt = parse_epg_datetime(start_raw)
stop_dt = parse_epg_datetime(stop_raw)
except Exception:
continue
programs.append({
'title': title_node.text if title_node is not None else 'No title',
'desc': desc_node.text if desc_node is not None else '',
'start': start_dt.isoformat(),
'stop': stop_dt.isoformat()
})
programs.sort(key=lambda p: p['start'])
return web.json_response({'programs': programs})
async def download_epg(request):
data = dict(request.query)
url = data['url']
filename = data['filename']
filepath = os.path.join(DOWNLOAD_EPG_DIR, filename)
print('download file :', url, flush=True)
response = requests.get(url)
with open(filepath, 'wb') as f:
f.write(response.content)
return web.json_response({'status':'ok'})
# Routage
BASE_DIR = pathlib.Path(__file__).parent
app.router.add_static('/static/', path=BASE_DIR / 'static', name='static')
app.router.add_static('/logos/', LOGO_DIR)
app.router.add_get('/', index)
app.router.add_get('/playlists', playlist_download)
app.router.add_get('/init_db', init_db)
app.router.add_get('/get_list_playlists', get_list_playlists)
app.router.add_post('/upload', handle_upload)
app.router.add_post('/upload_snapshot', upload_snapshot)
app.router.add_post('/analyse_m3u_file', analyse_m3u_file)
app.router.add_post('/import_xtream_codes', import_xtream_codes)
app.router.add_post('/update_playlist_comment', update_playlist_comment)
app.router.add_get('/delete_playlist', delete_playlist)
app.router.add_get('/cleanup_xtream_duplicates', cleanup_xtream_duplicates)
app.router.add_post('/list_channel', list_channel)
app.router.add_get('/get_logo', get_logo)
app.router.add_get('/list_logo', list_logo)
app.router.add_get('/list_devices', list_devices)
app.router.add_get('/recents', get_recents)
app.router.add_get('/clean_recents', clean_recents)
app.router.add_get('/activation_playlist', activation_playlist)
app.router.add_get('/delete_device', delete_device)
app.router.add_get('/count_channels', count_channels)
app.router.add_post('/delete_recent', delete_recent)
app.router.add_get('/get_source', get_source)
app.router.add_get('/most_recents', most_recents)
app.router.add_get('/epg_url', get_epg_url)
app.router.add_post('/epg_url', set_epg_url)
app.router.add_get('/gemini_key', get_gemini_key)
app.router.add_post('/gemini_key', set_gemini_key)
app.router.add_post('/gemini_search', gemini_search)
app.router.add_get('/channels_epg', channels_epg)
app.router.add_get('/epg', epg_programs)
# SIO FUNCTIONS
def messageReceived():
print('messageReceived')
@sio.event
def connect(sid, environ):
print("connect ", sid, flush=True)
@sio.event
async def communication(fromSid, data):
print("search request from ", fromSid, flush=True)
toSidDevice = data['toSidDevice']
print("send signal to ", toSidDevice, flush=True)
await sio.emit('communication', data, room=toSidDevice ,callback=messageReceived)
@sio.event
def disconnect(sid):
print('disconnect ', sid, flush=True)
@sio.event
async def polling(sid, data):
global devices
device_name = data['name']
IdDevice = data['IdDevice']
mediaRunning = data['mediaRunning']
device_last_datetime = time.time()
data_running = data['data_running']
auto_snapshot_enabled = data.get('auto_snapshot_enabled', True)
snapshot_interval_seconds = data.get('snapshot_interval_seconds', 10)
print('polling from device :', IdDevice, flush=True)
print('data_running :', data_running, flush=True)
try:
vpn_status = data['vpn_status']
except:
vpn_status = 'unknown'
device_data = devices_db.find_one({ 'id_device':IdDevice})
if device_data:
query = { 'id_device':IdDevice}
update = { "$set": { "last_datetime": device_last_datetime ,"name": device_name, "sid":sid, "mediaRunning":mediaRunning, 'vpn_status':vpn_status, 'data_running':data_running, 'auto_snapshot_enabled': auto_snapshot_enabled, 'snapshot_interval_seconds': snapshot_interval_seconds} }
devices_db.update_one(query, update)
else:
new_device = { 'id_device': IdDevice, 'name' : device_name, 'last_datetime' : device_last_datetime, 'sid':sid ,'mediaRunning' : mediaRunning, 'vpn_status':vpn_status, 'data_running':data_running, 'auto_snapshot_enabled': auto_snapshot_enabled, 'snapshot_interval_seconds': snapshot_interval_seconds}
devices_db.insert_one(new_device)
if data_running:
new_record = {
"id_device": IdDevice,
"id_channel": data_running['id_channel'],
"id_playlist": data_running['id_playlist'],
"datetime": datetime.utcnow(),
"name_channel": data_running['name_channel'],
"url_channel":data_running['url_channel'],
"progress": data_running['progress'],
"percentage": data_running['percentage']
}
latest_recent_record = recents.find_one(
{"id_device": IdDevice},
sort=[("datetime", -1)]
)
if latest_recent_record:
same_channel = (
latest_recent_record.get("id_channel") == data_running['id_channel']
and latest_recent_record.get("id_playlist") == data_running['id_playlist']
)
if same_channel:
recents.find_one_and_replace(
{"_id": latest_recent_record["_id"]},
new_record,
return_document=True
)
else:
recents.insert_one(new_record)
else:
recents.insert_one(new_record)
test_exist_device = False
for device in devices:
if device['IdDevice'] == IdDevice:
test_exist_device = True
device['last_datetime'] = device_last_datetime
device['name'] = device_name
device['sid'] = sid
device['mediaRunning'] = mediaRunning
try:
device['vpn_status'] = vpn_status
except:
device['vpn_status'] = 'unknown'
if not test_exist_device:
new_device = { 'IdDevice': IdDevice, 'name' : device_name, 'last_datetime' : device_last_datetime, 'sid':sid ,'mediaRunning' : mediaRunning, 'vpn_status':vpn_status}
devices.append(new_device)
# Main
if __name__ == '__main__':
logger.info('Starting server on port 443')
web.run_app(app, ssl_context=ssl_context, port=443)