Add ScrapeOps monitoring, Add zilean filtered endpoint, Fix Prowlarr scraping with download link & Fix PikPak login error etc (#319)

* Refactor ZileanScraper to use parallel requests for searching and filtering streams with new endpoints

* Fix DLHD scraping & enable DLHD without MediaFlow

* Switch to httpx for async HTTP requests

* Implement caching for PikPak token to reduce login error

* Refactor torrent cleanup logic

* Add inactivity monitor extension to close idle spiders

Introduced `InactivityMonitor` to automatically close spiders that remain inactive for a specified time. This extension checks activity at regular intervals and uses configurable settings for check intervals and inactivity timeouts. If no items are scraped within the timeout period, the spider is closed to free resources.

* Handle TypeError in dynamic sorting of streams

* Refactor torrent info scraper to support pre-processing.

Introduced a pre-processing function mapping to handle specific indexer requirements before parsing the HTML. Added a custom pre-processing function for "TheRARBG" to handle URL adjustments, improving modularity and readability in the `get_torrent_info` function.

* Refine error logging and fix hash key in spider

* Fix Prowlarr not stop on max process limit

* #315: Integrate ScrapeOps logging into all scrapers & spiders

Added ScrapeOps logging to Zilean, Torrentio, Prowlarr, and Prowlarr Feed scrapers to enhance request tracking and error handling. Configured ScrapeOps API key in settings and updated Pipfile/Pipfile.lock with scrapeops-python-requests and scrapeops-scrapy dependencies.

* Refactor scraper cache status handler

* verify torrent before parsing on prowlarr & prioritize magnet on badass_torrents

* Add support for provide locally hosted mediaflow proxy public address and reduce the time leg on private ip address checking

* handle RD exception cases

* do not setup scrapeops when api key is none

* Enhanced dynamic sorting of torrent streams

Revised the dynamic_sort_key function to handle different key types more efficiently with match-case. Simplified error handling and improved logging to capture sorting data in the case of exceptions.

* add missing last update date for metadata

* update domain for nowmesports
This commit is contained in:
Mohamed Zumair
2024-10-14 05:58:23 +05:30
committed by GitHub
parent 95fdd6a3de
commit 52d8956034
38 changed files with 1394 additions and 732 deletions
+2
View File
@@ -47,6 +47,8 @@ parsett = "*"
tenacity = "*"
ratelimit = "*"
qrcode = "*"
scrapeops-scrapy = "*"
scrapeops-python-requests = "*"
[dev-packages]
pysocks = "*"
Generated
+676 -493
View File
File diff suppressed because it is too large Load Diff
+12 -1
View File
@@ -1,11 +1,13 @@
from apscheduler.schedulers.asyncio import AsyncIOScheduler
from apscheduler.triggers.cron import CronTrigger
from db.config import settings
from mediafusion_scrapy.task import run_spider
from scrapers.imdb_data import fetch_movie_ids_to_update
from scrapers.prowlarr_feed import run_prowlarr_feed_scraper
from scrapers.trackers import update_torrent_seeders
from scrapers.tv import validate_tv_streams_in_db
from scrapers.prowlarr_feed import run_prowlarr_feed_scraper
from scrapers.utils import cleanup_expired_scraper_task
def setup_scheduler(scheduler: AsyncIOScheduler):
@@ -218,3 +220,12 @@ def setup_scheduler(scheduler: AsyncIOScheduler):
"crontab_expression": settings.prowlarr_feed_scraper_crontab,
},
)
scheduler.add_job(
cleanup_expired_scraper_task.send,
CronTrigger.from_crontab(settings.cleanup_expired_scraper_task_crontab),
name="cleanup_expired_scraper_task",
kwargs={
"crontab_expression": settings.cleanup_expired_scraper_task_crontab,
},
)
+9 -2
View File
@@ -23,10 +23,12 @@ class Settings(BaseSettings):
# External Service URLs
scraper_proxy_url: str | None = None
zilean_url: str = "http://zilean.zilean:8181"
playwright_cdp_url: str = "ws://browserless:3000?blockAds=true&stealth=true"
flaresolverr_url: str = "http://flaresolverr:8191/v1"
# External Service API Keys
scrapeops_api_key: str | None = None
# Prowlarr Settings
prowlarr_url: str = "http://prowlarr-service:9696"
prowlarr_api_key: str | None = None
@@ -42,6 +44,10 @@ class Settings(BaseSettings):
torrentio_search_interval_days: int = 3
torrentio_url: str = "https://torrentio.strem.fun"
# Zilean Settings
zilean_search_interval_hour: int = 24
zilean_url: str = "http://zilean.zilean:8181"
# Premiumize Settings
premiumize_oauth_client_id: str | None = None
premiumize_oauth_client_secret: str | None = None
@@ -93,7 +99,7 @@ class Settings(BaseSettings):
streamed_scheduler_crontab: str = "*/30 * * * *"
disable_streamed_scheduler: bool = False
streambtw_scheduler_crontab: str = "*/15 * * * *"
disable_streambtw_scheduler: bool = False
disable_streambtw_scheduler: bool = True
dlhd_scheduler_crontab: str = "25 * * * *"
disable_dlhd_scheduler: bool = False
update_imdb_data_crontab: str = "0 2 * * *"
@@ -108,6 +114,7 @@ class Settings(BaseSettings):
disable_ufc_tgx_scheduler: bool = False
prowlarr_feed_scraper_crontab: str = "0 */3 * * *"
disable_prowlarr_feed_scraper: bool = False
cleanup_expired_scraper_task_crontab: str = "0 * * * *"
@model_validator(mode="after")
def default_poster_host_url(self) -> "Settings":
+1
View File
@@ -154,6 +154,7 @@ class MediaFusionMetaData(Document):
runtime: Optional[str] = None
website: Optional[str] = None
genres: Optional[list[str]] = Field(default_factory=list)
last_updated_at: datetime = Field(default_factory=datetime.now)
class Settings:
is_root = True
+1
View File
@@ -109,6 +109,7 @@ class QBittorrentConfig(BaseModel):
class MediaFlowConfig(BaseModel):
proxy_url: str | None = Field(alias="pu")
api_password: str | None = Field(alias="ap")
public_ip: str | None = Field(alias="pip")
proxy_live_streams: bool = Field(default=False, alias="pls")
proxy_debrid_streams: bool = Field(default=False, alias="pds")
+2 -2
View File
@@ -2,7 +2,7 @@ version: '3.8'
services:
mediafusion:
image: mhdzumair/mediafusion:v4.0.1
image: mhdzumair/mediafusion:v4.0.2
ports:
- "8000:8000"
env_file:
@@ -41,7 +41,7 @@ services:
- "6379:6379"
dramatiq-worker:
image: mhdzumair/mediafusion:v4.0.1
image: mhdzumair/mediafusion:v4.0.2
command: ["pipenv", "run", "dramatiq", "api.task", "-p", "1", "-t", "1", "--queues", "scrapy"]
env_file:
- .env
+2 -2
View File
@@ -19,7 +19,7 @@ spec:
spec:
containers:
- name: mediafusion
image: mhdzumair/mediafusion:v4.0.1
image: mhdzumair/mediafusion:v4.0.2
ports:
- containerPort: 8000
resources:
@@ -124,7 +124,7 @@ spec:
spec:
containers:
- name: dramatiq-worker
image: mhdzumair/mediafusion:v4.0.1
image: mhdzumair/mediafusion:v4.0.2
command: ["pipenv", "run", "dramatiq", "api.task", "-p", "1", "-t", "1"]
env:
- name: HOST_URL
+1 -1
View File
@@ -1,5 +1,5 @@
<?xml version="1.0" encoding="UTF-8" standalone="yes"?>
<addon id="plugin.video.mediafusion" version="4.0.1" name="MediaFusion" provider-name="Mohamed Zumair">
<addon id="plugin.video.mediafusion" version="4.0.2" name="MediaFusion" provider-name="Mohamed Zumair">
<requires>
<import addon="xbmc.python" version="3.0.0"/>
<import addon="xbmc.metadata" version="2.1.0"/>
+63
View File
@@ -0,0 +1,63 @@
import time
from datetime import datetime, timedelta
from scrapy import signals
from scrapy.exceptions import NotConfigured, CloseSpider
from scrapy.utils.log import logger
from twisted.internet import task
class InactivityMonitor:
"""Monitor spider for inactivity and close if inactive for too long"""
def __init__(self, crawler, interval=60.0, inactivity_timeout=15):
self.crawler = crawler
self.stats = crawler.stats
self.interval = interval
self.inactivity_timeout = inactivity_timeout
self.task = None
self.last_scraped_time = time.monotonic()
@classmethod
def from_crawler(cls, crawler):
interval = crawler.settings.getfloat("INACTIVITY_CHECK_INTERVAL", 60.0)
inactivity_timeout = crawler.settings.getint("INACTIVITY_TIMEOUT_MINUTES", 15)
if not interval or not inactivity_timeout:
raise NotConfigured
ext = cls(crawler, interval, inactivity_timeout)
crawler.signals.connect(ext.spider_opened, signal=signals.spider_opened)
crawler.signals.connect(ext.spider_closed, signal=signals.spider_closed)
crawler.signals.connect(ext.item_scraped, signal=signals.item_scraped)
return ext
def spider_opened(self, spider):
self.last_scraped_time = time.monotonic()
self.task = task.LoopingCall(self.check_inactivity, spider)
self.task.start(self.interval)
def check_inactivity(self, spider):
time_since_last_scrape_seconds = time.monotonic() - self.last_scraped_time
time_since_last_scrape = timedelta(seconds=time_since_last_scrape_seconds)
logger.debug(
"Checking inactivity for spider %s, last scraped %s ago",
spider.name,
time_since_last_scrape,
)
if time_since_last_scrape > timedelta(minutes=self.inactivity_timeout):
msg = f"No items scraped in the last {self.inactivity_timeout} minutes. Closing spider."
logger.info(msg, extra={"spider": spider})
self.crawler.engine.close_spider(
spider, reason=f"Inactivity timeout: {msg}"
)
raise CloseSpider(f"Inactivity timeout: {msg}")
def item_scraped(self, item, response, spider):
self.last_scraped_time = datetime.now()
def spider_closed(self, spider, reason):
if self.task and self.task.running:
self.task.stop()
@@ -12,6 +12,7 @@ class LiveStreamResolverPipeline:
adapter = ItemAdapter(item)
stream_url = adapter.get("stream_url")
stream_headers = adapter.get("stream_headers")
response_headers = adapter.get("response_headers", {})
if not stream_headers:
referer = adapter.get("referer")
stream_headers = {"Referer": referer} if referer else {}
@@ -30,7 +31,9 @@ class LiveStreamResolverPipeline:
)
)
)
content_type = response.headers.get("Content-Type", b"").decode().lower()
content_type = response_headers.get(
"Content-Type", response.headers.get("Content-Type", b"").decode().lower()
)
if response.status == 200 and content_type in const.IPTV_VALID_CONTENT_TYPES:
stream_headers.update(
@@ -49,6 +52,7 @@ class LiveStreamResolverPipeline:
"notWebReady": True,
"proxyHeaders": {
"request": stream_headers,
"response": response_headers,
},
},
}
@@ -22,7 +22,7 @@ class SportVideoParserPipeline:
def process_item(self, item, spider):
adapter = ItemAdapter(item)
if "title" not in adapter or "torrent_name" not in adapter:
raise DropItem(f"title not found in item: {item}")
raise DropItem(f"title or torrent_name not found in item: {item}")
match = self.title_regex.search(adapter["title"]) or self.title_regex.search(
adapter["torrent_name"]
+21 -5
View File
@@ -51,16 +51,30 @@ SPIDER_MIDDLEWARES = {
# Enable or disable downloader middlewares
# See https://docs.scrapy.org/en/latest/topics/downloader-middleware.html
DOWNLOADER_MIDDLEWARES = {
"mediafusion_scrapy.middlewares.FlaresolverrMiddleware": 542,
"mediafusion_scrapy.middlewares.TooManyRequestsRetryMiddleware": 543,
"scrapy.downloadermiddlewares.retry.RetryMiddleware": 550,
}
# Enable or disable extensions
# See https://docs.scrapy.org/en/latest/topics/extensions.html
# EXTENSIONS = {
# "scrapy.extensions.telnet.TelnetConsole": None,
# }
EXTENSIONS = {
"mediafusion_scrapy.extensions.InactivityMonitor": 100,
}
if settings.scrapeops_api_key:
SCRAPEOPS_API_KEY = settings.scrapeops_api_key
DOWNLOADER_MIDDLEWARES.update(
{
"scrapeops_scrapy.middleware.retry.RetryMiddleware": 550,
"scrapy.downloadermiddlewares.retry.RetryMiddleware": None,
}
)
EXTENSIONS.update(
{
"scrapeops_scrapy.extension.ScrapeOpsMonitor": 500, # ScrapeOps Monitor
}
)
# Configure item pipelines
# See https://docs.scrapy.org/en/latest/topics/item-pipeline.html
@@ -114,3 +128,5 @@ RETRY_HTTP_CODES = [
RETRY_TIMES = 5
FLARESOLVERR_URL = settings.flaresolverr_url
INACTIVITY_TIMEOUT_MINUTES = 15
+5 -22
View File
@@ -96,33 +96,16 @@ class DaddyLiveHDSpider(scrapy.Spider):
item_copy.update(
{
"stream_name": channel["channel_name"],
"stream_url": m3u8_url,
"stream_headers": {
"Referer": self.referer,
"Origin": self.referer.rstrip("/"),
},
"response_headers": {
"Content-Type": "application/vnd.apple.mpegurl",
},
"channel_id": channel["channel_id"],
}
)
yield scrapy.Request(
m3u8_url,
self.parse_stream_link,
meta={
"item": item_copy,
"dont_redirect": True,
"handle_httpstatus_list": [301],
},
headers={"Referer": self.referer},
dont_filter=True,
)
def parse_stream_link(self, response):
item = response.meta["item"]
stream_url = response.headers.get("Location", b"").decode("utf-8")
if not stream_url:
self.logger.error(
f"Failed to get stream URL for {item['stream_name']} channel_id: {item['channel_id']}"
)
return
item["stream_url"] = stream_url
yield item
yield item_copy
+1 -1
View File
@@ -327,7 +327,7 @@ class FormulaTgxSpider(TgxSpider):
logo_image = "https://i.postimg.cc/Sqf4V8tj/f1logo.png?dl=1"
keyword_patterns = re.compile(r"formula[ .+]*[1234e]+", re.IGNORECASE)
scraped_info_hash_key = "formula_tgx_scraped_info_hash3"
scraped_info_hash_key = "formula_tgx_scraped_info_hash"
custom_settings = {
"ITEM_PIPELINES": {
Binary file not shown.
Binary file not shown.
Binary file not shown.
+9 -1
View File
@@ -509,7 +509,7 @@
<div class="mb-2">
<label for="mediaflow_proxy_url">MediaFlow Proxy URL:</label>
<input type="text" class="form-control" id="mediaflow_proxy_url" name="mediaflow_proxy_url"
placeholder="https://your-mediaflow-proxy-url.com"
placeholder="https://your-mediaflow-proxy-url.com or http://127.0.0.1:8888"
value="{{ user_data.mediaflow_config.proxy_url if user_data.mediaflow_config else '' }}">
<div class="invalid-feedback">
Please enter a valid MediaFlow Proxy URL.
@@ -528,6 +528,14 @@
Please enter the MediaFlow API password.
</div>
</div>
<div class="mb-3">
<label for="mediaflow_public_ip">MediaFlow Public IP (Optional):</label>
<i class="bi bi-question-circle" data-bs-toggle="tooltip" data-bs-placement="top"
title="Configure this only when running MediaFlow locally with a proxy service. Leave empty if MediaFlow is configured locally without a proxy server or if it's hosted on
a remote server."></i>
<input type="text" class="form-control" id="mediaflow_public_ip" placeholder="Enter public IP address. (Optional, See tooltip for details)"
value="{{ user_data.mediaflow_config.public_ip if user_data.mediaflow_config else '' }}">
</div>
<div class="form-check mb-2">
<input class="form-check-input" type="checkbox" id="proxy_live_streams" name="proxy_live_streams"
{% if user_data.mediaflow_config and user_data.mediaflow_config.proxy_live_streams %}checked{% endif %}>
+1
View File
@@ -313,6 +313,7 @@ function getUserData() {
mediaflowConfig = {
proxy_url: document.getElementById('mediaflow_proxy_url').value,
api_password: document.getElementById('mediaflow_api_password').value,
mediaflow_public_ip: document.getElementById('mediaflow_public_ip').value,
proxy_live_streams: document.getElementById('proxy_live_streams').checked,
proxy_debrid_streams: document.getElementById('proxy_debrid_streams').checked
};
+4 -3
View File
@@ -1,6 +1,6 @@
{
"start_urls": {
"nowsports": "https://nowmesports.nl/",
"nowsports": "https://nowmesport.com/",
"nowmetv": "https://nowmelive.com/",
"tamilultra": "https://tamilultra.tv/",
"tamilbulb": "https://tamilbulb.tv/",
@@ -283,8 +283,9 @@
},
"dlhd": {
"schedule_url": "https://dlhd.so/schedule/schedule-generated.json",
"m3u8_base_url": "https://webhdrunns.mizhls.ru/lb/premium{channel_id}/index.m3u8",
"referer": "https://cookiewebplay.xyz/",
"m3u8_base_url": "https://xyzdddd.mizhls.ru/lb/premium{channel_id}/index.m3u8",
"key_url": "https://key.keylocking.ru",
"referer": "https://ilovetoplay.xyz/",
"category_mapping": {
"Tv Shows": "Other Sports",
"Soccer": "Football",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"id": "stremio.addons",
"version": "4.0.1",
"version": "4.0.2",
"name": "Media Fusion",
"contactEmail": "mhdzumair@gmail.com",
"description": "Universal Stremio Add-on for Movies, Series, Live TV & Sports Events. Source: https://github.com/mhdzumair/MediaFusion",
+27 -9
View File
@@ -1,5 +1,6 @@
import abc
import logging
import time
from datetime import timedelta
from functools import wraps
from typing import List, Any, Dict
@@ -40,7 +41,7 @@ class BaseScraper(abc.ABC):
@staticmethod
def cache(ttl: int = 3600):
"""
Decorator for caching the results of a method.
Decorator for caching the scraping status using Redis Sorted Sets with timestamps.
:param ttl: Time to live for the cache in seconds
"""
@@ -48,12 +49,23 @@ class BaseScraper(abc.ABC):
@wraps(func)
async def wrapper(self, *args, **kwargs):
cache_key = self.get_cache_key(*args, **kwargs)
cached_result = await REDIS_ASYNC_CLIENT.get(cache_key)
if cached_result:
return []
current_time = int(time.time())
# Check if the item has been scraped recently
score = await REDIS_ASYNC_CLIENT.zscore(
self.cache_key_prefix, cache_key
)
if score and current_time - score < ttl:
return [] # Item has been scraped recently, no need to scrape again
result = await func(self, *args, **kwargs)
await REDIS_ASYNC_CLIENT.set(cache_key, "True", ex=ttl)
# Mark the item as scraped with the current timestamp
await REDIS_ASYNC_CLIENT.zadd(
self.cache_key_prefix, {cache_key: current_time}
)
return result
return wrapper
@@ -145,11 +157,9 @@ class BaseScraper(abc.ABC):
:return: Cache key string
"""
if catalog_type == "movie":
return f"{self.cache_key_prefix}:{catalog_type}:{metadata.id}"
return f"{catalog_type}:{metadata.id}"
return (
f"{self.cache_key_prefix}:{catalog_type}:{metadata.id}:{season}:{episode}"
)
return f"{catalog_type}:{metadata.id}:{season}:{episode}"
def validate_title_and_year(
self,
@@ -223,3 +233,11 @@ class BaseScraper(abc.ABC):
from db.crud import store_new_torrent_streams
await store_new_torrent_streams(streams)
@staticmethod
async def remove_expired_items(scraper_prefix: str, ttl: int = 3600):
"""
Remove expired items from the cache.
"""
current_time = int(time.time())
await REDIS_ASYNC_CLIENT.zremrangebyscore(scraper_prefix, 0, current_time - ttl)
+71 -7
View File
@@ -5,6 +5,7 @@ from typing import List, Dict, Any, AsyncGenerator, Literal, AsyncIterable
import PTT
import dramatiq
import httpx
from scrapeops_python_requests.scrapeops_requests import ScrapeOpsRequests
from torf import Magnet, MagnetError
from db.config import settings
@@ -21,7 +22,7 @@ from scrapers.base_scraper import BaseScraper
from scrapers.imdb_data import get_episode_by_date, get_season_episodes
from utils.network import CircuitBreaker, batch_process_with_circuit_breaker
from utils.parser import is_contain_18_plus_keywords
from utils.runtime_const import REDIS_ASYNC_CLIENT
from utils.runtime_const import REDIS_ASYNC_CLIENT, PROWLARR_SEARCH_TTL
from utils.torrent import extract_torrent_metadata
from utils.wrappers import minimum_run_interval
@@ -47,6 +48,10 @@ allowlist_keywords = [
# fmt: on
class MaxProcessLimitReached(Exception):
pass
class ProwlarrScraper(BaseScraper):
MOVIE_SEARCH_QUERY_TEMPLATES = [
"{title} ({year})", # Exact match with year
@@ -60,16 +65,17 @@ class ProwlarrScraper(BaseScraper):
"{title} S{season:02d}", # Short season search with leading zeros
"{title}", # Title-only fallback
]
cache_key_prefix = "prowlarr"
def __init__(self):
super().__init__(
cache_key_prefix="prowlarr", logger_name=self.__class__.__name__
cache_key_prefix=self.cache_key_prefix, logger_name=self.__class__.__name__
)
self.base_url = f"{settings.prowlarr_url}/api/v1/search"
self.scrapeops_logger = None
self.scrape_response = None
@BaseScraper.cache(
ttl=int(timedelta(hours=settings.prowlarr_search_interval_hour).total_seconds())
)
@BaseScraper.cache(ttl=PROWLARR_SEARCH_TTL)
@BaseScraper.rate_limit(calls=5, period=timedelta(seconds=1))
async def scrape_and_parse(
self,
@@ -80,6 +86,16 @@ class ProwlarrScraper(BaseScraper):
) -> List[TorrentStreams]:
results = []
processed_info_hashes: set[str] = set()
job_name = f"{metadata.title}:{metadata.id}"
if catalog_type == "series":
job_name += f":{season}:{episode}"
if settings.scrapeops_api_key:
self.scrapeops_logger = ScrapeOpsRequests(
scrapeops_api_key=settings.scrapeops_api_key,
spider_name="Prowlarr Scraper",
job_name=job_name,
)
try:
async for stream in self._scrape_and_parse(
@@ -90,6 +106,11 @@ class ProwlarrScraper(BaseScraper):
episode,
):
results.append(stream)
if settings.scrapeops_api_key:
self.scrapeops_logger.item_scraped(
item=stream.model_dump(include={"id"}),
response=self.scrape_response,
)
except httpx.ReadTimeout:
self.logger.warning("Timeout while fetching search results")
except httpx.HTTPStatusError as e:
@@ -98,6 +119,9 @@ class ProwlarrScraper(BaseScraper):
)
except Exception as e:
self.logger.exception(f"An error occurred during scraping: {str(e)}")
finally:
if settings.scrapeops_api_key:
self.scrapeops_logger.logger.close_sdk()
self.logger.info(
f"Returning {len(results)} scraped streams for {metadata.title}"
@@ -260,7 +284,7 @@ class ProwlarrScraper(BaseScraper):
if max_process and streams_processed >= max_process:
self.logger.info(f"Reached max process limit of {max_process}")
break
raise MaxProcessLimitReached("Max process limit reached")
try:
async with asyncio.timeout(max_process_time):
@@ -279,6 +303,10 @@ class ProwlarrScraper(BaseScraper):
f"Stream processing timed out after {max_process_time} seconds. "
f"Processed {streams_processed} streams"
)
except MaxProcessLimitReached:
self.logger.info(
f"Stream processing cancelled after reaching max process limit of {max_process}"
)
except Exception as e:
self.logger.error(f"An error occurred during stream processing: {e}")
self.logger.info(
@@ -295,6 +323,7 @@ class ProwlarrScraper(BaseScraper):
params=params,
timeout=settings.prowlarr_search_query_timeout,
)
self.scrape_response = response
response.raise_for_status()
return response.json()
@@ -655,7 +684,12 @@ class ProwlarrScraper(BaseScraper):
redirect_url = response.headers.get("Location")
return await self.get_torrent_data(redirect_url, indexer)
response.raise_for_status()
return extract_torrent_metadata(response.content, is_parse_ptt=False), True
if response.headers.get("Content-Type") == "application/x-bittorrent":
return (
extract_torrent_metadata(response.content, is_parse_ptt=False),
True,
)
return {}, False
@staticmethod
def parse_title_data(title: str) -> dict:
@@ -713,6 +747,13 @@ async def background_movie_title_search(
if not metadata:
return
if settings.scrapeops_api_key:
scraper.scrapeops_logger = ScrapeOpsRequests(
scrapeops_api_key=settings.scrapeops_api_key,
spider_name="Prowlarr Scraper",
job_name=f"background:{metadata.title}:{metadata.id}",
)
title_streams_generators = [
scraper.scrape_movie_by_title(
processed_info_hashes,
@@ -725,6 +766,11 @@ async def background_movie_title_search(
try:
async for stream in scraper.process_streams(*title_streams_generators):
await scraper.store_streams([stream])
if settings.scrapeops_api_key:
scraper.scrapeops_logger.item_scraped(
item=stream.model_dump(include={"id"}),
response=scraper.scrape_response,
)
except httpx.ReadTimeout:
scraper.logger.warning(
f"Timeout while fetching search results for movie {metadata.title} ({metadata.year}), retrying later"
@@ -738,6 +784,9 @@ async def background_movie_title_search(
scraper.logger.error(
f"Error fetching search results: {e.response.text}, status code: {e.response.status_code}"
)
finally:
if settings.scrapeops_api_key:
scraper.scrapeops_logger.logger.close_sdk()
scraper.logger.info(
f"Background title search completed for {metadata.title} ({metadata.year})"
@@ -762,6 +811,13 @@ async def background_series_title_search(
if not metadata:
return
if settings.scrapeops_api_key:
scraper.scrapeops_logger = ScrapeOpsRequests(
scrapeops_api_key=settings.scrapeops_api_key,
spider_name="Prowlarr Scraper",
job_name=f"background:{metadata.title}:{metadata.id}:{season}:{episode}",
)
title_streams_generators = [
scraper.scrape_series_by_title(
processed_info_hashes,
@@ -778,6 +834,11 @@ async def background_series_title_search(
try:
async for stream in scraper.process_streams(*title_streams_generators):
await scraper.store_streams([stream])
if settings.scrapeops_api_key:
scraper.scrapeops_logger.item_scraped(
item=stream.model_dump(include={"id"}),
response=scraper.scrape_response,
)
except httpx.ReadTimeout:
scraper.logger.warning(
f"Timeout while fetching search results for {metadata.title} S{season}E{episode}, retrying later"
@@ -792,6 +853,9 @@ async def background_series_title_search(
scraper.logger.error(
f"Error fetching search results: {e.response.text}, status code: {e.response.status_code}"
)
finally:
if settings.scrapeops_api_key:
scraper.scrapeops_logger.logger.close_sdk()
scraper.logger.info(
f"Background title search completed for {metadata.title} S{season}E{episode}"
+46 -1
View File
@@ -2,6 +2,7 @@ import logging
import dramatiq
import httpx
from scrapeops_python_requests.scrapeops_requests import ScrapeOpsRequests
from db.config import settings
from db.crud import (
@@ -40,6 +41,13 @@ async def mark_item_as_processed(item_id: str):
async def scrape_prowlarr_feed():
scraper = ProwlarrScraper()
if settings.scrapeops_api_key:
scraper.scrapeops_logger = ScrapeOpsRequests(
scrapeops_api_key=settings.scrapeops_api_key,
spider_name="Prowlarr Scraper",
job_name="Prowlarr Feed Scraper",
)
params = {
"type": "search",
"categories": [2000, 5000, 8000], # Movies, TV, and Other categories
@@ -70,6 +78,9 @@ async def scrape_prowlarr_feed():
except Exception as e:
logger.exception(f"Error scraping Prowlarr feed: {e}")
finally:
if scraper.scrapeops_logger:
scraper.scrapeops_logger.logger.close_sdk()
async def process_feed_item(item: dict, scraper: ProwlarrScraper):
@@ -83,6 +94,12 @@ async def process_feed_item(item: dict, scraper: ProwlarrScraper):
parsed_title_data = scraper.parse_title_data(item["title"])
if is_contain_18_plus_keywords(item["title"]):
logger.warning(f"Item {item['title']} contains black listed keywords")
if settings.scrapeops_api_key:
scraper.scrapeops_logger.logger.item_dropped(
item={"info_hash": item_id},
response=scraper.scrape_response,
message="Contains blacklisted keywords",
)
return item_id
# Determine media type
@@ -97,10 +114,22 @@ async def process_feed_item(item: dict, scraper: ProwlarrScraper):
logger.warning(
f"Category 8000 item {item['title']} does not match expected format"
)
if settings.scrapeops_api_key:
scraper.scrapeops_logger.logger.item_dropped(
item={"info_hash": item_id},
response=scraper.scrape_response,
message="Does not match expected format",
)
return item_id
media_type = "series" if parsed_title_data.get("seasons") else "movie"
else:
logger.warning(f"Unsupported category {category_ids} for item {item['title']}")
if settings.scrapeops_api_key:
scraper.scrapeops_logger.logger.item_dropped(
item={"info_hash": item_id},
response=scraper.scrape_response,
message="Unsupported category",
)
return item_id
# Fetch or create metadata
@@ -111,6 +140,12 @@ async def process_feed_item(item: dict, scraper: ProwlarrScraper):
if not metadata:
logger.warning(f"Unable to find or create metadata for {item['title']}")
if settings.scrapeops_api_key:
scraper.scrapeops_logger.logger.item_dropped(
item={"info_hash": item_id},
response=scraper.scrape_response,
message="Unable to find or create metadata",
)
return item_id
# Process the stream
@@ -119,9 +154,19 @@ async def process_feed_item(item: dict, scraper: ProwlarrScraper):
)
if stream:
await store_new_torrent_streams([stream])
if settings.scrapeops_api_key:
scraper.scrapeops_logger.item_scraped(
item={"info_hash": item_id}, response=scraper.scrape_response
)
return item_id
else:
logger.warning(f"Failed to process stream for {item['title']}")
if settings.scrapeops_api_key:
scraper.scrapeops_logger.logger.item_dropped(
item={"info_hash": item_id},
response=scraper.scrape_response,
message="Failed to process stream",
)
return item_id
@@ -150,6 +195,6 @@ async def search_and_create_metadata(metadata: dict, media_type: str):
max_backoff=3600000,
priority=50,
)
async def run_prowlarr_feed_scraper():
async def run_prowlarr_feed_scraper(**kwargs):
logger.info("Running Prowlarr feed scraper")
await scrape_prowlarr_feed()
+13 -5
View File
@@ -7,9 +7,10 @@ from utils.torrent import get_info_hash_from_magnet
async def get_torrent_info(url: str, indexer: str) -> dict:
torrent_info = {}
pre_process_func = pre_process_map.get(indexer, get_page_bs4)
parse_func = parser_function_map.get(indexer, parse_common_torrents)
soup = await get_page_bs4(url)
soup = await pre_process_func(url)
if not soup:
return torrent_info
@@ -43,7 +44,7 @@ def parse_torrent_downloads(soup, torrent_info, url):
info_hash = (
info_hash_span.find_parent("p").text.replace("Infohash:", "").strip()
)
torrent_info["infoHash"] = info_hash.strip()
torrent_info["infoHash"] = info_hash
return torrent_info
@@ -74,10 +75,8 @@ def parse_badass_torrents(soup, torrent_info, url):
torrent_link = soup.find("a", string="Torrent Download")
if torrent_link:
torrent_info["downloadUrl"] = urljoin(url, torrent_link.get("href"))
torrent_info["magnetUrl"] = None
elif magnet_link:
if magnet_link:
torrent_info["magnetUrl"] = magnet_link.get("href")
torrent_info["downloadUrl"] = None
return torrent_info
@@ -116,6 +115,15 @@ def parse_gktorrent(soup, torrent_info, url):
return torrent_info
async def pre_process_therarbg(url):
url = url.replace("?format=json", "")
return await get_page_bs4(url)
pre_process_map = {
"TheRARBG": pre_process_therarbg,
}
parser_function_map = {
"1337x": parse_1337x,
"TheRARBG": parse_common_torrents,
+20 -15
View File
@@ -5,6 +5,7 @@ from os import path
from typing import List, Dict, Any
import PTT
from scrapeops_python_requests.scrapeops_requests import ScrapeOpsRequests
from tenacity import RetryError
from db.config import settings
@@ -14,22 +15,21 @@ from utils.parser import (
convert_size_to_bytes,
is_contain_18_plus_keywords,
)
from utils.runtime_const import TORRENTIO_SEARCH_TTL
from utils.validation_helper import is_video_file
from scrapeops_python_requests.scrapeops_requests import ScrapeOpsRequests
class TorrentioScraper(BaseScraper):
cache_key_prefix = "torrentio"
def __init__(self):
super().__init__(
cache_key_prefix="torrentio", logger_name=self.__class__.__name__
cache_key_prefix=self.cache_key_prefix, logger_name=self.__class__.__name__
)
self.base_url = settings.torrentio_url
self.semaphore = asyncio.Semaphore(10)
@BaseScraper.cache(
ttl=int(timedelta(days=settings.torrentio_search_interval_days).total_seconds())
)
@BaseScraper.cache(ttl=TORRENTIO_SEARCH_TTL)
@BaseScraper.rate_limit(calls=5, period=timedelta(seconds=1))
async def scrape_and_parse(
self,
@@ -44,14 +44,17 @@ class TorrentioScraper(BaseScraper):
url = f"{self.base_url}/stream/{catalog_type}/{metadata.id}:{season}:{episode}.json"
job_name += f":{season}:{episode}"
scrapeops_logger = ScrapeOpsRequests(
scrapeops_api_key=settings.scrapeops_api_key,
spider_name="Torrentio Scraper",
job_name=job_name,
)
scrapeops_logger = None
if settings.scrapeops_api_key:
scrapeops_logger = ScrapeOpsRequests(
scrapeops_api_key=settings.scrapeops_api_key,
spider_name="Torrentio Scraper",
job_name=job_name,
)
try:
response = await self.make_request(url)
response.raise_for_status()
data = response.json()
if not self.validate_response(data):
@@ -62,9 +65,10 @@ class TorrentioScraper(BaseScraper):
data, metadata, catalog_type, season, episode
)
for stream in stream_data:
scrapeops_logger.item_scraped(
item=stream.model_dump(include={"id", "meta_id"}), response=response
)
if scrapeops_logger:
scrapeops_logger.item_scraped(
item=stream.model_dump(include={"id"}), response=response
)
return stream_data
except (ScraperError, RetryError):
return []
@@ -72,7 +76,8 @@ class TorrentioScraper(BaseScraper):
self.logger.exception(f"Error occurred while fetching {url}: {e}")
return []
finally:
scrapeops_logger.logger.close_sdk()
if scrapeops_logger:
scrapeops_logger.logger.close_sdk()
def validate_response(self, response: Dict[str, Any]) -> bool:
return "streams" in response and isinstance(response["streams"], list)
+24
View File
@@ -1,11 +1,19 @@
import asyncio
import logging
import dramatiq
from db.config import settings
from db.models import TorrentStreams, MediaFusionMetaData
from scrapers.base_scraper import BaseScraper
from scrapers.prowlarr import ProwlarrScraper
from scrapers.torrentio import TorrentioScraper
from scrapers.zilean import ZileanScraper
from utils.runtime_const import (
ZILEAN_SEARCH_TTL,
TORRENTIO_SEARCH_TTL,
PROWLARR_SEARCH_TTL,
)
async def run_scrapers(
@@ -46,3 +54,19 @@ async def run_scrapers(
unique_streams = set(scraped_streams)
logging.info(f"Scraped {len(scraped_streams)} streams for {metadata.title}")
return unique_streams
@dramatiq.actor(
time_limit=5 * 60 * 1000, # 5 minutes
priority=20,
)
async def cleanup_expired_scraper_task(**kwargs):
await BaseScraper.remove_expired_items(
ProwlarrScraper.cache_key_prefix, PROWLARR_SEARCH_TTL
)
await BaseScraper.remove_expired_items(
TorrentioScraper.cache_key_prefix, TORRENTIO_SEARCH_TTL
)
await BaseScraper.remove_expired_items(
ZileanScraper.cache_key_prefix, ZILEAN_SEARCH_TTL
)
+84 -14
View File
@@ -3,25 +3,28 @@ from datetime import timedelta
from typing import List, Dict, Any
import PTT
from httpx import Response
from tenacity import RetryError
from scrapeops_python_requests.scrapeops_requests import ScrapeOpsRequests
from db.config import settings
from db.models import TorrentStreams, Season, Episode, MediaFusionMetaData
from scrapers.base_scraper import BaseScraper, ScraperError
from utils.parser import (
is_contain_18_plus_keywords,
)
from utils.runtime_const import ZILEAN_SEARCH_TTL
class ZileanScraper(BaseScraper):
cache_key_prefix = "zilean"
def __init__(self):
super().__init__(cache_key_prefix="zilean", logger_name=self.__class__.__name__)
self.base_url = f"{settings.zilean_url}/dmm/search"
super().__init__(
cache_key_prefix=self.cache_key_prefix, logger_name=self.__class__.__name__
)
self.semaphore = asyncio.Semaphore(10)
@BaseScraper.cache(
ttl=int(timedelta(hours=settings.prowlarr_search_interval_hour).total_seconds())
)
@BaseScraper.cache(ttl=ZILEAN_SEARCH_TTL)
@BaseScraper.rate_limit(calls=5, period=timedelta(seconds=1))
async def scrape_and_parse(
self,
@@ -30,22 +33,86 @@ class ZileanScraper(BaseScraper):
season: int = None,
episode: int = None,
) -> List[TorrentStreams]:
try:
stream_response = await self.make_request(
self.base_url,
job_name = f"{metadata.title}:{metadata.id}"
if catalog_type == "series":
job_name += f":{season}:{episode}"
scrapeops_logger = None
if settings.scrapeops_api_key:
scrapeops_logger = ScrapeOpsRequests(
scrapeops_api_key=settings.scrapeops_api_key,
spider_name="Zilean Scraper",
job_name=job_name,
)
search_task = asyncio.create_task(
self.make_request(
f"{settings.zilean_url}/dmm/search",
method="POST",
json={"queryText": metadata.title},
timeout=10,
)
)
stream_data = stream_response.json()
if not self.validate_response(stream_data):
self.logger.warning(f"Invalid response received for {metadata.title}")
return []
if metadata.type == "movie":
params = {
"Query": metadata.title,
"Year": metadata.year,
}
else:
params = {
"Query": metadata.title,
"Season": season,
"Episode": episode,
}
return await self.parse_response(
filtered_task = asyncio.create_task(
self.make_request(
f"{settings.zilean_url}/dmm/filtered",
method="GET",
params=params,
timeout=10,
)
)
search_response, filtered_response = await asyncio.gather(
search_task, filtered_task, return_exceptions=True
)
stream_data = []
response = None
if isinstance(search_response, Response):
response = search_response
stream_data.extend(search_response.json())
else:
self.logger.error(
f"Error occurred while search {metadata.title}: {search_response}"
)
if isinstance(filtered_response, Response):
response = filtered_response
stream_data.extend(filtered_response.json())
else:
self.logger.error(
f"Error occurred while filtering {metadata.title}: {filtered_response}"
)
if not self.validate_response(stream_data):
self.logger.error(f"No valid streams found for {metadata.title}")
if scrapeops_logger:
scrapeops_logger.logger.close_sdk()
return []
try:
streams = await self.parse_response(
stream_data, metadata, catalog_type, season, episode
)
if scrapeops_logger:
for stream in streams:
scrapeops_logger.item_scraped(
item=stream.model_dump(include={"id"}),
response=response,
)
return streams
except (ScraperError, RetryError):
return []
except Exception as e:
@@ -53,6 +120,9 @@ class ZileanScraper(BaseScraper):
f"Error occurred while fetching {metadata.title}: {e}"
)
return []
finally:
if scrapeops_logger:
scrapeops_logger.logger.close_sdk()
async def parse_response(
self,
+34
View File
@@ -10,6 +10,8 @@ from db.models import TorrentStreams
from db.schemas import UserData
from streaming_providers.exceptions import ProviderException
from streaming_providers.parser import select_file_index_from_torrent
from utils import crypto
from utils.runtime_const import REDIS_ASYNC_CLIENT
async def get_torrent_file_by_info_hash(
@@ -141,6 +143,22 @@ async def find_file_in_folder_tree(
async def initialize_pikpak(user_data: UserData):
cache_key = f"pikpak:{crypto.get_text_hash(user_data.streaming_provider.email + user_data.streaming_provider.password, full_hash=True)}"
if pikpak_encrypted_token := await REDIS_ASYNC_CLIENT.get(cache_key):
pikpak_encoded_token = crypto.decrypt_text(
pikpak_encrypted_token, user_data.streaming_provider.password
)
pikpak = PikPakApi(
encoded_token=pikpak_encoded_token,
httpx_client_args={
"transport": httpx.AsyncHTTPTransport(retries=3),
"timeout": 10,
},
token_refresh_callback=store_pikpak_token_in_cache,
token_refresh_callback_kwargs={"user_data": user_data},
)
return pikpak
pikpak = PikPakApi(
username=user_data.streaming_provider.email,
password=user_data.streaming_provider.password,
@@ -148,7 +166,10 @@ async def initialize_pikpak(user_data: UserData):
"transport": httpx.AsyncHTTPTransport(retries=3),
"timeout": 10,
},
token_refresh_callback=store_pikpak_token_in_cache,
token_refresh_callback_kwargs={"user_data": user_data},
)
try:
await pikpak.login()
except PikpakException as error:
@@ -166,9 +187,22 @@ async def initialize_pikpak(user_data: UserData):
"Failed to connect to PikPak. Please try again later.",
"debrid_service_down_error.mp4",
)
await store_pikpak_token_in_cache(pikpak, user_data)
return pikpak
async def store_pikpak_token_in_cache(pikpak: PikPakApi, user_data: UserData):
cache_key = f"pikpak:{crypto.get_text_hash(user_data.streaming_provider.email + user_data.streaming_provider.password, full_hash=True)}"
await REDIS_ASYNC_CLIENT.set(
cache_key,
crypto.encrypt_text(
pikpak.encoded_token, user_data.streaming_provider.password
),
ex=7 * 24 * 60 * 60, # Store for 7 days
)
async def handle_torrent_status(
pikpak: PikPakApi,
info_hash: str,
+20 -7
View File
@@ -18,13 +18,26 @@ class RealDebrid(DebridClient):
super().__init__(token)
def _handle_service_specific_errors(self, error):
if (
error.response.status_code == 403
and error.response.json().get("error_code") == 9
):
raise ProviderException(
"Real-Debrid Permission denied for free account", "need_premium.mp4"
)
if error.response.status_code == 403:
error_code = error.response.json().get("error_code")
match error_code:
case 9:
raise ProviderException(
"Real-Debrid Permission denied for free account",
"need_premium.mp4",
)
case 22:
raise ProviderException(
"IP address not allowed", "ip_not_allowed.mp4"
)
case 34:
raise ProviderException(
"Too many requests", "too_many_requests.mp4"
)
case 35:
raise ProviderException(
"Content marked as infringing", "content_infringing.mp4"
)
def _make_request(
self,
+95 -54
View File
@@ -1,6 +1,7 @@
import asyncio
import sys
import logging
import sys
from datetime import datetime, timezone, timedelta
from typing import Dict
import PTT
@@ -12,10 +13,13 @@ from db.models import (
TorrentStreams,
MediaFusionSeriesMetaData,
MediaFusionMovieMetaData,
MediaFusionMetaData,
)
from utils.parser import calculate_max_similarity_ratio, is_contain_18_plus_keywords
BATCH_SIZE = 1000 # Adjust the batch size as needed
SKIP_TORRENTS = 0 # Skip the first n torrents
LAST_UPDATED = datetime.now(tz=timezone.utc) - timedelta(days=1)
# Configure logging
logging.basicConfig(
@@ -26,49 +30,63 @@ logger = logging.getLogger(__name__)
async def cleanup_torrents(dry_run: bool = True) -> Dict[str, int]:
metrics = {
"total_torrents": 0,
"removed_torrents": 0,
"updated_torrents": 0,
"total": 0,
"removed": 0,
"no_metadata": 0,
"valid": 0,
"ratio_error": 0,
"adult_content": 0,
"ratio": 0,
"adult": 0,
"year_mismatch": 0,
"removed_non_imdb": 0,
"non_movie": 0,
}
lock = asyncio.Lock()
total_torrents = await TorrentStreams.count()
progress_bar = tqdm(
range(0, total_torrents, BATCH_SIZE), desc="Processing torrents"
# Create a cursor for efficient pagination
cursor = TorrentStreams.find(TorrentStreams.updated_at < LAST_UPDATED).sort(
-TorrentStreams.updated_at
)
total_torrents = await cursor.count()
total_torrents = max(0, total_torrents - SKIP_TORRENTS)
progress_bar = tqdm(total=total_torrents, desc="Processing torrents")
torrent_bulk_writer = BulkWriter()
movies_bulk_writer = BulkWriter()
series_bulk_writer = BulkWriter()
for skip in progress_bar:
torrents = await TorrentStreams.find().skip(skip).limit(BATCH_SIZE).to_list()
tasks = [
process_torrent(
torrent,
dry_run,
metrics,
lock,
torrent_bulk_writer,
movies_bulk_writer,
series_bulk_writer,
)
for torrent in torrents
]
await asyncio.gather(*tasks)
await torrent_bulk_writer.commit()
await movies_bulk_writer.commit()
await series_bulk_writer.commit()
torrent_bulk_writer.operations.clear()
movies_bulk_writer.operations.clear()
series_bulk_writer.operations.clear()
cursor = cursor.skip(SKIP_TORRENTS)
processed_count = 0
async for torrent in cursor:
await process_torrent(
torrent,
dry_run,
metrics,
lock,
torrent_bulk_writer,
movies_bulk_writer,
series_bulk_writer,
)
processed_count += 1
if processed_count >= BATCH_SIZE:
await torrent_bulk_writer.commit()
await movies_bulk_writer.commit()
await series_bulk_writer.commit()
torrent_bulk_writer.operations.clear()
movies_bulk_writer.operations.clear()
series_bulk_writer.operations.clear()
processed_count = 0
progress_bar.update(1)
progress_bar.set_postfix(metrics)
# Commit any remaining operations
await torrent_bulk_writer.commit()
await movies_bulk_writer.commit()
await series_bulk_writer.commit()
progress_bar.close()
return metrics
@@ -82,20 +100,28 @@ async def process_torrent(
series_bulk_writer: BulkWriter,
) -> str:
async with lock:
metrics["total_torrents"] += 1
metrics["total"] += 1
meta_data = None
if any("series" in catalog for catalog in torrent.catalog):
meta_data = await MediaFusionSeriesMetaData.get(torrent.meta_id)
meta_data_bulk_writer = series_bulk_writer
else:
meta_data = await MediaFusionMovieMetaData.get(torrent.meta_id)
meta_data_bulk_writer = movies_bulk_writer
meta_data = await MediaFusionMetaData.get_motor_collection().find_one(
{"_id": torrent.meta_id}
)
if not meta_data:
await torrent.delete(bulk_writer=torrent_bulk_writer)
async with lock:
metrics["no_metadata"] += 1
return "no_metadata"
metrics["removed"] += 1
logger.debug(
f"No metadata found for torrent: {torrent.torrent_name}, {torrent.meta_id}"
)
return "removed"
if meta_data["type"] == "movie":
meta_data = MediaFusionMovieMetaData(**meta_data)
meta_data_bulk_writer = movies_bulk_writer
else:
meta_data = MediaFusionSeriesMetaData(**meta_data)
meta_data_bulk_writer = series_bulk_writer
if (
not dry_run
@@ -109,16 +135,17 @@ async def process_torrent(
await torrent.delete(bulk_writer=torrent_bulk_writer)
await meta_data.delete(bulk_writer=meta_data_bulk_writer)
async with lock:
metrics["removed_torrents"] += 1
metrics["removed"] += 1
metrics["removed_non_imdb"] += 1
logger.debug(f"Removed non-IMDb torrent: {torrent.torrent_name}")
return "removed"
if is_contain_18_plus_keywords(torrent.torrent_name):
if not dry_run:
await torrent.delete(bulk_writer=torrent_bulk_writer)
async with lock:
metrics["removed_torrents"] += 1
metrics["adult_content"] += 1
metrics["removed"] += 1
metrics["adult"] += 1
logger.debug(f"Removed torrent due to adult content: {torrent.torrent_name}")
return "removed"
@@ -128,26 +155,41 @@ async def process_torrent(
parsed_data.get("title", ""), meta_data.title, meta_data.aka_titles
)
expected_ratio = 70 if torrent.source in ["TamilMV", "TamilBlasters"] else 85
if max_ratio < expected_ratio and torrent.meta_id.startswith("tt"):
if (
max_ratio < expected_ratio
and torrent.meta_id.startswith("tt")
and not any(x in torrent.catalog for x in ["wwe_tgx", "ufc_tgx"])
):
if not dry_run:
await torrent.delete(bulk_writer=torrent_bulk_writer)
async with lock:
metrics["removed_torrents"] += 1
metrics["ratio_error"] += 1
metrics["removed"] += 1
metrics["ratio"] += 1
logger.debug(
f"Removed torrent due to low similarity ratio: {parsed_data.get('title')} != {meta_data.title} (ratio: {max_ratio}%) (full title: {torrent.torrent_name})"
)
return "removed"
if "year" in parsed_data and torrent.meta_id.startswith("tt"):
if meta_data.type == "movie" and parsed_data["year"] != meta_data.year:
if torrent.meta_id.startswith("tt") and meta_data.type == "movie":
if parsed_data.get("seasons"):
if not dry_run:
await torrent.delete(bulk_writer=torrent_bulk_writer)
async with lock:
metrics["removed_torrents"] += 1
metrics["removed"] += 1
metrics["non_movie"] += 1
logger.debug(
f"Removed non-movie torrent: {torrent.torrent_name}, found season: {parsed_data['seasons']}. Expected movie."
)
return "removed"
if meta_data.type == "movie" and parsed_data.get("year") != meta_data.year:
if not dry_run:
await torrent.delete(bulk_writer=torrent_bulk_writer)
async with lock:
metrics["removed"] += 1
metrics["year_mismatch"] += 1
logger.debug(
f"Removed torrent due to year mismatch: {torrent.torrent_name} (torrent year: {parsed_data['year']}, meta year: {meta_data.year})"
f"Removed torrent due to year mismatch: {torrent.torrent_name} ({parsed_data.get('year')} != {meta_data.year}) ({meta_data.id}) ({meta_data.type})"
)
return "removed"
@@ -163,14 +205,13 @@ async def process_torrent(
torrent.audio = (
parsed_data.get("audio")[0] if parsed_data.get("audio") else None
)
torrent.updated_at = datetime.now(tz=timezone.utc)
await torrent.save(bulk_writer=torrent_bulk_writer)
async with lock:
metrics["updated_torrents"] += 1
return "updated"
async with lock:
metrics["valid"] += 1
return "valid"
torrent.updated_at = datetime.now(tz=timezone.utc)
await torrent.save(bulk_writer=torrent_bulk_writer)
return "updated"
# Initialize the database connection and start the cleanup process
+29 -14
View File
@@ -14,6 +14,32 @@ from db.schemas import UserData
from utils.runtime_const import SECRET_KEY
def encrypt_text(text: str, secret_key: str | bytes) -> str:
iv = get_random_bytes(16)
if isinstance(secret_key, str):
secret_key = secret_key.encode("utf-8")
cipher = AES.new(secret_key.ljust(32)[:32], AES.MODE_CBC, iv)
encoded_text = text.encode("utf-8")
encrypted_data = cipher.encrypt(
encoded_text + b"\0" * (16 - len(encoded_text) % 16)
)
compressed_data = zlib.compress(iv + encrypted_data)
encrypted_str = urlsafe_b64encode(compressed_data).decode("utf-8")
return encrypted_str
def decrypt_text(secret_str: str, secret_key: str | bytes) -> str:
decoded_data = urlsafe_b64decode(secret_str)
encrypted_data = zlib.decompress(decoded_data)
iv = encrypted_data[:16]
if isinstance(secret_key, str):
secret_key = secret_key.encode("utf-8")
cipher = AES.new(secret_key.ljust(32)[:32], AES.MODE_CBC, iv)
decrypted_data = cipher.decrypt(encrypted_data[16:])
decrypted_data = decrypted_data.rstrip(b"\0")
return decrypted_data.decode("utf-8")
def encrypt_user_data(user_data: UserData) -> str:
data = user_data.model_dump_json(
exclude_none=True,
@@ -22,26 +48,15 @@ def encrypt_user_data(user_data: UserData) -> str:
round_trip=True,
by_alias=True,
)
dump_data = data.encode("utf-8")
iv = get_random_bytes(16)
cipher = AES.new(SECRET_KEY, AES.MODE_CBC, iv)
encrypted_data = cipher.encrypt(dump_data + b"\0" * (16 - len(dump_data) % 16))
compressed_data = zlib.compress(iv + encrypted_data)
encrypted_str = urlsafe_b64encode(compressed_data).decode("utf-8")
return encrypted_str
return encrypt_text(data, SECRET_KEY)
def decrypt_user_data(secret_str: str | None = None) -> UserData:
if not secret_str:
return UserData()
try:
decoded_data = urlsafe_b64decode(secret_str)
encrypted_data = zlib.decompress(decoded_data)
iv = encrypted_data[:16]
cipher = AES.new(SECRET_KEY, AES.MODE_CBC, iv)
decrypted_data = cipher.decrypt(encrypted_data[16:])
decrypted_data = decrypted_data.rstrip(b"\0")
user_data = UserData.model_validate_json(decrypted_data.decode("utf-8"))
decrypted_data = decrypt_text(secret_str, SECRET_KEY)
user_data = UserData.model_validate_json(decrypted_data)
except Exception:
user_data = UserData()
return user_data
+15 -11
View File
@@ -2,7 +2,7 @@ import asyncio
import logging
from typing import Callable, AsyncGenerator, Any, Tuple
from urllib import parse
from urllib.parse import urlencode
from urllib.parse import urlencode, urlparse
import httpx
from fastapi.requests import Request
@@ -169,14 +169,21 @@ def get_client_ip(request: Request) -> str | None:
return request.client.host if request.client else "127.0.0.1"
async def get_mediaflow_proxy_public_ip(
mediaflow_proxy_url: str, api_password
) -> str | None:
async def get_mediaflow_proxy_public_ip(mediaflow_config) -> str | None:
"""
Get the public IP address of the MediaFlow proxy server.
"""
if mediaflow_config.public_ip:
return mediaflow_config.public_ip
parsed_url = urlparse(mediaflow_config.mediaflow_proxy_url)
if PRIVATE_CIDR.match(parsed_url.netloc):
# MediaFlow proxy URL is a private IP address
return None
cache_key = crypto.get_text_hash(
f"{mediaflow_proxy_url}:{api_password}", full_hash=True
f"{mediaflow_config.mediaflow_proxy_url}:{mediaflow_config.api_password}",
full_hash=True,
)
if public_ip := await REDIS_ASYNC_CLIENT.getex(cache_key, ex=300):
return public_ip.decode()
@@ -184,8 +191,8 @@ async def get_mediaflow_proxy_public_ip(
try:
async with httpx.AsyncClient() as client:
response = await client.get(
parse.urljoin(mediaflow_proxy_url, "/proxy/ip"),
params={"api_password": api_password},
parse.urljoin(mediaflow_config.mediaflow_proxy_url, "/proxy/ip"),
params={"api_password": mediaflow_config.api_password},
timeout=10,
)
response.raise_for_status()
@@ -211,10 +218,7 @@ async def get_user_public_ip(
and user_data.mediaflow_config
and user_data.mediaflow_config.proxy_debrid_streams
):
public_ip = await get_mediaflow_proxy_public_ip(
user_data.mediaflow_config.proxy_url,
user_data.mediaflow_config.api_password,
)
public_ip = await get_mediaflow_proxy_public_ip(user_data.mediaflow_config)
if public_ip:
return public_ip
# Get the user's public IP address
+50 -32
View File
@@ -1,7 +1,8 @@
import asyncio
import functools
import logging
from typing import Optional, List
from datetime import datetime
from typing import Optional, List, Any
import math
import re
@@ -13,6 +14,7 @@ from db.models import TorrentStreams, TVStreams
from db.schemas import Stream, UserData
from streaming_providers import mapper
from utils import const
from utils.config import config_manager
from utils.const import STREAMING_PROVIDERS_SHORT_NAMES
from utils.network import encode_mediaflow_proxy_url
from utils.runtime_const import ADULT_CONTENT_KEYWORDS, TRACKERS
@@ -94,13 +96,29 @@ async def filter_and_sort_streams(
)
# Step 3: Dynamically sort streams based on user preferences
def dynamic_sort_key(stream):
return tuple(
(
const.RESOLUTION_RANKING.get(stream.filtered_resolution, 0)
if key == "resolution"
else (
-min(
def dynamic_sort_key(stream: TorrentStreams) -> tuple:
def key_value(key: str) -> Any:
match key:
case "cached":
return stream.cached or False
case "resolution":
return const.RESOLUTION_RANKING.get(stream.filtered_resolution, 0)
case "quality":
return const.QUALITY_RANKING.get(stream.filtered_quality, 0)
case "size":
return stream.size
case "seeders":
return stream.seeders or 0
case "created_at":
created_at = stream.created_at
if isinstance(created_at, datetime):
return created_at
elif isinstance(created_at, (int, float)):
return datetime.fromtimestamp(created_at)
else:
return datetime.min
case "language":
return -min(
(
user_data.language_sorting.index(lang)
for lang in stream.filtered_languages
@@ -108,28 +126,22 @@ async def filter_and_sort_streams(
),
default=len(user_data.language_sorting),
)
if key == "language"
else (
const.QUALITY_RANKING.get(stream.filtered_quality, 0)
if key == "quality"
else (
getattr(stream, key, 0)
if key in stream.model_fields_set
else 0
)
)
)
)
for key in user_data.torrent_sorting_priority
case _ if key in stream.model_fields_set:
return getattr(stream, key, 0)
case _:
return 0
return tuple(key_value(key) for key in user_data.torrent_sorting_priority)
try:
dynamically_sorted_streams = sorted(
filtered_streams, key=dynamic_sort_key, reverse=True
)
def safe_sort_key(stream):
raw_key = dynamic_sort_key(stream)
return tuple(0 if item is None else item for item in raw_key)
dynamically_sorted_streams = sorted(
filtered_streams, key=safe_sort_key, reverse=True
)
except (TypeError, Exception):
logging.exception(
f"torrent_sorting_priority: {user_data.torrent_sorting_priority}: sort data: {[dynamic_sort_key(stream) for stream in filtered_streams]}"
)
dynamically_sorted_streams = filtered_streams
# Step 4: Limit streams per resolution based on user preference, after dynamic sorting
limited_streams = []
@@ -390,7 +402,7 @@ async def process_stream(
stream_url, behavior_hints = stream.url, stream.behaviorHints
behavior_hints = behavior_hints if behavior_hints else {}
if stream.drm_key or "dlhd" in stream.source:
if stream.drm_key:
if not is_mediaflow_proxy_enabled:
return "MEDIAFLOW_NEEDED"
stream_url = get_proxy_url(stream, mediaflow_config)
@@ -418,7 +430,10 @@ def get_proxy_url(stream: TVStreams, mediaflow_config) -> str:
if stream.drm_key:
query_params = {"key_id": stream.drm_key_id, "key": stream.drm_key}
elif "dlhd" in stream.source:
query_params = {"use_request_proxy": False}
query_params = {
"use_request_proxy": False,
"key_url": config_manager.get_scraper_config("dlhd", "key_url"),
}
return encode_mediaflow_proxy_url(
mediaflow_config.proxy_url,
@@ -426,6 +441,9 @@ def get_proxy_url(stream: TVStreams, mediaflow_config) -> str:
stream.url,
query_params=query_params,
request_headers=stream.behaviorHints.get("proxyHeaders", {}).get("request", {}),
response_headers=stream.behaviorHints.get("proxyHeaders", {}).get(
"response", {}
),
encryption_api_password=mediaflow_config.api_password,
)
@@ -460,7 +478,7 @@ async def fetch_downloaded_info_hashes(
return downloaded_info_hashes
except Exception as error:
logging.error(
logging.exception(
f"Failed to fetch downloaded info hashes for {user_data.streaming_provider.service}: {error}"
)
pass
+12
View File
@@ -1,4 +1,5 @@
import re
from datetime import timedelta
from fastapi.templating import Jinja2Templates
import redis
@@ -41,3 +42,14 @@ REDIS_SYNC_CLIENT: redis.Redis = redis.Redis(
REDIS_ASYNC_CLIENT: redis.asyncio.Redis = redis.asyncio.Redis(
connection_pool=redis.asyncio.ConnectionPool.from_url(settings.redis_url)
)
PROWLARR_SEARCH_TTL = int(
timedelta(hours=settings.prowlarr_search_interval_hour).total_seconds()
)
TORRENTIO_SEARCH_TTL = int(
timedelta(days=settings.torrentio_search_interval_days).total_seconds()
)
ZILEAN_SEARCH_TTL = int(
timedelta(hours=settings.zilean_search_interval_hour).total_seconds()
)
+1 -1
View File
@@ -80,7 +80,7 @@ def extract_torrent_metadata(content: bytes, is_parse_ptt: bool = True) -> dict:
metadata.update(PTT.parse_title(torrent_name, True))
return metadata
except Exception as e:
logging.error(f"Error occurred: {e}")
logging.exception(f"Error occurred: {e}")
return {}
+36 -26
View File
@@ -4,8 +4,7 @@ import logging
from urllib import parse
from urllib.parse import urlparse
import aiohttp
from aiohttp import ClientError
import httpx
from db import schemas
from utils import const
@@ -18,14 +17,17 @@ def is_valid_url(url: str) -> bool:
async def does_url_exist(url: str) -> bool:
async with aiohttp.ClientSession() as session:
async with httpx.AsyncClient() as client:
try:
async with session.head(
url, allow_redirects=True, timeout=10, headers=const.UA_HEADER
) as response:
logging.info("URL: %s, Status: %s", url, response.status)
return response.status == 200
except (ClientError, asyncio.TimeoutError) as err:
response = await client.head(
url, timeout=10, headers=const.UA_HEADER, follow_redirects=True
)
response.raise_for_status()
return response.status_code == 200
except httpx.HTTPStatusError as err:
logging.error("URL: %s, Status: %s", url, err.response.status_code)
return False
except (httpx.RequestError, httpx.TimeoutException) as err:
logging.error("URL: %s, Status: %s", url, err)
return False
@@ -41,20 +43,29 @@ async def validate_live_stream_url(
return False
headers = behaviour_hint.get("proxyHeaders", {}).get("request", {})
async with aiohttp.ClientSession() as session:
async with httpx.AsyncClient() as client:
try:
async with session.head(
url,
allow_redirects=True,
headers=headers,
timeout=aiohttp.ClientTimeout(total=20),
) as response:
content_type = response.headers.get("content-type", "").lower()
is_valid = content_type in const.IPTV_VALID_CONTENT_TYPES
return is_valid
except (ClientError, asyncio.TimeoutError) as err:
response = await client.head(
url, timeout=10, headers=headers, follow_redirects=True
)
response.raise_for_status()
content_type = (
behaviour_hint.get("proxyHeaders", {})
.get("response", {})
.get("Content-Type", response.headers.get("content-type", "").lower())
)
is_valid = content_type in const.IPTV_VALID_CONTENT_TYPES
return is_valid
except (
httpx.RequestError,
httpx.TimeoutException,
httpx.HTTPStatusError,
) as err:
logging.error(err)
return False
except Exception as e:
logging.exception(e)
return False
async def validate_m3u8_or_mpd_url_with_cache(url: str, behaviour_hint: dict):
@@ -78,13 +89,12 @@ class ValidationError(Exception):
async def validate_yt_id(yt_id: str) -> bool:
image_url = f"https://img.youtube.com/vi/{yt_id}/mqdefault.jpg"
async with aiohttp.ClientSession() as session:
async with httpx.AsyncClient() as client:
try:
async with session.head(
image_url, allow_redirects=True, timeout=aiohttp.ClientTimeout(total=5)
) as response:
return response.status == 200
except (ClientError, asyncio.TimeoutError):
response = await client.head(image_url, timeout=10, follow_redirects=True)
response.raise_for_status()
return response.status_code == 200
except (httpx.HTTPStatusError, httpx.TimeoutException, Exception):
return False