"""Leaderboard read model over materialized RCON/AdminLog match stats.""" from __future__ import annotations import argparse import json import os from contextlib import closing from datetime import datetime, timedelta, timezone from pathlib import Path from typing import Literal from .config import get_historical_weekly_fallback_min_matches, 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", ) 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) -> Path: """Create ranking snapshot tables used by weekly/monthly public ranking reads.""" resolved_path = initialize_rcon_materialized_storage(db_path=db_path) if use_postgres_rcon_storage(explicit_sqlite_path=db_path): 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 when snapshot is missing.""" normalized = os.getenv( "HLL_BACKEND_RANKING_RUNTIME_FALLBACK_ENABLED", "true", ).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, 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) connection_scope = _connect_write_scope(resolved_path, db_path=db_path) 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 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)) resolved_path = initialize_ranking_snapshot_storage(db_path=db_path) connection_scope = _connect_scope(resolved_path, db_path=db_path) with connection_scope as connection: snapshot = _find_latest_snapshot( connection=connection, timeframe=normalized_timeframe, server_key=normalized_server_key, metric=normalized_metric, ) if snapshot is None: return { "snapshot_status": "missing", "timeframe": normalized_timeframe, "server_id": normalized_server_key, "metric": normalized_metric, "limit": normalized_limit, "requested_limit": normalized_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": [], } 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, ) 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, } 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, 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) connection_scope = _connect_scope(resolved_path, db_path=db_path) 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: for index, row in enumerate(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 _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): if use_postgres_rcon_storage(explicit_sqlite_path=db_path): from .postgres_rcon_storage import connect_postgres_compat return connect_postgres_compat() return closing(connect_sqlite_readonly(resolved_path)) def _connect_write_scope(resolved_path: Path, *, db_path: Path | None): if use_postgres_rcon_storage(explicit_sqlite_path=db_path): from .postgres_rcon_storage import connect_postgres_compat return connect_postgres_compat() 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 _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", ) parser.set_defaults(command="generate-ranking-snapshot", 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)) return 0 parser.print_help() return 2 if __name__ == "__main__": raise SystemExit(_main())