Files
comunidadhll/backend/app/rcon_historical_leaderboards.py
2026-06-11 10:22:37 +02:00

1539 lines
52 KiB
Python

"""Leaderboard read model over materialized RCON/AdminLog match stats."""
from __future__ import annotations
import argparse
import json
import os
import sqlite3
from contextlib import closing
from contextlib import contextmanager
from datetime import date, datetime, timedelta, timezone
from pathlib import Path
from typing import Literal
from .config import (
get_historical_weekly_fallback_min_matches,
get_storage_path,
use_postgres_rcon_storage,
)
from .historical_storage import ALL_SERVERS_SLUG
from .rcon_admin_log_materialization import (
MATCH_RESULT_SOURCE,
initialize_rcon_materialized_storage,
)
from .sqlite_utils import connect_sqlite_readonly, connect_sqlite_writer
LeaderboardTimeframe = Literal["weekly", "monthly"]
LeaderboardMetric = Literal[
"kills",
"deaths",
"teamkills",
"matches_considered",
"kd_ratio",
"kills_per_match",
"matches_over_100_kills",
"support",
]
SNAPSHOT_GENERATOR_TIMEFRAMES = ("weekly", "monthly")
SNAPSHOT_GENERATOR_SERVER_KEYS = (
ALL_SERVERS_SLUG,
"comunidad-hispana-01",
"comunidad-hispana-02",
)
SNAPSHOT_GENERATOR_METRICS = (
"kills",
"deaths",
"teamkills",
"matches_considered",
"kd_ratio",
"kills_per_match",
)
DEFAULT_RANKING_SNAPSHOT_REFRESH_LIMIT = 30
RANKING_SNAPSHOT_SQLITE_SCHEMA_SQL = """
CREATE TABLE IF NOT EXISTS ranking_snapshots (
id INTEGER PRIMARY KEY AUTOINCREMENT,
timeframe TEXT NOT NULL,
server_id TEXT NOT NULL,
metric TEXT NOT NULL,
window_start TEXT NOT NULL,
window_end TEXT NOT NULL,
generated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
source TEXT NOT NULL DEFAULT 'rcon-materialized-admin-log',
snapshot_status TEXT NOT NULL DEFAULT 'ready',
item_count INTEGER NOT NULL DEFAULT 0,
limit_size INTEGER NOT NULL DEFAULT 20,
source_matches_count INTEGER NOT NULL DEFAULT 0,
freshness TEXT NOT NULL DEFAULT 'fresh',
window_kind TEXT,
window_label TEXT,
error_message TEXT,
UNIQUE(timeframe, server_id, metric, window_start, window_end)
);
CREATE INDEX IF NOT EXISTS idx_ranking_snapshots_lookup
ON ranking_snapshots(timeframe, server_id, metric, snapshot_status, window_end DESC, generated_at DESC);
CREATE TABLE IF NOT EXISTS ranking_snapshot_items (
id INTEGER PRIMARY KEY AUTOINCREMENT,
snapshot_id INTEGER NOT NULL REFERENCES ranking_snapshots(id) ON DELETE CASCADE,
ranking_position INTEGER NOT NULL,
player_id TEXT NOT NULL,
player_name TEXT NOT NULL,
metric_value REAL NOT NULL DEFAULT 0,
matches_considered INTEGER NOT NULL DEFAULT 0,
kills INTEGER NOT NULL DEFAULT 0,
deaths INTEGER NOT NULL DEFAULT 0,
teamkills INTEGER NOT NULL DEFAULT 0,
kd_ratio REAL NOT NULL DEFAULT 0.0,
kills_per_match REAL NOT NULL DEFAULT 0.0,
UNIQUE(snapshot_id, ranking_position),
UNIQUE(snapshot_id, player_id)
);
CREATE INDEX IF NOT EXISTS idx_ranking_snapshot_items_snapshot
ON ranking_snapshot_items(snapshot_id, ranking_position);
CREATE INDEX IF NOT EXISTS idx_ranking_snapshot_items_player
ON ranking_snapshot_items(snapshot_id, player_id);
"""
def initialize_ranking_snapshot_storage(
*,
db_path: Path | None = None,
ensure_storage: bool = True,
) -> Path:
"""Create ranking snapshot tables used by weekly/monthly public ranking reads."""
resolved_path = initialize_rcon_materialized_storage(
db_path=db_path,
ensure_storage=ensure_storage,
)
if use_postgres_rcon_storage(explicit_sqlite_path=db_path):
if ensure_storage:
from .postgres_rcon_storage import initialize_postgres_rcon_storage
initialize_postgres_rcon_storage()
return resolved_path
with closing(connect_sqlite_writer(resolved_path)) as connection:
with connection:
connection.executescript(RANKING_SNAPSHOT_SQLITE_SCHEMA_SQL)
return resolved_path
def is_ranking_runtime_fallback_enabled() -> bool:
"""Return whether `/api/ranking` may fall back to runtime reads in explicit internal mode."""
normalized = os.getenv(
"HLL_BACKEND_RANKING_RUNTIME_FALLBACK_ENABLED",
"false",
).strip().lower()
return normalized in {"1", "true", "yes", "on"}
def generate_ranking_snapshot(
*,
timeframe: str,
server_key: str | None,
metric: str,
limit: int,
replace_existing: bool = True,
ensure_storage: bool = True,
now: datetime | None = None,
db_path: Path | None = None,
) -> dict[str, object]:
"""Generate and persist one weekly/monthly ranking snapshot."""
normalized_timeframe = _normalize_generator_timeframe(timeframe)
normalized_server_key = _normalize_generator_server_key(server_key)
normalized_metric = _normalize_generator_metric(metric)
normalized_limit = _normalize_limit(limit)
anchor = _as_utc(now or datetime.now(timezone.utc))
resolved_path = initialize_ranking_snapshot_storage(
db_path=db_path,
ensure_storage=ensure_storage,
)
connection_scope = _connect_write_scope(
resolved_path,
db_path=db_path,
initialize=ensure_storage,
)
with connection_scope as connection:
window = select_leaderboard_window(
connection=connection,
server_key=normalized_server_key,
timeframe=normalized_timeframe,
now=anchor,
)
window_start = _to_iso(window["start"])
window_end = _to_iso(window["end"])
existing_snapshot_id = _find_existing_snapshot_for_window(
connection=connection,
timeframe=normalized_timeframe,
server_key=normalized_server_key,
metric=normalized_metric,
window_start=window_start,
window_end=window_end,
)
if existing_snapshot_id is not None and not replace_existing:
snapshot = _get_snapshot_record(connection=connection, snapshot_id=existing_snapshot_id)
items = _list_snapshot_items(connection=connection, snapshot_id=existing_snapshot_id)
return {
"status": "ok",
"snapshot": snapshot,
"items": items,
"source_matches_count": int(snapshot.get("source_matches_count") or 0)
if snapshot
else 0,
"ranked_players": len(items),
"skipped_regeneration": True,
}
ranking_rows = _fetch_leaderboard_rows(
connection,
server_key=normalized_server_key,
metric=normalized_metric,
limit=normalized_limit,
window_start=window["start"],
window_end=window["end"],
)
source_matches_count = _count_matches(
connection,
server_key=normalized_server_key,
start=window["start"],
end=window["end"],
)
if existing_snapshot_id is not None:
_delete_snapshot(connection=connection, snapshot_id=existing_snapshot_id)
snapshot_id = _insert_snapshot_record(
connection=connection,
timeframe=normalized_timeframe,
server_key=normalized_server_key,
metric=normalized_metric,
limit=normalized_limit,
source_matches_count=source_matches_count,
window_start=window_start,
window_end=window_end,
generated_at=_to_iso(datetime.now(timezone.utc)),
window_kind=str(window["kind"]),
window_label=str(window["label"]),
item_count=len(ranking_rows),
)
_insert_snapshot_items(
connection=connection,
snapshot_id=snapshot_id,
rows=ranking_rows,
limit=normalized_limit,
)
snapshot = _get_snapshot_record(connection=connection, snapshot_id=snapshot_id)
items = _list_snapshot_items(connection=connection, snapshot_id=snapshot_id)
return {
"status": "ok",
"snapshot": snapshot,
"items": items,
"source_matches_count": source_matches_count,
"ranked_players": len(items),
"skipped_regeneration": False,
}
def refresh_ranking_snapshots(
*,
limit: int = DEFAULT_RANKING_SNAPSHOT_REFRESH_LIMIT,
replace_existing: bool = True,
timeframes: tuple[str, ...] | None = None,
ensure_storage: bool = True,
now: datetime | None = None,
db_path: Path | None = None,
) -> dict[str, object]:
"""Generate the full weekly/monthly ranking snapshot matrix with partial-failure reporting."""
normalized_limit = _normalize_limit(limit)
anchor = _as_utc(now or datetime.now(timezone.utc))
normalized_timeframes = tuple(
_normalize_generator_timeframe(timeframe)
for timeframe in (timeframes or SNAPSHOT_GENERATOR_TIMEFRAMES)
)
combinations = [
(timeframe, server_key, metric)
for timeframe in normalized_timeframes
for server_key in SNAPSHOT_GENERATOR_SERVER_KEYS
for metric in SNAPSHOT_GENERATOR_METRICS
]
results: list[dict[str, object]] = []
succeeded = 0
failed = 0
skipped_regeneration = 0
for timeframe, server_key, metric in combinations:
try:
payload = generate_ranking_snapshot(
timeframe=timeframe,
server_key=server_key,
metric=metric,
limit=normalized_limit,
replace_existing=replace_existing,
ensure_storage=ensure_storage if not results else False,
now=anchor,
db_path=db_path,
)
snapshot = payload.get("snapshot") if isinstance(payload, dict) else {}
skipped = bool(payload.get("skipped_regeneration")) if isinstance(payload, dict) else False
if skipped:
skipped_regeneration += 1
succeeded += 1
results.append(
{
"status": "ok",
"timeframe": timeframe,
"server_key": server_key,
"metric": metric,
"limit": normalized_limit,
"snapshot_id": snapshot.get("id") if isinstance(snapshot, dict) else None,
"snapshot_status": snapshot.get("snapshot_status")
if isinstance(snapshot, dict)
else None,
"window_start": snapshot.get("window_start") if isinstance(snapshot, dict) else None,
"window_end": snapshot.get("window_end") if isinstance(snapshot, dict) else None,
"generated_at": snapshot.get("generated_at") if isinstance(snapshot, dict) else None,
"ranked_players": int(payload.get("ranked_players") or 0),
"source_matches_count": int(payload.get("source_matches_count") or 0),
"skipped_regeneration": skipped,
}
)
except Exception as exc: # noqa: BLE001 - bulk refresh must report per-combination failures
failed += 1
results.append(
{
"status": "error",
"timeframe": timeframe,
"server_key": server_key,
"metric": metric,
"limit": normalized_limit,
"error_type": type(exc).__name__,
"error": str(exc),
}
)
overall_status = "ok"
if failed and succeeded:
overall_status = "partial"
elif failed:
overall_status = "error"
return {
"status": overall_status,
"generated_at": _to_iso(anchor),
"timeframes": list(normalized_timeframes),
"limit": normalized_limit,
"replace_existing": replace_existing,
"combinations_expected": len(combinations),
"totals": {
"combinations_expected": len(combinations),
"succeeded": succeeded,
"failed": failed,
"skipped_regeneration": skipped_regeneration,
},
"results": results,
}
def get_latest_ranking_snapshot(
*,
server_key: str | None = None,
timeframe: str = "weekly",
metric: str = "kills",
limit: int = 10,
db_path: Path | None = None,
) -> dict[str, object]:
"""Return the latest ready weekly/monthly ranking snapshot for the requested scope."""
normalized_server_key = _normalize_snapshot_server_key(server_key)
normalized_timeframe = _normalize_timeframe(timeframe)
normalized_metric = _normalize_metric(metric)
normalized_limit = max(1, int(limit or 10))
try:
with _open_ranking_snapshot_read_connection(db_path=db_path) as connection:
snapshot = _find_latest_snapshot(
connection=connection,
timeframe=normalized_timeframe,
server_key=normalized_server_key,
metric=normalized_metric,
)
if snapshot is None:
return _build_missing_ranking_snapshot_result(
timeframe=normalized_timeframe,
server_key=normalized_server_key,
metric=normalized_metric,
limit=normalized_limit,
)
snapshot_limit = max(1, int(snapshot.get("limit_size") or normalized_limit))
item_count = int(snapshot.get("item_count") or 0)
effective_limit = max(0, min(normalized_limit, snapshot_limit, item_count))
items = _list_snapshot_items(
connection=connection,
snapshot_id=int(snapshot["id"]),
limit=effective_limit if effective_limit > 0 else None,
)
except (FileNotFoundError, sqlite3.OperationalError):
return _build_missing_ranking_snapshot_result(
timeframe=normalized_timeframe,
server_key=normalized_server_key,
metric=normalized_metric,
limit=normalized_limit,
)
return {
"snapshot_status": "ready",
"timeframe": normalized_timeframe,
"server_id": normalized_server_key,
"metric": normalized_metric,
"limit": effective_limit,
"requested_limit": normalized_limit,
"effective_limit": effective_limit,
"snapshot_limit": snapshot_limit,
"item_count": item_count,
"generated_at": snapshot.get("generated_at"),
"window_start": snapshot.get("window_start"),
"window_end": snapshot.get("window_end"),
"window_kind": snapshot.get("window_kind"),
"window_label": snapshot.get("window_label"),
"source": snapshot.get("source") or "ranking-snapshot",
"freshness": snapshot.get("freshness") or "fresh",
"source_matches_count": int(snapshot.get("source_matches_count") or 0),
"items": items,
}
@contextmanager
def _open_ranking_snapshot_read_connection(*, db_path: Path | None = None):
if use_postgres_rcon_storage(explicit_sqlite_path=db_path):
from .postgres_rcon_storage import PostgresCompatConnection, connect_postgres
with connect_postgres() as connection:
yield PostgresCompatConnection(connection)
return
resolved_path = _resolve_ranking_snapshot_sqlite_path(db_path=db_path)
if not resolved_path.exists():
raise FileNotFoundError(resolved_path)
with closing(connect_sqlite_readonly(resolved_path)) as connection:
yield connection
def _resolve_ranking_snapshot_sqlite_path(*, db_path: Path | None = None) -> Path:
return db_path if db_path is not None else get_storage_path()
def _build_missing_ranking_snapshot_result(
*,
timeframe: str,
server_key: str,
metric: str,
limit: int,
) -> dict[str, object]:
return {
"snapshot_status": "missing",
"timeframe": timeframe,
"server_id": server_key,
"metric": metric,
"limit": limit,
"requested_limit": limit,
"effective_limit": 0,
"snapshot_limit": None,
"item_count": 0,
"generated_at": None,
"window_start": None,
"window_end": None,
"window_kind": None,
"window_label": None,
"source": "ranking-snapshot",
"freshness": "missing",
"source_matches_count": 0,
"items": [],
}
def build_rcon_materialized_leaderboard_snapshot_payload(
*,
server_id: str | None = None,
timeframe: str = "weekly",
metric: str = "kills",
limit: int = 10,
) -> dict[str, object]:
"""Return an API payload for RCON-backed leaderboard snapshots.
This is a runtime fast read over the materialized AdminLog tables. It intentionally
avoids the old public-scoreboard fallback because the UI is running in RCON mode.
"""
normalized_timeframe = _normalize_timeframe(timeframe)
normalized_metric = _normalize_metric(metric)
result = list_rcon_materialized_leaderboard(
server_key=server_id,
timeframe=normalized_timeframe,
metric=normalized_metric,
limit=limit,
)
items = list(result.get("items") or [])[:limit]
return {
"status": "ok",
"data": {
"title": _build_title(
metric=normalized_metric,
timeframe=normalized_timeframe,
server_id=server_id,
),
"context": f"historical-{normalized_timeframe}-leaderboard-snapshot",
"source": "rcon-materialized-admin-log-leaderboard",
"server_slug": server_id,
"timeframe": normalized_timeframe,
"metric": normalized_metric,
"found": True,
"snapshot_status": "ready",
"missing_reason": None,
"request_path_policy": "runtime-rcon-materialized-fast-path",
"generation_policy": "runtime-materialized-read",
"generated_at": _to_iso(datetime.now(timezone.utc)),
"source_range_start": result.get("source_range_start"),
"source_range_end": result.get("source_range_end"),
"is_stale": False,
"freshness": "runtime",
"window_days": result.get("window_days"),
"window_start": result.get("window_start"),
"window_end": result.get("window_end"),
"window_kind": result.get("window_kind"),
"window_label": result.get("window_label"),
"uses_fallback": False,
"selection_reason": result.get("selection_reason"),
"current_week_start": result.get("current_week_start"),
"current_week_closed_matches": result.get("current_week_closed_matches"),
"previous_week_closed_matches": result.get("previous_week_closed_matches"),
"current_month_start": result.get("current_month_start"),
"selected_month_start": result.get("selected_month_start"),
"selected_month_end": result.get("selected_month_end"),
"current_month_closed_matches": result.get("current_month_closed_matches"),
"previous_month_closed_matches": result.get("previous_month_closed_matches"),
"sufficient_sample": result.get("sufficient_sample"),
"snapshot_limit": result.get("limit"),
"limit": limit,
"runtime_enrichment": {
"applied": False,
"reason": None,
},
"primary_source": "rcon",
"selected_source": "rcon",
"fallback_used": False,
"fallback_reason": None,
"source_attempts": [
{
"source": "rcon",
"role": "primary",
"status": "success",
"reason": "leaderboard-served-by-rcon-materialized-admin-log",
"message": None,
}
],
"items": items,
},
}
def list_rcon_materialized_leaderboard(
*,
server_key: str | None = None,
timeframe: str = "weekly",
metric: str = "kills",
limit: int = 10,
ensure_storage: bool = True,
db_path: Path | None = None,
now: datetime | None = None,
) -> dict[str, object]:
"""Return a leaderboard built from materialized RCON/AdminLog player stats.
RCON/AdminLog materialization currently has reliable kill/death/teamkill counters,
but not public-scoreboard support points. For support, return an explicitly empty
supported payload rather than falling back to unrelated public scoreboard storage.
"""
normalized_timeframe = _normalize_timeframe(timeframe)
normalized_metric = _normalize_metric(metric)
normalized_limit = max(1, int(limit or 10))
anchor = _as_utc(now or datetime.now(timezone.utc))
resolved_path = initialize_rcon_materialized_storage(
db_path=db_path,
ensure_storage=ensure_storage,
)
connection_scope = _connect_scope(
resolved_path,
db_path=db_path,
initialize=ensure_storage,
)
with connection_scope as connection:
window = select_leaderboard_window(
connection=connection,
server_key=server_key,
timeframe=normalized_timeframe,
now=anchor,
)
if normalized_metric == "support":
return _empty_payload(
server_key=server_key,
timeframe=normalized_timeframe,
metric=normalized_metric,
limit=normalized_limit,
window=window,
reason="rcon-materialized-stats-do-not-include-support-score",
)
rows = _fetch_leaderboard_rows(
connection,
server_key=server_key,
metric=normalized_metric,
limit=normalized_limit,
window_start=window["start"],
window_end=window["end"],
)
source_range = _fetch_source_range(
connection,
server_key=server_key,
window_start=window["start"],
window_end=window["end"],
)
items = [_build_item(row, index=index + 1) for index, row in enumerate(rows)]
return {
"source": "rcon-materialized-admin-log-leaderboard",
"server_key": server_key,
"metric": normalized_metric,
"limit": normalized_limit,
"window_days": window["days"],
"window_start": _to_iso(window["start"]),
"window_end": _to_iso(window["end"]),
"window_kind": window["kind"],
"window_label": window["label"],
"uses_fallback": False,
"selection_reason": window["selection_reason"],
"current_week_start": _to_iso(window["current_week_start"]),
"current_week_closed_matches": window["current_week_closed_matches"],
"previous_week_closed_matches": window["previous_week_closed_matches"],
"current_month_start": _to_iso(window["current_month_start"]),
"selected_month_start": _to_iso(window["selected_month_start"]),
"selected_month_end": _to_iso(window["selected_month_end"]),
"current_month_closed_matches": window["current_month_closed_matches"],
"previous_month_closed_matches": window["previous_month_closed_matches"],
"sufficient_sample": window["sufficient_sample"],
"source_range_start": _to_iso(source_range[0]) if source_range[0] else None,
"source_range_end": _to_iso(source_range[1]) if source_range[1] else None,
"items": items,
}
def _find_latest_snapshot(
*,
connection: object,
timeframe: str,
server_key: str,
metric: str,
) -> dict[str, object] | None:
row = connection.execute(
"""
SELECT
id,
timeframe,
server_id,
metric,
window_start,
window_end,
generated_at,
source,
snapshot_status,
item_count,
limit_size,
source_matches_count,
freshness,
window_kind,
window_label
FROM ranking_snapshots
WHERE timeframe = ?
AND server_id = ?
AND metric = ?
AND snapshot_status = 'ready'
ORDER BY window_end DESC, generated_at DESC
LIMIT 1
""",
[timeframe, server_key, metric],
).fetchone()
return dict(row) if row else None
def _find_existing_snapshot_for_window(
*,
connection: object,
timeframe: str,
server_key: str,
metric: str,
window_start: str,
window_end: str,
) -> int | None:
row = connection.execute(
"""
SELECT id
FROM ranking_snapshots
WHERE timeframe = ?
AND server_id = ?
AND metric = ?
AND window_start = ?
AND window_end = ?
LIMIT 1
""",
[timeframe, server_key, metric, window_start, window_end],
).fetchone()
return int(row["id"]) if row else None
def _get_snapshot_record(
*,
connection: object,
snapshot_id: int,
) -> dict[str, object] | None:
row = connection.execute(
"""
SELECT
id,
timeframe,
server_id,
metric,
window_start,
window_end,
generated_at,
source,
snapshot_status,
item_count,
limit_size,
source_matches_count,
freshness,
window_kind,
window_label,
error_message
FROM ranking_snapshots
WHERE id = ?
LIMIT 1
""",
[snapshot_id],
).fetchone()
return dict(row) if row else None
def _list_snapshot_items(
*,
connection: object,
snapshot_id: int,
limit: int | None = None,
) -> list[dict[str, object]]:
params: list[object] = [snapshot_id]
query = """
SELECT
ranking_position,
player_id,
player_name,
metric_value,
matches_considered,
kills,
deaths,
teamkills,
kd_ratio,
kills_per_match
FROM ranking_snapshot_items
WHERE snapshot_id = ?
ORDER BY ranking_position ASC
"""
if limit is not None:
query = f"{query}\n LIMIT ?"
params.append(limit)
rows = connection.execute(query, params).fetchall()
return [dict(row) for row in rows]
def _delete_snapshot(*, connection: object, snapshot_id: int) -> None:
connection.execute(
"DELETE FROM ranking_snapshot_items WHERE snapshot_id = ?",
[snapshot_id],
)
connection.execute(
"DELETE FROM ranking_snapshots WHERE id = ?",
[snapshot_id],
)
def _insert_snapshot_record(
*,
connection: object,
timeframe: str,
server_key: str,
metric: str,
limit: int,
source_matches_count: int,
window_start: str,
window_end: str,
generated_at: str,
window_kind: str,
window_label: str,
item_count: int,
) -> int:
cursor = connection.execute(
"""
INSERT INTO ranking_snapshots (
timeframe,
server_id,
metric,
window_start,
window_end,
generated_at,
source,
snapshot_status,
item_count,
limit_size,
source_matches_count,
freshness,
window_kind,
window_label,
error_message
) VALUES (?, ?, ?, ?, ?, ?, ?, 'ready', ?, ?, ?, 'fresh', ?, ?, NULL)
RETURNING id
""",
[
timeframe,
server_key,
metric,
window_start,
window_end,
generated_at,
"rcon-materialized-admin-log",
item_count,
limit,
source_matches_count,
window_kind,
window_label,
],
)
row = cursor.fetchone()
if row is None or row["id"] is None:
raise RuntimeError("Unable to resolve ranking snapshot id after insert.")
return int(row["id"])
def _insert_snapshot_items(
*,
connection: object,
snapshot_id: int,
rows: list[dict[str, object]],
limit: int,
) -> None:
deduplicated_rows = _dedupe_snapshot_rows(rows)
for index, row in enumerate(deduplicated_rows[:limit], start=1):
item = _build_item(row, index=index)
connection.execute(
"""
INSERT INTO ranking_snapshot_items (
snapshot_id,
ranking_position,
player_id,
player_name,
metric_value,
matches_considered,
kills,
deaths,
teamkills,
kd_ratio,
kills_per_match
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
[
snapshot_id,
int(item["ranking_position"]),
str(item["player_id"]),
str(item["player_name"]),
float(item["metric_value"]),
int(item["matches_considered"]),
int(item["kills"]),
int(item["deaths"]),
int(item["teamkills"]),
float(item["kd_ratio"]),
float(item["kills_per_match"]),
],
)
def _dedupe_snapshot_rows(rows: list[dict[str, object]]) -> list[dict[str, object]]:
"""Keep the first ranked row for each player id to preserve deterministic snapshot order."""
deduplicated_rows: list[dict[str, object]] = []
seen_player_ids: set[str] = set()
for row in rows:
player_id = str(row.get("player_id") or "").strip()
if not player_id or player_id in seen_player_ids:
continue
seen_player_ids.add(player_id)
deduplicated_rows.append(row)
return deduplicated_rows
def _fetch_leaderboard_rows(
connection: object,
*,
server_key: str | None,
metric: str,
limit: int,
window_start: datetime,
window_end: datetime,
) -> list[dict[str, object]]:
scope_sql, scope_params = _build_scope_sql(server_key)
metric_sql, having_sql, order_by_sql = _resolve_metric_sql(metric)
params: list[object] = [
_to_iso(window_start),
_to_iso(window_end),
*scope_params,
limit,
]
rows = connection.execute(
f"""
SELECT
stats.player_id,
stats.player_name,
{metric_sql} AS metric_value,
COUNT(DISTINCT stats.match_key) AS matches_considered,
SUM(COALESCE(stats.kills, 0)) AS kills,
SUM(COALESCE(stats.deaths, 0)) AS deaths,
SUM(COALESCE(stats.teamkills, 0)) AS teamkills
FROM rcon_match_player_stats AS stats
INNER JOIN rcon_materialized_matches AS matches
ON matches.target_key = stats.target_key
AND matches.match_key = stats.match_key
WHERE matches.source_basis = ?
AND COALESCE(CAST(matches.ended_at AS TEXT), CAST(matches.started_at AS TEXT)) >= ?
AND COALESCE(CAST(matches.ended_at AS TEXT), CAST(matches.started_at AS TEXT)) <= ?
{scope_sql}
AND TRIM(COALESCE(stats.player_name, '')) != ''
GROUP BY stats.player_id, stats.player_name
{having_sql}
ORDER BY {order_by_sql}
LIMIT ?
""",
[MATCH_RESULT_SOURCE, *params],
).fetchall()
return [dict(row) for row in rows]
def _fetch_match_counts(
connection: object,
*,
server_key: str | None,
timeframe: str,
window_start: datetime,
window_end: datetime,
) -> dict[str, int]:
current_week_start = _week_start(window_end)
previous_week_start = current_week_start - timedelta(days=7)
current_month_start = _month_start(window_end)
previous_month_start = _previous_month_start(current_month_start)
return {
"current_week_closed_matches": _count_matches(
connection,
server_key=server_key,
start=current_week_start,
end=window_end,
),
"previous_week_closed_matches": _count_matches(
connection,
server_key=server_key,
start=previous_week_start,
end=current_week_start,
),
"current_month_closed_matches": _count_matches(
connection,
server_key=server_key,
start=current_month_start,
end=window_end,
),
"previous_month_closed_matches": _count_matches(
connection,
server_key=server_key,
start=previous_month_start,
end=current_month_start,
),
}
def select_leaderboard_window(
*,
connection: object,
server_key: str | None,
timeframe: str,
now: datetime | None = None,
) -> dict[str, object]:
"""Select the RCON leaderboard window using weekly/monthly fallback policy."""
anchor = _as_utc(now or datetime.now(timezone.utc))
current_week_start = _week_start(anchor)
previous_week_start = current_week_start - timedelta(days=7)
current_month_start = _month_start(anchor)
previous_month_start = _previous_month_start(current_month_start)
minimum_week_matches = get_historical_weekly_fallback_min_matches()
current_week_count = _count_matches(
connection,
server_key=server_key,
start=current_week_start,
end=anchor,
)
previous_week_count = _count_matches(
connection,
server_key=server_key,
start=previous_week_start,
end=current_week_start,
)
current_month_count = _count_matches(
connection,
server_key=server_key,
start=current_month_start,
end=anchor,
)
previous_month_count = _count_matches(
connection,
server_key=server_key,
start=previous_month_start,
end=current_month_start,
)
if timeframe == "monthly":
use_previous_month = anchor.day <= 7
start = previous_month_start if use_previous_month else current_month_start
end = current_month_start if use_previous_month else anchor
return {
"start": start,
"end": end,
"days": max(1, (end.date() - start.date()).days),
"kind": "previous-month" if use_previous_month else "current-month",
"label": "Mes anterior" if use_previous_month else "Mes actual",
"selection_reason": (
"monthly-uses-previous-month-until-day-8"
if use_previous_month
else "monthly-uses-current-month-after-day-7"
),
"current_week_start": current_week_start,
"current_week_closed_matches": current_week_count,
"previous_week_closed_matches": previous_week_count,
"current_month_start": current_month_start,
"selected_month_start": start,
"selected_month_end": end,
"current_month_closed_matches": current_month_count,
"previous_month_closed_matches": previous_month_count,
"sufficient_sample": {
"minimum_closed_matches": 1,
"current_month_closed_matches": current_month_count,
"previous_month_closed_matches": previous_month_count,
"current_month_has_sufficient_sample": current_month_count >= 1,
"uses_previous_month_until_day": 7,
},
}
current_week_has_sample = current_week_count >= minimum_week_matches
start = current_week_start if current_week_has_sample else previous_week_start
end = anchor if current_week_has_sample else current_week_start
return {
"start": start,
"end": end,
"days": max(1, (end.date() - start.date()).days),
"kind": "current-week" if current_week_has_sample else "previous-week",
"label": "Semana actual" if current_week_has_sample else "Semana anterior",
"selection_reason": (
"weekly-current-week-has-sufficient-closed-matches"
if current_week_has_sample
else "weekly-fallback-previous-week-insufficient-current-week-data"
),
"current_week_start": current_week_start,
"current_week_closed_matches": current_week_count,
"previous_week_closed_matches": previous_week_count,
"current_month_start": current_month_start,
"selected_month_start": current_month_start,
"selected_month_end": anchor,
"current_month_closed_matches": current_month_count,
"previous_month_closed_matches": previous_month_count,
"sufficient_sample": {
"minimum_closed_matches": minimum_week_matches,
"current_week_closed_matches": current_week_count,
"current_week_has_sufficient_sample": current_week_has_sample,
"previous_week_closed_matches": previous_week_count,
},
}
def _fetch_source_range(
connection: object,
*,
server_key: str | None,
window_start: datetime,
window_end: datetime,
) -> tuple[datetime | None, datetime | None]:
scope_sql, scope_params = _build_scope_sql(server_key, table_alias="matches")
row = connection.execute(
f"""
SELECT
MIN(COALESCE(CAST(matches.ended_at AS TEXT), CAST(matches.started_at AS TEXT))) AS source_range_start,
MAX(COALESCE(CAST(matches.ended_at AS TEXT), CAST(matches.started_at AS TEXT))) AS source_range_end
FROM rcon_materialized_matches AS matches
WHERE matches.source_basis = ?
AND COALESCE(CAST(matches.ended_at AS TEXT), CAST(matches.started_at AS TEXT)) >= ?
AND COALESCE(CAST(matches.ended_at AS TEXT), CAST(matches.started_at AS TEXT)) <= ?
{scope_sql}
""",
[MATCH_RESULT_SOURCE, _to_iso(window_start), _to_iso(window_end), *scope_params],
).fetchone()
if not row:
return None, None
return _parse_datetime(row["source_range_start"]), _parse_datetime(row["source_range_end"])
def _count_matches(
connection: object,
*,
server_key: str | None,
start: datetime,
end: datetime,
) -> int:
scope_sql, scope_params = _build_scope_sql(server_key, table_alias="matches")
row = connection.execute(
f"""
SELECT COUNT(*) AS count
FROM rcon_materialized_matches AS matches
WHERE matches.source_basis = ?
AND COALESCE(CAST(matches.ended_at AS TEXT), CAST(matches.started_at AS TEXT)) >= ?
AND COALESCE(CAST(matches.ended_at AS TEXT), CAST(matches.started_at AS TEXT)) < ?
{scope_sql}
""",
[MATCH_RESULT_SOURCE, _to_iso(start), _to_iso(end), *scope_params],
).fetchone()
return int(row["count"] or 0) if row else 0
def _build_item(row: dict[str, object], *, index: int) -> dict[str, object]:
kills = _coerce_int(row.get("kills"))
deaths = _coerce_int(row.get("deaths"))
matches_considered = _coerce_int(row.get("matches_considered"))
kd_ratio = round(kills / deaths, 2) if deaths else float(kills)
kills_per_match = round(kills / matches_considered, 2) if matches_considered else 0.0
return {
"ranking_position": index,
"player": {
"id": row.get("player_id"),
"name": row.get("player_name"),
},
"player_id": row.get("player_id"),
"player_name": row.get("player_name"),
"metric_value": _coerce_metric_value(row.get("metric_value")),
"matches_considered": matches_considered,
"kills": kills,
"deaths": deaths,
"teamkills": _coerce_int(row.get("teamkills")),
"kd_ratio": kd_ratio,
"kills_per_match": kills_per_match,
}
def _build_scope_sql(
server_key: str | None,
*,
table_alias: str = "matches",
) -> tuple[str, list[object]]:
if not server_key or server_key == ALL_SERVERS_SLUG:
return "", []
return f"AND ({table_alias}.target_key = ? OR {table_alias}.external_server_id = ?)", [
server_key,
server_key,
]
def _normalize_snapshot_server_key(server_key: str | None) -> str:
normalized = str(server_key or "").strip()
normalized_lower = normalized.lower()
if not normalized or normalized_lower in {ALL_SERVERS_SLUG, "all"}:
return ALL_SERVERS_SLUG
return normalized
def _connect_scope(resolved_path: Path, *, db_path: Path | None, initialize: bool = True):
if use_postgres_rcon_storage(explicit_sqlite_path=db_path):
from .postgres_rcon_storage import connect_postgres_compat
return connect_postgres_compat(initialize=initialize)
return closing(connect_sqlite_readonly(resolved_path))
def _connect_write_scope(
resolved_path: Path,
*,
db_path: Path | None,
initialize: bool = True,
):
if use_postgres_rcon_storage(explicit_sqlite_path=db_path):
from .postgres_rcon_storage import connect_postgres_compat
return connect_postgres_compat(initialize=initialize)
return connect_sqlite_writer(resolved_path)
def _empty_payload(
*,
server_key: str | None,
timeframe: str,
metric: str,
limit: int,
window: dict[str, object],
reason: str,
) -> dict[str, object]:
return {
"source": "rcon-materialized-admin-log-leaderboard",
"server_key": server_key,
"metric": metric,
"limit": limit,
"window_days": window["days"],
"window_start": _to_iso(window["start"]),
"window_end": _to_iso(window["end"]),
"window_kind": window["kind"],
"window_label": window["label"],
"uses_fallback": False,
"selection_reason": reason,
"current_week_start": _to_iso(window["current_week_start"]),
"current_week_closed_matches": window["current_week_closed_matches"],
"previous_week_closed_matches": window["previous_week_closed_matches"],
"current_month_start": _to_iso(window["current_month_start"]),
"selected_month_start": _to_iso(window["selected_month_start"]),
"selected_month_end": _to_iso(window["selected_month_end"]),
"current_month_closed_matches": window["current_month_closed_matches"],
"previous_month_closed_matches": window["previous_month_closed_matches"],
"sufficient_sample": window["sufficient_sample"],
"source_range_start": None,
"source_range_end": None,
"items": [],
}
def _build_window(timeframe: str) -> dict[str, object]:
now = datetime.now(timezone.utc)
if timeframe == "monthly":
start = _month_start(now)
return {
"start": start,
"end": now,
"days": max(1, (now.date() - start.date()).days + 1),
"kind": "current-month",
"label": "Mes actual",
}
start = _week_start(now)
return {
"start": start,
"end": now,
"days": max(1, (now.date() - start.date()).days + 1),
"kind": "current-week",
"label": "Semana actual",
}
def _as_utc(value: datetime) -> datetime:
if value.tzinfo is None:
return value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
def _week_start(value: datetime) -> datetime:
point = value.astimezone(timezone.utc)
start = point - timedelta(days=point.weekday())
return start.replace(hour=0, minute=0, second=0, microsecond=0)
def _month_start(value: datetime) -> datetime:
point = value.astimezone(timezone.utc)
return point.replace(day=1, hour=0, minute=0, second=0, microsecond=0)
def _previous_month_start(current_month_start: datetime) -> datetime:
previous_month_end = current_month_start - timedelta(days=1)
return _month_start(previous_month_end)
def _normalize_timeframe(value: str) -> LeaderboardTimeframe:
return "monthly" if str(value or "").strip().lower() == "monthly" else "weekly"
def _normalize_generator_timeframe(value: str) -> str:
normalized = str(value or "").strip().lower()
if normalized not in SNAPSHOT_GENERATOR_TIMEFRAMES:
raise ValueError(
f"timeframe must be one of: {', '.join(SNAPSHOT_GENERATOR_TIMEFRAMES)}"
)
return normalized
def _normalize_generator_server_key(server_key: str | None) -> str:
normalized = _normalize_snapshot_server_key(server_key)
if normalized not in SNAPSHOT_GENERATOR_SERVER_KEYS:
raise ValueError(
"server_key must be one of: all, all-servers, comunidad-hispana-01, comunidad-hispana-02"
)
return normalized
def _normalize_generator_metric(metric: str) -> str:
normalized = str(metric or "").strip().lower()
if normalized not in SNAPSHOT_GENERATOR_METRICS:
raise ValueError(
f"metric must be one of: {', '.join(SNAPSHOT_GENERATOR_METRICS)}"
)
return normalized
def _normalize_limit(limit: object, *, maximum: int = 100) -> int:
normalized_limit = int(limit or 1)
if normalized_limit < 1:
raise ValueError("limit must be greater than zero")
return min(normalized_limit, maximum)
def _normalize_metric(value: str) -> LeaderboardMetric:
normalized = str(value or "kills").strip().lower()
if normalized in {
"kills",
"deaths",
"teamkills",
"matches_considered",
"kd_ratio",
"kills_per_match",
"matches_over_100_kills",
"support",
}:
return normalized # type: ignore[return-value]
return "kills"
def _build_title(*, metric: str, timeframe: str, server_id: str | None) -> str:
timeframe_label = "mensual" if timeframe == "monthly" else "semanal"
scope = "totales" if server_id == ALL_SERVERS_SLUG else "por servidor"
metric_label = {
"kills": "Top kills",
"deaths": "Top muertes",
"teamkills": "Top teamkills",
"matches_considered": "Top partidas jugadas",
"kd_ratio": "Top K/D",
"kills_per_match": "Top kills por partida",
"matches_over_100_kills": "Partidas 100+ kills",
"support": "Top soporte",
}.get(metric, "Top kills")
return f"Snapshot {metric_label} {timeframe_label} {scope}"
def _coerce_int(value: object) -> int:
try:
return int(value or 0)
except (TypeError, ValueError):
return 0
def _coerce_metric_value(value: object) -> int | float:
numeric = _coerce_float(value)
if numeric.is_integer():
return int(numeric)
return round(numeric, 2)
def _coerce_float(value: object) -> float:
try:
return float(value or 0)
except (TypeError, ValueError):
return 0.0
def _resolve_metric_sql(metric: str) -> tuple[str, str, str]:
metric_sql_by_metric = {
"kills": "SUM(COALESCE(stats.kills, 0))",
"deaths": "SUM(COALESCE(stats.deaths, 0))",
"teamkills": "SUM(COALESCE(stats.teamkills, 0))",
"matches_considered": "COUNT(DISTINCT stats.match_key)",
"kd_ratio": (
"CASE "
"WHEN SUM(COALESCE(stats.deaths, 0)) > 0 "
"THEN ROUND(CAST(SUM(COALESCE(stats.kills, 0)) AS NUMERIC) / "
"CAST(SUM(COALESCE(stats.deaths, 0)) AS NUMERIC), 2) "
"ELSE CAST(SUM(COALESCE(stats.kills, 0)) AS NUMERIC) "
"END"
),
"kills_per_match": (
"CASE "
"WHEN COUNT(DISTINCT stats.match_key) > 0 "
"THEN ROUND(CAST(SUM(COALESCE(stats.kills, 0)) AS NUMERIC) / "
"CAST(COUNT(DISTINCT stats.match_key) AS NUMERIC), 2) "
"ELSE CAST(0 AS NUMERIC) "
"END"
),
"matches_over_100_kills": "SUM(CASE WHEN COALESCE(stats.kills, 0) >= 100 THEN 1 ELSE 0 END)",
}
having_sql_by_metric = {
"kills": "HAVING SUM(COALESCE(stats.kills, 0)) > 0",
"deaths": "HAVING SUM(COALESCE(stats.deaths, 0)) > 0",
"teamkills": "HAVING SUM(COALESCE(stats.teamkills, 0)) > 0",
"matches_considered": "HAVING COUNT(DISTINCT stats.match_key) > 0",
"kd_ratio": "HAVING SUM(COALESCE(stats.kills, 0)) > 0",
"kills_per_match": "HAVING COUNT(DISTINCT stats.match_key) > 0 AND SUM(COALESCE(stats.kills, 0)) > 0",
"matches_over_100_kills": "HAVING SUM(CASE WHEN COALESCE(stats.kills, 0) >= 100 THEN 1 ELSE 0 END) > 0",
}
order_by_sql_by_metric = {
"kills": "metric_value DESC, matches_considered DESC, stats.player_name ASC",
"deaths": "metric_value DESC, matches_considered DESC, stats.player_name ASC",
"teamkills": "metric_value DESC, matches_considered DESC, stats.player_name ASC",
"matches_considered": "metric_value DESC, kills DESC, stats.player_name ASC",
"kd_ratio": "metric_value DESC, kills DESC, matches_considered DESC, stats.player_name ASC",
"kills_per_match": "metric_value DESC, kills DESC, matches_considered DESC, stats.player_name ASC",
"matches_over_100_kills": "metric_value DESC, matches_considered DESC, stats.player_name ASC",
}
if metric == "support":
raise ValueError("support is handled separately")
return (
metric_sql_by_metric[metric],
having_sql_by_metric[metric],
order_by_sql_by_metric[metric],
)
def _parse_datetime(value: object) -> datetime | None:
if isinstance(value, datetime):
parsed = value
elif isinstance(value, str) and value.strip():
try:
parsed = datetime.fromisoformat(value.strip().replace("Z", "+00:00"))
except ValueError:
return None
else:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc)
def _to_iso(value: object) -> str:
parsed = _parse_datetime(value)
if parsed is None:
parsed = datetime.now(timezone.utc)
return parsed.astimezone(timezone.utc).isoformat().replace("+00:00", "Z")
def _json_default(value: object) -> str:
if isinstance(value, datetime):
return _to_iso(value)
if isinstance(value, date):
return value.isoformat()
raise TypeError(f"Object of type {type(value).__name__} is not JSON serializable")
def _main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(
description="Read and generate weekly/monthly ranking snapshots.",
)
subparsers = parser.add_subparsers(dest="command")
generate_parser = subparsers.add_parser("generate-ranking-snapshot")
generate_parser.add_argument(
"--timeframe",
required=True,
choices=SNAPSHOT_GENERATOR_TIMEFRAMES,
)
generate_parser.add_argument("--server-key", default=None)
generate_parser.add_argument(
"--metric",
required=True,
choices=SNAPSHOT_GENERATOR_METRICS,
)
generate_parser.add_argument("--limit", type=int, default=20)
generate_parser.add_argument(
"--sqlite-path",
type=Path,
default=None,
help="explicit local SQLite override; default operational mode uses PostgreSQL when configured",
)
generate_parser.add_argument(
"--no-replace-existing",
action="store_false",
dest="replace_existing",
help="keep an existing snapshot for the selected exact window and metric scope",
)
refresh_parser = subparsers.add_parser("refresh-ranking-snapshots")
refresh_parser.add_argument(
"--limit",
type=int,
default=DEFAULT_RANKING_SNAPSHOT_REFRESH_LIMIT,
)
refresh_parser.add_argument(
"--timeframe",
choices=SNAPSHOT_GENERATOR_TIMEFRAMES,
action="append",
dest="timeframes",
help="Optional timeframe filter. Repeat to refresh more than one timeframe.",
)
refresh_parser.add_argument(
"--sqlite-path",
type=Path,
default=None,
help="explicit local SQLite override; default operational mode uses PostgreSQL when configured",
)
refresh_parser.add_argument(
"--no-replace-existing",
action="store_false",
dest="replace_existing",
help="keep an existing snapshot for the selected exact window and metric scope",
)
parser.set_defaults(replace_existing=True)
args = parser.parse_args(argv)
if args.command == "generate-ranking-snapshot":
payload = generate_ranking_snapshot(
timeframe=args.timeframe,
server_key=args.server_key,
metric=args.metric,
limit=args.limit,
replace_existing=args.replace_existing,
db_path=args.sqlite_path,
)
print(
json.dumps(
{"status": "ok", "data": payload},
ensure_ascii=True,
indent=2,
default=_json_default,
)
)
return 0
if args.command == "refresh-ranking-snapshots":
payload = refresh_ranking_snapshots(
limit=args.limit,
replace_existing=args.replace_existing,
timeframes=tuple(args.timeframes) if args.timeframes else None,
db_path=args.sqlite_path,
)
print(
json.dumps(
{"status": "ok", "data": payload},
ensure_ascii=True,
indent=2,
default=_json_default,
)
)
return 0
parser.print_help()
return 2
if __name__ == "__main__":
raise SystemExit(_main())