From 9c0d39146e2665f86a4fe31c35d79e70d2fe736d Mon Sep 17 00:00:00 2001 From: guigerdts Date: Tue, 25 Aug 2026 16:36:34 +0000 Subject: [PATCH] refactor(backend): reduce numbers pipeline to stats->features->gen and decouple gen from MetaSelection --- backend/src/backend/app/api/v1/pipeline.py | 4 +- .../src/backend/app/services/gen_service.py | 35 +- .../backend/app/services/pipeline_service.py | 359 +----------------- 3 files changed, 40 insertions(+), 358 deletions(-) diff --git a/backend/src/backend/app/api/v1/pipeline.py b/backend/src/backend/app/api/v1/pipeline.py index c658bb2..9c40868 100644 --- a/backend/src/backend/app/api/v1/pipeline.py +++ b/backend/src/backend/app/api/v1/pipeline.py @@ -29,10 +29,10 @@ summary="Run the canonical numbers chain and return the per-stage report", ) def run_numbers(payload: PipelineRunRequest, db: DbSession) -> SuccessEnvelope[PipelineRunResult]: - """Heal-and-run stats→features→ml→dl→bt→rank→select→gen in one call (R1). + """Heal-and-run stats→features→gen in one call (R1). Missing/stale prerequisites are repaired by running exactly the deficient - stages (R2); every response carries the ordered eight-stage report (R3); + stages (R2); every response carries the ordered three-stage report (R3); unchanged inputs reuse stored fingerprints with zero side effects (R4). """ outcome = PipelineService(db).run( diff --git a/backend/src/backend/app/services/gen_service.py b/backend/src/backend/app/services/gen_service.py index 7feeebb..165e8d7 100644 --- a/backend/src/backend/app/services/gen_service.py +++ b/backend/src/backend/app/services/gen_service.py @@ -33,9 +33,9 @@ from backend.app.generators.version import GENERATOR_VERSION from backend.app.generators.weighting import build_weights from backend.app.models.gen_snapshot import GenSnapshot -from backend.app.services.probability_service import _classify_coverage from backend.app.repositories.stat_payload_repository import StatPayloadRepository from backend.app.services.errors import GenServiceError +from backend.app.services.probability_service import _classify_coverage from backend.app.statistics.engine import frequency DEFAULT_COUNT: int = 10 @@ -316,8 +316,14 @@ def _resolve_selection(self, lottery_id: int, selection_id: int | None) -> Any: """Resolve the selection to weight from; ``GEN_NO_SELECTION`` when absent. Without an override the active selection is used; an override must belong - to the target lottery (GEN-003). + to the target lottery (GEN-003). When no active MetaSelection exists, + returns a lightweight fallback object with id=0 and a deterministic + fingerprint derived from the probability snapshot checksum + lottery_id. """ + import hashlib + import json as _json + from types import SimpleNamespace + from backend.app.models.meta_selection import MetaSelection if selection_id is not None: @@ -341,10 +347,29 @@ def _resolve_selection(self, lottery_id: int, selection_id: int | None) -> Any: ) selection = self._session.execute(stmt).scalar_one_or_none() if selection is None: - raise GenServiceError( - GenServiceError.GEN_NO_SELECTION, - f"no active selection for lottery {lottery_id}", + # Fallback: build lightweight selection-like object (no MetaSelection needed) + from backend.app.models.prob_snapshot import ProbSnapshot + + prob_stmt = ( + select(ProbSnapshot) + .where(ProbSnapshot.lottery_id == lottery_id, ProbSnapshot.status == "active") + .order_by(ProbSnapshot.version.desc()) + .limit(1) + ) + prob = self._session.execute(prob_stmt).scalar_one_or_none() + if prob is None: + raise GenServiceError( + GenServiceError.GEN_NO_DISTRIBUTION, + f"no active probability distribution for lottery {lottery_id}", + ) + # Deterministic fingerprint from prob snapshot checksum + lottery_id + canonical = _json.dumps( + {"checksum": prob.checksum, "lottery_id": lottery_id}, + sort_keys=True, + separators=(",", ":"), ) + fp = hashlib.sha256(canonical.encode()).hexdigest() + return SimpleNamespace(id=0, fingerprint=fp) return selection def _load_distribution(self, lottery_id: int) -> dict[int, float]: diff --git a/backend/src/backend/app/services/pipeline_service.py b/backend/src/backend/app/services/pipeline_service.py index 0100369..b896346 100644 --- a/backend/src/backend/app/services/pipeline_service.py +++ b/backend/src/backend/app/services/pipeline_service.py @@ -1,25 +1,14 @@ """Numbers pipeline orchestrator service (S2, R1-R4). Single entry point that executes the canonical chain -``stats → features(+probability) → ml → dl → bt → rank → select → gen`` +``stats → features(+probability) → gen`` in order (R1), heals missing/stale prerequisites by running exactly the deficient stages (R2), reports per-stage status (R3), and stays zero-side-effect on unchanged inputs via fingerprint reuse (R4). -Classification strategy: every stage service in the chain is fingerprint- -idempotent EXCEPT two — ``MlService.train`` always persists a new snapshot -version and ``BtSnapshotStore.create_active`` delete+recreates rows. Those -two stages are therefore *gated*: they are invoked only when their active -artifact is missing or a dependency wrote during this run; all other stages -are invoked unconditionally and classified by comparing the active-artifact -fingerprint before/after the call (design §S2, task 5.1). - -bt-before-rank (D8): the ranking context is derived post-bt via -``resolve_context_vector(lottery, "backtesting")`` + ``compute_context_hash`` -and passed explicitly to ``MetaService.select(context_hash=...)`` — retiring -the hardcoded ``meta_service.py:242`` coupling without touching META logic. -A stale active ranking (created_at ≤ newest active BtSnapshot) triggers -exactly ONE rerank attempt, then ``PIPE_STAGE_FAILED(rank)``. +The features stage covers both feature and probability snapshots (D9). +Stages ml, dl, bt, rank, and select have been removed from the numbers +path; backtesting retains its own independent pipeline. """ from __future__ import annotations @@ -34,27 +23,9 @@ STAGE_ORDER: tuple[str, ...] = ( "stats", "features", - "ml", - "dl", - "bt", - "rank", - "select", "gen", ) -# Canonical backtest parameters for the chain (CLI bt-run defaults). -BT_STRATEGY_ID = "ml-core-5" - -# Stage dependencies within the chain. ``stats`` proxies draw coverage: its -# checksum-based idempotency rewrites iff the imported draw set changed. -_DEPS: dict[str, tuple[str, ...]] = { - "ml": ("stats", "features"), - "bt": ("stats",), -} - -# Stage services that are NOT no-op on identical inputs (see module docstring). -_GATED_STAGES: frozenset[str] = frozenset({"ml", "bt"}) - @dataclass(frozen=True) class StageRecord: @@ -85,25 +56,12 @@ class _Artifact: class _RunState: - """Mutable per-run state (ctx hash + generator output).""" + """Mutable per-run state (generator output).""" def __init__(self) -> None: - self.context_hash: str | None = None self.gen_result: object | None = None -def derive_context_hash(db: Session, lottery_id: int) -> str: - """Derive the backtesting context hash from the executed bt run (D8). - - Raises ``ValueError`` when no active BtSnapshot exists — the orchestrator - only calls this after the bt stage has run or been reused. - """ - from backend.app.meta.context import compute_context_hash, resolve_context_vector - - vector = resolve_context_vector(lottery_id, "backtesting", db) - return compute_context_hash(vector) - - class PipelineService: """Orchestrates the canonical eight-stage numbers chain (R1-R4).""" @@ -123,23 +81,7 @@ def run(self, *, lottery_id: int, count: int | None = None, seed: int | None = N state = _RunState() for name in STAGE_ORDER: - if name in ("rank", "select") and state.context_hash is None: - state.context_hash = self._safe_context_hash(lottery_id) - - before = self._read_artifact(lottery_id, name, state.context_hash) - - if self._gated_skip(name, before, changed): - report.append( - StageRecord( - name=name, - status="skipped", - snapshot_id=before.snapshot_id, - fingerprint=before.fingerprint, - detail="active artifact reused", - ) - ) - changed[name] = False - continue + before = self._read_artifact(lottery_id, name) try: self._execute_stage(name, lottery_id, count, seed, state) @@ -158,7 +100,7 @@ def run(self, *, lottery_id: int, count: int | None = None, seed: int | None = N error.stages = report # R3: failed entry travels with the error raise error from exc - after = self._read_artifact(lottery_id, name, state.context_hash) + after = self._read_artifact(lottery_id, name) wrote = after != before and after.fingerprint is not None report.append( StageRecord( @@ -199,32 +141,6 @@ def _execute_stage( FeatureEngineService(db).generate(lottery_id=lottery_id) ProbabilityService(db).generate(lottery_id=lottery_id) - elif name == "ml": - from backend.app.services.ml_service import MlService - - MlService(db, _MlDrawAdapter(db), _MlFeatureAdapter(db)).train(lottery_id) - elif name == "dl": - from backend.app.services.dl_service import DlService - - ml_draws = _MlDrawAdapter(db) - ml_features = _MlFeatureAdapter(db) - DlService( - db, - _DlDrawAdapter(ml_draws), - _DlFeatureAdapter(ml_features), - ).train(lottery_id) - elif name == "bt": - from backend.app.services.bt_service import BtService - - BtService(db).run(lottery_id=lottery_id, strategy_id=BT_STRATEGY_ID) - elif name == "rank": - self._run_rank(lottery_id, state) - elif name == "select": - from backend.app.services.meta_service import MetaService - - if state.context_hash is None: - state.context_hash = derive_context_hash(db, lottery_id) - MetaService(db).select(lottery_id=lottery_id, context_hash=state.context_hash) elif name == "gen": from backend.app.services.gen_service import GenService @@ -234,75 +150,12 @@ def _execute_stage( else: # pragma: no cover - STAGE_ORDER is closed raise ValueError(f"unknown stage {name!r}") - def _run_rank(self, lottery_id: int, state: _RunState) -> None: - """Rank, then detect-and-rerank once against bt freshness (D8).""" - from backend.app.services.meta_service import MetaService - - meta = MetaService(self._db) - meta.rank(lottery_id=lottery_id) - - if state.context_hash is None: - state.context_hash = derive_context_hash(self._db, lottery_id) - - if self._ranking_stale(lottery_id, state.context_hash): - meta.rank(lottery_id=lottery_id) # exactly ONE repair attempt (D8) - if self._ranking_stale(lottery_id, state.context_hash): - raise RuntimeError("ranking stale for backtest context after one rerank") - - def _ranking_stale(self, lottery_id: int, context_hash: str) -> bool: - """True when no active ranking exists for ctx or it predates bt (D8).""" - from backend.app.models.bt_snapshot import BtSnapshot - from backend.app.models.meta_ranking import MetaRanking - - bt = ( - self._db.execute( - select(BtSnapshot) - .where(BtSnapshot.lottery_id == lottery_id, BtSnapshot.status == "active") - .order_by(BtSnapshot.created_at.desc()) - .limit(1) - ) - .scalars() - .first() - ) - ranking = ( - self._db.execute( - select(MetaRanking) - .where( - MetaRanking.lottery_id == lottery_id, - MetaRanking.context_hash == context_hash, - MetaRanking.status == "active", - ) - .limit(1) - ) - .scalars() - .first() - ) - if bt is None or ranking is None or bt.created_at is None or ranking.created_at is None: - return True - return _naive(ranking.created_at) <= _naive(bt.created_at) - - def _gated_skip(self, name: str, before: _Artifact, changed: dict[str, bool]) -> bool: - """Gate non-idempotent writers on missing/stale artifacts (D12).""" - if name not in _GATED_STAGES: - return False - if before.fingerprint is None: - return False - deps_changed = any(changed.get(dep, False) for dep in _DEPS.get(name, ())) - return not deps_changed - - def _safe_context_hash(self, lottery_id: int) -> str | None: - try: - return derive_context_hash(self._db, lottery_id) - except ValueError: - return None - # ------------------------------------------------------------------ # Active artifact resolution (skip-vs-run classification) # ------------------------------------------------------------------ - def _read_artifact(self, lottery_id: int, stage: str, context_hash: str | None) -> _Artifact: + def _read_artifact(self, lottery_id: int, stage: str) -> _Artifact: """Return the active artifact reference for a stage (absent → empty).""" - db = self._db if stage == "stats": from backend.app.models.stat_snapshot import StatSnapshot @@ -349,85 +202,6 @@ def _read_artifact(self, lottery_id: int, stage: str, context_hash: str | None) fp = f"{parts[0]}|{parts[1]}" ids = [r.id for r in (row, prob) if r is not None] return _Artifact(ids[0] if ids else None, fp) - if stage == "ml": - from backend.app.ml.registry import MODEL_SET_CORE_5 - from backend.app.models.ml_snapshot import MlSnapshot - - row = self._latest( - select(MlSnapshot) - .where( - MlSnapshot.lottery_id == lottery_id, - MlSnapshot.model_set == MODEL_SET_CORE_5, - MlSnapshot.status == "active", - ) - .order_by(MlSnapshot.version.desc()) - ) - return _Artifact(row.id if row else None, row.input_fingerprint if row else None) - if stage == "dl": - from backend.app.dl.registry import MODEL_SET_CORE_3 - from backend.app.models.dl_snapshot import DlSnapshot - - row = self._latest( - select(DlSnapshot) - .where( - DlSnapshot.lottery_id == lottery_id, - DlSnapshot.model_set == MODEL_SET_CORE_3, - DlSnapshot.status == "active", - ) - .order_by(DlSnapshot.version.desc()) - ) - return _Artifact(row.id if row else None, row.input_fingerprint if row else None) - if stage == "bt": - from backend.app.models.bt_snapshot import BtSnapshot - - row = self._latest( - select(BtSnapshot) - .where( - BtSnapshot.lottery_id == lottery_id, - BtSnapshot.strategy_id == BT_STRATEGY_ID, - BtSnapshot.status == "active", - ) - .order_by(BtSnapshot.created_at.desc()) - ) - return _Artifact(row.id if row else None, row.fingerprint if row else None) - if stage == "rank": - from backend.app.models.meta_ranking import MetaRanking - - if context_hash is None: - return _Artifact() - row = ( - db.execute( - select(MetaRanking) - .where( - MetaRanking.lottery_id == lottery_id, - MetaRanking.context_hash == context_hash, - MetaRanking.status == "active", - ) - .limit(1) - ) - .scalars() - .first() - ) - return _Artifact(row.id if row else None, row.fingerprint if row else None) - if stage == "select": - from backend.app.models.meta_selection import MetaSelection - - if context_hash is None: - return _Artifact() - row = ( - db.execute( - select(MetaSelection) - .where( - MetaSelection.lottery_id == lottery_id, - MetaSelection.context_hash == context_hash, - MetaSelection.status == "active", - ) - .limit(1) - ) - .scalars() - .first() - ) - return _Artifact(row.id if row else None, row.fingerprint if row else None) if stage == "gen": from backend.app.models.gen_snapshot import GenSnapshot @@ -442,120 +216,3 @@ def _read_artifact(self, lottery_id: int, stage: str, context_hash: str | None) def _latest(self, stmt): # noqa: ANN001 - private helper over a typed select """Return the single row produced by a typed select statement.""" return self._db.execute(stmt.limit(1)).scalars().first() - - -def _naive(value: object) -> object: - """Strip tzinfo so SQLite-loaded datetimes compare uniformly.""" - return value.replace(tzinfo=None) if getattr(value, "tzinfo", None) else value - - -# ---------------------------------------------------------------------- -# Provider adapters (mirror api/v1/ml.py + cli.py composition seams) -# ---------------------------------------------------------------------- - - -class _MlDrawAdapter: - """Minimal ML DrawHistoryProvider adapter over the draw tables.""" - - def __init__(self, session: Session) -> None: - self._session = session - - def iter_draws(self, lottery_id: int, *, after_draw_number: int | None = None): - """Yield ML ``DrawRow`` carriers in ascending draw-number order.""" - from sqlalchemy import select - - from backend.app.ml.providers import DrawRow - from backend.app.models.draw import Draw - from backend.app.models.draw_number import DrawNumber - - stmt = select(Draw).where(Draw.lottery_id == lottery_id).order_by(Draw.draw_number) - if after_draw_number is not None: - stmt = stmt.where(Draw.draw_number > after_draw_number) - for draw in self._session.execute(stmt).scalars().all(): - nums_stmt = ( - select(DrawNumber.number) - .where(DrawNumber.draw_id == draw.id) - .order_by(DrawNumber.position) - ) - numbers = tuple(self._session.execute(nums_stmt).scalars().all()) - yield DrawRow(draw_number=draw.draw_number, numbers=numbers) - - -class _MlFeatureAdapter: - """Minimal ML FeatureSnapshotProvider adapter over feature_values.""" - - def __init__(self, session: Session) -> None: - self._session = session - - def active_snapshot_id(self, lottery_id: int) -> int | None: - """Return the newest active ML snapshot id for the lottery, if any.""" - from sqlalchemy import select - - from backend.app.models.feature_snapshot import FeatureSnapshot - - stmt = ( - select(FeatureSnapshot) - .where( - FeatureSnapshot.lottery_id == lottery_id, - FeatureSnapshot.status == "active", - ) - .order_by(FeatureSnapshot.version.desc()) - .limit(1) - ) - snap = self._session.execute(stmt).scalar_one_or_none() - return snap.id if snap is not None else None - - def feature_rows(self, snapshot_id: int): - """Yield ML feature-value carriers ordered by draw then feature.""" - from sqlalchemy import select - - from backend.app.ml.feature_reader import FeatureValueRow - from backend.app.models.feature_value import FeatureValue - - stmt = ( - select(FeatureValue) - .where(FeatureValue.snapshot_id == snapshot_id) - .order_by(FeatureValue.draw_number, FeatureValue.feature_id) - ) - for fv in self._session.execute(stmt).scalars().all(): - yield FeatureValueRow( - feature_id=fv.feature_id, - draw_number=fv.draw_number, - value=float(fv.value), - ) - - -class _DlDrawAdapter: - """Converts ML draw carriers to DL carriers at the composition root (DLE-13).""" - - def __init__(self, inner: _MlDrawAdapter) -> None: - self._inner = inner - - def iter_draws(self, lottery_id: int, *, after_draw_number: int | None = None): - """Yield DL ``DrawRow`` carriers in ascending draw-number order.""" - from backend.app.dl.providers import DrawRow - - for row in self._inner.iter_draws(lottery_id, after_draw_number=after_draw_number): - yield DrawRow(draw_number=row.draw_number, numbers=tuple(row.numbers)) - - -class _DlFeatureAdapter: - """Converts ML feature carriers to DL carriers at the composition root (DLE-13).""" - - def __init__(self, inner: _MlFeatureAdapter) -> None: - self._inner = inner - - def active_snapshot_id(self, lottery_id: int) -> int | None: - """Return the newest active DL snapshot id for the lottery, if any.""" - return self._inner.active_snapshot_id(lottery_id) - - def feature_rows(self, snapshot_id: int): - """Yield DL feature carriers ordered by draw then feature.""" - from backend.app.dl.providers import FeatureRow - - for row in self._inner.feature_rows(snapshot_id): - yield FeatureRow( - feature_id=row.feature_id, - draw_number=row.draw_number, - value=float(row.value), - )