mirror of
https://github.com/g0ldyy/comet.git
synced 2026-01-12 01:16:12 +01:00
@@ -221,6 +221,8 @@ TORRENT_DISABLED_STREAM_URL=https://comet.fast # Optional URL included in the pl
|
||||
# Content Filtering #
|
||||
# ============================== #
|
||||
REMOVE_ADULT_CONTENT=False
|
||||
DIGITAL_RELEASE_FILTER=False # Filter unreleased content
|
||||
TMDB_READ_ACCESS_TOKEN= # Optional: Provide your own TMDB Read Access Token to avoid using the shared default key
|
||||
|
||||
# ============================== #
|
||||
# UI Customization #
|
||||
|
||||
@@ -16,6 +16,7 @@ from comet.background_scraper.worker import background_scraper
|
||||
from comet.core.database import (cleanup_expired_locks,
|
||||
cleanup_expired_sessions, setup_database,
|
||||
teardown_database)
|
||||
from comet.core.execution import setup_executor, shutdown_executor
|
||||
from comet.core.logger import logger
|
||||
from comet.core.models import settings
|
||||
from comet.services.anime import anime_mapper
|
||||
@@ -49,6 +50,7 @@ async def lifespan(app: FastAPI):
|
||||
# loop.set_debug(True)
|
||||
|
||||
await setup_database()
|
||||
setup_executor()
|
||||
await download_best_trackers()
|
||||
|
||||
# Load anime ID mapping for enhanced metadata and anime detection
|
||||
@@ -107,6 +109,7 @@ async def lifespan(app: FastAPI):
|
||||
await torrent_update_queue.stop()
|
||||
|
||||
await teardown_database()
|
||||
shutdown_executor()
|
||||
|
||||
|
||||
tags_metadata = [
|
||||
|
||||
@@ -8,7 +8,9 @@ from fastapi import APIRouter, BackgroundTasks, Request
|
||||
from comet.core.config_validation import config_check
|
||||
from comet.core.logger import logger
|
||||
from comet.core.models import database, settings, trackers
|
||||
from comet.debrid.exceptions import DebridAuthError
|
||||
from comet.debrid.manager import get_debrid_extension
|
||||
from comet.metadata.filter import release_filter
|
||||
from comet.metadata.manager import MetadataScraper
|
||||
from comet.services.debrid import DebridService
|
||||
from comet.services.lock import DistributedLock, is_scrape_in_progress
|
||||
@@ -135,6 +137,9 @@ async def stream(
|
||||
b64config: str = None,
|
||||
chilllink: bool = False,
|
||||
):
|
||||
if media_type not in ["movie", "series"]:
|
||||
return {"streams": []}
|
||||
|
||||
if "tmdb:" in media_id:
|
||||
return {"streams": []}
|
||||
|
||||
@@ -167,8 +172,27 @@ async def stream(
|
||||
async with aiohttp.ClientSession(connector=connector) as session:
|
||||
metadata_scraper = MetadataScraper(session)
|
||||
|
||||
# First, check if metadata is already cached
|
||||
id, season, episode = parse_media_id(media_type, media_id)
|
||||
|
||||
# Digital Release Filter
|
||||
if settings.DIGITAL_RELEASE_FILTER:
|
||||
is_released = await release_filter.check_is_released(
|
||||
session, media_type, media_id, season, episode
|
||||
)
|
||||
|
||||
if not is_released:
|
||||
logger.log("FILTER", f"🚫 {media_id} is not released yet. Skipping.")
|
||||
return {
|
||||
"streams": [
|
||||
{
|
||||
"name": "[🚫] Comet",
|
||||
"description": "Content not digitally released yet.",
|
||||
"url": "https://comet.fast",
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
# Check if metadata is already cached
|
||||
cached_metadata = await metadata_scraper.get_cached(
|
||||
id, season if "kitsu" not in media_id else 1, episode
|
||||
)
|
||||
@@ -390,14 +414,25 @@ async def stream(
|
||||
and debrid_service != "torrent"
|
||||
):
|
||||
logger.log("SCRAPER", "🔄 Checking availability on debrid service...")
|
||||
await debrid_service_instance.get_and_cache_availability(
|
||||
session,
|
||||
torrent_manager.torrents,
|
||||
media_id,
|
||||
media_only_id,
|
||||
season,
|
||||
episode,
|
||||
)
|
||||
try:
|
||||
await debrid_service_instance.get_and_cache_availability(
|
||||
session,
|
||||
torrent_manager.torrents,
|
||||
media_id,
|
||||
media_only_id,
|
||||
season,
|
||||
episode,
|
||||
)
|
||||
except DebridAuthError as e:
|
||||
return {
|
||||
"streams": [
|
||||
{
|
||||
"name": "[❌] Comet",
|
||||
"description": e.display_message,
|
||||
"url": "https://comet.fast",
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
if debrid_service != "torrent":
|
||||
cached_count = sum(
|
||||
|
||||
+27
-2
@@ -294,10 +294,10 @@ async def setup_database():
|
||||
)
|
||||
|
||||
await database.execute(
|
||||
f"""
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS bandwidth_stats (
|
||||
id INTEGER PRIMARY KEY,
|
||||
total_bytes {"INTEGER" if settings.DATABASE_TYPE == "sqlite" else "BIGINT"},
|
||||
total_bytes BIGINT,
|
||||
last_updated INTEGER
|
||||
)
|
||||
"""
|
||||
@@ -361,6 +361,16 @@ async def setup_database():
|
||||
"""
|
||||
)
|
||||
|
||||
await database.execute(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS digital_release_cache (
|
||||
media_id TEXT PRIMARY KEY,
|
||||
release_date BIGINT,
|
||||
timestamp INTEGER
|
||||
)
|
||||
"""
|
||||
)
|
||||
|
||||
# =============================================================================
|
||||
# TORRENTS TABLE INDEXES - Most critical for performance
|
||||
# =============================================================================
|
||||
@@ -617,6 +627,13 @@ async def setup_database():
|
||||
"""
|
||||
)
|
||||
|
||||
await database.execute(
|
||||
"""
|
||||
CREATE INDEX IF NOT EXISTS idx_digital_release_timestamp
|
||||
ON digital_release_cache (timestamp)
|
||||
"""
|
||||
)
|
||||
|
||||
if settings.DATABASE_TYPE == "sqlite":
|
||||
await database.execute("PRAGMA busy_timeout=30000") # 30 seconds timeout
|
||||
await database.execute("PRAGMA journal_mode=WAL")
|
||||
@@ -697,6 +714,14 @@ async def _run_startup_cleanup():
|
||||
{"cache_ttl": settings.DEBRID_CACHE_TTL, "current_time": current_time},
|
||||
)
|
||||
|
||||
await database.execute(
|
||||
"""
|
||||
DELETE FROM digital_release_cache
|
||||
WHERE timestamp + :cache_ttl < :current_time;
|
||||
""",
|
||||
{"cache_ttl": settings.METADATA_CACHE_TTL, "current_time": current_time},
|
||||
)
|
||||
|
||||
await database.execute("DELETE FROM download_links_cache")
|
||||
|
||||
await database.execute(
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
import atexit
|
||||
import multiprocessing
|
||||
import os
|
||||
import signal
|
||||
from concurrent.futures import ProcessPoolExecutor
|
||||
|
||||
_mp_context = None
|
||||
try:
|
||||
_mp_context = multiprocessing.get_context("forkserver")
|
||||
except ValueError:
|
||||
_mp_context = multiprocessing.get_context("spawn")
|
||||
|
||||
app_executor = None
|
||||
|
||||
|
||||
def worker_initializer():
|
||||
signal.signal(signal.SIGINT, signal.SIG_IGN)
|
||||
|
||||
|
||||
def setup_executor(max_workers: int | None = None):
|
||||
global app_executor
|
||||
|
||||
if max_workers is None:
|
||||
cpu_count = os.cpu_count() or 1
|
||||
max_workers = min(cpu_count, 4)
|
||||
|
||||
app_executor = ProcessPoolExecutor(
|
||||
max_workers=max_workers, mp_context=_mp_context, initializer=worker_initializer
|
||||
)
|
||||
|
||||
|
||||
def shutdown_executor():
|
||||
global app_executor
|
||||
if app_executor:
|
||||
app_executor.shutdown(wait=True, cancel_futures=True)
|
||||
app_executor = None
|
||||
|
||||
|
||||
atexit.register(shutdown_executor)
|
||||
|
||||
|
||||
def get_executor():
|
||||
return app_executor
|
||||
@@ -62,6 +62,12 @@ CUSTOM_LOG_LEVELS = {
|
||||
"loguru_color": "<fg #E24A90>",
|
||||
"no": 20,
|
||||
},
|
||||
"FILTER": {
|
||||
"color": "#FFD700",
|
||||
"icon": "🛡️",
|
||||
"loguru_color": "<fg #FFD700>",
|
||||
"no": 35,
|
||||
},
|
||||
}
|
||||
|
||||
ALL_LOG_LEVELS = {**STANDARD_LOG_LEVELS, **CUSTOM_LOG_LEVELS}
|
||||
|
||||
@@ -390,6 +390,13 @@ def log_startup_info(settings):
|
||||
)
|
||||
|
||||
logger.log("COMET", f"Remove Adult Content: {bool(settings.REMOVE_ADULT_CONTENT)}")
|
||||
logger.log(
|
||||
"COMET", f"Digital Release Filter: {bool(settings.DIGITAL_RELEASE_FILTER)}"
|
||||
)
|
||||
logger.log(
|
||||
"COMET",
|
||||
f"TMDB Read Access Token: {settings.TMDB_READ_ACCESS_TOKEN if settings.TMDB_READ_ACCESS_TOKEN else 'Shared'}",
|
||||
)
|
||||
logger.log("COMET", f"Custom Header HTML: {bool(settings.CUSTOM_HEADER_HTML)}")
|
||||
|
||||
background_scraper_display = (
|
||||
|
||||
@@ -129,6 +129,8 @@ class AppSettings(BaseSettings):
|
||||
BACKGROUND_SCRAPER_MAX_SERIES_PER_RUN: Optional[int] = 100
|
||||
ANIME_MAPPING_SOURCE: Optional[str] = "remote"
|
||||
ANIME_MAPPING_REFRESH_INTERVAL: Optional[int] = 86400
|
||||
DIGITAL_RELEASE_FILTER: Optional[bool] = False
|
||||
TMDB_READ_ACCESS_TOKEN: Optional[str] = None
|
||||
|
||||
@field_validator("INDEXER_MANAGER_TYPE")
|
||||
def set_indexer_manager_type(cls, v, values):
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
class DebridError(Exception):
|
||||
"""Base exception for debrid-related errors."""
|
||||
|
||||
def __init__(self, message: str, display_message: str = None):
|
||||
self.message = message
|
||||
self.display_message = display_message or message
|
||||
super().__init__(self.message)
|
||||
|
||||
|
||||
class DebridAuthError(DebridError):
|
||||
"""Raised when debrid authentication fails (not premium, invalid API key, etc.)."""
|
||||
|
||||
def __init__(self, debrid_name: str, message: str = None):
|
||||
self.debrid_name = debrid_name
|
||||
default_message = f"{debrid_name}: Authentication failed or not premium"
|
||||
display_message = (
|
||||
message
|
||||
or f"{debrid_name}: Invalid API key or no active subscription.\nPlease check your debrid account."
|
||||
)
|
||||
super().__init__(default_message, display_message)
|
||||
+48
-11
@@ -4,13 +4,19 @@ from urllib.parse import quote, unquote
|
||||
import aiohttp
|
||||
from RTN import parse, title_match
|
||||
|
||||
from comet.core.execution import get_executor
|
||||
from comet.core.logger import logger
|
||||
from comet.core.models import settings
|
||||
from comet.debrid.exceptions import DebridAuthError
|
||||
from comet.services.debrid_cache import cache_availability
|
||||
from comet.services.torrent_manager import torrent_update_queue
|
||||
from comet.utils.parsing import is_video
|
||||
|
||||
|
||||
def batch_parse(filenames):
|
||||
return [parse(f) for f in filenames]
|
||||
|
||||
|
||||
class StremThru:
|
||||
def __init__(
|
||||
self,
|
||||
@@ -43,17 +49,29 @@ class StremThru:
|
||||
|
||||
async def check_premium(self):
|
||||
try:
|
||||
user = await self.session.get(
|
||||
response = await self.session.get(
|
||||
f"{self.base_url}/user?client_ip={self.client_ip}"
|
||||
)
|
||||
user = await user.json()
|
||||
return user["data"]["subscription_status"] == "premium"
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
f"Exception while checking premium status on {self.name}: {e}"
|
||||
)
|
||||
user = await response.json()
|
||||
|
||||
return False
|
||||
if "data" not in user:
|
||||
raise DebridAuthError(
|
||||
self.name,
|
||||
f"{self.name}: Invalid API key.\nPlease check your configuration.",
|
||||
)
|
||||
|
||||
if user["data"]["subscription_status"] != "premium":
|
||||
raise DebridAuthError(
|
||||
self.name,
|
||||
f"{self.name}: No active subscription.\nPlease renew your debrid account.",
|
||||
)
|
||||
except DebridAuthError:
|
||||
raise
|
||||
except Exception as e:
|
||||
raise DebridAuthError(
|
||||
self.name,
|
||||
f"{self.name}: Failed to check account status.\n{e}",
|
||||
)
|
||||
|
||||
async def get_instant(self, magnets: list):
|
||||
try:
|
||||
@@ -72,8 +90,7 @@ class StremThru:
|
||||
tracker_map: dict,
|
||||
sources_map: dict,
|
||||
):
|
||||
if not await self.check_premium():
|
||||
return []
|
||||
await self.check_premium()
|
||||
|
||||
chunk_size = 50
|
||||
chunks = [
|
||||
@@ -95,6 +112,26 @@ class StremThru:
|
||||
|
||||
is_offcloud = self.real_debrid_name == "offcloud"
|
||||
|
||||
filenames_to_parse = []
|
||||
if not is_offcloud:
|
||||
for result in availability:
|
||||
for torrent in result:
|
||||
if torrent["status"] != "cached":
|
||||
continue
|
||||
for file in torrent["files"]:
|
||||
filename = file["name"].split("/")[-1]
|
||||
if not is_video(filename) or "sample" in filename.lower():
|
||||
continue
|
||||
filenames_to_parse.append(filename)
|
||||
|
||||
parsed_iter = iter([])
|
||||
if filenames_to_parse:
|
||||
loop = asyncio.get_running_loop()
|
||||
parsed_results = await loop.run_in_executor(
|
||||
get_executor(), batch_parse, filenames_to_parse
|
||||
)
|
||||
parsed_iter = iter(parsed_results)
|
||||
|
||||
files = []
|
||||
cached_count = 0
|
||||
for result in availability:
|
||||
@@ -127,7 +164,7 @@ class StremThru:
|
||||
if not is_video(filename) or "sample" in filename.lower():
|
||||
continue
|
||||
|
||||
filename_parsed = parse(filename)
|
||||
filename_parsed = next(parsed_iter)
|
||||
|
||||
season = (
|
||||
filename_parsed.seasons[0]
|
||||
|
||||
+11
-13
@@ -44,20 +44,18 @@ def run_with_uvicorn():
|
||||
workers=settings.FASTAPI_WORKERS,
|
||||
log_config=None,
|
||||
)
|
||||
server = Server(config=config)
|
||||
server = uvicorn.Server(config=config)
|
||||
|
||||
with server.run_in_thread():
|
||||
log_startup_info(settings)
|
||||
try:
|
||||
while True:
|
||||
time.sleep(1) # Keep the main thread alive
|
||||
except KeyboardInterrupt:
|
||||
logger.log("COMET", "Server stopped by user")
|
||||
except Exception as e:
|
||||
logger.error(f"Unexpected error: {e}")
|
||||
logger.exception(traceback.format_exc())
|
||||
finally:
|
||||
logger.log("COMET", "Server Shutdown")
|
||||
log_startup_info(settings)
|
||||
try:
|
||||
server.run()
|
||||
except KeyboardInterrupt:
|
||||
logger.log("COMET", "Server stopped by user")
|
||||
except Exception as e:
|
||||
logger.error(f"Unexpected error: {e}")
|
||||
logger.exception(traceback.format_exc())
|
||||
finally:
|
||||
logger.log("COMET", "Server Shutdown")
|
||||
|
||||
|
||||
def run_with_gunicorn():
|
||||
|
||||
@@ -0,0 +1,100 @@
|
||||
import time
|
||||
from datetime import datetime
|
||||
|
||||
from comet.core.logger import logger
|
||||
from comet.core.models import database, settings
|
||||
from comet.metadata.tmdb import TMDBApi
|
||||
|
||||
|
||||
class DigitalReleaseFilter:
|
||||
async def check_is_released(
|
||||
self,
|
||||
session,
|
||||
media_type: str,
|
||||
media_id: str,
|
||||
season: int = None,
|
||||
episode: int = None,
|
||||
):
|
||||
try:
|
||||
cached_date = await database.fetch_val(
|
||||
"""
|
||||
SELECT release_date FROM digital_release_cache
|
||||
WHERE media_id = :media_id
|
||||
AND timestamp + :cache_ttl >= :current_time
|
||||
""",
|
||||
{
|
||||
"media_id": media_id,
|
||||
"cache_ttl": settings.METADATA_CACHE_TTL,
|
||||
"current_time": time.time(),
|
||||
},
|
||||
)
|
||||
|
||||
if cached_date is not None:
|
||||
return self._is_released(cached_date)
|
||||
|
||||
tmdb_id = None
|
||||
|
||||
tmdb = TMDBApi(session)
|
||||
|
||||
if media_id.startswith("tt"):
|
||||
imdb_id = media_id.split(":")[0]
|
||||
tmdb_id = await tmdb.get_tmdb_id_from_imdb(imdb_id)
|
||||
else:
|
||||
# Other formats (e.g. kitsu) are not supported
|
||||
return True
|
||||
|
||||
if not tmdb_id:
|
||||
logger.warning(
|
||||
f"DigitalReleaseFilter: Could not resolve {media_id} to TMDB ID. Allowing search."
|
||||
)
|
||||
return True
|
||||
|
||||
release_date_str = None
|
||||
if media_type == "movie":
|
||||
release_date_str = await tmdb.get_upcoming_movie_release_date(tmdb_id)
|
||||
elif media_type == "series":
|
||||
release_date_str = await tmdb.get_episode_air_date(
|
||||
tmdb_id, season, episode
|
||||
)
|
||||
|
||||
cache_timestamp = int(time.time())
|
||||
if release_date_str is None:
|
||||
# Not found, treat as released in far future to block
|
||||
release_date_timestamp = 253402300799 # 9999-12-31
|
||||
# Cache for only 1 day (86400s) to recheck later
|
||||
if settings.METADATA_CACHE_TTL > 86400:
|
||||
cache_timestamp = int(
|
||||
time.time() - settings.METADATA_CACHE_TTL + 86400
|
||||
)
|
||||
else:
|
||||
release_date_timestamp = int(
|
||||
datetime.strptime(release_date_str, "%Y-%m-%d").timestamp()
|
||||
)
|
||||
|
||||
await database.execute(
|
||||
"""
|
||||
INSERT INTO digital_release_cache (media_id, release_date, timestamp)
|
||||
VALUES (:media_id, :release_date, :timestamp)
|
||||
ON CONFLICT (media_id) DO UPDATE SET release_date = :release_date, timestamp = :timestamp
|
||||
""",
|
||||
{
|
||||
"media_id": media_id,
|
||||
"release_date": release_date_timestamp,
|
||||
"timestamp": cache_timestamp,
|
||||
},
|
||||
)
|
||||
|
||||
return self._is_released(release_date_timestamp)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
f"DigitalReleaseFilter: Error checking release status for {media_id}: {e}"
|
||||
)
|
||||
return True
|
||||
|
||||
def _is_released(self, release_timestamp: float):
|
||||
if release_timestamp is None:
|
||||
return True
|
||||
return release_timestamp <= time.time()
|
||||
|
||||
|
||||
release_filter = DigitalReleaseFilter()
|
||||
@@ -0,0 +1,79 @@
|
||||
import aiohttp
|
||||
|
||||
from comet.core.logger import logger
|
||||
from comet.core.models import settings
|
||||
|
||||
DEFAULT_TMDB_READ_ACCESS_TOKEN = "eyJhbGciOiJIUzI1NiJ9.eyJhdWQiOiJlNTkxMmVmOWFhM2IxNzg2Zjk3ZTE1NWY1YmQ3ZjY1MSIsInN1YiI6IjY1M2NjNWUyZTg5NGE2MDBmZjE2N2FmYyIsInNjb3BlcyI6WyJhcGlfcmVhZCJdLCJ2ZXJzaW9uIjoxfQ.xrIXsMFJpI1o1j5g2QpQcFP1X3AfRjFA5FlBFO5Naw8"
|
||||
|
||||
|
||||
class TMDBApi:
|
||||
def __init__(self, session: aiohttp.ClientSession):
|
||||
self.session = session
|
||||
self.base_url = "https://api.themoviedb.org/3"
|
||||
self.headers = {
|
||||
"Authorization": f"Bearer {settings.TMDB_READ_ACCESS_TOKEN if settings.TMDB_READ_ACCESS_TOKEN else DEFAULT_TMDB_READ_ACCESS_TOKEN}",
|
||||
"Content-Type": "application/json",
|
||||
}
|
||||
|
||||
async def get_upcoming_movie_release_date(self, tmdb_id: str):
|
||||
try:
|
||||
url = f"{self.base_url}/movie/{tmdb_id}/release_dates"
|
||||
async with self.session.get(url, headers=self.headers) as response:
|
||||
if response.status != 200:
|
||||
return None
|
||||
|
||||
data = await response.json()
|
||||
|
||||
release_dates = []
|
||||
for result in data.get("results", []):
|
||||
for release in result.get("release_dates", []):
|
||||
if release.get("type") in [4, 5]: # Digital or Physical
|
||||
date_str = release.get("release_date", "").split("T")[0]
|
||||
if date_str:
|
||||
release_dates.append(date_str)
|
||||
|
||||
if release_dates:
|
||||
return min(release_dates)
|
||||
|
||||
return None
|
||||
except Exception as e:
|
||||
logger.error(f"TMDB: Error getting movie release date for {tmdb_id}: {e}")
|
||||
return None
|
||||
|
||||
async def get_episode_air_date(self, tmdb_id: str, season: int, episode: int):
|
||||
try:
|
||||
url = f"{self.base_url}/tv/{tmdb_id}/season/{season}/episode/{episode}"
|
||||
async with self.session.get(url, headers=self.headers) as response:
|
||||
if response.status != 200:
|
||||
return None
|
||||
|
||||
data = await response.json()
|
||||
return data.get("air_date")
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
f"TMDB: Error getting episode air date for {tmdb_id} S{season}E{episode}: {e}"
|
||||
)
|
||||
return None
|
||||
|
||||
async def get_tmdb_id_from_imdb(self, imdb_id: str):
|
||||
try:
|
||||
url = f"{self.base_url}/find/{imdb_id}?external_source=imdb_id"
|
||||
async with self.session.get(url, headers=self.headers) as response:
|
||||
if response.status != 200:
|
||||
text = await response.text()
|
||||
logger.error(
|
||||
f"TMDB: Failed to get TMDB ID from IMDB ID {imdb_id}: {text}"
|
||||
)
|
||||
return None
|
||||
|
||||
data = await response.json()
|
||||
|
||||
if data.get("movie_results"):
|
||||
return str(data["movie_results"][0]["id"])
|
||||
if data.get("tv_results"):
|
||||
return str(data["tv_results"][0]["id"])
|
||||
|
||||
return None
|
||||
except Exception as e:
|
||||
logger.error(f"TMDB: Error converting IMDB ID {imdb_id}: {e}")
|
||||
return None
|
||||
@@ -5,12 +5,13 @@ import aiohttp
|
||||
import orjson
|
||||
from RTN import DefaultRanking, ParsedData
|
||||
|
||||
from comet.core.execution import get_executor
|
||||
from comet.core.logger import logger
|
||||
from comet.core.models import CometSettingsModel, database, settings
|
||||
from comet.scrapers.manager import scraper_manager
|
||||
from comet.services.filtering import filter_worker
|
||||
from comet.services.ranking import rank_worker
|
||||
from comet.utils.parsing import default_dump
|
||||
from comet.services.torrent_manager import torrent_update_queue
|
||||
|
||||
|
||||
class TorrentManager:
|
||||
@@ -72,7 +73,7 @@ class TorrentManager:
|
||||
async for scraper_name, results in scraper_manager.scrape_all(request, session):
|
||||
await self.filter_manager(scraper_name, results)
|
||||
|
||||
await self.cache_torrents()
|
||||
asyncio.create_task(self.cache_torrents())
|
||||
|
||||
for torrent in self.ready_to_cache:
|
||||
season = torrent["parsed"].seasons[0] if torrent["parsed"].seasons else None
|
||||
@@ -128,41 +129,24 @@ class TorrentManager:
|
||||
}
|
||||
|
||||
async def cache_torrents(self):
|
||||
current_time = time.time()
|
||||
values = [
|
||||
{
|
||||
"media_id": self.media_only_id,
|
||||
for torrent in self.ready_to_cache:
|
||||
file_info = {
|
||||
"info_hash": torrent["infoHash"],
|
||||
"file_index": int(torrent["fileIndex"])
|
||||
if torrent["fileIndex"] is not None
|
||||
else None,
|
||||
"index": torrent["fileIndex"],
|
||||
"title": torrent["title"],
|
||||
"size": torrent["size"],
|
||||
"season": torrent["parsed"].seasons[0]
|
||||
if torrent["parsed"].seasons
|
||||
else self.season,
|
||||
"episode": torrent["parsed"].episodes[0]
|
||||
if torrent["parsed"].episodes
|
||||
else None,
|
||||
"title": torrent["title"],
|
||||
"seeders": int(torrent["seeders"])
|
||||
if torrent["seeders"] is not None
|
||||
else None,
|
||||
"size": int(torrent["size"]) if torrent["size"] is not None else None,
|
||||
"parsed": torrent["parsed"],
|
||||
"seeders": torrent["seeders"],
|
||||
"tracker": torrent["tracker"],
|
||||
"sources": orjson.dumps(torrent["sources"]).decode("utf-8"),
|
||||
"parsed": orjson.dumps(torrent["parsed"], default_dump).decode("utf-8"),
|
||||
"timestamp": current_time,
|
||||
"sources": torrent["sources"],
|
||||
}
|
||||
for torrent in self.ready_to_cache
|
||||
]
|
||||
|
||||
query = f"""
|
||||
INSERT {"OR REPLACE " if settings.DATABASE_TYPE == "sqlite" else ""}
|
||||
INTO torrents
|
||||
VALUES (:media_id, :info_hash, :file_index, :season, :episode, :title, :seeders, :size, :tracker, :sources, :parsed, :timestamp)
|
||||
{" ON CONFLICT DO NOTHING" if settings.DATABASE_TYPE == "postgresql" else ""}
|
||||
"""
|
||||
|
||||
await database.execute_many(query, values)
|
||||
await torrent_update_queue.add_torrent_info(file_info, self.media_only_id)
|
||||
|
||||
async def filter_manager(self, scraper_name: str, torrents: list):
|
||||
if len(torrents) == 0:
|
||||
@@ -188,10 +172,10 @@ class TorrentManager:
|
||||
return
|
||||
|
||||
loop = asyncio.get_running_loop()
|
||||
chunk_size = 50
|
||||
chunk_size = 20
|
||||
tasks = [
|
||||
loop.run_in_executor(
|
||||
None,
|
||||
get_executor(),
|
||||
filter_worker,
|
||||
new_torrents[i : i + chunk_size],
|
||||
self.title,
|
||||
@@ -217,7 +201,7 @@ class TorrentManager:
|
||||
):
|
||||
loop = asyncio.get_running_loop()
|
||||
self.ranked_torrents = await loop.run_in_executor(
|
||||
None,
|
||||
get_executor(),
|
||||
rank_worker,
|
||||
self.torrents,
|
||||
self.debrid_service,
|
||||
|
||||
@@ -1,11 +1,10 @@
|
||||
import asyncio
|
||||
import base64
|
||||
import hashlib
|
||||
import html
|
||||
import re
|
||||
import time
|
||||
from collections import defaultdict
|
||||
from urllib.parse import parse_qs, urlparse
|
||||
from urllib.parse import unquote
|
||||
|
||||
import aiohttp
|
||||
import anyio
|
||||
@@ -20,15 +19,14 @@ from comet.core.logger import logger
|
||||
from comet.core.models import database, settings
|
||||
from comet.utils.parsing import default_dump, is_video
|
||||
|
||||
TRACKER_PATTERN = re.compile(r"[&?]tr=([^&]+)")
|
||||
INFO_HASH_PATTERN = re.compile(r"btih:([a-fA-F0-9]{40}|[a-zA-Z0-9]{32})")
|
||||
|
||||
|
||||
def extract_trackers_from_magnet(magnet_uri: str):
|
||||
try:
|
||||
decoded_uri = html.unescape(magnet_uri)
|
||||
parsed = urlparse(decoded_uri)
|
||||
params = parse_qs(parsed.query)
|
||||
return params.get("tr", [])
|
||||
trackers = TRACKER_PATTERN.findall(magnet_uri)
|
||||
return [unquote(tracker) for tracker in trackers]
|
||||
except Exception as e:
|
||||
logger.warning(f"Failed to extract trackers from magnet URI: {e}")
|
||||
return []
|
||||
@@ -247,7 +245,7 @@ add_torrent_queue = AddTorrentQueue()
|
||||
|
||||
|
||||
class TorrentUpdateQueue:
|
||||
def __init__(self, batch_size: int = 100, flush_interval: float = 5.0):
|
||||
def __init__(self, batch_size: int = 1000, flush_interval: float = 5.0):
|
||||
self.queue = asyncio.Queue()
|
||||
self.batch_size = batch_size
|
||||
self.flush_interval = flush_interval
|
||||
|
||||
Reference in New Issue
Block a user