From 99155271bbfc3301e0374e8d28072ddaece516f6 Mon Sep 17 00:00:00 2001 From: KiraStream Date: Tue, 19 May 2026 15:42:37 +0000 Subject: [PATCH] =?UTF-8?q?feat:=20domain=20health=20check=20autom=C3=A1ti?= =?UTF-8?q?co=20con=20stickiness=20de=20dominio?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - health_checker: usa TCP connect en lugar de Xtream API para dominios CDN (correcto para dominios que solo sirven streams, no la API de credenciales) Intervalo reducido a 3 min; detecta recuperación automática y lo registra en log - pool: _build_stream_urls excluye dominios con status='error'; si todos están caídos usa todos como fallback para no dejar el stream sin opciones - restream: stickiness de dominio — en renovación de sesión normal reconecta al mismo dominio siempre; solo rota al siguiente en fallos reales (<25s) - restream: _mark_domain_down marca el dominio como error en BD inmediatamente al detectar un fallo rápido, sin esperar al próximo ciclo del health checker - database: usa ADD COLUMN IF NOT EXISTS para la migración (evita abortar transacción) Co-Authored-By: Claude Sonnet 4.6 --- backend/core/health_checker.py | 91 +++++++++++++++++++++--- backend/core/pool.py | 33 ++++++++- backend/core/restream.py | 125 ++++++++++++++++++++------------- backend/database.py | 9 +-- 4 files changed, 189 insertions(+), 69 deletions(-) diff --git a/backend/core/health_checker.py b/backend/core/health_checker.py index 66ec71b..8f4bc22 100644 --- a/backend/core/health_checker.py +++ b/backend/core/health_checker.py @@ -1,4 +1,11 @@ -"""Background health checker for provider URLs.""" +"""Background health checker for provider URLs. + +Two types of checks: +- _check_domain_reachable: TCP connectivity (DNS + port open). Used by the background + loop for CDN/stream domains — no auth needed, just "can we reach this host?". +- _check_xtream_api: Full Xtream API call with credentials. Used by the admin API + endpoint to verify that provider account credentials are valid. +""" import asyncio import logging import time @@ -12,10 +19,43 @@ from ..models.provider import ProviderAccount, ProviderUrl logger = logging.getLogger(__name__) -CHECK_INTERVAL = 300 # 5 minutes +CHECK_INTERVAL = 180 # 3 minutes between full sweeps +CONNECT_TIMEOUT = 6.0 # seconds for TCP connect -async def _check_url(base_url: str, username: str, password: str) -> tuple[str, int | None]: +async def _check_domain_reachable(url: str) -> tuple[str, int | None]: + """DNS resolution + TCP connect to port 80/443. + + Returns ("ok"|"error"|"timeout", response_ms). + Does NOT need credentials — works for any CDN stream domain. + A domain is "ok" if we can open a TCP connection, regardless of + what the HTTP layer returns. "error" means DNS failure or refused. + """ + from urllib.parse import urlparse + parsed = urlparse(url) + host = parsed.hostname or url + port = parsed.port or (443 if parsed.scheme == "https" else 80) + start = time.time() + try: + _, writer = await asyncio.wait_for( + asyncio.open_connection(host, port), + timeout=CONNECT_TIMEOUT, + ) + writer.close() + try: + await writer.wait_closed() + except Exception: + pass + return "ok", int((time.time() - start) * 1000) + except asyncio.TimeoutError: + return "timeout", int(CONNECT_TIMEOUT * 1000) + except OSError: + # Covers: DNS NXDOMAIN, connection refused, network unreachable + return "error", int((time.time() - start) * 1000) + + +async def _check_xtream_api(base_url: str, username: str, password: str) -> tuple[str, int | None]: + """Full Xtream Codes API check with credentials (for admin UI only).""" url = f"{base_url.rstrip('/')}/player_api.php" params = {"username": username, "password": password, "action": "user_info"} start = time.time() @@ -39,7 +79,11 @@ async def _check_url(base_url: str, username: str, password: str) -> tuple[str, async def check_provider_urls(provider_account_id: int) -> list[dict]: - """Check all URLs for a specific provider. Returns list of result dicts.""" + """Check all URLs for a specific provider (admin API endpoint). + + Uses the Xtream API check for the primary URL and TCP check for extras, + so the admin panel shows meaningful status for both types. + """ async with AsyncSessionLocal() as db: result = await db.execute( select(ProviderUrl.id, ProviderUrl.url, ProviderAccount.username, ProviderAccount.password) @@ -48,13 +92,20 @@ async def check_provider_urls(provider_account_id: int) -> list[dict]: ProviderUrl.provider_account_id == provider_account_id, ProviderUrl.is_active == True, # noqa: E712 ) + .order_by(ProviderUrl.priority) ) rows = result.all() if not rows: return [] - tasks = [_check_url(url, username, password) for _, url, username, password in rows] + # Primary (priority=0) → Xtream API check; extras → TCP check + tasks = [] + for i, (_, url, username, password) in enumerate(rows): + if i == 0: + tasks.append(_check_xtream_api(url, username, password)) + else: + tasks.append(_check_domain_reachable(url)) results = await asyncio.gather(*tasks, return_exceptions=True) now = datetime.now(timezone.utc) @@ -77,14 +128,19 @@ async def check_provider_urls(provider_account_id: int) -> list[dict]: async def health_check_loop() -> None: - """Runs forever, checking all active provider URLs every CHECK_INTERVAL seconds.""" - logger.info("Provider URL health checker started") + """Runs forever, checking all active provider URLs every CHECK_INTERVAL seconds. + + Uses TCP connectivity check — works for CDN stream domains without credentials. + Domains previously marked 'error' that are now reachable get restored to 'ok' + automatically. + """ + logger.info(f"Provider URL health checker started (interval={CHECK_INTERVAL}s)") while True: await asyncio.sleep(CHECK_INTERVAL) try: async with AsyncSessionLocal() as db: result = await db.execute( - select(ProviderUrl.id, ProviderUrl.url, ProviderAccount.username, ProviderAccount.password) + select(ProviderUrl.id, ProviderUrl.url) .join(ProviderAccount) .where(ProviderUrl.is_active == True) # noqa: E712 ) @@ -93,23 +149,36 @@ async def health_check_loop() -> None: if not rows: continue - tasks = [_check_url(url, username, password) for _, url, username, password in rows] + tasks = [_check_domain_reachable(url) for _, url in rows] check_results = await asyncio.gather(*tasks, return_exceptions=True) now = datetime.now(timezone.utc) + recovered, degraded = [], [] async with AsyncSessionLocal() as db: - for (url_id, _, _, _), check_result in zip(rows, check_results): + for (url_id, url), check_result in zip(rows, check_results): if isinstance(check_result, Exception): status, response_ms = "error", None else: status, response_ms = check_result pu = await db.get(ProviderUrl, url_id) if pu: + prev = pu.status pu.status = status pu.response_ms = response_ms pu.last_checked_at = now + if prev == "error" and status == "ok": + recovered.append(url) + elif prev == "ok" and status == "error": + degraded.append(url) await db.commit() - logger.info(f"Health check done: {len(rows)} URLs checked") + if recovered: + logger.info(f"Health check: domains RECOVERED → {recovered}") + if degraded: + logger.warning(f"Health check: domains DOWN → {degraded}") + logger.info(f"Health check done: {len(rows)} URLs checked" + f" ({sum(1 for r in check_results if not isinstance(r, Exception) and r[0]=='ok')} ok," + f" {sum(1 for r in check_results if not isinstance(r, Exception) and r[0]=='error')} error," + f" {sum(1 for r in check_results if not isinstance(r, Exception) and r[0]=='timeout')} timeout)") except Exception as e: logger.error(f"Health check loop error: {e}") diff --git a/backend/core/pool.py b/backend/core/pool.py index 54909f1..67b7d50 100644 --- a/backend/core/pool.py +++ b/backend/core/pool.py @@ -241,7 +241,13 @@ class ProviderPool: return preferred + rest async def _build_stream_urls(self, account_id: int, primary_stream_url: str) -> list[str]: - """Build list of stream URLs for failover: one per active ProviderUrl, substituting base URL.""" + """Build list of stream URLs for multi-domain failover. + + Only includes domains with status != 'error' (healthy or not-yet-checked). + If ALL domains are currently marked 'error', falls back to using all of them + so the stream isn't stranded — the health checker will restore them when + reachable again. + """ from urllib.parse import urlparse, urlunparse from ..models.provider import ProviderUrl from ..database import AsyncSessionLocal @@ -250,19 +256,39 @@ class ProviderPool: parsed = urlparse(primary_stream_url) async with AsyncSessionLocal() as db: + # Prefer healthy/unknown domains only result = await db.execute( select(ProviderUrl) .where( ProviderUrl.provider_account_id == account_id, - ProviderUrl.is_active == True, # noqa: E712 + ProviderUrl.is_active == True, # noqa: E712 + ProviderUrl.status != "error", ) .order_by(ProviderUrl.priority) ) provider_urls = result.scalars().all() + if not provider_urls: + # Fallback: all active domains even if all are 'error' + result = await db.execute( + select(ProviderUrl) + .where( + ProviderUrl.provider_account_id == account_id, + ProviderUrl.is_active == True, # noqa: E712 + ) + .order_by(ProviderUrl.priority) + ) + provider_urls = result.scalars().all() + if provider_urls: + logger.warning( + f"provider {account_id}: all domains marked 'error' — " + f"using all {len(provider_urls)} as fallback" + ) + if not provider_urls: return [primary_stream_url] + skipped = [] urls = [] for pu in provider_urls: alt = urlparse(pu.url.rstrip('/')) @@ -275,6 +301,9 @@ class ProviderPool: parsed.fragment, )) urls.append(reconstructed) + + if skipped: + logger.debug(f"provider {account_id}: skipped error domains {skipped}") return urls async def _find_free_slot(self, provider_maps: list[dict]) -> tuple[int | None, str | None]: diff --git a/backend/core/restream.py b/backend/core/restream.py index 2f72494..be38bb8 100644 --- a/backend/core/restream.py +++ b/backend/core/restream.py @@ -426,13 +426,43 @@ class BroadcastGroup: except Exception as e: logger.debug(f"probe_info failed for ch{self.channel_id}: {e}") + async def _mark_domain_down(self, url: str) -> None: + """Update ProviderUrl.status='error' immediately so other streams skip this domain.""" + from urllib.parse import urlparse + from datetime import datetime, timezone + from ..database import AsyncSessionLocal + from ..models.provider import ProviderUrl + from sqlalchemy import select + netloc = urlparse(url).netloc + if not netloc: + return + try: + async with AsyncSessionLocal() as db: + result = await db.execute( + select(ProviderUrl).where( + ProviderUrl.provider_account_id == self.provider_account_id, + ProviderUrl.url.ilike(f"%{netloc}%"), + ) + ) + pu = result.scalar_one_or_none() + if pu and pu.status != "error": + pu.status = "error" + pu.last_checked_at = datetime.now(timezone.utc) + await db.commit() + logger.info(f"[ch{self.channel_id}] Dominio {netloc} marcado como error (fallo rápido)") + except Exception as exc: + logger.debug(f"[ch{self.channel_id}] _mark_domain_down failed: {exc}") + async def _pump_with_retry(self) -> None: - # Failure-mode tracking: a "real failure" is any pump that exits in - # less than PROVIDER_SESSION_MIN seconds. Longer pumps are normal - # provider session timeouts (~2 min) and are invisible to the user. - in_failure_mode = False # True once at least one real failure has occurred + # Domain stickiness: on normal provider session timeouts (>PROVIDER_SESSION_MIN) + # we reconnect to the same domain — providers rotate DNS/CDN themselves, so + # staying on the same host avoids unnecessary domain cycling. + # Only on real quick failures (25 s) — reconnect to same domain + self._session_count += 1 + preferred_idx = url_idx # remember: this domain works + url_idx = preferred_idx # stay on it + + if in_failure_mode: + outage_secs = int(time.time() - (failure_start or time.time())) + asyncio.create_task(self._log_stream_event( + "recuperado", _url_netloc(current_url), + f"Restablecido tras {failure_attempts} intento(s) y {outage_secs}s de interrupción", + failure_attempts, outage_secs, + )) + logger.info( + f"[ch{self.channel_id}] RECUPERADO — tras {failure_attempts} " + f"intento(s) en {outage_secs}s, dominio estable: {_url_netloc(current_url)}" + ) + in_failure_mode = False + failure_start = None + failure_attempts = 0 + else: + logger.info( + f"[ch{self.channel_id}] Sesión del proveedor cerrada tras " + f"{pump_duration:.1f}s (normal, sesión #{self._session_count}, " + f"dominio: {_url_netloc(current_url)})" + ) + + self._status = "reconnecting" + # Reconnect immediately — no backoff for normal session rotations finally: self._status = "stopped" diff --git a/backend/database.py b/backend/database.py index 77c75cb..2e58336 100644 --- a/backend/database.py +++ b/backend/database.py @@ -32,12 +32,9 @@ async def init_db(): async with engine.begin() as conn: await conn.run_sync(Base.metadata.create_all) # Migrate: add pump_duration_secs column if it doesn't exist yet - try: - await conn.execute(text( - "ALTER TABLE stream_events ADD COLUMN pump_duration_secs INTEGER" - )) - except Exception: - pass # column already exists + await conn.execute(text( + "ALTER TABLE stream_events ADD COLUMN IF NOT EXISTS pump_duration_secs INTEGER" + )) # Migrate: create provider_urls entries from existing base_url if not already done await conn.execute(text(""" INSERT INTO provider_urls (provider_account_id, url, priority, is_active, status)