Files
KiraStream 99155271bb feat: domain health check automático con stickiness de dominio
- 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 <noreply@anthropic.com>
2026-05-19 15:42:37 +00:00

185 lines
7.3 KiB
Python

"""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
from datetime import datetime, timezone
import aiohttp
from sqlalchemy import select
from ..database import AsyncSessionLocal
from ..models.provider import ProviderAccount, ProviderUrl
logger = logging.getLogger(__name__)
CHECK_INTERVAL = 180 # 3 minutes between full sweeps
CONNECT_TIMEOUT = 6.0 # seconds for TCP connect
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()
try:
timeout = aiohttp.ClientTimeout(total=8)
async with aiohttp.ClientSession(timeout=timeout) as session:
async with session.get(url, params=params, ssl=False) as resp:
elapsed_ms = int((time.time() - start) * 1000)
if resp.status == 200:
try:
data = await resp.json(content_type=None)
if "user_info" in data:
return "ok", elapsed_ms
except Exception:
pass
return "error", elapsed_ms
except asyncio.TimeoutError:
return "timeout", None
except Exception:
return "error", None
async def check_provider_urls(provider_account_id: int) -> list[dict]:
"""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)
.join(ProviderAccount)
.where(
ProviderUrl.provider_account_id == provider_account_id,
ProviderUrl.is_active == True, # noqa: E712
)
.order_by(ProviderUrl.priority)
)
rows = result.all()
if not rows:
return []
# 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)
output = []
async with AsyncSessionLocal() as db:
for (url_id, url, _, _), check_result in zip(rows, 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:
pu.status = status
pu.response_ms = response_ms
pu.last_checked_at = now
output.append({"id": url_id, "url": url, "status": status, "response_ms": response_ms})
await db.commit()
return output
async def health_check_loop() -> None:
"""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)
.join(ProviderAccount)
.where(ProviderUrl.is_active == True) # noqa: E712
)
rows = result.all()
if not rows:
continue
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, 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()
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}")