Files
KiraStream 4bfc7592a6 fix: distinguir sesiones normales del proveedor de fallos reales en logs y dashboard
- Solo registra "caída" cuando ffmpeg sale en <25s (error real: DNS, 403, auth, stall)
- Timeouts de sesión del proveedor (~2 min) NO se registran como fallos — son normales
- "Recuperado" se registra cuando el stream vuelve a funcionar tras un fallo real
- Añade clasificación automática de causa: conexión rechazada, DNS, 403, timeout...
- Dashboard: "reinicios" ahora distingue sesiones del proveedor (gris, normal) vs fallos (rojo)
- Registro: nueva columna "Duración fallo", proveedor+dominio juntos, filas con color de fondo
- StreamEvent: campo pump_duration_secs para persistir duración del fallo o interrupción

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-19 15:12:04 +00:00

197 lines
6.2 KiB
Python

"""Connection log API — query and stats."""
from datetime import datetime
from fastapi import APIRouter, Depends, Query
from sqlalchemy import select, func, desc, distinct
from sqlalchemy.ext.asyncio import AsyncSession
from ...database import get_db
from ...models.log import ConnectionLog
from ...models.stream_event import StreamEvent
from ..auth import get_current_admin
router = APIRouter(prefix="/logs", tags=["logs"])
@router.get("")
async def get_logs(
username: str | None = Query(None),
channel: str | None = Query(None),
content_type: str | None = Query(None),
date_from: str | None = Query(None),
date_to: str | None = Query(None),
page: int = Query(1, ge=1),
limit: int = Query(50, ge=1, le=200),
db: AsyncSession = Depends(get_db),
_=Depends(get_current_admin),
):
q = select(ConnectionLog).order_by(desc(ConnectionLog.started_at))
if username:
q = q.where(ConnectionLog.username.ilike(f"%{username}%"))
if channel:
q = q.where(ConnectionLog.channel_name.ilike(f"%{channel}%"))
if content_type:
q = q.where(ConnectionLog.content_type == content_type)
if date_from:
try:
q = q.where(ConnectionLog.started_at >= datetime.fromisoformat(date_from))
except ValueError:
pass
if date_to:
try:
q = q.where(ConnectionLog.started_at <= datetime.fromisoformat(date_to))
except ValueError:
pass
total_result = await db.execute(select(func.count()).select_from(q.subquery()))
total = total_result.scalar() or 0
q = q.offset((page - 1) * limit).limit(limit)
result = await db.execute(q)
logs = result.scalars().all()
return {
"total": total,
"page": page,
"limit": limit,
"items": [_serialize(log) for log in logs],
}
@router.get("/stats")
async def get_stats(db: AsyncSession = Depends(get_db), _=Depends(get_current_admin)):
top_users_q = await db.execute(
select(
ConnectionLog.username,
func.count().label("sessions"),
func.sum(ConnectionLog.bytes_transferred).label("bytes"),
)
.group_by(ConnectionLog.username)
.order_by(desc("sessions"))
.limit(10)
)
top_channels_q = await db.execute(
select(ConnectionLog.channel_name, func.count().label("sessions"))
.where(ConnectionLog.channel_name.isnot(None))
.group_by(ConnectionLog.channel_name)
.order_by(desc("sessions"))
.limit(10)
)
totals_q = await db.execute(
select(
func.count().label("total_sessions"),
func.sum(ConnectionLog.bytes_transferred).label("total_bytes"),
func.avg(ConnectionLog.duration_seconds).label("avg_duration"),
func.count(distinct(ConnectionLog.username)).label("unique_users"),
)
)
t = totals_q.one()
return {
"total_sessions": t.total_sessions or 0,
"total_bytes": t.total_bytes or 0,
"avg_duration_seconds": round(t.avg_duration or 0),
"unique_users": t.unique_users or 0,
"top_users": [
{"username": r.username, "sessions": r.sessions, "bytes": r.bytes or 0}
for r in top_users_q
],
"top_channels": [
{"channel_name": r.channel_name, "sessions": r.sessions}
for r in top_channels_q
],
}
@router.delete("", status_code=204)
async def clear_logs(db: AsyncSession = Depends(get_db), _=Depends(get_current_admin)):
"""Delete all logs (admin use)."""
from sqlalchemy import delete
await db.execute(delete(ConnectionLog))
await db.commit()
@router.get("/events")
async def get_events(
channel: str | None = Query(None),
event_type: str | None = Query(None),
date_from: str | None = Query(None),
date_to: str | None = Query(None),
page: int = Query(1, ge=1),
limit: int = Query(50, ge=1, le=200),
db: AsyncSession = Depends(get_db),
_=Depends(get_current_admin),
):
q = select(StreamEvent).order_by(desc(StreamEvent.created_at))
if channel:
q = q.where(StreamEvent.channel_name.ilike(f"%{channel}%"))
if event_type:
q = q.where(StreamEvent.event_type == event_type)
if date_from:
try:
q = q.where(StreamEvent.created_at >= datetime.fromisoformat(date_from))
except ValueError:
pass
if date_to:
try:
q = q.where(StreamEvent.created_at <= datetime.fromisoformat(date_to))
except ValueError:
pass
total_result = await db.execute(select(func.count()).select_from(q.subquery()))
total = total_result.scalar() or 0
q = q.offset((page - 1) * limit).limit(limit)
result = await db.execute(q)
events = result.scalars().all()
return {
"total": total,
"page": page,
"limit": limit,
"items": [_serialize_event(ev) for ev in events],
}
@router.delete("/events", status_code=204)
async def clear_events(db: AsyncSession = Depends(get_db), _=Depends(get_current_admin)):
from sqlalchemy import delete
await db.execute(delete(StreamEvent))
await db.commit()
def _serialize_event(ev: StreamEvent) -> dict:
return {
"id": ev.id,
"channel_id": ev.channel_id,
"channel_name": ev.channel_name,
"provider_name": ev.provider_name,
"event_type": ev.event_type,
"domain": ev.domain,
"error_message": ev.error_message,
"attempt_number": ev.attempt_number,
"pump_duration_secs": ev.pump_duration_secs,
"created_at": ev.created_at.isoformat() if ev.created_at else None,
}
def _serialize(log: ConnectionLog) -> dict:
return {
"id": log.id,
"user_id": log.user_id,
"username": log.username,
"channel_id": log.channel_id,
"channel_name": log.channel_name,
"content_type": log.content_type,
"stream_id": log.stream_id,
"started_at": log.started_at.isoformat() if log.started_at else None,
"ended_at": log.ended_at.isoformat() if log.ended_at else None,
"duration_seconds": log.duration_seconds,
"bytes_transferred": log.bytes_transferred,
"client_ip": log.client_ip,
}