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>
This commit is contained in:
KiraStream
2026-05-19 15:42:37 +00:00
parent 4bfc7592a6
commit 99155271bb
4 changed files with 189 additions and 69 deletions
+80 -11
View File
@@ -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}")
+30 -1
View File
@@ -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,6 +256,20 @@ 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.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(
@@ -259,10 +279,16 @@ class ProviderPool:
.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]:
+75 -50
View File
@@ -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 (<PROVIDER_SESSION_MIN) do we advance to the
# next domain and immediately mark the failed one as 'error' in the DB.
in_failure_mode = False
failure_start: float | None = None
failure_attempts = 0 # consecutive real failures in this failure episode
failure_attempts = 0
preferred_idx = 0 # index of last domain that served a full session
base_wait = 0.5
url_idx = 0
n = len(self._stream_urls) or 1
@@ -444,7 +474,7 @@ class BroadcastGroup:
pump_start = time.time()
try:
await self._pump_once(current_url)
url_idx = 0 # reset to primary after successful pump
url_idx = 0
break
except asyncio.CancelledError:
raise
@@ -457,75 +487,70 @@ class BroadcastGroup:
is_real_failure = pump_duration < PROVIDER_SESSION_MIN
if is_real_failure:
# Quick exit → real provider/network error
# Quick exit → real error (DNS, refused, HTTP 4xx, stall)
self._failure_count += 1
failure_attempts += 1
cause = _classify_error(e, pump_duration)
if not in_failure_mode:
# First real failure of this episode → log caída
in_failure_mode = True
failure_start = time.time()
asyncio.create_task(self._log_stream_event(
"caida",
_url_netloc(current_url),
cause,
1,
int(pump_duration),
"caida", _url_netloc(current_url), cause, 1, int(pump_duration),
))
logger.warning(
f"[ch{self.channel_id}] CAÍDA — {cause} "
f"(pump={pump_duration:.1f}s)"
)
logger.warning(f"[ch{self.channel_id}] CAÍDA — {cause} (pump={pump_duration:.1f}s)")
else:
# Still failing — don't spam the log table
logger.warning(
f"[ch{self.channel_id}] Fallo continuado intento {failure_attempts} — "
f"{cause} (pump={pump_duration:.1f}s)"
)
else:
# Long pump → normal provider session timeout
self._session_count += 1
if in_failure_mode:
# Was in failure mode, now ran long → stream recovered
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"
)
in_failure_mode = False
failure_start = None
failure_attempts = 0
else:
# Healthy session timeout — normal, no event logged
logger.info(
f"[ch{self.channel_id}] Sesión del proveedor cerrada tras "
f"{pump_duration:.1f}s (normal, sesión #{self._session_count})"
)
# Mark domain as error immediately so other streams skip it
asyncio.create_task(self._mark_domain_down(current_url))
next_idx = (url_idx + 1) % n
url_idx = next_idx
self._status = "reconnecting"
# Rotate to next domain
url_idx = (url_idx + 1) % n
if is_real_failure:
cycle_done = (failure_attempts % n == 0)
self._status = "reconnecting"
if cycle_done:
# All domains tried — back off before next cycle
wait = min(base_wait * (2 ** (failure_attempts // n - 1)), _BACKOFF_CAP)
logger.warning(
f"[ch{self.channel_id}] All {n} domain(s) failed "
f"(attempt {failure_attempts}), retrying in {wait:.1f}s"
)
await asyncio.sleep(wait)
# Normal session timeout: reconnect immediately (no sleep needed)
# else: try next domain immediately
else:
# Normal provider session timeout (>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"
+1 -4
View File
@@ -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"
"ALTER TABLE stream_events ADD COLUMN IF NOT EXISTS pump_duration_secs INTEGER"
))
except Exception:
pass # column already exists
# 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)