feat: materialize rcon matches and backend detail data
This commit is contained in:
@@ -388,8 +388,14 @@ def build_recent_historical_matches_payload(
|
||||
)
|
||||
|
||||
if not bool(rcon_source_policy.get("fallback_used")):
|
||||
if 0 < len(items) < limit:
|
||||
fallback_items = list_recent_historical_matches(limit=limit, server_slug=server_slug)
|
||||
if 0 < len(items) < limit and not _recent_items_include_rcon_results(items):
|
||||
fallback_items = [
|
||||
_with_recent_result_source(item, "public-scoreboard-fallback")
|
||||
for item in list_recent_historical_matches(
|
||||
limit=limit,
|
||||
server_slug=server_slug,
|
||||
)
|
||||
]
|
||||
merged_items = _merge_recent_match_items(
|
||||
primary_items=items,
|
||||
fallback_items=fallback_items,
|
||||
@@ -453,7 +459,10 @@ def build_recent_historical_matches_payload(
|
||||
"capabilities": capabilities,
|
||||
},
|
||||
}
|
||||
items = list_recent_historical_matches(limit=limit, server_slug=server_slug)
|
||||
items = [
|
||||
_with_recent_result_source(item, "public-scoreboard-fallback")
|
||||
for item in list_recent_historical_matches(limit=limit, server_slug=server_slug)
|
||||
]
|
||||
return {
|
||||
"status": "ok",
|
||||
"data": {
|
||||
@@ -1828,6 +1837,23 @@ def _merge_recent_match_items(
|
||||
return merged[:limit]
|
||||
|
||||
|
||||
def _with_recent_result_source(
|
||||
item: dict[str, object],
|
||||
result_source: str,
|
||||
) -> dict[str, object]:
|
||||
enriched = dict(item)
|
||||
enriched.setdefault("result_source", result_source)
|
||||
return enriched
|
||||
|
||||
|
||||
def _recent_items_include_rcon_results(items: list[dict[str, object]]) -> bool:
|
||||
return any(
|
||||
item.get("result_source") in {"admin-log-match-ended", "rcon-session"}
|
||||
for item in items
|
||||
if isinstance(item, dict)
|
||||
)
|
||||
|
||||
|
||||
def _build_recent_match_dedupe_key(item: dict[str, object]) -> str:
|
||||
server = item.get("server") if isinstance(item.get("server"), dict) else {}
|
||||
map_payload = item.get("map") if isinstance(item.get("map"), dict) else {}
|
||||
|
||||
725
backend/app/rcon_admin_log_materialization.py
Normal file
725
backend/app/rcon_admin_log_materialization.py
Normal file
@@ -0,0 +1,725 @@
|
||||
"""Materialize RCON AdminLog events into match and player-stat read models."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import sqlite3
|
||||
from collections import Counter
|
||||
from collections.abc import Iterable
|
||||
from contextlib import closing
|
||||
from pathlib import Path
|
||||
|
||||
from .config import get_storage_path
|
||||
from .normalizers import normalize_map_name
|
||||
from .rcon_admin_log_storage import initialize_rcon_admin_log_storage
|
||||
from .rcon_historical_storage import list_rcon_historical_competitive_windows
|
||||
from .sqlite_utils import connect_sqlite_readonly, connect_sqlite_writer
|
||||
|
||||
|
||||
MATCH_RESULT_SOURCE = "admin-log-match-ended"
|
||||
SESSION_RESULT_SOURCE = "rcon-session"
|
||||
|
||||
|
||||
def initialize_rcon_materialized_storage(*, db_path: Path | None = None) -> Path:
|
||||
"""Create SQLite structures used by the materialized RCON match pipeline."""
|
||||
resolved_path = initialize_rcon_admin_log_storage(db_path=db_path)
|
||||
with closing(connect_sqlite_writer(resolved_path)) as connection:
|
||||
with connection:
|
||||
connection.executescript(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS rcon_materialized_matches (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
target_key TEXT NOT NULL,
|
||||
external_server_id TEXT,
|
||||
match_key TEXT NOT NULL,
|
||||
map_name TEXT,
|
||||
map_pretty_name TEXT,
|
||||
game_mode TEXT,
|
||||
started_server_time INTEGER,
|
||||
ended_server_time INTEGER,
|
||||
started_at TEXT,
|
||||
ended_at TEXT,
|
||||
allied_score INTEGER,
|
||||
axis_score INTEGER,
|
||||
winner TEXT,
|
||||
confidence_mode TEXT NOT NULL,
|
||||
source_basis TEXT NOT NULL,
|
||||
created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
UNIQUE(target_key, match_key)
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_rcon_materialized_matches_recent
|
||||
ON rcon_materialized_matches(target_key, ended_at DESC, ended_server_time DESC);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS rcon_match_player_stats (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
target_key TEXT NOT NULL,
|
||||
match_key TEXT NOT NULL,
|
||||
player_id TEXT NOT NULL,
|
||||
player_name TEXT NOT NULL,
|
||||
team TEXT,
|
||||
kills INTEGER NOT NULL DEFAULT 0,
|
||||
deaths INTEGER NOT NULL DEFAULT 0,
|
||||
teamkills INTEGER NOT NULL DEFAULT 0,
|
||||
deaths_by_teamkill INTEGER NOT NULL DEFAULT 0,
|
||||
weapons_json TEXT NOT NULL DEFAULT '{}',
|
||||
death_by_weapons_json TEXT NOT NULL DEFAULT '{}',
|
||||
most_killed_json TEXT NOT NULL DEFAULT '{}',
|
||||
death_by_json TEXT NOT NULL DEFAULT '{}',
|
||||
first_seen_server_time INTEGER,
|
||||
last_seen_server_time INTEGER,
|
||||
created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
UNIQUE(target_key, match_key, player_id)
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_rcon_match_player_stats_match
|
||||
ON rcon_match_player_stats(target_key, match_key);
|
||||
"""
|
||||
)
|
||||
return resolved_path
|
||||
|
||||
|
||||
def materialize_rcon_admin_log(*, db_path: Path | None = None) -> dict[str, object]:
|
||||
"""Materialize matches and player stats from stored AdminLog events."""
|
||||
resolved_path = initialize_rcon_materialized_storage(db_path=db_path)
|
||||
errors: list[str] = []
|
||||
matches_seen = 0
|
||||
matches_materialized = 0
|
||||
matches_updated = 0
|
||||
player_stats_seen = 0
|
||||
player_stats_materialized = 0
|
||||
player_stats_updated = 0
|
||||
|
||||
with closing(connect_sqlite_writer(resolved_path)) as connection:
|
||||
with connection:
|
||||
try:
|
||||
match_rows = _derive_admin_log_matches(connection)
|
||||
matches_seen = len(match_rows)
|
||||
for row in match_rows:
|
||||
outcome = _upsert_match(connection, row)
|
||||
matches_materialized += int(outcome == "inserted")
|
||||
matches_updated += int(outcome == "updated")
|
||||
session_rows = _derive_session_fallback_matches(connection, db_path=resolved_path)
|
||||
matches_seen += len(session_rows)
|
||||
for row in session_rows:
|
||||
outcome = _upsert_match(connection, row)
|
||||
matches_materialized += int(outcome == "inserted")
|
||||
matches_updated += int(outcome == "updated")
|
||||
|
||||
persisted_matches = _list_materialized_matches(connection)
|
||||
for match in persisted_matches:
|
||||
stats = _derive_player_stats_for_match(connection, match)
|
||||
player_stats_seen += len(stats)
|
||||
connection.execute(
|
||||
"""
|
||||
DELETE FROM rcon_match_player_stats
|
||||
WHERE target_key = ? AND match_key = ?
|
||||
""",
|
||||
(match["target_key"], match["match_key"]),
|
||||
)
|
||||
for stat in stats:
|
||||
_insert_player_stat(connection, stat)
|
||||
player_stats_materialized += 1
|
||||
except sqlite3.Error as error:
|
||||
errors.append(str(error))
|
||||
|
||||
return {
|
||||
"matches_seen": matches_seen,
|
||||
"matches_materialized": matches_materialized,
|
||||
"matches_updated": matches_updated,
|
||||
"player_stats_seen": player_stats_seen,
|
||||
"player_stats_materialized": player_stats_materialized,
|
||||
"player_stats_updated": player_stats_updated,
|
||||
"errors": errors,
|
||||
}
|
||||
|
||||
|
||||
def list_materialized_rcon_matches(
|
||||
*,
|
||||
target_key: str | None = None,
|
||||
only_ended: bool = False,
|
||||
limit: int = 20,
|
||||
db_path: Path | None = None,
|
||||
) -> list[dict[str, object]]:
|
||||
"""Return recent materialized RCON matches."""
|
||||
resolved_path = initialize_rcon_materialized_storage(db_path=db_path)
|
||||
clauses: list[str] = []
|
||||
params: list[object] = []
|
||||
if target_key:
|
||||
clauses.append("(target_key = ? OR external_server_id = ?)")
|
||||
params.extend([target_key, target_key])
|
||||
if only_ended:
|
||||
clauses.append("source_basis = ?")
|
||||
params.append(MATCH_RESULT_SOURCE)
|
||||
where = "WHERE " + " AND ".join(clauses) if clauses else ""
|
||||
params.append(limit)
|
||||
with closing(connect_sqlite_readonly(resolved_path)) as connection:
|
||||
rows = connection.execute(
|
||||
f"""
|
||||
SELECT *
|
||||
FROM rcon_materialized_matches
|
||||
{where}
|
||||
ORDER BY COALESCE(ended_at, started_at) DESC,
|
||||
COALESCE(ended_server_time, started_server_time) DESC
|
||||
LIMIT ?
|
||||
""",
|
||||
params,
|
||||
).fetchall()
|
||||
return [dict(row) for row in rows]
|
||||
|
||||
|
||||
def get_materialized_rcon_match_detail(
|
||||
*,
|
||||
server_key: str,
|
||||
match_key: str,
|
||||
db_path: Path | None = None,
|
||||
) -> dict[str, object] | None:
|
||||
"""Return one materialized match with player stats."""
|
||||
resolved_path = initialize_rcon_materialized_storage(db_path=db_path)
|
||||
with closing(connect_sqlite_readonly(resolved_path)) as connection:
|
||||
match = connection.execute(
|
||||
"""
|
||||
SELECT *
|
||||
FROM rcon_materialized_matches
|
||||
WHERE match_key = ?
|
||||
AND (target_key = ? OR external_server_id = ?)
|
||||
LIMIT 1
|
||||
""",
|
||||
(match_key, server_key, server_key),
|
||||
).fetchone()
|
||||
if match is None:
|
||||
return None
|
||||
stat_rows = connection.execute(
|
||||
"""
|
||||
SELECT *
|
||||
FROM rcon_match_player_stats
|
||||
WHERE target_key = ? AND match_key = ?
|
||||
ORDER BY kills DESC, deaths ASC, player_name ASC
|
||||
""",
|
||||
(match["target_key"], match["match_key"]),
|
||||
).fetchall()
|
||||
timeline_rows = connection.execute(
|
||||
"""
|
||||
SELECT event_type, COUNT(*) AS event_count
|
||||
FROM rcon_admin_log_events
|
||||
WHERE target_key = ?
|
||||
AND server_time IS NOT NULL
|
||||
AND (? IS NULL OR server_time >= ?)
|
||||
AND (? IS NULL OR server_time <= ?)
|
||||
GROUP BY event_type
|
||||
ORDER BY event_count DESC, event_type ASC
|
||||
""",
|
||||
(
|
||||
match["target_key"],
|
||||
match["started_server_time"],
|
||||
match["started_server_time"],
|
||||
match["ended_server_time"],
|
||||
match["ended_server_time"],
|
||||
),
|
||||
).fetchall()
|
||||
|
||||
return {
|
||||
"match": dict(match),
|
||||
"players": [dict(row) for row in stat_rows],
|
||||
"timeline": [dict(row) for row in timeline_rows],
|
||||
}
|
||||
|
||||
|
||||
def summarize_rcon_materialization_status(*, db_path: Path | None = None) -> dict[str, object]:
|
||||
"""Return a small diagnostic summary for stored RCON materialization state."""
|
||||
resolved_path = initialize_rcon_materialized_storage(db_path=db_path)
|
||||
with closing(connect_sqlite_readonly(resolved_path)) as connection:
|
||||
match_count = connection.execute(
|
||||
"SELECT COUNT(*) AS count FROM rcon_materialized_matches"
|
||||
).fetchone()["count"]
|
||||
stats_match_count = connection.execute(
|
||||
"""
|
||||
SELECT COUNT(*) AS count
|
||||
FROM (
|
||||
SELECT 1
|
||||
FROM rcon_match_player_stats
|
||||
GROUP BY target_key, match_key
|
||||
)
|
||||
"""
|
||||
).fetchone()["count"]
|
||||
ranges = connection.execute(
|
||||
"""
|
||||
SELECT target_key, MIN(server_time) AS first_server_time, MAX(server_time) AS last_server_time
|
||||
FROM rcon_admin_log_events
|
||||
GROUP BY target_key
|
||||
ORDER BY target_key ASC
|
||||
"""
|
||||
).fetchall()
|
||||
event_counts = connection.execute(
|
||||
"""
|
||||
SELECT target_key, event_type, COUNT(*) AS event_count
|
||||
FROM rcon_admin_log_events
|
||||
GROUP BY target_key, event_type
|
||||
ORDER BY target_key ASC, event_count DESC
|
||||
"""
|
||||
).fetchall()
|
||||
return {
|
||||
"materialized_matches": int(match_count or 0),
|
||||
"matches_with_player_stats": int(stats_match_count or 0),
|
||||
"server_time_ranges": [dict(row) for row in ranges],
|
||||
"event_counts": [dict(row) for row in event_counts],
|
||||
}
|
||||
|
||||
|
||||
def _derive_admin_log_matches(connection: sqlite3.Connection) -> list[dict[str, object]]:
|
||||
rows = connection.execute(
|
||||
"""
|
||||
SELECT *
|
||||
FROM rcon_admin_log_events
|
||||
WHERE event_type IN ('match_start', 'match_end')
|
||||
ORDER BY target_key ASC, server_time ASC, id ASC
|
||||
"""
|
||||
).fetchall()
|
||||
matches: list[dict[str, object]] = []
|
||||
open_by_target: dict[str, sqlite3.Row] = {}
|
||||
for row in rows:
|
||||
target_key = row["target_key"]
|
||||
payload = _json_object(row["parsed_payload_json"])
|
||||
if row["event_type"] == "match_start":
|
||||
if target_key in open_by_target:
|
||||
matches.append(_build_match_row(open_by_target.pop(target_key), None))
|
||||
open_by_target[target_key] = row
|
||||
continue
|
||||
start_row = open_by_target.pop(target_key, None)
|
||||
matches.append(_build_match_row(start_row, row, end_payload=payload))
|
||||
for start_row in open_by_target.values():
|
||||
matches.append(_build_match_row(start_row, None))
|
||||
return matches
|
||||
|
||||
|
||||
def _derive_session_fallback_matches(
|
||||
connection: sqlite3.Connection,
|
||||
*,
|
||||
db_path: Path,
|
||||
) -> list[dict[str, object]]:
|
||||
rows: list[dict[str, object]] = []
|
||||
existing = {
|
||||
(row["target_key"], normalize_map_name(row["map_pretty_name"] or row["map_name"]))
|
||||
for row in connection.execute(
|
||||
"""
|
||||
SELECT target_key, map_name, map_pretty_name
|
||||
FROM rcon_materialized_matches
|
||||
WHERE source_basis = ?
|
||||
""",
|
||||
(MATCH_RESULT_SOURCE,),
|
||||
).fetchall()
|
||||
}
|
||||
for window in list_rcon_historical_competitive_windows(limit=100, db_path=db_path):
|
||||
target_key = str(window.get("target_key") or "")
|
||||
map_name = window.get("map_pretty_name") or window.get("map_name")
|
||||
if (target_key, normalize_map_name(map_name)) in existing:
|
||||
continue
|
||||
session_key = str(window.get("session_key") or "").strip()
|
||||
if not target_key or not session_key:
|
||||
continue
|
||||
rows.append(
|
||||
{
|
||||
"target_key": target_key,
|
||||
"external_server_id": window.get("external_server_id"),
|
||||
"match_key": f"session:{session_key}",
|
||||
"map_name": window.get("map_name"),
|
||||
"map_pretty_name": normalize_map_name(map_name),
|
||||
"game_mode": None,
|
||||
"started_server_time": None,
|
||||
"ended_server_time": None,
|
||||
"started_at": window.get("first_seen_at"),
|
||||
"ended_at": window.get("last_seen_at"),
|
||||
"allied_score": _nested_int(window.get("latest_payload"), "allied_score"),
|
||||
"axis_score": _nested_int(window.get("latest_payload"), "axis_score"),
|
||||
"winner": _resolve_winner(
|
||||
_nested_int(window.get("latest_payload"), "allied_score"),
|
||||
_nested_int(window.get("latest_payload"), "axis_score"),
|
||||
),
|
||||
"confidence_mode": "partial",
|
||||
"source_basis": SESSION_RESULT_SOURCE,
|
||||
}
|
||||
)
|
||||
return rows
|
||||
|
||||
|
||||
def _build_match_row(
|
||||
start_row: sqlite3.Row | None,
|
||||
end_row: sqlite3.Row | None,
|
||||
*,
|
||||
end_payload: dict[str, object] | None = None,
|
||||
) -> dict[str, object]:
|
||||
start_payload = _json_object(start_row["parsed_payload_json"]) if start_row else {}
|
||||
end_payload = end_payload or (_json_object(end_row["parsed_payload_json"]) if end_row else {})
|
||||
target_key = str((end_row or start_row)["target_key"])
|
||||
external_server_id = (end_row or start_row)["external_server_id"]
|
||||
started_server_time = start_row["server_time"] if start_row else None
|
||||
ended_server_time = end_row["server_time"] if end_row else None
|
||||
map_name = end_payload.get("map_name") or start_payload.get("map_name")
|
||||
match_key = _build_match_key(
|
||||
target_key=target_key,
|
||||
started_server_time=started_server_time,
|
||||
ended_server_time=ended_server_time,
|
||||
map_name=map_name,
|
||||
)
|
||||
return {
|
||||
"target_key": target_key,
|
||||
"external_server_id": external_server_id,
|
||||
"match_key": match_key,
|
||||
"map_name": map_name,
|
||||
"map_pretty_name": normalize_map_name(map_name),
|
||||
"game_mode": start_payload.get("game_mode"),
|
||||
"started_server_time": started_server_time,
|
||||
"ended_server_time": ended_server_time,
|
||||
"started_at": start_row["event_timestamp"] if start_row else None,
|
||||
"ended_at": end_row["event_timestamp"] if end_row else None,
|
||||
"allied_score": _coerce_int(end_payload.get("allied_score")),
|
||||
"axis_score": _coerce_int(end_payload.get("axis_score")),
|
||||
"winner": end_payload.get("winner")
|
||||
or _resolve_winner(
|
||||
_coerce_int(end_payload.get("allied_score")),
|
||||
_coerce_int(end_payload.get("axis_score")),
|
||||
),
|
||||
"confidence_mode": "exact" if end_row else "partial",
|
||||
"source_basis": MATCH_RESULT_SOURCE if end_row else "admin-log-match-start",
|
||||
}
|
||||
|
||||
|
||||
def _upsert_match(connection: sqlite3.Connection, row: dict[str, object]) -> str:
|
||||
existing = connection.execute(
|
||||
"""
|
||||
SELECT id
|
||||
FROM rcon_materialized_matches
|
||||
WHERE target_key = ? AND match_key = ?
|
||||
""",
|
||||
(row["target_key"], row["match_key"]),
|
||||
).fetchone()
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO rcon_materialized_matches (
|
||||
target_key, external_server_id, match_key, map_name, map_pretty_name, game_mode,
|
||||
started_server_time, ended_server_time, started_at, ended_at,
|
||||
allied_score, axis_score, winner, confidence_mode, source_basis
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(target_key, match_key) DO UPDATE SET
|
||||
external_server_id = excluded.external_server_id,
|
||||
map_name = excluded.map_name,
|
||||
map_pretty_name = excluded.map_pretty_name,
|
||||
game_mode = excluded.game_mode,
|
||||
started_server_time = excluded.started_server_time,
|
||||
ended_server_time = excluded.ended_server_time,
|
||||
started_at = excluded.started_at,
|
||||
ended_at = excluded.ended_at,
|
||||
allied_score = excluded.allied_score,
|
||||
axis_score = excluded.axis_score,
|
||||
winner = excluded.winner,
|
||||
confidence_mode = excluded.confidence_mode,
|
||||
source_basis = excluded.source_basis,
|
||||
updated_at = CURRENT_TIMESTAMP
|
||||
""",
|
||||
(
|
||||
row["target_key"],
|
||||
row.get("external_server_id"),
|
||||
row["match_key"],
|
||||
row.get("map_name"),
|
||||
row.get("map_pretty_name"),
|
||||
row.get("game_mode"),
|
||||
row.get("started_server_time"),
|
||||
row.get("ended_server_time"),
|
||||
row.get("started_at"),
|
||||
row.get("ended_at"),
|
||||
row.get("allied_score"),
|
||||
row.get("axis_score"),
|
||||
row.get("winner"),
|
||||
row["confidence_mode"],
|
||||
row["source_basis"],
|
||||
),
|
||||
)
|
||||
return "updated" if existing else "inserted"
|
||||
|
||||
|
||||
def _list_materialized_matches(connection: sqlite3.Connection) -> list[dict[str, object]]:
|
||||
rows = connection.execute(
|
||||
"""
|
||||
SELECT *
|
||||
FROM rcon_materialized_matches
|
||||
WHERE started_server_time IS NOT NULL OR ended_server_time IS NOT NULL
|
||||
ORDER BY target_key ASC, COALESCE(started_server_time, ended_server_time) ASC
|
||||
"""
|
||||
).fetchall()
|
||||
return [dict(row) for row in rows]
|
||||
|
||||
|
||||
def _derive_player_stats_for_match(
|
||||
connection: sqlite3.Connection,
|
||||
match: dict[str, object],
|
||||
) -> list[dict[str, object]]:
|
||||
lower = match.get("started_server_time")
|
||||
upper = match.get("ended_server_time")
|
||||
if lower is None and upper is None:
|
||||
return []
|
||||
clauses = ["target_key = ?", "server_time IS NOT NULL"]
|
||||
params: list[object] = [match["target_key"]]
|
||||
if lower is not None:
|
||||
clauses.append("server_time >= ?")
|
||||
params.append(lower)
|
||||
if upper is not None:
|
||||
clauses.append("server_time <= ?")
|
||||
params.append(upper)
|
||||
rows = connection.execute(
|
||||
f"""
|
||||
SELECT *
|
||||
FROM rcon_admin_log_events
|
||||
WHERE {" AND ".join(clauses)}
|
||||
AND event_type IN ('kill', 'team_switch', 'connected', 'disconnected', 'chat')
|
||||
ORDER BY server_time ASC, id ASC
|
||||
""",
|
||||
params,
|
||||
).fetchall()
|
||||
|
||||
players: dict[str, dict[str, object]] = {}
|
||||
team_by_player: dict[str, str] = {}
|
||||
for row in rows:
|
||||
payload = _json_object(row["parsed_payload_json"])
|
||||
server_time = _coerce_int(row["server_time"])
|
||||
event_type = row["event_type"]
|
||||
if event_type == "kill":
|
||||
killer_key = _player_key(payload.get("killer_id"), payload.get("killer_name"))
|
||||
victim_key = _player_key(payload.get("victim_id"), payload.get("victim_name"))
|
||||
killer = _ensure_player(
|
||||
players,
|
||||
player_id=killer_key,
|
||||
player_name=payload.get("killer_name"),
|
||||
team=payload.get("killer_team") or team_by_player.get(killer_key),
|
||||
server_time=server_time,
|
||||
)
|
||||
victim = _ensure_player(
|
||||
players,
|
||||
player_id=victim_key,
|
||||
player_name=payload.get("victim_name"),
|
||||
team=payload.get("victim_team") or team_by_player.get(victim_key),
|
||||
server_time=server_time,
|
||||
)
|
||||
team_by_player[killer_key] = str(payload.get("killer_team") or killer.get("team") or "")
|
||||
team_by_player[victim_key] = str(payload.get("victim_team") or victim.get("team") or "")
|
||||
weapon = str(payload.get("weapon") or "Unknown")
|
||||
same_team = payload.get("killer_team") and payload.get("killer_team") == payload.get("victim_team")
|
||||
if same_team:
|
||||
killer["teamkills"] = int(killer["teamkills"]) + 1
|
||||
victim["deaths_by_teamkill"] = int(victim["deaths_by_teamkill"]) + 1
|
||||
else:
|
||||
killer["kills"] = int(killer["kills"]) + 1
|
||||
victim["deaths"] = int(victim["deaths"]) + 1
|
||||
_counter(killer, "weapons")[weapon] += 1
|
||||
_counter(victim, "death_by_weapons")[weapon] += 1
|
||||
_counter(killer, "most_killed")[str(victim["player_name"])] += 1
|
||||
_counter(victim, "death_by")[str(killer["player_name"])] += 1
|
||||
_touch_player(killer, server_time)
|
||||
_touch_player(victim, server_time)
|
||||
continue
|
||||
|
||||
if event_type == "team_switch" and not payload.get("player_id"):
|
||||
continue
|
||||
player_id = _player_key(payload.get("player_id"), payload.get("player_name"))
|
||||
team = payload.get("to_team") or payload.get("chat_team") or team_by_player.get(player_id)
|
||||
player = _ensure_player(
|
||||
players,
|
||||
player_id=player_id,
|
||||
player_name=payload.get("player_name"),
|
||||
team=team,
|
||||
server_time=server_time,
|
||||
)
|
||||
if team:
|
||||
player["team"] = team
|
||||
team_by_player[player_id] = str(team)
|
||||
_touch_player(player, server_time)
|
||||
|
||||
stats = []
|
||||
for player in players.values():
|
||||
stats.append(
|
||||
{
|
||||
"target_key": match["target_key"],
|
||||
"match_key": match["match_key"],
|
||||
"player_id": player["player_id"],
|
||||
"player_name": player["player_name"],
|
||||
"team": player.get("team"),
|
||||
"kills": player["kills"],
|
||||
"deaths": player["deaths"],
|
||||
"teamkills": player["teamkills"],
|
||||
"deaths_by_teamkill": player["deaths_by_teamkill"],
|
||||
"weapons_json": _dump_counter(player["weapons"]),
|
||||
"death_by_weapons_json": _dump_counter(player["death_by_weapons"]),
|
||||
"most_killed_json": _dump_counter(player["most_killed"]),
|
||||
"death_by_json": _dump_counter(player["death_by"]),
|
||||
"first_seen_server_time": player.get("first_seen_server_time"),
|
||||
"last_seen_server_time": player.get("last_seen_server_time"),
|
||||
}
|
||||
)
|
||||
return stats
|
||||
|
||||
|
||||
def _insert_player_stat(connection: sqlite3.Connection, stat: dict[str, object]) -> None:
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO rcon_match_player_stats (
|
||||
target_key, match_key, player_id, player_name, team,
|
||||
kills, deaths, teamkills, deaths_by_teamkill,
|
||||
weapons_json, death_by_weapons_json, most_killed_json, death_by_json,
|
||||
first_seen_server_time, last_seen_server_time
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
stat["target_key"],
|
||||
stat["match_key"],
|
||||
stat["player_id"],
|
||||
stat["player_name"],
|
||||
stat.get("team"),
|
||||
stat["kills"],
|
||||
stat["deaths"],
|
||||
stat["teamkills"],
|
||||
stat["deaths_by_teamkill"],
|
||||
stat["weapons_json"],
|
||||
stat["death_by_weapons_json"],
|
||||
stat["most_killed_json"],
|
||||
stat["death_by_json"],
|
||||
stat.get("first_seen_server_time"),
|
||||
stat.get("last_seen_server_time"),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _ensure_player(
|
||||
players: dict[str, dict[str, object]],
|
||||
*,
|
||||
player_id: str,
|
||||
player_name: object,
|
||||
team: object,
|
||||
server_time: int | None,
|
||||
) -> dict[str, object]:
|
||||
if player_id not in players:
|
||||
players[player_id] = {
|
||||
"player_id": player_id,
|
||||
"player_name": str(player_name or player_id),
|
||||
"team": team,
|
||||
"kills": 0,
|
||||
"deaths": 0,
|
||||
"teamkills": 0,
|
||||
"deaths_by_teamkill": 0,
|
||||
"weapons": Counter(),
|
||||
"death_by_weapons": Counter(),
|
||||
"most_killed": Counter(),
|
||||
"death_by": Counter(),
|
||||
"first_seen_server_time": server_time,
|
||||
"last_seen_server_time": server_time,
|
||||
}
|
||||
player = players[player_id]
|
||||
if player_name:
|
||||
player["player_name"] = str(player_name)
|
||||
if team:
|
||||
player["team"] = team
|
||||
_touch_player(player, server_time)
|
||||
return player
|
||||
|
||||
|
||||
def _touch_player(player: dict[str, object], server_time: int | None) -> None:
|
||||
if server_time is None:
|
||||
return
|
||||
first_seen = _coerce_int(player.get("first_seen_server_time"))
|
||||
last_seen = _coerce_int(player.get("last_seen_server_time"))
|
||||
player["first_seen_server_time"] = server_time if first_seen is None else min(first_seen, server_time)
|
||||
player["last_seen_server_time"] = server_time if last_seen is None else max(last_seen, server_time)
|
||||
|
||||
|
||||
def _counter(player: dict[str, object], key: str) -> Counter[str]:
|
||||
value = player[key]
|
||||
if isinstance(value, Counter):
|
||||
return value
|
||||
counter: Counter[str] = Counter()
|
||||
player[key] = counter
|
||||
return counter
|
||||
|
||||
|
||||
def _player_key(player_id: object, player_name: object) -> str:
|
||||
raw_id = str(player_id or "").strip()
|
||||
if raw_id:
|
||||
return raw_id
|
||||
return f"name:{str(player_name or 'unknown').strip().lower()}"
|
||||
|
||||
|
||||
def _build_match_key(
|
||||
*,
|
||||
target_key: str,
|
||||
started_server_time: object,
|
||||
ended_server_time: object,
|
||||
map_name: object,
|
||||
) -> str:
|
||||
map_part = "".join(character.lower() for character in str(map_name or "unknown") if character.isalnum())
|
||||
start_part = "missing" if started_server_time is None else str(started_server_time)
|
||||
end_part = "open" if ended_server_time is None else str(ended_server_time)
|
||||
return f"{target_key}:{start_part}:{end_part}:{map_part}"
|
||||
|
||||
|
||||
def _json_object(raw_value: object) -> dict[str, object]:
|
||||
if not isinstance(raw_value, str) or not raw_value.strip():
|
||||
return {}
|
||||
try:
|
||||
parsed = json.loads(raw_value)
|
||||
except json.JSONDecodeError:
|
||||
return {}
|
||||
return parsed if isinstance(parsed, dict) else {}
|
||||
|
||||
|
||||
def _dump_counter(counter: Counter[str]) -> str:
|
||||
ordered = dict(sorted(counter.items(), key=lambda item: (-item[1], item[0])))
|
||||
return json.dumps(ordered, ensure_ascii=False, separators=(",", ":"))
|
||||
|
||||
|
||||
def _coerce_int(value: object) -> int | None:
|
||||
if value is None:
|
||||
return None
|
||||
try:
|
||||
return int(value)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def _nested_int(payload: object, key: str) -> int | None:
|
||||
if not isinstance(payload, dict):
|
||||
return None
|
||||
return _coerce_int(payload.get(key))
|
||||
|
||||
|
||||
def _resolve_winner(allied_score: int | None, axis_score: int | None) -> str | None:
|
||||
if allied_score is None or axis_score is None:
|
||||
return None
|
||||
if allied_score > axis_score:
|
||||
return "allied"
|
||||
if axis_score > allied_score:
|
||||
return "axis"
|
||||
return "draw"
|
||||
|
||||
|
||||
def _main(argv: Iterable[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(description="Materialize stored RCON AdminLog events.")
|
||||
parser.add_argument(
|
||||
"command",
|
||||
nargs="?",
|
||||
choices=("materialize", "status"),
|
||||
default="materialize",
|
||||
)
|
||||
parser.add_argument("--db-path", type=Path, default=None)
|
||||
args = parser.parse_args(list(argv) if argv is not None else None)
|
||||
db_path = args.db_path or get_storage_path()
|
||||
payload = (
|
||||
summarize_rcon_materialization_status(db_path=db_path)
|
||||
if args.command == "status"
|
||||
else materialize_rcon_admin_log(db_path=db_path)
|
||||
)
|
||||
print(json.dumps({"status": "ok", "data": payload}, ensure_ascii=False, indent=2))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(_main())
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from .historical_storage import ALL_SERVERS_SLUG
|
||||
@@ -14,6 +15,9 @@ from .rcon_historical_storage import (
|
||||
list_rcon_historical_competitive_windows,
|
||||
)
|
||||
|
||||
MATCH_RESULT_SOURCE = "admin-log-match-ended"
|
||||
SESSION_RESULT_SOURCE = "rcon-session"
|
||||
|
||||
|
||||
def list_rcon_historical_server_summaries(
|
||||
*,
|
||||
@@ -41,9 +45,23 @@ def list_rcon_historical_recent_activity(
|
||||
limit: int = 20,
|
||||
) -> list[dict[str, object]]:
|
||||
"""Return recent RCON-backed competitive windows for one or all targets."""
|
||||
from .rcon_admin_log_materialization import list_materialized_rcon_matches
|
||||
|
||||
normalized_server_key = None if server_key == ALL_SERVERS_SLUG else server_key
|
||||
items = list_rcon_historical_competitive_windows(target_key=normalized_server_key, limit=limit)
|
||||
return [
|
||||
materialized_items = list_materialized_rcon_matches(
|
||||
target_key=normalized_server_key,
|
||||
only_ended=True,
|
||||
limit=limit,
|
||||
)
|
||||
primary_items = [_build_materialized_recent_item(item) for item in materialized_items]
|
||||
if primary_items:
|
||||
return primary_items[:limit]
|
||||
|
||||
session_items = list_rcon_historical_competitive_windows(
|
||||
target_key=normalized_server_key,
|
||||
limit=limit,
|
||||
)
|
||||
fallback_items = [
|
||||
{
|
||||
"server": {
|
||||
"slug": item["target_key"],
|
||||
@@ -52,6 +70,7 @@ def list_rcon_historical_recent_activity(
|
||||
"region": item["region"],
|
||||
},
|
||||
"match_id": item["session_key"],
|
||||
"internal_detail_match_id": item["session_key"],
|
||||
"started_at": item["first_seen_at"],
|
||||
"ended_at": item["last_seen_at"],
|
||||
"closed_at": item["last_seen_at"],
|
||||
@@ -66,11 +85,13 @@ def list_rcon_historical_recent_activity(
|
||||
"sample_count": item.get("sample_count"),
|
||||
"duration_seconds": item.get("duration_seconds"),
|
||||
"capture_basis": "rcon-competitive-window",
|
||||
"result_source": SESSION_RESULT_SOURCE,
|
||||
"capabilities": item.get("capabilities"),
|
||||
"minutes_since_capture": _minutes_since_timestamp(item.get("last_seen_at")),
|
||||
}
|
||||
for item in items
|
||||
for item in session_items
|
||||
]
|
||||
return _merge_recent_items(primary_items, fallback_items, limit=limit)
|
||||
|
||||
|
||||
def describe_rcon_historical_read_model() -> dict[str, object]:
|
||||
@@ -95,11 +116,11 @@ def describe_rcon_historical_read_model() -> dict[str, object]:
|
||||
],
|
||||
"capabilities": {
|
||||
"server_summary": "exact",
|
||||
"recent_matches": "approximate",
|
||||
"recent_matches": "exact-when-admin-log-match-ended",
|
||||
"competitive_quality": "partial",
|
||||
"result": "session-score",
|
||||
"gamestate": "session",
|
||||
"player_stats": "unavailable",
|
||||
"result": "admin-log-match-ended",
|
||||
"gamestate": "session-fallback",
|
||||
"player_stats": "admin-log-derived",
|
||||
},
|
||||
"limitations": [
|
||||
"No retroactive backfill of closed matches.",
|
||||
@@ -130,6 +151,12 @@ def get_rcon_historical_match_detail(
|
||||
match_id: str,
|
||||
) -> dict[str, object] | None:
|
||||
"""Return one RCON competitive window as a match-detail compatible payload."""
|
||||
from .rcon_admin_log_materialization import get_materialized_rcon_match_detail
|
||||
|
||||
materialized = get_materialized_rcon_match_detail(server_key=server_key, match_key=match_id)
|
||||
if materialized is not None:
|
||||
return _build_materialized_detail_item(materialized)
|
||||
|
||||
item = get_rcon_historical_competitive_window_by_session(
|
||||
server_key=server_key,
|
||||
session_key=match_id,
|
||||
@@ -161,6 +188,9 @@ def get_rcon_historical_match_detail(
|
||||
"sample_count": item.get("sample_count"),
|
||||
"players": [],
|
||||
"capture_basis": "rcon-competitive-window",
|
||||
"confidence": item.get("confidence_mode"),
|
||||
"source_basis": "rcon-session",
|
||||
"result_source": SESSION_RESULT_SOURCE,
|
||||
"capabilities": item.get("capabilities"),
|
||||
"match_url": resolve_rcon_scoreboard_match_url(
|
||||
server_slug=server_slug,
|
||||
@@ -174,6 +204,137 @@ def get_rcon_historical_match_detail(
|
||||
}
|
||||
|
||||
|
||||
def _build_materialized_recent_item(item: dict[str, object]) -> dict[str, object]:
|
||||
server_slug = item.get("external_server_id") or item.get("target_key")
|
||||
return {
|
||||
"server": {
|
||||
"slug": item.get("target_key"),
|
||||
"name": _server_display_name(item.get("external_server_id") or item.get("target_key")),
|
||||
"external_server_id": item.get("external_server_id"),
|
||||
"region": None,
|
||||
},
|
||||
"match_id": item.get("match_key"),
|
||||
"internal_detail_match_id": item.get("match_key"),
|
||||
"started_at": item.get("started_at"),
|
||||
"ended_at": item.get("ended_at"),
|
||||
"closed_at": item.get("ended_at") or item.get("started_at"),
|
||||
"map": {
|
||||
"name": item.get("map_name"),
|
||||
"pretty_name": item.get("map_pretty_name") or normalize_map_name(item.get("map_name")),
|
||||
},
|
||||
"game_mode": item.get("game_mode"),
|
||||
"result": {
|
||||
"allied_score": item.get("allied_score"),
|
||||
"axis_score": item.get("axis_score"),
|
||||
"winner": item.get("winner"),
|
||||
},
|
||||
"winner": item.get("winner"),
|
||||
"player_count": None,
|
||||
"peak_players": None,
|
||||
"sample_count": None,
|
||||
"duration_seconds": _calculate_match_duration_seconds(item),
|
||||
"capture_basis": "rcon-materialized-admin-log",
|
||||
"confidence": item.get("confidence_mode"),
|
||||
"source_basis": item.get("source_basis"),
|
||||
"result_source": (
|
||||
MATCH_RESULT_SOURCE
|
||||
if item.get("source_basis") == MATCH_RESULT_SOURCE
|
||||
else SESSION_RESULT_SOURCE
|
||||
),
|
||||
"match_url": resolve_rcon_scoreboard_match_url(
|
||||
server_slug=server_slug,
|
||||
map_name=item.get("map_pretty_name") or item.get("map_name"),
|
||||
started_at=item.get("started_at"),
|
||||
ended_at=item.get("ended_at"),
|
||||
duration_seconds=_calculate_match_duration_seconds(item),
|
||||
),
|
||||
"capabilities": describe_rcon_historical_read_model()["capabilities"],
|
||||
}
|
||||
|
||||
|
||||
def _build_materialized_detail_item(materialized: dict[str, object]) -> dict[str, object]:
|
||||
match = materialized["match"]
|
||||
recent_item = _build_materialized_recent_item(match)
|
||||
players = [_build_player_row(row) for row in materialized["players"]]
|
||||
return {
|
||||
**recent_item,
|
||||
"match_id": match["match_key"],
|
||||
"game_mode": match.get("game_mode"),
|
||||
"winner": match.get("winner"),
|
||||
"confidence": match.get("confidence_mode"),
|
||||
"source_basis": match.get("source_basis"),
|
||||
"players": players,
|
||||
"timeline": {
|
||||
"event_counts": materialized.get("timeline", []),
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def _build_player_row(row: dict[str, object]) -> dict[str, object]:
|
||||
kills = _coerce_optional_int(row.get("kills")) or 0
|
||||
deaths = _coerce_optional_int(row.get("deaths")) or 0
|
||||
return {
|
||||
"player_name": row.get("player_name"),
|
||||
"team": row.get("team"),
|
||||
"kills": kills,
|
||||
"deaths": deaths,
|
||||
"teamkills": _coerce_optional_int(row.get("teamkills")) or 0,
|
||||
"kd_ratio": round(kills / deaths, 2) if deaths else float(kills),
|
||||
"top_weapons": _top_counter(row.get("weapons_json")),
|
||||
"most_killed": _top_counter(row.get("most_killed_json")),
|
||||
"death_by": _top_counter(row.get("death_by_json")),
|
||||
}
|
||||
|
||||
|
||||
def _top_counter(raw_value: object, *, limit: int = 5) -> list[dict[str, object]]:
|
||||
if not isinstance(raw_value, str) or not raw_value.strip():
|
||||
return []
|
||||
try:
|
||||
payload = json.loads(raw_value)
|
||||
except (NameError, ValueError, TypeError):
|
||||
return []
|
||||
if not isinstance(payload, dict):
|
||||
return []
|
||||
rows = [
|
||||
{"name": str(name), "count": int(count)}
|
||||
for name, count in payload.items()
|
||||
if _coerce_optional_int(count) is not None
|
||||
]
|
||||
rows.sort(key=lambda item: (-int(item["count"]), str(item["name"])))
|
||||
return rows[:limit]
|
||||
|
||||
|
||||
def _merge_recent_items(
|
||||
primary_items: list[dict[str, object]],
|
||||
fallback_items: list[dict[str, object]],
|
||||
*,
|
||||
limit: int,
|
||||
) -> list[dict[str, object]]:
|
||||
merged: list[dict[str, object]] = []
|
||||
seen: set[tuple[object, object]] = set()
|
||||
for item in primary_items + fallback_items:
|
||||
map_payload = item.get("map") if isinstance(item.get("map"), dict) else {}
|
||||
key = (
|
||||
item.get("server", {}).get("slug") if isinstance(item.get("server"), dict) else None,
|
||||
normalize_map_name(map_payload.get("pretty_name") or map_payload.get("name")),
|
||||
)
|
||||
if key in seen:
|
||||
continue
|
||||
seen.add(key)
|
||||
merged.append(item)
|
||||
merged.sort(key=lambda row: str(row.get("closed_at") or row.get("ended_at") or row.get("started_at") or ""), reverse=True)
|
||||
return merged[:limit]
|
||||
|
||||
|
||||
def _server_display_name(server_slug: object) -> str:
|
||||
slug = str(server_slug or "").strip()
|
||||
if slug == "comunidad-hispana-01":
|
||||
return "Comunidad Hispana #01"
|
||||
if slug == "comunidad-hispana-02":
|
||||
return "Comunidad Hispana #02"
|
||||
return slug or "RCON"
|
||||
|
||||
|
||||
def _build_rcon_result(latest_payload: object) -> dict[str, object]:
|
||||
payload = latest_payload if isinstance(latest_payload, dict) else {}
|
||||
allied_score = _coerce_optional_int(payload.get("allied_score"))
|
||||
@@ -341,3 +502,26 @@ def _calculate_coverage_hours(
|
||||
last_point = last_point.replace(tzinfo=timezone.utc)
|
||||
delta = last_point.astimezone(timezone.utc) - first_point.astimezone(timezone.utc)
|
||||
return round(delta.total_seconds() / 3600, 2)
|
||||
|
||||
|
||||
def _calculate_duration_seconds(first_seen_at: object, last_seen_at: object) -> int | None:
|
||||
if not isinstance(first_seen_at, str) or not isinstance(last_seen_at, str):
|
||||
return None
|
||||
first_point = datetime.fromisoformat(first_seen_at.replace("Z", "+00:00"))
|
||||
last_point = datetime.fromisoformat(last_seen_at.replace("Z", "+00:00"))
|
||||
if first_point.tzinfo is None:
|
||||
first_point = first_point.replace(tzinfo=timezone.utc)
|
||||
if last_point.tzinfo is None:
|
||||
last_point = last_point.replace(tzinfo=timezone.utc)
|
||||
return max(0, int((last_point.astimezone(timezone.utc) - first_point.astimezone(timezone.utc)).total_seconds()))
|
||||
|
||||
|
||||
def _calculate_match_duration_seconds(item: dict[str, object]) -> int | None:
|
||||
duration = _calculate_duration_seconds(item.get("started_at"), item.get("ended_at"))
|
||||
if duration:
|
||||
return duration
|
||||
started_server_time = _coerce_optional_int(item.get("started_server_time"))
|
||||
ended_server_time = _coerce_optional_int(item.get("ended_server_time"))
|
||||
if started_server_time is None or ended_server_time is None:
|
||||
return duration
|
||||
return max(0, ended_server_time - started_server_time)
|
||||
|
||||
@@ -696,7 +696,8 @@ def get_rcon_historical_competitive_window_by_session(
|
||||
windows.total_players,
|
||||
windows.peak_players,
|
||||
windows.confidence_mode,
|
||||
windows.capabilities_json
|
||||
windows.capabilities_json,
|
||||
windows.latest_payload_json
|
||||
FROM rcon_historical_competitive_windows AS windows
|
||||
INNER JOIN rcon_historical_targets AS targets
|
||||
ON targets.id = windows.target_id
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
import re
|
||||
from urllib.parse import urlparse
|
||||
|
||||
|
||||
@@ -32,6 +33,8 @@ TRUSTED_PUBLIC_SCOREBOARD_ORIGINS = (
|
||||
),
|
||||
)
|
||||
|
||||
_TRUSTED_GAME_PATH_RE = re.compile(r"^/games/\d+/?$")
|
||||
|
||||
|
||||
def list_trusted_public_scoreboard_origins() -> tuple[TrustedScoreboardOrigin, ...]:
|
||||
"""Return trusted public scoreboard origins for active default servers."""
|
||||
@@ -71,6 +74,8 @@ def resolve_trusted_scoreboard_match_url(
|
||||
return None
|
||||
if candidate_parts.username or candidate_parts.password:
|
||||
return None
|
||||
if not candidate_parts.path.startswith("/games/"):
|
||||
if not _TRUSTED_GAME_PATH_RE.match(candidate_parts.path):
|
||||
return None
|
||||
if candidate_parts.params or candidate_parts.query or candidate_parts.fragment:
|
||||
return None
|
||||
return candidate
|
||||
|
||||
221
backend/tests/test_rcon_materialization_pipeline.py
Normal file
221
backend/tests/test_rcon_materialization_pipeline.py
Normal file
@@ -0,0 +1,221 @@
|
||||
"""Regression tests for the materialized RCON AdminLog pipeline."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import gc
|
||||
import os
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
|
||||
from app.historical_storage import upsert_historical_match
|
||||
from app.payloads import build_recent_historical_matches_payload
|
||||
from app.rcon_admin_log_materialization import (
|
||||
get_materialized_rcon_match_detail,
|
||||
materialize_rcon_admin_log,
|
||||
summarize_rcon_materialization_status,
|
||||
)
|
||||
from app.rcon_admin_log_storage import persist_rcon_admin_log_entries
|
||||
from app.rcon_historical_read_model import (
|
||||
get_rcon_historical_match_detail,
|
||||
list_rcon_historical_recent_activity,
|
||||
)
|
||||
from app.scoreboard_origins import resolve_trusted_scoreboard_match_url
|
||||
|
||||
|
||||
class RconMaterializationPipelineTests(unittest.TestCase):
|
||||
def test_materializes_match_result_and_player_stats_idempotently(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmpdir:
|
||||
db_path = Path(tmpdir) / "historical.sqlite3"
|
||||
_persist_admin_log_fixture(db_path)
|
||||
|
||||
first = materialize_rcon_admin_log(db_path=db_path)
|
||||
second = materialize_rcon_admin_log(db_path=db_path)
|
||||
detail = get_materialized_rcon_match_detail(
|
||||
server_key="comunidad-hispana-01",
|
||||
match_key="comunidad-hispana-01:100:500:stmariedumontwarfare",
|
||||
db_path=db_path,
|
||||
)
|
||||
status = summarize_rcon_materialization_status(db_path=db_path)
|
||||
|
||||
self.assertEqual(first["matches_materialized"], 1)
|
||||
self.assertEqual(second["matches_materialized"], 0)
|
||||
self.assertEqual(second["matches_updated"], 1)
|
||||
self.assertIsNotNone(detail)
|
||||
match = detail["match"]
|
||||
self.assertEqual(match["allied_score"], 5)
|
||||
self.assertEqual(match["axis_score"], 0)
|
||||
self.assertEqual(match["winner"], "allied")
|
||||
players = {row["player_name"]: row for row in detail["players"]}
|
||||
self.assertEqual(players["Alpha"]["kills"], 1)
|
||||
self.assertEqual(players["Alpha"]["teamkills"], 1)
|
||||
self.assertEqual(players["Bravo"]["deaths"], 1)
|
||||
self.assertEqual(players["Charlie"]["deaths_by_teamkill"], 1)
|
||||
self.assertEqual(status["materialized_matches"], 1)
|
||||
self.assertEqual(status["matches_with_player_stats"], 1)
|
||||
gc.collect()
|
||||
|
||||
def test_match_detail_read_model_hides_raw_player_ids(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmpdir:
|
||||
db_path = Path(tmpdir) / "historical.sqlite3"
|
||||
previous_storage_path = os.environ.get("HLL_BACKEND_STORAGE_PATH")
|
||||
os.environ["HLL_BACKEND_STORAGE_PATH"] = str(db_path)
|
||||
try:
|
||||
_persist_admin_log_fixture(db_path)
|
||||
materialize_rcon_admin_log(db_path=db_path)
|
||||
detail = get_rcon_historical_match_detail(
|
||||
server_key="comunidad-hispana-01",
|
||||
match_id="comunidad-hispana-01:100:500:stmariedumontwarfare",
|
||||
)
|
||||
finally:
|
||||
_restore_env("HLL_BACKEND_STORAGE_PATH", previous_storage_path)
|
||||
|
||||
self.assertIsNotNone(detail)
|
||||
self.assertEqual(detail["result_source"], "admin-log-match-ended")
|
||||
self.assertEqual(detail["result"]["allied_score"], 5)
|
||||
self.assertNotIn("player_id", detail["players"][0])
|
||||
self.assertIn("kd_ratio", detail["players"][0])
|
||||
gc.collect()
|
||||
|
||||
def test_recent_matches_prefer_materialized_rcon_over_scoreboard_fallback(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmpdir:
|
||||
db_path = Path(tmpdir) / "historical.sqlite3"
|
||||
previous_storage_path = os.environ.get("HLL_BACKEND_STORAGE_PATH")
|
||||
os.environ["HLL_BACKEND_STORAGE_PATH"] = str(db_path)
|
||||
try:
|
||||
_persist_admin_log_fixture(db_path)
|
||||
materialize_rcon_admin_log(db_path=db_path)
|
||||
_persist_scoreboard_match(db_path)
|
||||
|
||||
payload = build_recent_historical_matches_payload(
|
||||
limit=5,
|
||||
server_slug="comunidad-hispana-01",
|
||||
)
|
||||
recent = list_rcon_historical_recent_activity(
|
||||
server_key="comunidad-hispana-01",
|
||||
limit=5,
|
||||
)
|
||||
finally:
|
||||
_restore_env("HLL_BACKEND_STORAGE_PATH", previous_storage_path)
|
||||
|
||||
self.assertEqual(payload["data"]["selected_source"], "rcon")
|
||||
self.assertEqual(payload["data"]["items"][0]["result_source"], "admin-log-match-ended")
|
||||
self.assertEqual(recent[0]["result_source"], "admin-log-match-ended")
|
||||
self.assertNotEqual(payload["data"]["selected_source"], "public-scoreboard")
|
||||
gc.collect()
|
||||
|
||||
def test_public_scoreboard_fallback_used_only_without_rcon_activity(self) -> None:
|
||||
with tempfile.TemporaryDirectory() as tmpdir:
|
||||
db_path = Path(tmpdir) / "historical.sqlite3"
|
||||
previous_storage_path = os.environ.get("HLL_BACKEND_STORAGE_PATH")
|
||||
os.environ["HLL_BACKEND_STORAGE_PATH"] = str(db_path)
|
||||
try:
|
||||
_persist_scoreboard_match(db_path)
|
||||
payload = build_recent_historical_matches_payload(
|
||||
limit=5,
|
||||
server_slug="comunidad-hispana-01",
|
||||
)
|
||||
finally:
|
||||
_restore_env("HLL_BACKEND_STORAGE_PATH", previous_storage_path)
|
||||
|
||||
self.assertTrue(payload["data"]["fallback_used"])
|
||||
self.assertEqual(payload["data"]["selected_source"], "public-scoreboard")
|
||||
self.assertEqual(payload["data"]["items"][0]["result_source"], "public-scoreboard-fallback")
|
||||
gc.collect()
|
||||
|
||||
def test_safe_scoreboard_match_url_allowlist_for_active_origins(self) -> None:
|
||||
self.assertEqual(
|
||||
resolve_trusted_scoreboard_match_url(
|
||||
"https://scoreboard.comunidadhll.es/games/1561515",
|
||||
"comunidad-hispana-01",
|
||||
),
|
||||
"https://scoreboard.comunidadhll.es/games/1561515",
|
||||
)
|
||||
self.assertEqual(
|
||||
resolve_trusted_scoreboard_match_url(
|
||||
"https://scoreboard.comunidadhll.es:5443/games/222",
|
||||
"comunidad-hispana-02",
|
||||
),
|
||||
"https://scoreboard.comunidadhll.es:5443/games/222",
|
||||
)
|
||||
self.assertIsNone(
|
||||
resolve_trusted_scoreboard_match_url(
|
||||
"https://example.com/games/222",
|
||||
"comunidad-hispana-02",
|
||||
)
|
||||
)
|
||||
self.assertIsNone(
|
||||
resolve_trusted_scoreboard_match_url(
|
||||
"https://scoreboard.comunidadhll.es:5443/admin/222",
|
||||
"comunidad-hispana-02",
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
def _persist_admin_log_fixture(db_path: Path) -> None:
|
||||
persist_rcon_admin_log_entries(
|
||||
target={
|
||||
"target_key": "comunidad-hispana-01",
|
||||
"external_server_id": "comunidad-hispana-01",
|
||||
},
|
||||
entries=[
|
||||
{
|
||||
"timestamp": "2026-05-01T10:00:00Z",
|
||||
"message": "[1 min (100)] MATCH START ST MARIE DU MONT Warfare",
|
||||
},
|
||||
{
|
||||
"timestamp": "2026-05-01T10:05:00Z",
|
||||
"message": "[6 min (150)] CONNECTED Alpha (76561198000000001)",
|
||||
},
|
||||
{
|
||||
"timestamp": "2026-05-01T10:06:00Z",
|
||||
"message": "[7 min (160)] TEAMSWITCH Alpha (None > Allies)",
|
||||
},
|
||||
{
|
||||
"timestamp": "2026-05-01T10:10:00Z",
|
||||
"message": (
|
||||
"[11 min (200)] KILL: Alpha(Allies/76561198000000001) -> "
|
||||
"Bravo(Axis/76561198000000002) with M1 Garand"
|
||||
),
|
||||
},
|
||||
{
|
||||
"timestamp": "2026-05-01T10:12:00Z",
|
||||
"message": (
|
||||
"[13 min (220)] KILL: Alpha(Allies/76561198000000001) -> "
|
||||
"Charlie(Allies/nonsteam-local) with M1 Garand"
|
||||
),
|
||||
},
|
||||
{
|
||||
"timestamp": "2026-05-01T11:20:00Z",
|
||||
"message": "[81 min (500)] MATCH ENDED `ST MARIE DU MONT Warfare` ALLIED (5 - 0) AXIS",
|
||||
},
|
||||
],
|
||||
db_path=db_path,
|
||||
)
|
||||
|
||||
|
||||
def _persist_scoreboard_match(db_path: Path) -> None:
|
||||
upsert_historical_match(
|
||||
server_slug="comunidad-hispana-01",
|
||||
match_payload={
|
||||
"id": "1561515",
|
||||
"creation_time": "2026-05-01T10:00:00Z",
|
||||
"start": "2026-05-01T10:00:00Z",
|
||||
"end": "2026-05-01T11:20:00Z",
|
||||
"map": {"name": "ST MARIE DU MONT Warfare"},
|
||||
"result": {"allied": 2, "axis": 3},
|
||||
"player_stats": [],
|
||||
},
|
||||
db_path=db_path,
|
||||
)
|
||||
|
||||
|
||||
def _restore_env(name: str, previous_value: str | None) -> None:
|
||||
if previous_value is None:
|
||||
os.environ.pop(name, None)
|
||||
else:
|
||||
os.environ[name] = previous_value
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user