diff --git a/.env-sample b/.env-sample index 8cd1adf..8422879 100644 --- a/.env-sample +++ b/.env-sample @@ -213,6 +213,7 @@ PROXY_DEBRID_STREAM_PASSWORD=CHANGE_ME PROXY_DEBRID_STREAM_MAX_CONNECTIONS=-1 # -1 to disable connection limits PROXY_DEBRID_STREAM_DEBRID_DEFAULT_SERVICE=realdebrid PROXY_DEBRID_STREAM_DEBRID_DEFAULT_APIKEY=CHANGE_ME +PROXY_DEBRID_STREAM_INACTIVITY_THRESHOLD=300 # Seconds to wait before cleaning up an inactive connection. Set to 0 to disable. # ============================== # # Torrent Stream Policy # diff --git a/comet/core/logger.py b/comet/core/logger.py index 11cbf69..9a8f75b 100644 --- a/comet/core/logger.py +++ b/comet/core/logger.py @@ -373,7 +373,7 @@ def log_startup_info(settings): ) debrid_stream_proxy_display = ( - f" - Password: {settings.PROXY_DEBRID_STREAM_PASSWORD} - Max Connections: {settings.PROXY_DEBRID_STREAM_MAX_CONNECTIONS} - Default Debrid Service: {settings.PROXY_DEBRID_STREAM_DEBRID_DEFAULT_SERVICE} - Default Debrid API Key: {settings.PROXY_DEBRID_STREAM_DEBRID_DEFAULT_APIKEY}" + f" - Password: {settings.PROXY_DEBRID_STREAM_PASSWORD} - Max Connections: {settings.PROXY_DEBRID_STREAM_MAX_CONNECTIONS} - Inactivity Threshold: {settings.PROXY_DEBRID_STREAM_INACTIVITY_THRESHOLD}s - Default Debrid Service: {settings.PROXY_DEBRID_STREAM_DEBRID_DEFAULT_SERVICE} - Default Debrid API Key: {settings.PROXY_DEBRID_STREAM_DEBRID_DEFAULT_APIKEY}" if settings.PROXY_DEBRID_STREAM else "" ) diff --git a/comet/core/models.py b/comet/core/models.py index 67e21ed..a5a1a29 100644 --- a/comet/core/models.py +++ b/comet/core/models.py @@ -114,6 +114,7 @@ class AppSettings(BaseSettings): PROXY_DEBRID_STREAM_MAX_CONNECTIONS: Optional[int] = -1 PROXY_DEBRID_STREAM_DEBRID_DEFAULT_SERVICE: Optional[str] = "realdebrid" PROXY_DEBRID_STREAM_DEBRID_DEFAULT_APIKEY: Optional[str] = None + PROXY_DEBRID_STREAM_INACTIVITY_THRESHOLD: Optional[int] = 300 STREMTHRU_URL: Optional[str] = "https://stremthru.13377001.xyz" DISABLE_TORRENT_STREAMS: Optional[bool] = False TORRENT_DISABLED_STREAM_NAME: Optional[str] = "[INFO] Comet" diff --git a/comet/services/bandwidth.py b/comet/services/bandwidth.py index bdeb9e9..19eef0f 100644 --- a/comet/services/bandwidth.py +++ b/comet/services/bandwidth.py @@ -4,7 +4,7 @@ import time from dataclasses import dataclass, field from comet.core.logger import logger -from comet.core.models import database +from comet.core.models import database, settings @dataclass @@ -67,7 +67,10 @@ class BandwidthMonitor: pass # Start background tasks - self._cleanup_task = asyncio.create_task(self._cleanup_inactive_connections()) + if settings.PROXY_DEBRID_STREAM_INACTIVITY_THRESHOLD > 0: + self._cleanup_task = asyncio.create_task( + self._cleanup_inactive_connections() + ) self._db_sync_task = asyncio.create_task(self._sync_to_database()) self._initialized = True @@ -150,21 +153,50 @@ class BandwidthMonitor: return f"{bytes_per_second / (1024**3):.2f} GB/s" async def _cleanup_inactive_connections(self): + threshold = settings.PROXY_DEBRID_STREAM_INACTIVITY_THRESHOLD + # min 5s, max 60s, or 1/5th of threshold + sleep_interval = max(5, min(60, threshold // 5)) + while True: try: - await asyncio.sleep(30) # Check every 30 seconds + await asyncio.sleep(sleep_interval) + current_time = time.time() with self._lock: inactive_connections = [ conn_id for conn_id, metrics in self._connections.items() - if current_time - metrics.last_update > 60 # 60 seconds timeout + if current_time - metrics.last_update > threshold ] - # Remove inactive connections - for conn_id in inactive_connections: - await self.end_connection(conn_id) + for conn_id in inactive_connections: + self._connections.pop(conn_id, None) + + if inactive_connections: + self._global_stats["active_connections"] = len( + self._connections + ) + + if not inactive_connections: + continue + + try: + placeholders = ", ".join( + f":id_{i}" for i in range(len(inactive_connections)) + ) + params = { + f"id_{i}": conn_id + for i, conn_id in enumerate(inactive_connections) + } + await database.execute( + f"DELETE FROM active_connections WHERE id IN ({placeholders})", + params, + ) + except Exception as e: + logger.warning( + f"Error batch cleaning {len(inactive_connections)} inactive connections from DB: {e}" + ) except Exception as e: logger.warning(f"Error in bandwidth monitor cleanup: {e}")