mirror of
https://github.com/Viren070/MediaFusion.git
synced 2025-12-01 23:21:11 +01:00
enhance prowlarr scrapping for movies and series and refactoring helper functions
This commit is contained in:
+2
-1
@@ -16,7 +16,8 @@ class Settings(BaseSettings):
|
||||
premiumize_oauth_client_id: str | None = None
|
||||
premiumize_oauth_client_secret: str | None = None
|
||||
prowlarr_url: str = "http://prowlarr-service:9696"
|
||||
prowlarr_api_key: str
|
||||
prowlarr_api_key: str | None = None
|
||||
prowlarr_search_interval_hour: int = 6
|
||||
|
||||
class Config:
|
||||
env_file = ".env"
|
||||
|
||||
+38
-63
@@ -21,17 +21,15 @@ from db.models import (
|
||||
)
|
||||
from db.schemas import Stream, MetaIdProjection
|
||||
from scrappers import tamilmv
|
||||
from scrappers.prowlarr import scrap_streams_from_prowlarr
|
||||
from scrappers.torrentio import (
|
||||
scrap_streams_from_torrentio,
|
||||
scrap_series_streams_from_torrentio,
|
||||
)
|
||||
from scrappers.prowlarr import get_streams_from_prowlarr
|
||||
from scrappers.torrentio import get_streams_from_torrentio
|
||||
from utils.parser import (
|
||||
parse_stream_data,
|
||||
get_catalogs,
|
||||
search_imdb,
|
||||
parse_tv_stream_data,
|
||||
fetch_downloaded_info_hashes,
|
||||
get_imdb_data,
|
||||
)
|
||||
|
||||
|
||||
@@ -131,37 +129,19 @@ async def get_movie_streams(
|
||||
.to_list()
|
||||
)
|
||||
|
||||
|
||||
if settings.is_scrap_from_torrentio:
|
||||
last_torrentio_stream = next(
|
||||
(stream for stream in streams if stream.source == "Torrentio"), None
|
||||
streams = await get_streams_from_torrentio(
|
||||
user_data, streams, video_id, "movie", background_tasks
|
||||
)
|
||||
if settings.prowlarr_api_key:
|
||||
movie_metadata = await get_movie_data_by_id(video_id, False)
|
||||
if movie_metadata:
|
||||
title, year = movie_metadata.title, movie_metadata.year
|
||||
else:
|
||||
title, year = get_imdb_data(video_id)
|
||||
streams = await get_streams_from_prowlarr(
|
||||
user_data, streams, video_id, "movie", title, year, background_tasks
|
||||
)
|
||||
|
||||
if (
|
||||
video_id.startswith("tt")
|
||||
and "torrentio_streams" in user_data.selected_catalogs
|
||||
):
|
||||
if (
|
||||
last_torrentio_stream is None
|
||||
or last_torrentio_stream.updated_at < datetime.now() - timedelta(days=3)
|
||||
):
|
||||
streams.extend(
|
||||
await scrap_streams_from_torrentio(
|
||||
video_id, "movie", background_tasks
|
||||
)
|
||||
)
|
||||
|
||||
last_prowlarr_stream = next(
|
||||
(stream for stream in streams if stream.source == "Prowlarr"), None
|
||||
)
|
||||
if "prowlarr_streams" in user_data.selected_catalogs:
|
||||
if (
|
||||
last_prowlarr_stream is None
|
||||
or last_prowlarr_stream.updated_at < datetime.now() - timedelta(hours=6)
|
||||
):
|
||||
streams.extend(
|
||||
await scrap_streams_from_prowlarr(video_id, "movie", background_tasks)
|
||||
)
|
||||
|
||||
return await parse_stream_data(streams, user_data, secret_str)
|
||||
|
||||
@@ -184,32 +164,30 @@ async def get_series_streams(
|
||||
)
|
||||
|
||||
if settings.is_scrap_from_torrentio:
|
||||
last_torrentio_stream = next(
|
||||
(
|
||||
stream
|
||||
for stream in streams
|
||||
if stream.source == "Torrentio" and stream.get_episode(season, episode)
|
||||
),
|
||||
None,
|
||||
streams = await get_streams_from_torrentio(
|
||||
user_data, streams, video_id, "series", background_tasks, season, episode
|
||||
)
|
||||
if settings.prowlarr_api_key:
|
||||
series_metadata = await get_series_data_by_id(video_id, False)
|
||||
if series_metadata:
|
||||
title, year = series_metadata.title, series_metadata.year
|
||||
else:
|
||||
title, year = get_imdb_data(video_id)
|
||||
streams = await get_streams_from_prowlarr(
|
||||
user_data,
|
||||
streams,
|
||||
video_id,
|
||||
"series",
|
||||
title,
|
||||
year,
|
||||
background_tasks,
|
||||
season,
|
||||
episode,
|
||||
)
|
||||
|
||||
if (
|
||||
video_id.startswith("tt")
|
||||
and "torrentio_streams" in user_data.selected_catalogs
|
||||
):
|
||||
if (
|
||||
last_torrentio_stream is None
|
||||
or last_torrentio_stream.updated_at < datetime.now() - timedelta(days=3)
|
||||
):
|
||||
streams.extend(
|
||||
await scrap_series_streams_from_torrentio(
|
||||
video_id, "series", season, episode, background_tasks
|
||||
)
|
||||
)
|
||||
|
||||
matched_episode_streams = [
|
||||
stream for stream in streams if stream.get_episode(season, episode)
|
||||
]
|
||||
matched_episode_streams = filter(
|
||||
lambda stream: stream.get_episode(season, episode), streams
|
||||
)
|
||||
|
||||
return await parse_stream_data(
|
||||
matched_episode_streams, user_data, secret_str, season, episode
|
||||
@@ -334,9 +312,6 @@ async def save_movie_metadata(metadata: dict):
|
||||
background = existing_movie.background
|
||||
meta_id = existing_movie.id
|
||||
|
||||
# Determine file index for the main movie file (largest file)
|
||||
largest_file = max(metadata["file_data"], key=lambda x: x["size"])
|
||||
|
||||
if "language" in metadata:
|
||||
languages = (
|
||||
[metadata["language"]]
|
||||
@@ -352,8 +327,8 @@ async def save_movie_metadata(metadata: dict):
|
||||
torrent_name=metadata["torrent_name"],
|
||||
announce_list=metadata["announce_list"],
|
||||
size=metadata["total_size"],
|
||||
filename=largest_file["filename"],
|
||||
file_index=largest_file["index"],
|
||||
filename=metadata["largest_file"]["filename"],
|
||||
file_index=metadata["largest_file"]["index"],
|
||||
languages=languages,
|
||||
resolution=metadata.get("resolution"),
|
||||
codec=metadata.get("codec"),
|
||||
|
||||
+85
-18
@@ -1,5 +1,6 @@
|
||||
import json
|
||||
import logging
|
||||
from datetime import datetime
|
||||
|
||||
import cloudscraper
|
||||
import httpx
|
||||
@@ -8,8 +9,8 @@ from requests.adapters import HTTPAdapter
|
||||
from urllib3.util.retry import Retry
|
||||
|
||||
from db.config import settings
|
||||
from utils.torrent import extract_torrent_metadata
|
||||
|
||||
from db.models import TorrentStreams, Episode, Season
|
||||
from utils.torrent import extract_torrent_metadata, info_hashes_to_torrent_metadata
|
||||
|
||||
UA_HEADER = {
|
||||
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/118.0.0.0 Safari/537.36"
|
||||
@@ -53,19 +54,7 @@ async def get_page_content(page, url):
|
||||
return await page.content()
|
||||
|
||||
|
||||
async def download_and_save_torrent(
|
||||
torrent_element,
|
||||
metadata: dict,
|
||||
media_type: str,
|
||||
page_link: str,
|
||||
scraper=None,
|
||||
page=None,
|
||||
):
|
||||
from db import crud # Avoid circular import
|
||||
|
||||
torrent_link = torrent_element.get("href")
|
||||
logging.info(f"Downloading torrent: {torrent_link}")
|
||||
|
||||
async def download_torrent(torrent_link, scraper=None, page=None):
|
||||
if scraper:
|
||||
response = scraper.get(torrent_link)
|
||||
torrent_metadata = extract_torrent_metadata(response.content)
|
||||
@@ -82,18 +71,24 @@ async def download_and_save_torrent(
|
||||
|
||||
if not torrent_metadata:
|
||||
logging.error(f"Info hash not found for {torrent_link}")
|
||||
return False
|
||||
return None
|
||||
|
||||
return torrent_metadata
|
||||
|
||||
|
||||
async def save_torrent(torrent_metadata: dict, metadata: dict, media_type: str):
|
||||
metadata.update(torrent_metadata)
|
||||
|
||||
if not metadata.get("year"):
|
||||
logging.error(f"Year not found for {page_link}")
|
||||
logging.error("Year not found")
|
||||
return False
|
||||
|
||||
from db import crud # Avoid circular import
|
||||
|
||||
# Saving the metadata
|
||||
if media_type == "series":
|
||||
if not metadata.get("season"):
|
||||
logging.error(f"Season not found for {page_link}")
|
||||
logging.error("Season not found")
|
||||
return False
|
||||
await crud.save_series_metadata(metadata)
|
||||
else:
|
||||
@@ -105,6 +100,26 @@ async def download_and_save_torrent(
|
||||
return True
|
||||
|
||||
|
||||
async def download_and_save_torrent(
|
||||
torrent_element,
|
||||
metadata: dict,
|
||||
media_type: str,
|
||||
page_link: str,
|
||||
scraper=None,
|
||||
page=None,
|
||||
):
|
||||
torrent_link = torrent_element.get("href")
|
||||
logging.info(f"Downloading torrent: {torrent_link}")
|
||||
torrent_metadata = await download_torrent(torrent_link, scraper, page)
|
||||
if torrent_metadata is None:
|
||||
return False
|
||||
|
||||
result = await save_torrent(torrent_metadata, metadata, media_type)
|
||||
if not result:
|
||||
logging.error(f"Failed to save torrent for {page_link}")
|
||||
return result
|
||||
|
||||
|
||||
def get_scrapper_config(site_name: str, get_key: str) -> dict:
|
||||
with open("resources/json/scrapper_config.json") as file:
|
||||
config = json.load(file)
|
||||
@@ -123,3 +138,55 @@ async def add_to_bitsearch(magnet_link: str):
|
||||
logging.error(
|
||||
f"Failed to add magnet link {magnet_link}: {response.status_code}"
|
||||
)
|
||||
|
||||
|
||||
async def update_torrent_movie_streams_metadata(info_hashes: list[str]):
|
||||
"""Update torrent streams metadata."""
|
||||
streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, [])
|
||||
|
||||
for stream_metadata in streams_metadata:
|
||||
if not stream_metadata:
|
||||
continue
|
||||
|
||||
torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"])
|
||||
if torrent_stream:
|
||||
torrent_stream.torrent_name = stream_metadata["torrent_name"]
|
||||
torrent_stream.size = stream_metadata["total_size"]
|
||||
|
||||
torrent_stream.filename = stream_metadata["largest_file"]["filename"]
|
||||
torrent_stream.file_index = stream_metadata["largest_file"]["index"]
|
||||
torrent_stream.updated_at = datetime.now()
|
||||
await torrent_stream.save()
|
||||
|
||||
|
||||
async def update_torrent_series_streams_metadata(info_hashes: list[str]):
|
||||
"""Update torrent streams metadata."""
|
||||
streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, [])
|
||||
|
||||
for stream_metadata in streams_metadata:
|
||||
if not stream_metadata:
|
||||
continue
|
||||
|
||||
torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"])
|
||||
if torrent_stream:
|
||||
episodes = [
|
||||
Episode(
|
||||
episode_number=file["episode"],
|
||||
filename=file["filename"],
|
||||
size=file["size"],
|
||||
file_index=file["index"],
|
||||
)
|
||||
for file in stream_metadata["file_data"]
|
||||
if file["episode"]
|
||||
]
|
||||
torrent_stream.season = Season(
|
||||
season_number=stream_metadata["season"],
|
||||
episodes=episodes,
|
||||
)
|
||||
|
||||
torrent_stream.torrent_name = stream_metadata["torrent_name"]
|
||||
torrent_stream.size = stream_metadata["total_size"]
|
||||
|
||||
torrent_stream.updated_at = datetime.now()
|
||||
|
||||
await torrent_stream.save()
|
||||
|
||||
+367
-98
@@ -1,137 +1,406 @@
|
||||
import os
|
||||
import re
|
||||
from datetime import datetime
|
||||
import asyncio
|
||||
import logging
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
import PTN
|
||||
import httpx
|
||||
from fastapi import BackgroundTasks
|
||||
from torf import Magnet
|
||||
|
||||
from db import crud
|
||||
from db.config import settings
|
||||
from db.models import TorrentStreams
|
||||
from utils.parser import convert_size_to_bytes, get_imdb_title
|
||||
from utils.torrent import info_hashes_to_torrent_metadata
|
||||
from db.models import TorrentStreams, Season, Episode
|
||||
from db.schemas import UserData
|
||||
from scrappers.helpers import (
|
||||
update_torrent_series_streams_metadata,
|
||||
update_torrent_movie_streams_metadata,
|
||||
)
|
||||
from utils.torrent import extract_torrent_metadata
|
||||
|
||||
|
||||
async def fetch_stream_data(url: str, params: dict, headers: dict) -> dict:
|
||||
async def get_streams_from_prowlarr(
|
||||
user_data: UserData,
|
||||
streams: list[TorrentStreams],
|
||||
video_id: str,
|
||||
catalog_type: str,
|
||||
title: str,
|
||||
year: str,
|
||||
background_tasks: BackgroundTasks,
|
||||
season: int = None,
|
||||
episode: int = None,
|
||||
):
|
||||
last_stream = next(
|
||||
(stream for stream in streams if "prowlarr_streams" in stream.catalog), None
|
||||
)
|
||||
if video_id.startswith("tt") and "prowlarr_streams" in user_data.selected_catalogs:
|
||||
if last_stream is None or last_stream.updated_at < datetime.now() - timedelta(
|
||||
hours=settings.prowlarr_search_interval_hour
|
||||
):
|
||||
if catalog_type == "movie":
|
||||
streams.extend(
|
||||
await scrap_movies_streams_from_prowlarr(
|
||||
video_id, title, year, background_tasks
|
||||
)
|
||||
)
|
||||
elif catalog_type == "series":
|
||||
streams.extend(
|
||||
await scrap_series_streams_from_prowlarr(
|
||||
video_id, title, background_tasks, season, episode
|
||||
)
|
||||
)
|
||||
return streams
|
||||
|
||||
|
||||
async def fetch_stream_data(url: str, params: dict) -> dict:
|
||||
"""Fetch stream data asynchronously."""
|
||||
headers = {
|
||||
"accept": "application/json",
|
||||
"X-Api-Key": settings.prowlarr_api_key,
|
||||
}
|
||||
async with httpx.AsyncClient() as client:
|
||||
response = await client.get(url, timeout=100, params=params, headers=headers)
|
||||
response = await client.get(url, timeout=10, params=params, headers=headers)
|
||||
response.raise_for_status() # Will raise an exception for 4xx/5xx responses
|
||||
return response.json()
|
||||
|
||||
|
||||
async def scrap_streams_from_prowlarr(
|
||||
video_id: str, catalog_type: str, background_tasks: BackgroundTasks = None, season: int = None, episode: int = None
|
||||
async def scrap_movies_streams_from_prowlarr(
|
||||
video_id: str,
|
||||
title: str,
|
||||
year: str,
|
||||
background_tasks: BackgroundTasks = None,
|
||||
) -> list[TorrentStreams]:
|
||||
"""
|
||||
Get streams by IMDb ID from prowlarr.
|
||||
Get movie streams by IMDb ID from prowlarr.
|
||||
"""
|
||||
url = f"{settings.prowlarr_url}/api/v1/search"
|
||||
stream_data = []
|
||||
|
||||
prefixes = ["{" + f"ImdbId:{video_id}" + "} "]
|
||||
if catalog_type == "tv":
|
||||
prefixes = ["{" + f"Season:{season}" + "} ",
|
||||
"{" + f"Season:{season}" + "} {" + f"Episode:{episode}" + "} "]
|
||||
for prefix in prefixes:
|
||||
headers = {
|
||||
'accept': 'application/json',
|
||||
'X-Api-Key': settings.prowlarr_api_key,
|
||||
}
|
||||
params = {
|
||||
'query': f"{prefix}{get_imdb_title(video_id.removeprefix('tt'))}",
|
||||
'categories': [
|
||||
'2000', # Movies
|
||||
'5000', # TV
|
||||
'8000', # Other
|
||||
],
|
||||
'type': "tvsearch" if catalog_type == "tv" else "moviesearch",
|
||||
'limit': '20',
|
||||
'offset': '0',
|
||||
}
|
||||
try:
|
||||
stream_data.extend(await fetch_stream_data(url, params, headers))
|
||||
except (httpx.HTTPError, httpx.TimeoutException):
|
||||
return [] # Return empty list in case of HTTP errors or timeouts
|
||||
# Params for IMDb ID search
|
||||
params_imdb = {
|
||||
"query": f"{{ImdbId:{video_id}}}",
|
||||
"categories": [2000], # Movies
|
||||
"type": "movie",
|
||||
"limit": 20,
|
||||
"offset": 0,
|
||||
}
|
||||
|
||||
return await store_and_parse_stream_data(
|
||||
video_id, stream_data, background_tasks
|
||||
# Params for title search
|
||||
params_title = {
|
||||
"query": title,
|
||||
"categories": [2000], # Movies
|
||||
"type": "search",
|
||||
"limit": 20,
|
||||
"offset": 0,
|
||||
}
|
||||
|
||||
try:
|
||||
# Fetch data for both searches simultaneously
|
||||
imdb_search, title_search = await asyncio.gather(
|
||||
fetch_stream_data(url, params_imdb),
|
||||
fetch_stream_data(url, params_title),
|
||||
)
|
||||
stream_data.extend(imdb_search)
|
||||
stream_data.extend(title_search)
|
||||
except (httpx.HTTPError, httpx.TimeoutException):
|
||||
return [] # Return an empty list in case of HTTP errors or timeouts
|
||||
|
||||
return await store_and_parse_movie_stream_data(
|
||||
video_id, title, year, stream_data, background_tasks
|
||||
)
|
||||
|
||||
|
||||
async def store_and_parse_stream_data(
|
||||
video_id: str, stream_data: list, background_tasks: BackgroundTasks
|
||||
async def scrap_series_streams_from_prowlarr(
|
||||
video_id: str,
|
||||
title: str,
|
||||
background_tasks: BackgroundTasks = None,
|
||||
season: int = None,
|
||||
episode: int = None,
|
||||
) -> list[TorrentStreams]:
|
||||
"""
|
||||
Get series streams by IMDb ID from prowlarr.
|
||||
"""
|
||||
url = f"{settings.prowlarr_url}/api/v1/search"
|
||||
stream_data = []
|
||||
|
||||
# Params for IMDb ID, season, and episode search
|
||||
params_imdb = {
|
||||
"query": f"{{ImdbId:{video_id}}}{{Season:{season}}}{{Episode:{episode}}}",
|
||||
"categories": [5000], # TV
|
||||
"type": "tvsearch",
|
||||
"limit": 20,
|
||||
"offset": 0,
|
||||
}
|
||||
|
||||
# Params for title search
|
||||
params_title = {
|
||||
"query": title,
|
||||
"categories": [5000], # TV
|
||||
"type": "search",
|
||||
"limit": 20,
|
||||
"offset": 0,
|
||||
}
|
||||
|
||||
try:
|
||||
# Fetch data for both searches simultaneously
|
||||
imdb_search, title_search = await asyncio.gather(
|
||||
fetch_stream_data(url, params_imdb),
|
||||
fetch_stream_data(url, params_title),
|
||||
)
|
||||
stream_data.extend(imdb_search)
|
||||
stream_data.extend(title_search)
|
||||
except (httpx.HTTPError, httpx.TimeoutException):
|
||||
return [] # Return an empty list in case of HTTP errors or timeouts
|
||||
|
||||
return await store_and_parse_series_stream_data(
|
||||
video_id, title, season, episode, stream_data, background_tasks
|
||||
)
|
||||
|
||||
|
||||
async def get_torrent_data_from_prowlarr(download_url: str) -> tuple[dict, bool]:
|
||||
"""Get torrent data from prowlarr."""
|
||||
if not download_url:
|
||||
return {}, False
|
||||
if download_url.startswith("magnet:"):
|
||||
magnet = Magnet.from_string(download_url)
|
||||
return {"info_hash": magnet.infohash, "announce_list": magnet.tr}, False
|
||||
|
||||
async with httpx.AsyncClient() as client:
|
||||
response = await client.get(download_url, follow_redirects=False)
|
||||
|
||||
if response.status_code == 301:
|
||||
magnet_url = response.headers.get("Location")
|
||||
magnet = Magnet.from_string(magnet_url)
|
||||
return {"info_hash": magnet.infohash, "announce_list": magnet.tr}, False
|
||||
elif response.status_code == 200:
|
||||
return extract_torrent_metadata(response.content), True
|
||||
else:
|
||||
return {}, False
|
||||
|
||||
|
||||
async def prowlarr_data_parser(meta_data: dict) -> tuple[dict, bool]:
|
||||
"""Parse prowlarr data."""
|
||||
try:
|
||||
torrent_data, is_torrent_downloaded = await get_torrent_data_from_prowlarr(
|
||||
meta_data.get("downloadUrl") or meta_data.get("magnetUrl")
|
||||
)
|
||||
except Exception as e:
|
||||
logging.warning(f"Error parsing torrent data: {e} {e.__class__.__name__}")
|
||||
if meta_data.get("infoHash"):
|
||||
torrent_data = {
|
||||
"info_hash": meta_data.get("infoHash"),
|
||||
"announce_list": [],
|
||||
}
|
||||
is_torrent_downloaded = False
|
||||
else:
|
||||
return {}, False
|
||||
|
||||
torrent_data.update(
|
||||
{
|
||||
"seeders": meta_data.get("seeders"),
|
||||
"created_at": datetime.strptime(
|
||||
meta_data.get("publishDate"), "%Y-%m-%dT%H:%M:%SZ"
|
||||
),
|
||||
"source": meta_data.get("indexer"),
|
||||
"poster_url": meta_data.get("posterUrl"),
|
||||
}
|
||||
)
|
||||
if is_torrent_downloaded is False:
|
||||
torrent_data.update(
|
||||
{
|
||||
"torrent_name": meta_data.get("title"),
|
||||
"total_size": meta_data.get("size"),
|
||||
**PTN.parse(meta_data.get("title")),
|
||||
}
|
||||
)
|
||||
return torrent_data, is_torrent_downloaded
|
||||
|
||||
|
||||
async def store_and_parse_movie_stream_data(
|
||||
video_id: str,
|
||||
title: str,
|
||||
year: str,
|
||||
stream_data: list,
|
||||
background_tasks: BackgroundTasks,
|
||||
) -> list[TorrentStreams]:
|
||||
streams = []
|
||||
info_hashes = []
|
||||
for stream in stream_data:
|
||||
infohash = stream.get("infoHash", "").lower()
|
||||
if infohash:
|
||||
torrent_stream = None # TODO ADD ME LATER await TorrentStreams.get(infohash)
|
||||
if torrent_stream:
|
||||
# Update existing stream
|
||||
torrent_stream.seeders = stream["seeders"]
|
||||
torrent_stream.updated_at = datetime.now()
|
||||
# TODO ADD ME LATER
|
||||
await torrent_stream.save()
|
||||
else:
|
||||
title = stream["title"]
|
||||
metadata = PTN.parse(title)
|
||||
languages = []
|
||||
language = metadata.get("language", "")
|
||||
languages.extend(language if type(language) == list else [language])
|
||||
# Create new stream
|
||||
torrent_stream = TorrentStreams(
|
||||
id=infohash,
|
||||
torrent_name=title,
|
||||
announce_list=[],
|
||||
size=stream["size"],
|
||||
filename=None,
|
||||
file_index=stream.get("indexerId"),
|
||||
languages=languages,
|
||||
resolution=metadata.get("resolution"),
|
||||
codec=metadata.get("codec"),
|
||||
quality=metadata.get("quality"),
|
||||
audio=metadata.get("audio"),
|
||||
encoder=metadata.get("encoder"),
|
||||
source="Prowlarr",
|
||||
catalog=["prowlarr_streams"],
|
||||
updated_at=datetime.now(),
|
||||
seeders=stream["seeders"],
|
||||
meta_id=video_id,
|
||||
)
|
||||
# TODO ADD ME LATER
|
||||
await torrent_stream.save()
|
||||
parsed_data, _ = await prowlarr_data_parser(stream)
|
||||
info_hash = parsed_data.get("info_hash", "").lower()
|
||||
if not info_hash or not parsed_data.get("seeders"):
|
||||
logging.warning(
|
||||
f"Skipping {info_hash} due to missing info_hash or seeders: {parsed_data.get('seeders')}"
|
||||
)
|
||||
continue
|
||||
|
||||
streams.append(torrent_stream)
|
||||
if torrent_stream.filename is None:
|
||||
info_hashes.append(infohash)
|
||||
if not (
|
||||
parsed_data.get("title").lower() == title.lower()
|
||||
and parsed_data.get("year") == year
|
||||
) and (str(stream.get("imdbId", "")) not in video_id):
|
||||
logging.warning(
|
||||
f"Skipping {info_hash} due to title mismatch: '{parsed_data.get('title')}' != '{title}' or year mismatch: '{parsed_data.get('year')}' != '{year}'"
|
||||
)
|
||||
continue
|
||||
|
||||
# TODO ADD ME LATER
|
||||
background_tasks.add_task(update_torrent_streams_metadata, info_hashes)
|
||||
torrent_stream = await TorrentStreams.get(info_hash)
|
||||
|
||||
if torrent_stream:
|
||||
# Update existing stream
|
||||
torrent_stream.seeders = parsed_data["seeders"]
|
||||
torrent_stream.updated_at = datetime.now()
|
||||
await torrent_stream.save()
|
||||
else:
|
||||
# Create new stream
|
||||
torrent_stream = TorrentStreams(
|
||||
id=info_hash,
|
||||
torrent_name=parsed_data.get("torrent_name"),
|
||||
announce_list=parsed_data.get("announce_list"),
|
||||
size=parsed_data.get("total_size"),
|
||||
filename=parsed_data.get("largest_file", {}).get("file_name"),
|
||||
file_index=parsed_data.get("largest_file", {}).get("index"),
|
||||
languages=parsed_data.get("languages", []),
|
||||
resolution=parsed_data.get("resolution"),
|
||||
codec=parsed_data.get("codec"),
|
||||
quality=parsed_data.get("quality"),
|
||||
audio=parsed_data.get("audio"),
|
||||
encoder=parsed_data.get("encoder"),
|
||||
source=parsed_data.get("source"),
|
||||
catalog=[
|
||||
"prowlarr_streams",
|
||||
"prowlarr_movies",
|
||||
f"{parsed_data.get('source').lower()}_movies",
|
||||
],
|
||||
updated_at=datetime.now(),
|
||||
seeders=parsed_data.get("seeders"),
|
||||
created_at=parsed_data.get("created_at"),
|
||||
meta_id=video_id,
|
||||
)
|
||||
await torrent_stream.save()
|
||||
|
||||
streams.append(torrent_stream)
|
||||
if torrent_stream.filename is None:
|
||||
info_hashes.append(info_hash)
|
||||
|
||||
background_tasks.add_task(update_torrent_movie_streams_metadata, info_hashes)
|
||||
|
||||
return streams
|
||||
|
||||
|
||||
async def update_torrent_streams_metadata(info_hashes: list[str]):
|
||||
"""Update torrent streams metadata."""
|
||||
streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, [])
|
||||
|
||||
for stream_metadata in streams_metadata:
|
||||
if not stream_metadata:
|
||||
async def store_and_parse_series_stream_data(
|
||||
video_id: str,
|
||||
title: str,
|
||||
season: int,
|
||||
episode: int,
|
||||
stream_data: list,
|
||||
background_tasks: BackgroundTasks,
|
||||
) -> list[TorrentStreams]:
|
||||
streams = []
|
||||
info_hashes = []
|
||||
for stream in stream_data:
|
||||
parsed_data, is_torrent_downloaded = await prowlarr_data_parser(stream)
|
||||
info_hash = parsed_data.get("info_hash", "").lower()
|
||||
if not info_hash or not parsed_data.get("seeders"):
|
||||
logging.warning(
|
||||
f"Skipping {info_hash} due to missing info_hash or seeders: {parsed_data.get('seeders')}"
|
||||
)
|
||||
continue
|
||||
|
||||
torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"])
|
||||
if torrent_stream:
|
||||
torrent_stream.torrent_name = stream_metadata.get("torrent_name")
|
||||
torrent_stream.size = stream_metadata.get("total_size")
|
||||
if not (parsed_data.get("title").lower() == title.lower()) and (
|
||||
str(stream.get("imdbId", "")) not in video_id
|
||||
):
|
||||
logging.warning(
|
||||
f"Skipping {info_hash} due to title mismatch: '{parsed_data.get('title')}' != '{title}'"
|
||||
)
|
||||
continue
|
||||
|
||||
largest_file = max(stream_metadata.get("file_data"), key=lambda x: x["size"])
|
||||
torrent_stream.filename = largest_file["filename"]
|
||||
torrent_stream.file_index = largest_file["index"]
|
||||
if parsed_data.get("season"):
|
||||
if isinstance(parsed_data["season"], int):
|
||||
season_number = parsed_data["season"]
|
||||
else:
|
||||
# Skip This Stream due to multiple seasons in one torrent.
|
||||
# TODO: Handle this case later.
|
||||
# Need to refactor DB and how streaming provider works.
|
||||
continue
|
||||
else:
|
||||
season_number = season
|
||||
|
||||
if is_torrent_downloaded is True:
|
||||
episode_data = [
|
||||
Episode(
|
||||
episode_number=file["episode"],
|
||||
filename=file["filename"],
|
||||
size=file["size"],
|
||||
file_index=file["index"],
|
||||
)
|
||||
for file in parsed_data["file_data"]
|
||||
if file["episode"]
|
||||
]
|
||||
elif parsed_data.get("episode"):
|
||||
if isinstance(parsed_data["episode"], int):
|
||||
episode_data = [Episode(episode_number=parsed_data["episode"])]
|
||||
else:
|
||||
episode_data = [
|
||||
Episode(
|
||||
episode_number=episode_number,
|
||||
)
|
||||
for episode_number in parsed_data["episode"]
|
||||
]
|
||||
|
||||
else:
|
||||
episode_data = [Episode(episode_number=episode)]
|
||||
|
||||
torrent_stream = await TorrentStreams.get(info_hash)
|
||||
|
||||
if torrent_stream:
|
||||
torrent_stream.seeders = parsed_data["seeders"]
|
||||
torrent_stream.updated_at = datetime.now()
|
||||
torrent_stream.resolution = stream_metadata.get("resolution")
|
||||
torrent_stream.quality = stream_metadata.get("quality")
|
||||
torrent_stream.codec = stream_metadata.get("codec")
|
||||
episode_item = torrent_stream.get_episode(season, episode)
|
||||
if episode_item is None:
|
||||
if torrent_stream.season:
|
||||
torrent_stream.season.episodes.extend(episode_data)
|
||||
else:
|
||||
torrent_stream.season = Season(
|
||||
season_number=season_number,
|
||||
episodes=episode_data,
|
||||
)
|
||||
episode_item = torrent_stream.get_episode(season, episode)
|
||||
await torrent_stream.save()
|
||||
else:
|
||||
# Create new stream
|
||||
torrent_stream = TorrentStreams(
|
||||
id=info_hash,
|
||||
torrent_name=parsed_data.get("torrent_name"),
|
||||
announce_list=parsed_data.get("announce_list"),
|
||||
size=parsed_data.get("total_size"),
|
||||
filename=None,
|
||||
languages=parsed_data.get("languages", []),
|
||||
resolution=parsed_data.get("resolution"),
|
||||
codec=parsed_data.get("codec"),
|
||||
quality=parsed_data.get("quality"),
|
||||
audio=parsed_data.get("audio"),
|
||||
encoder=parsed_data.get("encoder"),
|
||||
source=parsed_data.get("source"),
|
||||
catalog=[
|
||||
"prowlarr_streams",
|
||||
"prowlarr_series",
|
||||
f"{parsed_data.get('source').lower()}_series",
|
||||
],
|
||||
updated_at=datetime.now(),
|
||||
seeders=parsed_data.get("seeders"),
|
||||
created_at=parsed_data.get("created_at"),
|
||||
meta_id=video_id,
|
||||
season=Season(
|
||||
season_number=season_number,
|
||||
episodes=episode_data,
|
||||
),
|
||||
)
|
||||
await torrent_stream.save()
|
||||
episode_item = torrent_stream.get_episode(season, episode)
|
||||
|
||||
streams.append(torrent_stream)
|
||||
|
||||
if episode_item and episode_item.size is None:
|
||||
info_hashes.append(info_hash)
|
||||
|
||||
background_tasks.add_task(update_torrent_series_streams_metadata, info_hashes)
|
||||
|
||||
return streams
|
||||
|
||||
+44
-62
@@ -1,6 +1,6 @@
|
||||
import logging
|
||||
import re
|
||||
from datetime import datetime
|
||||
from datetime import datetime, timedelta
|
||||
from os import path
|
||||
|
||||
import PTN
|
||||
@@ -9,12 +9,47 @@ from fastapi import BackgroundTasks
|
||||
|
||||
from db.config import settings
|
||||
from db.models import TorrentStreams, Season, Episode
|
||||
from scrappers.helpers import UA_HEADER
|
||||
from db.schemas import UserData
|
||||
from scrappers.helpers import (
|
||||
UA_HEADER,
|
||||
update_torrent_series_streams_metadata,
|
||||
update_torrent_movie_streams_metadata,
|
||||
)
|
||||
from utils.parser import convert_size_to_bytes
|
||||
from utils.torrent import info_hashes_to_torrent_metadata
|
||||
from utils.validation_helper import is_video_file
|
||||
|
||||
|
||||
async def get_streams_from_torrentio(
|
||||
user_data: UserData,
|
||||
streams: list[TorrentStreams],
|
||||
video_id: str,
|
||||
catalog_type: str,
|
||||
background_tasks: BackgroundTasks,
|
||||
season: int = None,
|
||||
episode: int = None,
|
||||
):
|
||||
last_stream = next(
|
||||
(stream for stream in streams if stream.source == "Torrentio"), None
|
||||
)
|
||||
if video_id.startswith("tt") and "torrentio_streams" in user_data.selected_catalogs:
|
||||
if last_stream is None or last_stream.updated_at < datetime.now() - timedelta(
|
||||
days=3
|
||||
):
|
||||
if catalog_type == "movie":
|
||||
streams.extend(
|
||||
await scrap_movie_streams_from_torrentio(
|
||||
video_id, catalog_type, background_tasks
|
||||
)
|
||||
)
|
||||
elif catalog_type == "series":
|
||||
streams.extend(
|
||||
await scrap_series_streams_from_torrentio(
|
||||
video_id, catalog_type, season, episode, background_tasks
|
||||
)
|
||||
)
|
||||
return streams
|
||||
|
||||
|
||||
async def fetch_stream_data(url: str) -> dict:
|
||||
"""Fetch stream data asynchronously."""
|
||||
async with httpx.AsyncClient(
|
||||
@@ -25,7 +60,7 @@ async def fetch_stream_data(url: str) -> dict:
|
||||
return response.json()
|
||||
|
||||
|
||||
async def scrap_streams_from_torrentio(
|
||||
async def scrap_movie_streams_from_torrentio(
|
||||
video_id: str, catalog_type: str, background_tasks: BackgroundTasks
|
||||
) -> list[TorrentStreams]:
|
||||
"""
|
||||
@@ -34,11 +69,11 @@ async def scrap_streams_from_torrentio(
|
||||
url = f"{settings.torrentio_url}/stream/{catalog_type}/{video_id}.json"
|
||||
try:
|
||||
stream_data = await fetch_stream_data(url)
|
||||
return await store_and_parse_stream_data(
|
||||
return await store_and_parse_movie_stream_data(
|
||||
video_id, stream_data.get("streams", []), background_tasks
|
||||
)
|
||||
except (httpx.HTTPError, httpx.TimeoutException):
|
||||
return [] # Return empty list in case of HTTP errors or timeouts
|
||||
return [] # Return an empty list in case of HTTP errors or timeouts
|
||||
except Exception as e:
|
||||
logging.error(f"Error while fetching stream data from torrentio: {e}")
|
||||
return []
|
||||
@@ -61,7 +96,7 @@ async def scrap_series_streams_from_torrentio(
|
||||
video_id, season, episode, stream_data.get("streams", []), background_tasks
|
||||
)
|
||||
except (httpx.HTTPError, httpx.TimeoutException):
|
||||
return [] # Return empty list in case of HTTP errors or timeouts
|
||||
return [] # Return an empty list in case of HTTP errors or timeouts
|
||||
except Exception as e:
|
||||
logging.error(f"Error while fetching stream data from torrentio: {e}")
|
||||
return []
|
||||
@@ -82,7 +117,7 @@ def parse_stream_title(stream: dict) -> dict:
|
||||
}
|
||||
|
||||
|
||||
async def store_and_parse_stream_data(
|
||||
async def store_and_parse_movie_stream_data(
|
||||
video_id: str, stream_data: list, background_tasks: BackgroundTasks
|
||||
) -> list[TorrentStreams]:
|
||||
streams = []
|
||||
@@ -126,7 +161,7 @@ async def store_and_parse_stream_data(
|
||||
if torrent_stream.filename is None:
|
||||
info_hashes.append(stream["infoHash"])
|
||||
|
||||
background_tasks.add_task(update_torrent_streams_metadata, info_hashes)
|
||||
background_tasks.add_task(update_torrent_movie_streams_metadata, info_hashes)
|
||||
|
||||
return streams
|
||||
|
||||
@@ -268,56 +303,3 @@ def extract_size_string(details: str) -> str:
|
||||
"""Extract the size string from the details."""
|
||||
size_match = re.search(r"💾 (\d+(?:\.\d+)?\s*(GB|MB))", details, re.IGNORECASE)
|
||||
return size_match.group(1) if size_match else ""
|
||||
|
||||
|
||||
async def update_torrent_streams_metadata(info_hashes: list[str]):
|
||||
"""Update torrent streams metadata."""
|
||||
streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, [])
|
||||
|
||||
for stream_metadata in streams_metadata:
|
||||
if not stream_metadata:
|
||||
continue
|
||||
|
||||
torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"])
|
||||
if torrent_stream:
|
||||
torrent_stream.torrent_name = stream_metadata["torrent_name"]
|
||||
torrent_stream.size = stream_metadata["total_size"]
|
||||
|
||||
largest_file = max(stream_metadata["file_data"], key=lambda x: x["size"])
|
||||
torrent_stream.filename = largest_file["filename"]
|
||||
torrent_stream.file_index = largest_file["index"]
|
||||
torrent_stream.updated_at = datetime.now()
|
||||
await torrent_stream.save()
|
||||
|
||||
|
||||
async def update_torrent_series_streams_metadata(info_hashes: list[str]):
|
||||
"""Update torrent streams metadata."""
|
||||
streams_metadata = await info_hashes_to_torrent_metadata(info_hashes, [])
|
||||
|
||||
for stream_metadata in streams_metadata:
|
||||
if not stream_metadata:
|
||||
continue
|
||||
|
||||
torrent_stream = await TorrentStreams.get(stream_metadata["info_hash"])
|
||||
if torrent_stream:
|
||||
episodes = [
|
||||
Episode(
|
||||
episode_number=file["episode"],
|
||||
filename=file["filename"],
|
||||
size=file["size"],
|
||||
file_index=file["index"],
|
||||
)
|
||||
for file in stream_metadata["file_data"]
|
||||
if file["episode"]
|
||||
]
|
||||
torrent_stream.season = Season(
|
||||
season_number=stream_metadata["season"],
|
||||
episodes=episodes,
|
||||
)
|
||||
|
||||
torrent_stream.torrent_name = stream_metadata["torrent_name"]
|
||||
torrent_stream.size = stream_metadata["total_size"]
|
||||
|
||||
torrent_stream.updated_at = datetime.now()
|
||||
|
||||
await torrent_stream.save()
|
||||
|
||||
+3
-3
@@ -221,9 +221,9 @@ def search_imdb(title: str, year: int, retry: int = 5) -> dict:
|
||||
return {}
|
||||
|
||||
|
||||
def get_imdb_title(video_id:str) -> str:
|
||||
movie = ia.get_movie(video_id)
|
||||
return movie.get("title")
|
||||
def get_imdb_data(video_id: str) -> tuple[str, str]:
|
||||
movie = ia.get_movie(video_id.removeprefix("tt"))
|
||||
return movie.get("title"), movie.get("year")
|
||||
|
||||
|
||||
def parse_tv_stream_data(tv_data: MediaFusionTVMetaData) -> list[Stream]:
|
||||
|
||||
@@ -70,6 +70,7 @@ def extract_torrent_metadata(content: Torrent | bytes) -> dict:
|
||||
"episode": parsed_data.get("episode"),
|
||||
}
|
||||
)
|
||||
largest_file = max(file_data, key=lambda x: x["size"])
|
||||
|
||||
return {
|
||||
**PTN.parse(torrent.name),
|
||||
@@ -78,6 +79,7 @@ def extract_torrent_metadata(content: Torrent | bytes) -> dict:
|
||||
"total_size": total_size,
|
||||
"file_data": file_data,
|
||||
"torrent_name": torrent.name,
|
||||
"largest_file": largest_file,
|
||||
}
|
||||
except Exception as e:
|
||||
logging.error(f"Error occurred: {e}")
|
||||
|
||||
Reference in New Issue
Block a user