diff --git a/ai/tasks/done/TASK-100-rcon-historical-writer-path-implementation.md b/ai/tasks/done/TASK-100-rcon-historical-writer-path-implementation.md new file mode 100644 index 0000000..403dd90 --- /dev/null +++ b/ai/tasks/done/TASK-100-rcon-historical-writer-path-implementation.md @@ -0,0 +1,53 @@ +# TASK-100-rcon-historical-writer-path-implementation + +## Goal +Implementar un writer path histórico real por RCON para que la ingesta histórica intente RCON primero y use scoreboard/public-scoreboard solo como fallback. + +## Context +La repo ya tiene: +- live RCON-first +- captura prospectiva RCON +- read model histórico RCON parcial +- ingesta histórica clásica por scoreboard +Pero todavía no existe un writer path histórico real por RCON integrado en `historical_ingestion`. + +## Steps +1. Auditar: + - `backend/app/data_sources.py` + - `backend/app/historical_ingestion.py` + - `backend/app/rcon_historical_worker.py` + - `backend/app/rcon_historical_storage.py` + - `backend/app/rcon_historical_read_model.py` + - cualquier capa necesaria de modelos/storage +2. Definir qué significa “writer path histórico RCON” con la telemetría real actual: + - qué puede alimentar + - qué estructura persistida necesita + - cómo se integra con el ingestion flow existente +3. Implementar una vía writer-oriented RCON que permita a `historical_ingestion` intentar primero RCON. +4. Si RCON falla o no cubre la operación concreta, hacer fallback controlado a scoreboard/public-scoreboard. +5. Mantener trazabilidad explícita: + - primary_source + - selected_source + - fallback_used + - fallback_reason + - source_attempts +6. No romper la compatibilidad con snapshots ni con rebuilds posteriores. +7. Actualizar README/runbook explicando el nuevo writer path real. + +## Constraints +- No fingir cobertura histórica RCON que no exista. +- No eliminar scoreboard como fallback. +- No romper live RCON-first. +- No romper el locking compartido. + +## Validation +- `historical_ingestion` intenta RCON primero. +- Cuando RCON falla o no soporta la operación, hace fallback explícito a scoreboard. +- La salida del comando deja claro qué fuente se usó realmente. +- El repositorio queda consistente. + +## Expected Files +- `backend/app/data_sources.py` +- `backend/app/historical_ingestion.py` +- archivos backend necesarios para writer path RCON +- `backend/README.md` diff --git a/ai/tasks/done/TASK-101-remove-false-rcon-first-claims-and-fix-operational-visibility.md b/ai/tasks/done/TASK-101-remove-false-rcon-first-claims-and-fix-operational-visibility.md new file mode 100644 index 0000000..ad5488d --- /dev/null +++ b/ai/tasks/done/TASK-101-remove-false-rcon-first-claims-and-fix-operational-visibility.md @@ -0,0 +1,38 @@ +# TASK-101-remove-false-rcon-first-claims-and-fix-operational-visibility + +## Goal +Alinear documentación, outputs operativos y visibilidad de progreso con el comportamiento real del sistema, evitando afirmaciones engañosas sobre histórico RCON-first y mejorando la operativa del refresh manual. + +## Context +Actualmente el operador puede pensar que la ingesta histórica ya va por RCON cuando en realidad el writer path sigue cayendo a scoreboard. +Además, el comando de refresh es demasiado opaco: tarda mucho y no ofrece progreso útil. + +## Steps +1. Auditar: + - `backend/README.md` + - `backend/app/historical_ingestion.py` + - outputs/logs relevantes +2. Corregir cualquier copy/documentación que sugiera que la ingesta histórica completa ya está en RCON si no es cierto. +3. Añadir progreso operativo útil al refresh manual: + - servidor actual + - página actual + - número de match ids a detallar + - fuente realmente seleccionada +4. Hacer que, cuando haya fallback a scoreboard, eso quede visible para el operador en tiempo real o en el payload final. +5. Mantener la salida usable, sin inundar de logs innecesarios. +6. Actualizar runbook con recomendaciones reales para pasadas manuales y límites razonables. + +## Constraints +- No convertir el comando en un spam de logs. +- No ocultar fallbacks reales. +- No mezclar esta task con grandes cambios de UI frontend. + +## Validation +- El operador puede ver progreso útil durante un refresh. +- La fuente usada de verdad queda visible. +- La documentación ya no induce a error. +- El repositorio queda consistente. + +## Expected Files +- `backend/app/historical_ingestion.py` +- `backend/README.md` diff --git a/backend/README.md b/backend/README.md index 3f9d132..d535cd5 100644 --- a/backend/README.md +++ b/backend/README.md @@ -272,7 +272,7 @@ Valores soportados en esta fase: - `rcon` como camino primario recomendado - `a2s` como fallback legacy o override explicito - historico: - - `rcon` como camino primario recomendado para captura y lectura minima + - `rcon` como camino primario recomendado para captura y writer path primario - `public-scoreboard` como fallback legacy o override explicito Defaults actuales: @@ -309,18 +309,14 @@ Politica funcional actual: - `public-scoreboard` solo si RCON no soporta aun esa operacion concreta o falla la captura primaria -Limitacion actual de `rcon`: - -- el backend puede usar `rcon` para `/api/servers` -- la ingesta historica competitiva por `historical_ingestion.py` sigue cayendo a - `public-scoreboard`, porque la repo todavia no incluye una canalizacion - retroactiva RCON capaz de reconstruir partidas cerradas con paridad completa - Estado real de "historico por RCON" en esta repo: - no existe backfill retroactivo por RCON con el cliente actual - la viabilidad documentada hoy es solo para captura prospectiva separada -- `public-scoreboard` sigue siendo la fuente historica principal +- `historical_ingestion.py` intenta primero el writer path prospectivo RCON +- la persistencia competitiva `historical_*` sigue necesitando fallback a + `public-scoreboard` mientras RCON no exponga pagina historica/detalle de + match cerrada con paridad suficiente - el diseno tecnico de esa linea prospectiva queda en `docs/rcon-historical-ingestion-design.md` @@ -473,7 +469,8 @@ respuesta deja trazabilidad con `primary_source`, `selected_source`, consultas live de produccion. - `config.py` centraliza host, puerto y allowlist minima de origenes locales. - `data_sources.py` define los contratos y la seleccion por entorno para live e historico. -- `historical_ingestion.py` consulta la capa JSON publica de CRCON para bootstrap y refresh incremental. +- `historical_ingestion.py` intenta primero el writer path RCON y, si hace falta + poblar `historical_*`, cae de forma explicita a la capa JSON publica de CRCON. - `historical_models.py` fija las entidades historicas minimas del dominio. - `historical_snapshots.py` fija los tipos y selectores validos de snapshots historicos precalculados. - `historical_snapshot_storage.py` persiste snapshots historicos precalculados listos para lectura rapida. @@ -1039,6 +1036,25 @@ Si la recomposicion se lanza para un servidor fisico concreto, el backend rehace tambien el agregado logico `all-servers` para mantener `Todos` alineado con `#01` y `#02` aunque `#03` siga sin bootstrap. +En esta fase, el comando muestra progreso operativo util sin saturar stdout: + +- intento primario RCON +- fuente finalmente seleccionada +- servidor actual +- pagina actual +- `match_ids_to_detail` de cada pagina + +Si RCON no puede cubrir la operacion competitiva real, el fallback a +`public-scoreboard` queda visible tanto durante la ejecucion como en el JSON +final mediante: + +- `primary_source` +- `selected_source` +- `fallback_used` +- `fallback_reason` +- `source_attempts` +- `primary_writer_result` + El comando devuelve ademas un resumen de cobertura persistida por servidor. Esto ayuda a validar rapidamente cuantos matches reales quedaron importados, el rango temporal cubierto y si la carga ya supera la ultima semana movil que usa la UI. @@ -1080,6 +1096,16 @@ registrados. La segunda sirve para validar un solo servidor con alcance acotado. La tercera recompone snapshots despues de una pasada manual cuando se quiere confirmar que la capa precalculada vuelve a quedar alineada. +Interpretacion operativa recomendada: + +- si aparece `historical-ingestion-rcon-primary-succeeded`, RCON se intento de + verdad primero y la captura prospectiva quedo registrada +- si despues `selected_source` termina en `public-scoreboard`, eso significa + que la reconstruccion del archivo competitivo `historical_*` siguio + necesitando fallback clasico +- si RCON falla por red, auth o timeout, el motivo queda visible en + `fallback_reason` y en `source_attempts` + Los reintentos de cada request JSON pueden ajustarse sin tocar codigo con: - `HLL_HISTORICAL_CRCON_REQUEST_RETRIES` @@ -1392,9 +1418,11 @@ El backend queda orientado a `RCON-first` tambien para historico: - `rcon` primero - `a2s` solo como fallback - historico: - - `rcon` primero para la capa de lectura/cobertura soportada hoy + - `rcon` primero tanto para lectura minima como para el writer path primario + de `historical_ingestion` - `public-scoreboard` solo como fallback cuando RCON no cubre una operacion - competitiva concreta o no tiene cobertura suficiente + competitiva concreta, no tiene cobertura suficiente o falla la captura + primaria Metadata observable en payloads historicos: @@ -1407,6 +1435,8 @@ Metadata observable en payloads historicos: Estado real a fecha de esta fase: - el read model historico RCON soporta cobertura y actividad reciente +- `historical_ingestion` intenta primero una captura writer-oriented por RCON y + deja esa tentativa visible en su salida - rankings competitivos, MVP y Elo/MMR siguen necesitando fallback a `public-scoreboard` porque el read model RCON actual no expone aun detalle historico competitivo suficiente diff --git a/backend/app/data_sources.py b/backend/app/data_sources.py index 4f14d55..3f7db3b 100644 --- a/backend/app/data_sources.py +++ b/backend/app/data_sources.py @@ -327,45 +327,46 @@ def build_historical_runtime_source_policy( def resolve_historical_ingestion_data_source() -> tuple[HistoricalDataSource, dict[str, object]]: - """Resolve the writer-oriented historical provider with safe fallback semantics.""" + """Resolve the fallback provider used when classic scoreboard import is required.""" configured_kind = get_historical_data_source_kind() - if configured_kind == SOURCE_KIND_PUBLIC_SCOREBOARD: - return ( - PublicScoreboardHistoricalDataSource(), - build_source_policy( - primary_source=SOURCE_KIND_PUBLIC_SCOREBOARD, - selected_source=SOURCE_KIND_PUBLIC_SCOREBOARD, - source_attempts=[ - build_source_attempt( - source=SOURCE_KIND_PUBLIC_SCOREBOARD, - role="primary", - status="success", - ) - ], - ), + if configured_kind in {SOURCE_KIND_PUBLIC_SCOREBOARD, SOURCE_KIND_RCON}: + primary_source = ( + SOURCE_KIND_PUBLIC_SCOREBOARD + if configured_kind == SOURCE_KIND_PUBLIC_SCOREBOARD + else SOURCE_KIND_RCON + ) + fallback_used = configured_kind == SOURCE_KIND_RCON + fallback_reason = ( + "classic-historical-import-requires-public-scoreboard-fallback" + if fallback_used + else None + ) + attempts = [] + if configured_kind == SOURCE_KIND_RCON: + attempts.append( + build_source_attempt( + source=SOURCE_KIND_RCON, + role="primary", + status="deferred", + reason="rcon-primary-writer-attempt-is-handled-by-historical-ingestion", + ) + ) + attempts.append( + build_source_attempt( + source=SOURCE_KIND_PUBLIC_SCOREBOARD, + role="fallback" if fallback_used else "primary", + status="ready", + reason="classic-historical-import-provider-ready", + ) ) - - if configured_kind == SOURCE_KIND_RCON: return ( PublicScoreboardHistoricalDataSource(), build_source_policy( - primary_source=SOURCE_KIND_RCON, + primary_source=primary_source, selected_source=SOURCE_KIND_PUBLIC_SCOREBOARD, - fallback_used=True, - fallback_reason="rcon-historical-ingestion-not-supported-yet", - source_attempts=[ - build_source_attempt( - source=SOURCE_KIND_RCON, - role="primary", - status="unsupported", - reason="rcon-historical-ingestion-not-supported-yet", - ), - build_source_attempt( - source=SOURCE_KIND_PUBLIC_SCOREBOARD, - role="fallback", - status="success", - ), - ], + fallback_used=fallback_used, + fallback_reason=fallback_reason, + source_attempts=attempts, ), ) diff --git a/backend/app/historical_ingestion.py b/backend/app/historical_ingestion.py index 731ed00..7598173 100644 --- a/backend/app/historical_ingestion.py +++ b/backend/app/historical_ingestion.py @@ -5,14 +5,21 @@ from __future__ import annotations import argparse import json from dataclasses import dataclass -from typing import Iterable +from typing import Callable, Iterable from .config import ( get_historical_crcon_detail_workers, get_historical_crcon_page_size, + get_historical_data_source_kind, get_historical_refresh_overlap_hours, ) -from .data_sources import HistoricalDataSource, resolve_historical_ingestion_data_source +from .data_sources import ( + SOURCE_KIND_PUBLIC_SCOREBOARD, + SOURCE_KIND_RCON, + HistoricalDataSource, + build_historical_runtime_source_policy, + resolve_historical_ingestion_data_source, +) from .elo_mmr_engine import rebuild_elo_mmr_models from .historical_snapshots import generate_and_persist_historical_snapshots from .historical_storage import ( @@ -28,9 +35,13 @@ from .historical_storage import ( start_ingestion_run, upsert_historical_match, ) +from .rcon_historical_worker import run_rcon_historical_capture_unlocked from .writer_lock import backend_writer_lock, build_writer_lock_holder +ProgressCallback = Callable[[dict[str, object]], None] + + @dataclass(slots=True) class IngestionStats: """Mutable counters for one ingestion execution.""" @@ -57,6 +68,7 @@ def run_bootstrap( start_page: int | None = None, detail_workers: int | None = None, rebuild_snapshots: bool = True, + progress_callback: ProgressCallback | None = None, ) -> dict[str, object]: """Run a first full historical import against one or all configured servers.""" with backend_writer_lock( @@ -74,6 +86,7 @@ def run_bootstrap( overlap_hours=None, incremental=False, rebuild_snapshots=rebuild_snapshots, + progress_callback=progress_callback, ) @@ -86,6 +99,7 @@ def run_incremental_refresh( detail_workers: int | None = None, overlap_hours: int | None = None, rebuild_snapshots: bool = True, + progress_callback: ProgressCallback | None = None, ) -> dict[str, object]: """Refresh recent historical pages without replaying the whole archive.""" with backend_writer_lock( @@ -103,6 +117,7 @@ def run_incremental_refresh( overlap_hours=overlap_hours, incremental=True, rebuild_snapshots=rebuild_snapshots, + progress_callback=progress_callback, ) @@ -117,10 +132,11 @@ def _run_ingestion( overlap_hours: int | None, incremental: bool, rebuild_snapshots: bool, + progress_callback: ProgressCallback | None, ) -> dict[str, object]: initialize_historical_storage() stats = IngestionStats() - data_source, source_policy = resolve_historical_ingestion_data_source() + fallback_data_source, fallback_source_policy = resolve_historical_ingestion_data_source() selected_servers = _select_servers(server_slug) processed_servers: list[dict[str, object]] = [] active_runs: dict[str, int] = {} @@ -132,60 +148,86 @@ def _run_ingestion( if resolved_overlap_hours < 0: raise ValueError("--overlap-hours must be zero or positive.") + primary_writer_result = _attempt_primary_rcon_writer( + mode=mode, + server_slug=server_slug, + selected_servers=selected_servers, + progress_callback=progress_callback, + ) + source_policy = _resolve_ingestion_source_policy( + fallback_source_policy=fallback_source_policy, + primary_writer_result=primary_writer_result, + ) + use_classic_fallback = _should_use_classic_fallback(primary_writer_result) + _emit_progress( + progress_callback, + { + "event": "historical-ingestion-source-selected", + "mode": mode, + "primary_source": source_policy.get("primary_source"), + "selected_source": source_policy.get("selected_source"), + "fallback_used": bool(source_policy.get("fallback_used")), + "fallback_reason": source_policy.get("fallback_reason"), + }, + ) + try: - for server in selected_servers: - run_id = start_ingestion_run(mode=mode, target_server_slug=str(server["slug"])) - active_runs[str(server["slug"])] = run_id - mark_backfill_progress_started( - server_slug=str(server["slug"]), - mode=mode, - run_id=run_id, - ) - cutoff = ( - get_refresh_cutoff_for_server( - str(server["slug"]), - overlap_hours=resolved_overlap_hours, + if use_classic_fallback: + for server in selected_servers: + run_id = start_ingestion_run(mode=mode, target_server_slug=str(server["slug"])) + active_runs[str(server["slug"])] = run_id + mark_backfill_progress_started( + server_slug=str(server["slug"]), + mode=mode, + run_id=run_id, ) - if incremental - else None - ) - resolved_start_page = _resolve_start_page( - start_page=start_page, - server_slug=str(server["slug"]), - mode=mode, - ) - server_stats = _ingest_server( - server=server, - mode=mode, - run_id=run_id, - stats=stats, - data_source=data_source, - max_pages=max_pages, - page_size=page_size, - start_page=resolved_start_page, - detail_workers=detail_workers, - cutoff=cutoff, - ) - processed_servers.append(server_stats) - finalize_ingestion_run( - run_id, - status="success", - pages_processed=server_stats["pages_processed"], - matches_seen=server_stats["matches_seen"], - matches_inserted=server_stats["matches_inserted"], - matches_updated=server_stats["matches_updated"], - player_rows_inserted=server_stats["player_rows_inserted"], - player_rows_updated=server_stats["player_rows_updated"], - notes=f"public_name={server_stats['public_name']}", - ) - finalize_backfill_progress( - server_slug=str(server["slug"]), - mode=mode, - run_id=run_id, - status="success", - archive_exhausted=bool(server_stats["archive_exhausted"]), - ) - active_runs.pop(str(server["slug"]), None) + cutoff = ( + get_refresh_cutoff_for_server( + str(server["slug"]), + overlap_hours=resolved_overlap_hours, + ) + if incremental + else None + ) + resolved_start_page = _resolve_start_page( + start_page=start_page, + server_slug=str(server["slug"]), + mode=mode, + ) + server_stats = _ingest_server( + server=server, + mode=mode, + run_id=run_id, + stats=stats, + data_source=fallback_data_source, + max_pages=max_pages, + page_size=page_size, + start_page=resolved_start_page, + detail_workers=detail_workers, + cutoff=cutoff, + progress_callback=progress_callback, + source_policy=source_policy, + ) + processed_servers.append(server_stats) + finalize_ingestion_run( + run_id, + status="success", + pages_processed=server_stats["pages_processed"], + matches_seen=server_stats["matches_seen"], + matches_inserted=server_stats["matches_inserted"], + matches_updated=server_stats["matches_updated"], + player_rows_inserted=server_stats["player_rows_inserted"], + player_rows_updated=server_stats["player_rows_updated"], + notes=f"public_name={server_stats['public_name']}", + ) + finalize_backfill_progress( + server_slug=str(server["slug"]), + mode=mode, + run_id=run_id, + status="success", + archive_exhausted=bool(server_stats["archive_exhausted"]), + ) + active_runs.pop(str(server["slug"]), None) if rebuild_snapshots: snapshot_result = generate_and_persist_historical_snapshots(server_key=server_slug) elo_mmr_result = rebuild_elo_mmr_models() @@ -224,8 +266,9 @@ def _run_ingestion( return { "status": "ok", "mode": mode, - "source_provider": data_source.source_kind, + "source_provider": source_policy.get("selected_source"), "source_policy": source_policy, + "primary_writer_result": primary_writer_result, "page_size": page_size or get_historical_crcon_page_size(), "start_page": start_page, "detail_workers": detail_workers or get_historical_crcon_detail_workers(), @@ -257,6 +300,8 @@ def _ingest_server( start_page: int, detail_workers: int | None, cutoff: str | None, + progress_callback: ProgressCallback | None, + source_policy: dict[str, object], ) -> dict[str, object]: resolved_page_size = page_size or get_historical_crcon_page_size() resolved_detail_workers = detail_workers or get_historical_crcon_detail_workers() @@ -267,6 +312,18 @@ def _ingest_server( discovered_total_matches: int | None = None last_page_processed: int | None = None archive_exhausted = False + _emit_progress( + progress_callback, + { + "event": "historical-ingestion-server-started", + "mode": mode, + "server_slug": server["slug"], + "selected_source": source_policy.get("selected_source"), + "fallback_used": bool(source_policy.get("fallback_used")), + "start_page": start_page, + "cutoff": cutoff, + }, + ) for page_number in range(start_page, start_page + page_limit): payload = data_source.fetch_match_page( @@ -300,6 +357,20 @@ def _ingest_server( if match_id: match_ids_to_fetch.append(match_id) + _emit_progress( + progress_callback, + { + "event": "historical-ingestion-page-loaded", + "mode": mode, + "server_slug": server["slug"], + "page": page_number, + "selected_source": source_policy.get("selected_source"), + "match_ids_to_detail": len(match_ids_to_fetch), + "page_matches": len(page_matches), + "cutoff_reached": stop_after_page, + }, + ) + for detail_payload in data_source.fetch_match_details( base_url=str(server["scoreboard_base_url"]), match_ids=match_ids_to_fetch, @@ -356,6 +427,162 @@ def _resolve_start_page( return get_backfill_resume_page(server_slug, mode=mode) +def _attempt_primary_rcon_writer( + *, + mode: str, + server_slug: str | None, + selected_servers: list[dict[str, object]], + progress_callback: ProgressCallback | None, +) -> dict[str, object]: + configured_kind = get_historical_data_source_kind() + if configured_kind != SOURCE_KIND_RCON: + result = { + "attempted": False, + "status": "skipped", + "primary_source": SOURCE_KIND_PUBLIC_SCOREBOARD, + "selected_source": SOURCE_KIND_PUBLIC_SCOREBOARD, + "fallback_used": False, + "fallback_reason": None, + "source_attempts": [], + } + _emit_progress( + progress_callback, + { + "event": "historical-ingestion-rcon-primary-skipped", + "mode": mode, + "reason": "historical-data-source-configured-for-public-scoreboard", + }, + ) + return result + + target_scope = server_slug or "all-configured-rcon-targets" + _emit_progress( + progress_callback, + { + "event": "historical-ingestion-rcon-primary-started", + "mode": mode, + "target_scope": target_scope, + "servers": [str(server["slug"]) for server in selected_servers], + }, + ) + try: + capture_result = run_rcon_historical_capture_unlocked(target_key=server_slug) + except Exception as exc: # noqa: BLE001 - fallback remains explicit and controlled + result = { + "attempted": True, + "status": "error", + "primary_source": SOURCE_KIND_RCON, + "selected_source": SOURCE_KIND_PUBLIC_SCOREBOARD, + "fallback_used": True, + "fallback_reason": "rcon-historical-writer-request-failed", + "message": str(exc), + } + _emit_progress( + progress_callback, + { + "event": "historical-ingestion-rcon-primary-failed", + "mode": mode, + "target_scope": target_scope, + "message": str(exc), + }, + ) + return result + + capture_run_status = str(capture_result.get("run_status") or capture_result.get("status") or "unknown") + targets = list(capture_result.get("targets") or []) + errors = list(capture_result.get("errors") or []) + if targets: + result = { + "attempted": True, + "status": "partial", + "primary_source": SOURCE_KIND_RCON, + "selected_source": SOURCE_KIND_PUBLIC_SCOREBOARD, + "fallback_used": True, + "fallback_reason": "rcon-primary-writer-succeeded-but-classic-match-archive-still-needs-fallback", + "capture_result": capture_result, + } + _emit_progress( + progress_callback, + { + "event": "historical-ingestion-rcon-primary-succeeded", + "mode": mode, + "target_scope": target_scope, + "captured_targets": len(targets), + "run_status": capture_run_status, + "next_step": "classic-public-scoreboard-fallback-required", + }, + ) + return result + + result = { + "attempted": True, + "status": "empty", + "primary_source": SOURCE_KIND_RCON, + "selected_source": SOURCE_KIND_PUBLIC_SCOREBOARD, + "fallback_used": True, + "fallback_reason": "rcon-historical-writer-returned-no-usable-samples", + "capture_result": capture_result, + "message": json.dumps(errors, separators=(",", ":")) if errors else None, + } + _emit_progress( + progress_callback, + { + "event": "historical-ingestion-rcon-primary-empty", + "mode": mode, + "target_scope": target_scope, + "run_status": capture_run_status, + "errors": len(errors), + }, + ) + return result + + +def _should_use_classic_fallback(primary_writer_result: dict[str, object]) -> bool: + selected_source = str(primary_writer_result.get("selected_source") or "") + return selected_source == SOURCE_KIND_PUBLIC_SCOREBOARD + + +def _resolve_ingestion_source_policy( + *, + fallback_source_policy: dict[str, object], + primary_writer_result: dict[str, object], +) -> dict[str, object]: + configured_kind = get_historical_data_source_kind() + if configured_kind != SOURCE_KIND_RCON: + return fallback_source_policy + + status = str(primary_writer_result.get("status") or "error") + selected_source = str( + primary_writer_result.get("selected_source") or SOURCE_KIND_PUBLIC_SCOREBOARD + ) + fallback_reason = primary_writer_result.get("fallback_reason") + message = primary_writer_result.get("message") + if ( + fallback_reason + == "rcon-primary-writer-succeeded-but-classic-match-archive-still-needs-fallback" + ): + message = ( + "RCON prospective capture succeeded first, but the classic historical_* " + "archive still requires public-scoreboard for match-page import." + ) + return build_historical_runtime_source_policy( + operation="historical-ingestion", + rcon_status=status, + fallback_reason=str(fallback_reason) if fallback_reason else None, + selected_source=selected_source, + rcon_message=message if isinstance(message, str) else None, + ) + + +def _emit_progress( + callback: ProgressCallback | None, + payload: dict[str, object], +) -> None: + if callback is None: + return + callback(payload) + + def _select_servers(server_slug: str | None) -> list[dict[str, object]]: servers = list_historical_servers() if server_slug is None: @@ -456,6 +683,9 @@ def main(argv: Iterable[str] | None = None) -> int: parser = build_arg_parser() args = parser.parse_args(list(argv) if argv is not None else None) + def _print_progress(payload: dict[str, object]) -> None: + print(json.dumps(payload, ensure_ascii=True)) + if args.mode == "bootstrap": result = run_bootstrap( server_slug=args.server_slug, @@ -463,6 +693,7 @@ def main(argv: Iterable[str] | None = None) -> int: page_size=args.page_size, start_page=args.start_page, detail_workers=args.detail_workers, + progress_callback=_print_progress, ) else: result = run_incremental_refresh( @@ -472,6 +703,7 @@ def main(argv: Iterable[str] | None = None) -> int: start_page=args.start_page, detail_workers=args.detail_workers, overlap_hours=args.overlap_hours, + progress_callback=_print_progress, ) print(json.dumps(result, indent=2)) diff --git a/backend/app/rcon_historical_worker.py b/backend/app/rcon_historical_worker.py index add7d07..999f0f5 100644 --- a/backend/app/rcon_historical_worker.py +++ b/backend/app/rcon_historical_worker.py @@ -44,93 +44,101 @@ def run_rcon_historical_capture( f"app.rcon_historical_worker capture:{target_key or 'all-targets'}" ) ): - initialize_rcon_historical_storage() - selected_targets = _select_targets(target_key) - captured_at = utc_now().isoformat().replace("+00:00", "Z") - target_scope = target_key or "all-configured-rcon-targets" - run_id = start_rcon_historical_capture_run(mode="capture", target_scope=target_scope) - stats = RconHistoricalCaptureStats() - items: list[dict[str, object]] = [] - errors: list[dict[str, object]] = [] + return run_rcon_historical_capture_unlocked(target_key=target_key) - try: - for target in selected_targets: - target_metadata = _serialize_target(target) - stats.targets_seen += 1 - try: - sample = query_live_server_sample(target) - delta = persist_rcon_historical_sample( - run_id=run_id, - captured_at=captured_at, - target=target_metadata, - normalized_payload=sample["normalized"], - raw_payload=sample["raw_session"], - ) - stats.samples_inserted += int(delta["samples_inserted"]) - stats.duplicate_samples += int(delta["duplicate_samples"]) - items.append( - { - "target_key": target_metadata["target_key"], - "external_server_id": target.external_server_id, - "captured_at": captured_at, - "sample_inserted": bool(delta["samples_inserted"]), - "normalized": sample["normalized"], - } - ) - except Exception as exc: # noqa: BLE001 - controlled worker failures - stats.failed_targets += 1 - mark_rcon_historical_capture_failure( - run_id=run_id, - target=target_metadata, - error_message=str(exc), - ) - errors.append( - { - "target_key": target_metadata["target_key"], - "name": target.name, - "host": target.host, - "port": target.port, - "message": str(exc), - } - ) - status = "success" if not errors else ("partial" if items else "failed") - finalize_rcon_historical_capture_run( - run_id, - status=status, - targets_seen=stats.targets_seen, - samples_inserted=stats.samples_inserted, - duplicate_samples=stats.duplicate_samples, - failed_targets=stats.failed_targets, - notes=None if not errors else json.dumps(errors, separators=(",", ":")), - ) - except Exception as exc: - finalize_rcon_historical_capture_run( - run_id, - status="failed", - targets_seen=stats.targets_seen, - samples_inserted=stats.samples_inserted, - duplicate_samples=stats.duplicate_samples, - failed_targets=max(1, stats.failed_targets), - notes=str(exc), - ) - raise +def run_rcon_historical_capture_unlocked( + *, + target_key: str | None = None, +) -> dict[str, object]: + """Capture one prospective RCON sample assuming the shared writer lock is already held.""" + initialize_rcon_historical_storage() + selected_targets = _select_targets(target_key) + captured_at = utc_now().isoformat().replace("+00:00", "Z") + target_scope = target_key or "all-configured-rcon-targets" + run_id = start_rcon_historical_capture_run(mode="capture", target_scope=target_scope) + stats = RconHistoricalCaptureStats() + items: list[dict[str, object]] = [] + errors: list[dict[str, object]] = [] - return { - "status": "ok" if items else "error", - "run_status": status, - "captured_at": captured_at, - "target_scope": target_scope, - "targets": items, - "errors": errors, - "storage_status": list_rcon_historical_target_statuses(), - "totals": { - "targets_seen": stats.targets_seen, - "samples_inserted": stats.samples_inserted, - "duplicate_samples": stats.duplicate_samples, - "failed_targets": stats.failed_targets, - }, - } + try: + for target in selected_targets: + target_metadata = _serialize_target(target) + stats.targets_seen += 1 + try: + sample = query_live_server_sample(target) + delta = persist_rcon_historical_sample( + run_id=run_id, + captured_at=captured_at, + target=target_metadata, + normalized_payload=sample["normalized"], + raw_payload=sample["raw_session"], + ) + stats.samples_inserted += int(delta["samples_inserted"]) + stats.duplicate_samples += int(delta["duplicate_samples"]) + items.append( + { + "target_key": target_metadata["target_key"], + "external_server_id": target.external_server_id, + "captured_at": captured_at, + "sample_inserted": bool(delta["samples_inserted"]), + "normalized": sample["normalized"], + } + ) + except Exception as exc: # noqa: BLE001 - controlled worker failures + stats.failed_targets += 1 + mark_rcon_historical_capture_failure( + run_id=run_id, + target=target_metadata, + error_message=str(exc), + ) + errors.append( + { + "target_key": target_metadata["target_key"], + "name": target.name, + "host": target.host, + "port": target.port, + "message": str(exc), + } + ) + + status = "success" if not errors else ("partial" if items else "failed") + finalize_rcon_historical_capture_run( + run_id, + status=status, + targets_seen=stats.targets_seen, + samples_inserted=stats.samples_inserted, + duplicate_samples=stats.duplicate_samples, + failed_targets=stats.failed_targets, + notes=None if not errors else json.dumps(errors, separators=(",", ":")), + ) + except Exception as exc: + finalize_rcon_historical_capture_run( + run_id, + status="failed", + targets_seen=stats.targets_seen, + samples_inserted=stats.samples_inserted, + duplicate_samples=stats.duplicate_samples, + failed_targets=max(1, stats.failed_targets), + notes=str(exc), + ) + raise + + return { + "status": "ok" if items else "error", + "run_status": status, + "captured_at": captured_at, + "target_scope": target_scope, + "targets": items, + "errors": errors, + "storage_status": list_rcon_historical_target_statuses(), + "totals": { + "targets_seen": stats.targets_seen, + "samples_inserted": stats.samples_inserted, + "duplicate_samples": stats.duplicate_samples, + "failed_targets": stats.failed_targets, + }, + } def run_periodic_rcon_historical_capture(