diff --git a/.dockerignore b/.dockerignore index 6840a92..93d7f7e 100644 --- a/.dockerignore +++ b/.dockerignore @@ -10,3 +10,9 @@ robustness_quick severity_outputs severity_quick ml_polish_outputs +archive +docs +presentation +train.csv +1_przebiegi_usterek.png +2_przebiegi_silnika.png diff --git a/.gitea/workflows/ci.yaml b/.gitea/workflows/ci.yaml index 7b8bfce..82747bc 100644 --- a/.gitea/workflows/ci.yaml +++ b/.gitea/workflows/ci.yaml @@ -23,42 +23,65 @@ jobs: run: | docker rm -f engin-ci-app 2>/dev/null || true docker image rm -f engin-console:ci 2>/dev/null || true + docker image rm -f engin-console:test 2>/dev/null || true - - name: Build production image - run: docker build --tag engin-console:ci . + - name: Build test image + run: docker build --target test --tag engin-console:test . - name: Run Ruff run: | docker run --rm \ --entrypoint sh \ - engin-console:ci \ + engin-console:test \ -c "python -m pip install --quiet \ --disable-pip-version-check \ --root-user-action=ignore \ ruff==0.16.4 && ruff check --no-cache \ - app.py engin tests final_pipeline.py ml_polish_benchmark.py" + app.py engin tests final_pipeline.py ml_polish_benchmark.py scripts" - name: Check installed dependencies run: | docker run --rm \ --entrypoint python \ - engin-console:ci \ + engin-console:test \ -m pip check - name: Run test suite run: | docker run --rm \ --entrypoint python \ - engin-console:ci \ + engin-console:test \ -m unittest discover -v + - name: Verify frozen model artifact + run: | + docker run --rm \ + --entrypoint python \ + engin-console:test \ + -m scripts.build_model_artifact --verify-only + - name: Validate Python syntax run: | docker run --rm \ --entrypoint python \ + engin-console:test \ + -m compileall -q app.py engin tests final_pipeline.py scripts + + - name: Build production image + run: docker build --target runtime --tag engin-console:ci . + + - name: Assert runtime excludes training assets + run: | + docker run --rm \ + --entrypoint sh \ engin-console:ci \ - -m compileall -q app.py engin tests final_pipeline.py + -c "test ! -e val.csv && + test ! -e train.csv && + test ! -e final_pipeline.py && + test ! -e severity_benchmark.py && + test -s artifacts/engin-2026.08.25-1/model.pkl && + test -s artifacts/engin-2026.08.25-1/manifest.json" - name: Start application run: | @@ -90,3 +113,4 @@ jobs: docker logs engin-ci-app 2>/dev/null || true docker rm -f engin-ci-app 2>/dev/null || true docker image rm -f engin-console:ci 2>/dev/null || true + docker image rm -f engin-console:test 2>/dev/null || true diff --git a/Dockerfile b/Dockerfile index 3865f5d..34f35f1 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,4 +1,5 @@ -FROM python:3.12-slim +ARG PYTHON_IMAGE=python:3.12-slim@sha256:7a8b475003c4fe15a2cd4e55e5cfc2f3560bdc9333d624f24cdd6d4340fd7a17 +FROM ${PYTHON_IMAGE} AS base ENV PYTHONDONTWRITEBYTECODE=1 \ PYTHONUNBUFFERED=1 \ @@ -9,8 +10,18 @@ WORKDIR /app COPY requirements-lock.txt ./ RUN pip install --no-cache-dir -r requirements-lock.txt +FROM base AS test COPY . . +FROM base AS runtime + +COPY app.py ./ +COPY engin/ ./engin/ +COPY assets/ ./assets/ +COPY .streamlit/ ./.streamlit/ +COPY artifacts/ ./artifacts/ +COPY test.csv predictions.csv prediction_diagnostics.csv ./ + EXPOSE 8501 HEALTHCHECK --interval=30s --timeout=5s --start-period=30s --retries=3 \ diff --git a/README.md b/README.md index ec10efa..5f687c3 100644 --- a/README.md +++ b/README.md @@ -4,7 +4,8 @@ Gotowy do demonstracji system diagnostyki cylindrów przemysłowych silników Di ## Uruchomienie -Wymagany jest Python 3.12 (3.11 również powinien działać). +Wymagany jest Python 3.12. Artefakt modelu sprawdza zgodność wersji Pythona +i bibliotek przed deserializacją. ```bash python -m venv .venv @@ -48,6 +49,8 @@ docker run --rm -p 8501:8501 engin-console Interfejs jest cienką warstwą prezentacji. Zależności są wstrzykiwane do `DiagnosticService`, dlatego odczyt pliku, walidator i model można niezależnie wymieniać oraz testować. +Proces produkcyjny nie trenuje modelu przy starcie. Aplikacja weryfikuje checksumę, schemat wejścia i wersje bibliotek, a następnie ładuje artefakt `engin-2026.08.25-1`. Uszkodzony lub niezgodny artefakt uruchamia ograniczony fallback demonstracyjny. + ```mermaid flowchart TD UI["Streamlit UI"] --> S["DiagnosticService"] @@ -64,6 +67,9 @@ app.py cienka warstwa Streamlit engin/io.py adapter odczytu CSV engin/validation.py kontrakt i walidacja danych engin/model.py adapter modelu i fallback demo +engin/artifact.py manifest, checksumy i ładowanie artefaktu +engin/features.py stabilne transformacje używane w train/inference +engin/inference.py runtime predykcji bez zależności od benchmarków engin/service.py przypadek użycia / dependency injection engin/explainability.py status, porządkowy triage i uzasadnienia engin/charts.py czyste, testowalne fabryki Plotly @@ -102,6 +108,17 @@ python final_pipeline.py Powstają `predictions.csv` (600/600 rekordów, bez duplikatów i braków) oraz `prediction_diagnostics.csv` używany przez aplikację do explainability. +Wersjonowany artefakt modelu buduje osobny krok release: + +```bash +python -m scripts.build_model_artifact \ + --source-revision "$(git rev-parse HEAD)" + +python -m scripts.build_model_artifact --verify-only +``` + +Druga komenda nie trenuje modelu. Ładuje artefakt, sprawdza manifest i checksumę oraz wymaga dokładnej zgodności wszystkich 600 predykcji z zamrożonym `predictions.csv`. + ## Kontrakt danych wejściowych Każdy wiersz reprezentuje cylinder. Wymagane kolumny: @@ -121,13 +138,14 @@ Pełna kontrola projektu: python -m unittest discover -v ``` -Pakiet zawiera 32 testy: +Pakiet zawiera 38 testów: - testy braku leakage, cech, OOD i kontraktu submission, - testy jednostkowe walidatora i błędnych CSV, - testy serwisu z fake reader/validator/model — weryfikują dependency injection, - testy explainability, rankingu i wykresów, - testy smoke Streamlit dla demo, pustego uploadu i synchronizacji mapy cylindrów ze szczegółami. +- testy artefaktu: checksumy, zgodność bibliotek, brak odczytu danych treningowych przy starcie i regresja predykcji 600/600. Szybka kontrola serwera: @@ -141,7 +159,8 @@ Oczekiwana odpowiedź: `ok`. ## Odporność operacyjna -- Model jest trenowany raz i przechowywany w cache procesu. +- Model jest trenowany poza aplikacją i ładowany raz z wersjonowanego artefaktu do cache procesu. +- Finalny obraz runtime nie zawiera `val.csv`, `train.csv` ani skryptów treningowych. - Nieoczekiwany błąd inicjalizacji modelu uruchamia read-only fallback dla danych demo. - Nieprawidłowe dane nie docierają do modelu. - Błędy użytkownika mają stabilne kody i nie pokazują tracebacków w interfejsie. diff --git a/app.py b/app.py index 6cea172..3e3b8e5 100644 --- a/app.py +++ b/app.py @@ -34,6 +34,8 @@ from engin.validation import SpectrumFrameValidator LOGGER = logging.getLogger("engin.app") BASE_DIR = Path(__file__).resolve().parent +MODEL_VERSION = "engin-2026.08.25-1" +MODEL_ARTIFACT_DIR = BASE_DIR / "artifacts" / MODEL_VERSION @dataclass(frozen=True) @@ -42,24 +44,28 @@ class AppDependencies: demo_payload: bytes live_inference: bool = True startup_warning: str | None = None + model_version: str = "unknown" def _reader() -> PandasCsvReader: return PandasCsvReader(AppConfig().max_upload_bytes) -@st.cache_resource(show_spinner="Przygotowanie modelu diagnostycznego…") +@st.cache_resource(show_spinner="Ładowanie modelu diagnostycznego…") def build_dependencies() -> AppDependencies: reader = _reader() validator = SpectrumFrameValidator() demo_path = BASE_DIR / "test.csv" demo_payload = demo_path.read_bytes() try: - reference = reader.read_path(BASE_DIR / "val.csv") - model = SklearnPredictionModel.train(reference, random_state=42, n_jobs=-1) + model = SklearnPredictionModel.from_artifact( + MODEL_ARTIFACT_DIR, + expected_model_version=MODEL_VERSION, + ) return AppDependencies( service=DiagnosticService(reader=reader, validator=validator, model=model), demo_payload=demo_payload, + model_version=MODEL_VERSION, ) except Exception: LOGGER.exception("Live model initialization failed; enabling demo fallback") @@ -71,9 +77,10 @@ def build_dependencies() -> AppDependencies: demo_payload=demo_payload, live_inference=False, startup_warning=( - "Model live nie został przygotowany. Aplikacja działa w bezpiecznym " + "Artefakt modelu nie został załadowany. Aplikacja działa w bezpiecznym " "trybie demonstracyjnym na zapisanych predykcjach." ), + model_version="demo-precomputed", ) @@ -236,7 +243,7 @@ def _render_technical(result: DiagnosisResult) -> None: c1.metric("Grouped Macro F1", "0.981") c2.metric("Severity accuracy", "0.930") c3.metric("Walidacyjne ML", "33.65 / 40") - c4.metric("Testy", "32/32") + c4.metric("Bramka jakości", "CI GREEN") st.markdown( """ - **Zero leakage:** każdy fold zawiera kompletne, wcześniej niewidziane silniki. @@ -316,7 +323,9 @@ def render_app(dependencies: AppDependencies | None = None) -> None: width="stretch", ) st.divider() - st.caption("Model ENGIN v1 · CPU inference · brak połączeń zewnętrznych") + st.caption( + f"Model {deps.model_version} · CPU inference · brak połączeń zewnętrznych" + ) analysis = analyze_engine(result, str(engine_id)) summary = summarize_engine(analysis) diff --git a/artifacts/engin-2026.08.25-1/manifest.json b/artifacts/engin-2026.08.25-1/manifest.json new file mode 100644 index 0000000..56d3ca0 --- /dev/null +++ b/artifacts/engin-2026.08.25-1/manifest.json @@ -0,0 +1,49 @@ +{ + "artifact_format_version": 1, + "created_at_utc": "2026-08-25T10:57:21.704599+00:00", + "input_columns": [ + "engine_id", + "cylinder", + "n_cylinders", + "mV_0", + "mV_1", + "mV_2", + "mV_3", + "mV_4", + "mV_5", + "mV_6", + "mV_7", + "mV_8", + "mV_9", + "mV_10", + "mV_11", + "mV_12", + "mV_13", + "mV_14", + "mV_15", + "mV_16", + "mV_17", + "mV_18", + "mV_19", + "mV_20" + ], + "label_model": "deviation_logistic_c10", + "model_sha256": "a649e96ab13dd49a69d73e12fab46f830d128dac49a93464998a4bdfe8cbec2f", + "model_version": "engin-2026.08.25-1", + "n_jobs": 1, + "ood_absolute_threshold_mv": 7.25, + "ood_ratio_threshold": 2.5, + "python_version": "3.12.13", + "random_state": 42, + "runtime_versions": { + "joblib": "1.5.3", + "numpy": "2.3.5", + "pandas": "2.2.3", + "scikit-learn": "1.8.0", + "scipy": "1.18.1", + "threadpoolctl": "3.6.0" + }, + "severity_model": "deviation_extra_trees_mf03", + "source_revision": "8981ae0a0086e1bd6d8ab6900ddd5a78bc92a304", + "training_data_sha256": "82120af46fa62eefc1cc62f7cd953c4af830207f1d5083dccff5900fb57737f1" +} diff --git a/artifacts/engin-2026.08.25-1/model.pkl b/artifacts/engin-2026.08.25-1/model.pkl new file mode 100644 index 0000000..e7cd682 Binary files /dev/null and b/artifacts/engin-2026.08.25-1/model.pkl differ diff --git a/benchmark_grouped.py b/benchmark_grouped.py index e43b445..52c50cf 100644 --- a/benchmark_grouped.py +++ b/benchmark_grouped.py @@ -33,13 +33,9 @@ from sklearn.model_selection import GroupKFold, StratifiedGroupKFold from sklearn.pipeline import Pipeline from sklearn.preprocessing import StandardScaler +from engin.config import FAULT_LABELS, FREQ_COLS, LABELS, NOT_APPLICABLE, SEVERITIES RANDOM_STATE = 42 -FREQ_COLS = [f"mV_{frequency}" for frequency in range(21)] -LABELS = ["ok", "zakoksowany", "lejacy", "pompa", "iglica", "unknown"] -FAULT_LABELS = ["zakoksowany", "lejacy", "pompa", "iglica"] -SEVERITIES = ["male", "srednie", "duze"] -NOT_APPLICABLE = "nie_dotyczy" FEATURE_SETS = ("raw", "relative", "combined") MODEL_NAMES = ("logistic", "extra_trees") diff --git a/engin/artifact.py b/engin/artifact.py new file mode 100644 index 0000000..ab85c3f --- /dev/null +++ b/engin/artifact.py @@ -0,0 +1,257 @@ +"""Versioned, checksummed model-artifact persistence for runtime loading.""" + +from __future__ import annotations + +import hashlib +import json +import os +import pickle +import platform +import tempfile +from dataclasses import asdict, dataclass +from datetime import UTC, datetime +from importlib.metadata import version +from pathlib import Path + +from .inference import ( + INFERENCE_COLUMNS, + OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD, + OOD_OK_TO_UNKNOWN_THRESHOLD_MV, + DiagnosticModels, +) + +ARTIFACT_FORMAT_VERSION = 1 +MODEL_FILENAME = "model.pkl" +MANIFEST_FILENAME = "manifest.json" +RUNTIME_PACKAGES = ( + "joblib", + "numpy", + "pandas", + "scikit-learn", + "scipy", + "threadpoolctl", +) + + +class ModelArtifactError(RuntimeError): + """Raised when a model artifact is missing, corrupt, or incompatible.""" + + +@dataclass(frozen=True) +class ArtifactMetadata: + artifact_format_version: int + model_version: str + created_at_utc: str + source_revision: str + training_data_sha256: str + model_sha256: str + input_columns: tuple[str, ...] + python_version: str + random_state: int + n_jobs: int + label_model: str + severity_model: str + ood_absolute_threshold_mv: float + ood_ratio_threshold: float + runtime_versions: dict[str, str] + + @classmethod + def from_dict(cls, payload: dict[str, object]) -> "ArtifactMetadata": + try: + return cls( + artifact_format_version=int(payload["artifact_format_version"]), + model_version=str(payload["model_version"]), + created_at_utc=str(payload["created_at_utc"]), + source_revision=str(payload["source_revision"]), + training_data_sha256=str(payload["training_data_sha256"]), + model_sha256=str(payload["model_sha256"]), + input_columns=tuple(str(value) for value in payload["input_columns"]), + python_version=str(payload["python_version"]), + random_state=int(payload["random_state"]), + n_jobs=int(payload["n_jobs"]), + label_model=str(payload["label_model"]), + severity_model=str(payload["severity_model"]), + ood_absolute_threshold_mv=float( + payload["ood_absolute_threshold_mv"] + ), + ood_ratio_threshold=float(payload["ood_ratio_threshold"]), + runtime_versions={ + str(name): str(package_version) + for name, package_version in dict( + payload["runtime_versions"] + ).items() + }, + ) + except (KeyError, TypeError, ValueError) as exc: + raise ModelArtifactError("Model manifest has an invalid schema.") from exc + + def to_dict(self) -> dict[str, object]: + payload = asdict(self) + payload["input_columns"] = list(self.input_columns) + return payload + + +@dataclass(frozen=True) +class LoadedModelArtifact: + models: DiagnosticModels + metadata: ArtifactMetadata + + +def sha256_file(path: Path) -> str: + digest = hashlib.sha256() + with path.open("rb") as handle: + for chunk in iter(lambda: handle.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def _runtime_versions() -> dict[str, str]: + return {package: version(package) for package in RUNTIME_PACKAGES} + + +def _atomic_write(path: Path, payload: bytes) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + descriptor, temporary_name = tempfile.mkstemp( + dir=path.parent, prefix=f".{path.name}.", suffix=".tmp" + ) + temporary_path = Path(temporary_name) + try: + with os.fdopen(descriptor, "wb") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) + os.replace(temporary_path, path) + except Exception: + temporary_path.unlink(missing_ok=True) + raise + + +def save_model_artifact( + models: DiagnosticModels, + directory: Path, + *, + model_version: str, + source_revision: str, + training_data_sha256: str, + random_state: int, + n_jobs: int, + label_model: str, + severity_model: str, +) -> ArtifactMetadata: + if not model_version.strip(): + raise ValueError("model_version cannot be empty.") + if len(training_data_sha256) != 64 or any( + character not in "0123456789abcdef" for character in training_data_sha256 + ): + raise ValueError("training_data_sha256 must be a SHA-256 hex digest.") + + envelope = { + "artifact_format_version": ARTIFACT_FORMAT_VERSION, + "model_version": model_version, + "models": models, + } + model_payload = pickle.dumps(envelope, protocol=pickle.HIGHEST_PROTOCOL) + model_sha256 = hashlib.sha256(model_payload).hexdigest() + metadata = ArtifactMetadata( + artifact_format_version=ARTIFACT_FORMAT_VERSION, + model_version=model_version, + created_at_utc=datetime.now(UTC).isoformat(), + source_revision=source_revision, + training_data_sha256=training_data_sha256, + model_sha256=model_sha256, + input_columns=tuple(INFERENCE_COLUMNS), + python_version=platform.python_version(), + random_state=random_state, + n_jobs=n_jobs, + label_model=label_model, + severity_model=severity_model, + ood_absolute_threshold_mv=OOD_OK_TO_UNKNOWN_THRESHOLD_MV, + ood_ratio_threshold=OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD, + runtime_versions=_runtime_versions(), + ) + + directory.mkdir(parents=True, exist_ok=True) + _atomic_write(directory / MODEL_FILENAME, model_payload) + manifest_payload = json.dumps( + metadata.to_dict(), ensure_ascii=False, indent=2, sort_keys=True + ).encode("utf-8") + b"\n" + _atomic_write(directory / MANIFEST_FILENAME, manifest_payload) + return metadata + + +def load_model_artifact( + directory: Path, *, expected_model_version: str | None = None +) -> LoadedModelArtifact: + manifest_path = directory / MANIFEST_FILENAME + model_path = directory / MODEL_FILENAME + try: + manifest_payload = json.loads(manifest_path.read_text(encoding="utf-8")) + except FileNotFoundError as exc: + raise ModelArtifactError(f"Model manifest is missing: {manifest_path}") from exc + except (OSError, json.JSONDecodeError) as exc: + raise ModelArtifactError("Model manifest cannot be read.") from exc + + if not isinstance(manifest_payload, dict): + raise ModelArtifactError("Model manifest must be a JSON object.") + metadata = ArtifactMetadata.from_dict(manifest_payload) + if metadata.artifact_format_version != ARTIFACT_FORMAT_VERSION: + raise ModelArtifactError( + "Unsupported model artifact format: " + f"{metadata.artifact_format_version}." + ) + if expected_model_version and metadata.model_version != expected_model_version: + raise ModelArtifactError( + f"Expected model {expected_model_version}, received {metadata.model_version}." + ) + if metadata.input_columns != tuple(INFERENCE_COLUMNS): + raise ModelArtifactError("Model input schema does not match this application.") + expected_python = metadata.python_version.split(".")[:2] + installed_python = platform.python_version().split(".")[:2] + if expected_python != installed_python: + raise ModelArtifactError( + "Model Python version does not match the runtime. " + f"expected={metadata.python_version}, installed={platform.python_version()}" + ) + + current_versions = _runtime_versions() + if metadata.runtime_versions != current_versions: + raise ModelArtifactError( + "Model runtime versions do not match installed dependencies. " + f"expected={metadata.runtime_versions}, installed={current_versions}" + ) + + try: + model_payload = model_path.read_bytes() + except OSError as exc: + raise ModelArtifactError(f"Model file cannot be read: {model_path}") from exc + actual_sha256 = hashlib.sha256(model_payload).hexdigest() + if actual_sha256 != metadata.model_sha256: + raise ModelArtifactError("Model checksum verification failed.") + + try: + envelope = pickle.loads(model_payload) # noqa: S301 - trusted, checksummed release artifact + except Exception as exc: + raise ModelArtifactError("Model payload cannot be deserialized.") from exc + if not isinstance(envelope, dict): + raise ModelArtifactError("Model payload has an invalid envelope.") + if envelope.get("artifact_format_version") != ARTIFACT_FORMAT_VERSION: + raise ModelArtifactError("Model payload format does not match the manifest.") + if envelope.get("model_version") != metadata.model_version: + raise ModelArtifactError("Model payload version does not match the manifest.") + models = envelope.get("models") + if not isinstance(models, DiagnosticModels): + raise ModelArtifactError("Model payload has an unexpected object type.") + return LoadedModelArtifact(models=models, metadata=metadata) + + +__all__ = [ + "ARTIFACT_FORMAT_VERSION", + "MANIFEST_FILENAME", + "MODEL_FILENAME", + "ArtifactMetadata", + "LoadedModelArtifact", + "ModelArtifactError", + "load_model_artifact", + "save_model_artifact", + "sha256_file", +] diff --git a/engin/config.py b/engin/config.py index 2cedbc7..ac1ee05 100644 --- a/engin/config.py +++ b/engin/config.py @@ -4,14 +4,11 @@ from __future__ import annotations from dataclasses import dataclass -from benchmark_grouped import ( - FAULT_LABELS, - FREQ_COLS, - LABELS, - NOT_APPLICABLE, - SEVERITIES, -) - +FREQ_COLS = [f"mV_{frequency}" for frequency in range(21)] +LABELS = ["ok", "zakoksowany", "lejacy", "pompa", "iglica", "unknown"] +FAULT_LABELS = ["zakoksowany", "lejacy", "pompa", "iglica"] +SEVERITIES = ["male", "srednie", "duze"] +NOT_APPLICABLE = "nie_dotyczy" ALLOWED_ENGINE_SIZES = (8, 12, 16) MAX_UPLOAD_BYTES = 10 * 1024 * 1024 MAX_MISSING_FRACTION_PER_CYLINDER = 0.50 diff --git a/engin/features.py b/engin/features.py new file mode 100644 index 0000000..c82c6f1 --- /dev/null +++ b/engin/features.py @@ -0,0 +1,171 @@ +"""Stable feature transformers shared by model training and runtime inference.""" + +from __future__ import annotations + +from typing import Iterable + +import numpy as np +import pandas as pd + +from .config import FAULT_LABELS, FREQ_COLS + +FEATURE_SETS = ("raw", "deviation", "all") + + +class SeverityFeatures: + """Create engine-relative spectral features without using labels. + + Engine-relative values use a leave-one-cylinder-out median reference. Input + to ``transform`` must therefore contain every cylinder of each engine. + ``diagnostic_label`` is optional and may contain only a model-predicted + label at validation or inference time. + """ + + def __init__(self, feature_set: str, include_fault_type: bool = False) -> None: + if feature_set not in FEATURE_SETS: + raise ValueError(f"Unknown feature_set={feature_set!r}") + self.feature_set = feature_set + self.include_fault_type = include_fault_type + + def fit( + self, X: pd.DataFrame, y: Iterable[str] | None = None + ) -> "SeverityFeatures": + self._validate(X) + spectra = X[FREQ_COLS].apply(pd.to_numeric, errors="coerce") + fallback = spectra.median(axis=0).to_numpy(dtype=float) + if np.isnan(fallback).any(): + raise ValueError("A frequency column is entirely NaN in training.") + self.fallback_medians_ = fallback + return self + + def transform(self, X: pd.DataFrame) -> np.ndarray: + self._validate(X) + if not hasattr(self, "fallback_medians_"): + raise RuntimeError("SeverityFeatures must be fitted before transform().") + raw = self._clean_spectra(X) + reference = self._leave_one_out_engine_median( + raw, X["engine_id"].to_numpy() + ) + relative = raw - reference + absolute_relative = np.abs(relative) + ratio_delta = raw / np.maximum(np.abs(reference), 1e-6) - 1.0 + + if self.feature_set == "raw": + features = raw + elif self.feature_set == "deviation": + features = np.hstack( + [ + relative, + absolute_relative, + ratio_delta, + self._summary_features(raw, relative, ratio_delta), + ] + ) + else: + features = np.hstack( + [ + raw, + relative, + absolute_relative, + ratio_delta, + self._summary_features(raw, relative, ratio_delta), + ] + ) + + if self.include_fault_type: + if "diagnostic_label" not in X.columns: + raise ValueError( + "diagnostic_label is required when include_fault_type=True" + ) + labels = X["diagnostic_label"].to_numpy(dtype=object) + one_hot = np.column_stack( + [labels == label for label in FAULT_LABELS] + ).astype(float) + features = np.hstack([features, one_hot]) + if np.isnan(features).any() or np.isinf(features).any(): + raise ValueError("Severity features contain NaN or infinity.") + return features + + def fit_transform( + self, X: pd.DataFrame, y: Iterable[str] | None = None + ) -> np.ndarray: + return self.fit(X).transform(X) + + def _clean_spectra(self, X: pd.DataFrame) -> np.ndarray: + spectra = X[FREQ_COLS].apply(pd.to_numeric, errors="coerce") + spectra = spectra.interpolate(axis=1, limit_direction="both") + spectra = spectra.fillna(pd.Series(self.fallback_medians_, index=FREQ_COLS)) + return spectra.to_numpy(dtype=float) + + @staticmethod + def _leave_one_out_engine_median( + raw: np.ndarray, groups: np.ndarray + ) -> np.ndarray: + reference = np.empty_like(raw) + for engine_id in pd.unique(groups): + positions = np.flatnonzero(groups == engine_id) + if len(positions) < 2: + raise ValueError( + f"Engine {engine_id!r} has fewer than two cylinders in transform()." + ) + engine_values = raw[positions] + for local_position, global_position in enumerate(positions): + keep = np.arange(len(positions)) != local_position + reference[global_position] = np.median(engine_values[keep], axis=0) + return reference + + @staticmethod + def _summary_features( + raw: np.ndarray, relative: np.ndarray, ratio_delta: np.ndarray + ) -> np.ndarray: + absolute_relative = np.abs(relative) + gradients = np.diff(raw, axis=1) + frequency = np.arange(raw.shape[1], dtype=float) + centered_frequency = frequency - frequency.mean() + slope = raw @ centered_frequency / np.sum(centered_frequency**2) + if hasattr(np, "trapezoid"): + spectral_auc = np.trapezoid(raw, axis=1) + else: # pragma: no cover - compatibility with NumPy 1.x + spectral_auc = np.trapz(raw, axis=1) + + columns = [ + raw.mean(axis=1), + raw.std(axis=1), + raw.min(axis=1), + raw.max(axis=1), + np.ptp(raw, axis=1), + spectral_auc, + raw.argmax(axis=1) / 20.0, + raw.argmin(axis=1) / 20.0, + slope, + np.abs(gradients).mean(axis=1), + np.abs(gradients).max(axis=1), + relative.mean(axis=1), + relative.std(axis=1), + relative.min(axis=1), + relative.max(axis=1), + absolute_relative.mean(axis=1), + absolute_relative.max(axis=1), + np.sqrt(np.mean(relative**2, axis=1)), + ratio_delta.mean(axis=1), + ratio_delta.std(axis=1), + ] + for start, stop in ((0, 5), (5, 10), (10, 15), (15, 21)): + columns.extend( + [ + raw[:, start:stop].mean(axis=1), + relative[:, start:stop].mean(axis=1), + absolute_relative[:, start:stop].mean(axis=1), + ] + ) + return np.column_stack(columns) + + @staticmethod + def _validate(X: pd.DataFrame) -> None: + required = {"engine_id", *FREQ_COLS} + missing = sorted(required.difference(X.columns)) + if missing: + raise ValueError(f"Missing columns: {missing}") + + +__all__ = ["FEATURE_SETS", "SeverityFeatures"] diff --git a/engin/inference.py b/engin/inference.py new file mode 100644 index 0000000..ffbc82c --- /dev/null +++ b/engin/inference.py @@ -0,0 +1,205 @@ +"""Runtime-only prediction contract for the frozen ENGIN model.""" + +from __future__ import annotations + +from dataclasses import dataclass + +import numpy as np +import pandas as pd + +from .config import FAULT_LABELS, FREQ_COLS, LABELS, NOT_APPLICABLE, SEVERITIES + +OOD_OK_TO_UNKNOWN_THRESHOLD_MV = 7.25 +OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD = 2.5 +SUBMISSION_COLUMNS = ["engine_id", "cylinder", "label", "severity"] +KEY_COLUMNS = ["engine_id", "cylinder"] +INFERENCE_COLUMNS = ["engine_id", "cylinder", "n_cylinders", *FREQ_COLS] + + +@dataclass +class DiagnosticModels: + label_pipeline: object + severity_transformer: object + severity_estimator: object + + +def apply_ood_override( + predicted_label: np.ndarray, + deviation_features: np.ndarray, + engine_ids: np.ndarray | pd.Series, + threshold_mv: float = OOD_OK_TO_UNKNOWN_THRESHOLD_MV, + threshold_ratio: float = OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD, +) -> tuple[np.ndarray, np.ndarray, np.ndarray, np.ndarray]: + """Turn an implausibly anomalous ``ok`` into ``unknown``.""" + + result = np.asarray(predicted_label, dtype=object).copy() + absolute_deviation = deviation_features[ + :, len(FREQ_COLS) : 2 * len(FREQ_COLS) + ] + anomaly_score = absolute_deviation.mean(axis=1) + ids = np.asarray(engine_ids, dtype=object) + if len(ids) != len(result): + raise ValueError("engine_ids must align with predicted labels.") + score_frame = pd.DataFrame({"engine_id": ids, "anomaly_score": anomaly_score}) + engine_median = ( + score_frame.groupby("engine_id", sort=False)["anomaly_score"] + .transform("median") + .to_numpy(dtype=float) + ) + anomaly_ratio = anomaly_score / np.maximum(engine_median, 1e-6) + override = ( + (result == "ok") + & (anomaly_score > threshold_mv) + & (anomaly_ratio > threshold_ratio) + ) + result[override] = "unknown" + return result, override, anomaly_score, anomaly_ratio + + +def validate_inference_data(df: pd.DataFrame) -> None: + required = set(INFERENCE_COLUMNS) + missing = sorted(required.difference(df.columns)) + if missing: + raise ValueError(f"Inference data is missing columns: {missing}") + if df.empty: + raise ValueError("Inference data is empty.") + if df[["engine_id", "cylinder", "n_cylinders"]].isna().any().any(): + raise ValueError("Inference identifiers cannot contain NaN.") + if df.duplicated(KEY_COLUMNS).any(): + raise ValueError("Duplicate engine_id + cylinder keys found in inference data.") + + engine_sizes = df.groupby("engine_id").agg( + rows=("cylinder", "size"), + expected=("n_cylinders", "first"), + size_values=("n_cylinders", "nunique"), + unique_cylinders=("cylinder", "nunique"), + ) + complete = ( + (engine_sizes["rows"] == engine_sizes["expected"]) + & (engine_sizes["unique_cylinders"] == engine_sizes["expected"]) + & (engine_sizes["size_values"] == 1) + & (engine_sizes["rows"] >= 2) + ) + if not complete.all(): + bad = engine_sizes.index[~complete].tolist() + raise ValueError(f"Incomplete or inconsistent inference engines: {bad}") + + +def _prepare_labeled_frame(df: pd.DataFrame, labels: np.ndarray) -> pd.DataFrame: + result = df.copy() + result["diagnostic_label"] = labels + return result + + +def predict_test( + models: DiagnosticModels, test: pd.DataFrame +) -> tuple[pd.DataFrame, pd.DataFrame]: + validate_inference_data(test) + + raw_model_label = models.label_pipeline.predict(test).astype(object) + label_features = models.label_pipeline.named_steps["features"].transform(test) + predicted_label, ood_override, anomaly_score, anomaly_ratio = apply_ood_override( + raw_model_label, + label_features, + test["engine_id"], + ) + label_probabilities = models.label_pipeline.predict_proba(test) + raw_model_confidence = label_probabilities.max(axis=1) + label_confidence = raw_model_confidence.copy() + label_confidence[ood_override] = np.nan + ordered_probabilities = np.sort(label_probabilities, axis=1) + label_margin = ordered_probabilities[:, -1] - ordered_probabilities[:, -2] + + test_with_labels = _prepare_labeled_frame(test, predicted_label) + severity_features = models.severity_transformer.transform(test_with_labels) + predicted_severity_all = models.severity_estimator.predict( + severity_features + ).astype(object) + + predicted_fault = np.isin(predicted_label, FAULT_LABELS) + emitted_severity = np.full(len(test), NOT_APPLICABLE, dtype=object) + emitted_severity[predicted_fault] = predicted_severity_all[predicted_fault] + + severity_probabilities = models.severity_estimator.predict_proba( + severity_features + ) + severity_confidence = np.full(len(test), np.nan, dtype=float) + severity_confidence[predicted_fault] = severity_probabilities.max(axis=1)[ + predicted_fault + ] + + absolute_deviation = severity_features[:, len(FREQ_COLS) : 2 * len(FREQ_COLS)] + top_frequency_indices = np.argsort(absolute_deviation, axis=1)[:, -3:][:, ::-1] + top_frequencies = [ + "|".join(str(int(index)) for index in row) for row in top_frequency_indices + ] + + submission = test[KEY_COLUMNS].copy() + submission["label"] = predicted_label + submission["severity"] = emitted_severity + + diagnostics = submission.copy() + diagnostics["raw_model_label"] = raw_model_label + diagnostics["decision_source"] = np.where( + ood_override, "ood_override", "classifier" + ) + diagnostics["label_confidence"] = label_confidence + diagnostics["raw_model_confidence"] = raw_model_confidence + diagnostics["label_margin"] = label_margin + diagnostics["severity_confidence"] = severity_confidence + diagnostics["anomaly_score_mean_abs_mv"] = anomaly_score + diagnostics["anomaly_ratio_to_engine_median"] = anomaly_ratio + diagnostics["ood_absolute_threshold_mv"] = OOD_OK_TO_UNKNOWN_THRESHOLD_MV + diagnostics["ood_ratio_threshold"] = OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD + diagnostics["top_anomalous_frequencies_khz"] = top_frequencies + diagnostics["missing_spectral_cells"] = ( + test[FREQ_COLS].isna().sum(axis=1).to_numpy(dtype=int) + ) + return submission, diagnostics + + +def validate_submission(submission: pd.DataFrame, test: pd.DataFrame) -> None: + if submission.columns.tolist() != SUBMISSION_COLUMNS: + raise ValueError( + f"Submission columns must be exactly {SUBMISSION_COLUMNS}; " + f"received={submission.columns.tolist()}" + ) + if len(submission) != len(test): + raise ValueError("Submission row count does not match test.csv.") + if submission.isna().any().any(): + raise ValueError("Submission contains NaN.") + if submission.duplicated(KEY_COLUMNS).any(): + raise ValueError("Submission contains duplicate engine_id + cylinder keys.") + if not submission[KEY_COLUMNS].reset_index(drop=True).equals( + test[KEY_COLUMNS].reset_index(drop=True) + ): + raise ValueError("Submission keys or row order do not match test.csv.") + + invalid_labels = sorted(set(submission["label"]).difference(LABELS)) + if invalid_labels: + raise ValueError(f"Submission contains invalid labels: {invalid_labels}") + + fault = submission["label"].isin(FAULT_LABELS) + invalid_fault_severity = sorted( + set(submission.loc[fault, "severity"]).difference(SEVERITIES) + ) + if invalid_fault_severity: + raise ValueError( + f"Fault predictions contain invalid severity: {invalid_fault_severity}" + ) + if not submission.loc[~fault, "severity"].eq(NOT_APPLICABLE).all(): + raise ValueError("ok/unknown predictions must use severity=nie_dotyczy.") + + +__all__ = [ + "DiagnosticModels", + "INFERENCE_COLUMNS", + "KEY_COLUMNS", + "OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD", + "OOD_OK_TO_UNKNOWN_THRESHOLD_MV", + "SUBMISSION_COLUMNS", + "apply_ood_override", + "predict_test", + "validate_inference_data", + "validate_submission", +] diff --git a/engin/model.py b/engin/model.py index 81a8b64..e1aa688 100644 --- a/engin/model.py +++ b/engin/model.py @@ -3,18 +3,14 @@ from __future__ import annotations from dataclasses import dataclass +from pathlib import Path from typing import Protocol import pandas as pd -from final_pipeline import ( - DiagnosticModels, - predict_test, - train_models, - validate_submission, -) - +from .artifact import ArtifactMetadata, load_model_artifact from .errors import InferenceError +from .inference import DiagnosticModels, predict_test, validate_submission @dataclass(frozen=True) @@ -28,31 +24,31 @@ class PredictionModel(Protocol): class SklearnPredictionModel: - """Thin adapter around the frozen, validated competition pipeline.""" + """Runtime adapter around a verified, pre-trained model artifact.""" - def __init__(self, models: DiagnosticModels) -> None: + def __init__( + self, models: DiagnosticModels, metadata: ArtifactMetadata | None = None + ) -> None: self._models = models + self.metadata = metadata @classmethod - def train( + def from_artifact( cls, - reference_frame: pd.DataFrame, + directory: Path, *, - random_state: int = 42, - n_jobs: int = -1, + expected_model_version: str | None = None, ) -> "SklearnPredictionModel": try: - models = train_models( - reference_frame, - random_state=random_state, - n_jobs=n_jobs, + artifact = load_model_artifact( + directory, expected_model_version=expected_model_version ) except Exception as exc: raise InferenceError( - "Nie udało się przygotować modelu referencyjnego.", - hint="Sprawdź kompletność val.csv oraz zgodność wersji zależności.", + "Nie udało się załadować zweryfikowanego artefaktu modelu.", + hint="Sprawdź manifest, checksumę oraz zgodność wersji zależności.", ) from exc - return cls(models) + return cls(artifact.models, artifact.metadata) def predict(self, frame: pd.DataFrame) -> PredictionBundle: try: @@ -63,6 +59,13 @@ class SklearnPredictionModel: "Model odrzucił pomiary podczas predykcji.", hint="Zweryfikuj kompletność silników oraz brak nietypowych wartości w widmie.", ) from exc + diagnostics = diagnostics.copy() + diagnostics["model_version"] = ( + self.metadata.model_version if self.metadata else "unversioned" + ) + diagnostics["model_artifact_sha256"] = ( + self.metadata.model_sha256 if self.metadata else "unavailable" + ) return PredictionBundle(submission, diagnostics) @@ -78,9 +81,9 @@ class PrecomputedPredictionModel: if not frame[keys].reset_index(drop=True).equals(self._submission[keys]): raise InferenceError( "Tryb awaryjny obsługuje wyłącznie dołączony zestaw demonstracyjny.", - hint="Przywróć val.csv i uruchom ponownie aplikację, aby diagnozować własne pliki.", + hint="Przywróć artefakt modelu i uruchom ponownie aplikację, aby diagnozować własne pliki.", ) - return PredictionBundle( - self._submission.copy(), - self._diagnostics.copy(), - ) + diagnostics = self._diagnostics.copy() + diagnostics["model_version"] = "demo-precomputed" + diagnostics["model_artifact_sha256"] = "unavailable" + return PredictionBundle(self._submission.copy(), diagnostics) diff --git a/final_pipeline.py b/final_pipeline.py index ce28fcd..d46e306 100644 --- a/final_pipeline.py +++ b/final_pipeline.py @@ -13,46 +13,47 @@ used because its benefit has not been established in grouped validation. from __future__ import annotations import argparse -from dataclasses import dataclass from pathlib import Path -import numpy as np import pandas as pd from sklearn.linear_model import LogisticRegression from sklearn.pipeline import Pipeline from sklearn.preprocessing import StandardScaler -from benchmark_grouped import ( - FAULT_LABELS, - FREQ_COLS, - LABELS, - NOT_APPLICABLE, - SEVERITIES, - validate_data, +from benchmark_grouped import validate_data +from engin.config import FAULT_LABELS, LABELS, NOT_APPLICABLE, SEVERITIES +from engin.features import SeverityFeatures +from engin.inference import ( + KEY_COLUMNS, + OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD, + OOD_OK_TO_UNKNOWN_THRESHOLD_MV, + DiagnosticModels, + apply_ood_override, + predict_test, + validate_inference_data, + validate_submission, ) from severity_benchmark import ( CANDIDATE_BY_ID, - SeverityFeatures, fit_candidate, - predict_candidate, prepare_labeled_frame, ) +__all__ = [ + "DiagnosticModels", + "OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD", + "OOD_OK_TO_UNKNOWN_THRESHOLD_MV", + "apply_ood_override", + "predict_test", + "run_pipeline", + "train_models", + "validate_inference_data", + "validate_submission", +] + LABEL_MODEL_NAME = "deviation_logistic_c10" LABEL_FEATURE_SET = "deviation" SEVERITY_CANDIDATE_ID = "deviation_extra_trees_mf03" -OOD_OK_TO_UNKNOWN_THRESHOLD_MV = 7.25 -OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD = 2.5 -SUBMISSION_COLUMNS = ["engine_id", "cylinder", "label", "severity"] -KEY_COLUMNS = ["engine_id", "cylinder"] - - -@dataclass -class DiagnosticModels: - label_pipeline: object - severity_candidate: object - severity_transformer: object - severity_estimator: object def make_final_label_pipeline(random_state: int = 42) -> Pipeline: @@ -75,45 +76,6 @@ def make_final_label_pipeline(random_state: int = 42) -> Pipeline: ) -def apply_ood_override( - predicted_label: np.ndarray, - deviation_features: np.ndarray, - engine_ids: np.ndarray | pd.Series, - threshold_mv: float = OOD_OK_TO_UNKNOWN_THRESHOLD_MV, - threshold_ratio: float = OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD, -) -> tuple[np.ndarray, np.ndarray, np.ndarray, np.ndarray]: - """Turn an implausibly anomalous ``ok`` into ``unknown``. - - The absolute threshold was frozen after repeated grouped validation on - clean spectra and spectra with 5% masked measurements. A second guard - requires the cylinder to be anomalous relative to the median anomaly level - of its own engine. This prevents a globally noisy unit from being treated - as a collection of isolated OOD cylinders. The rule never changes a named - fault into another class. - """ - - result = np.asarray(predicted_label, dtype=object).copy() - absolute_deviation = deviation_features[ - :, len(FREQ_COLS) : 2 * len(FREQ_COLS) - ] - anomaly_score = absolute_deviation.mean(axis=1) - ids = np.asarray(engine_ids, dtype=object) - if len(ids) != len(result): - raise ValueError("engine_ids must align with predicted labels.") - score_frame = pd.DataFrame({"engine_id": ids, "anomaly_score": anomaly_score}) - engine_median = score_frame.groupby("engine_id", sort=False)[ - "anomaly_score" - ].transform("median").to_numpy(dtype=float) - anomaly_ratio = anomaly_score / np.maximum(engine_median, 1e-6) - override = ( - (result == "ok") - & (anomaly_score > threshold_mv) - & (anomaly_ratio > threshold_ratio) - ) - result[override] = "unknown" - return result, override, anomaly_score, anomaly_ratio - - def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--val", type=Path, default=Path("val.csv")) @@ -133,35 +95,6 @@ def parse_args() -> argparse.Namespace: return parser.parse_args() -def validate_inference_data(df: pd.DataFrame) -> None: - required = {"engine_id", "cylinder", "n_cylinders", *FREQ_COLS} - missing = sorted(required.difference(df.columns)) - if missing: - raise ValueError(f"Inference data is missing columns: {missing}") - if df.empty: - raise ValueError("Inference data is empty.") - if df[["engine_id", "cylinder", "n_cylinders"]].isna().any().any(): - raise ValueError("Inference identifiers cannot contain NaN.") - if df.duplicated(KEY_COLUMNS).any(): - raise ValueError("Duplicate engine_id + cylinder keys found in inference data.") - - engine_sizes = df.groupby("engine_id").agg( - rows=("cylinder", "size"), - expected=("n_cylinders", "first"), - size_values=("n_cylinders", "nunique"), - unique_cylinders=("cylinder", "nunique"), - ) - complete = ( - (engine_sizes["rows"] == engine_sizes["expected"]) - & (engine_sizes["unique_cylinders"] == engine_sizes["expected"]) - & (engine_sizes["size_values"] == 1) - & (engine_sizes["rows"] >= 2) - ) - if not complete.all(): - bad = engine_sizes.index[~complete].tolist() - raise ValueError(f"Incomplete or inconsistent inference engines: {bad}") - - def validate_sample_keys(sample: pd.DataFrame, test: pd.DataFrame) -> None: missing = sorted(set(KEY_COLUMNS).difference(sample.columns)) if missing: @@ -210,119 +143,11 @@ def train_models( ) return DiagnosticModels( label_pipeline=label_pipeline, - severity_candidate=candidate, severity_transformer=transformer, severity_estimator=estimator, ) -def predict_test( - models: DiagnosticModels, test: pd.DataFrame -) -> tuple[pd.DataFrame, pd.DataFrame]: - validate_inference_data(test) - - raw_model_label = models.label_pipeline.predict(test).astype(object) - label_features = models.label_pipeline.named_steps["features"].transform(test) - predicted_label, ood_override, anomaly_score, anomaly_ratio = apply_ood_override( - raw_model_label, - label_features, - test["engine_id"], - ) - label_probabilities = models.label_pipeline.predict_proba(test) - raw_model_confidence = label_probabilities.max(axis=1) - label_confidence = raw_model_confidence.copy() - # OOD is a deterministic safety rule, not a probabilistic classifier. - # Emitting NaN is more honest than manufacturing a probability-like score. - label_confidence[ood_override] = np.nan - ordered_probabilities = np.sort(label_probabilities, axis=1) - label_margin = ordered_probabilities[:, -1] - ordered_probabilities[:, -2] - - test_with_labels = prepare_labeled_frame(test, predicted_label) - predicted_severity_all = predict_candidate( - models.severity_candidate, - models.severity_transformer, - models.severity_estimator, - test_with_labels, - predicted_label, - ).astype(object) - - predicted_fault = np.isin(predicted_label, FAULT_LABELS) - emitted_severity = np.full(len(test), NOT_APPLICABLE, dtype=object) - emitted_severity[predicted_fault] = predicted_severity_all[predicted_fault] - - severity_features = models.severity_transformer.transform(test_with_labels) - severity_probabilities = models.severity_estimator.predict_proba(severity_features) - severity_confidence = np.full(len(test), np.nan, dtype=float) - severity_confidence[predicted_fault] = severity_probabilities.max(axis=1)[ - predicted_fault - ] - - # The selected deviation representation starts with 21 signed deviations - # followed by their 21 absolute values. Reuse them for UI explainability. - absolute_deviation = severity_features[:, len(FREQ_COLS) : 2 * len(FREQ_COLS)] - top_frequency_indices = np.argsort(absolute_deviation, axis=1)[:, -3:][:, ::-1] - top_frequencies = [ - "|".join(str(int(index)) for index in row) - for row in top_frequency_indices - ] - - submission = test[KEY_COLUMNS].copy() - submission["label"] = predicted_label - submission["severity"] = emitted_severity - - diagnostics = submission.copy() - diagnostics["raw_model_label"] = raw_model_label - diagnostics["decision_source"] = np.where( - ood_override, "ood_override", "classifier" - ) - diagnostics["label_confidence"] = label_confidence - diagnostics["raw_model_confidence"] = raw_model_confidence - diagnostics["label_margin"] = label_margin - diagnostics["severity_confidence"] = severity_confidence - diagnostics["anomaly_score_mean_abs_mv"] = anomaly_score - diagnostics["anomaly_ratio_to_engine_median"] = anomaly_ratio - diagnostics["ood_absolute_threshold_mv"] = OOD_OK_TO_UNKNOWN_THRESHOLD_MV - diagnostics["ood_ratio_threshold"] = OOD_OK_TO_UNKNOWN_RATIO_THRESHOLD - diagnostics["top_anomalous_frequencies_khz"] = top_frequencies - diagnostics["missing_spectral_cells"] = ( - test[FREQ_COLS].isna().sum(axis=1).to_numpy(dtype=int) - ) - return submission, diagnostics - - -def validate_submission(submission: pd.DataFrame, test: pd.DataFrame) -> None: - if submission.columns.tolist() != SUBMISSION_COLUMNS: - raise ValueError( - f"Submission columns must be exactly {SUBMISSION_COLUMNS}; " - f"received={submission.columns.tolist()}" - ) - if len(submission) != len(test): - raise ValueError("Submission row count does not match test.csv.") - if submission.isna().any().any(): - raise ValueError("Submission contains NaN.") - if submission.duplicated(KEY_COLUMNS).any(): - raise ValueError("Submission contains duplicate engine_id + cylinder keys.") - if not submission[KEY_COLUMNS].reset_index(drop=True).equals( - test[KEY_COLUMNS].reset_index(drop=True) - ): - raise ValueError("Submission keys or row order do not match test.csv.") - - invalid_labels = sorted(set(submission["label"]).difference(LABELS)) - if invalid_labels: - raise ValueError(f"Submission contains invalid labels: {invalid_labels}") - - fault = submission["label"].isin(FAULT_LABELS) - invalid_fault_severity = sorted( - set(submission.loc[fault, "severity"]).difference(SEVERITIES) - ) - if invalid_fault_severity: - raise ValueError( - f"Fault predictions contain invalid severity: {invalid_fault_severity}" - ) - if not submission.loc[~fault, "severity"].eq(NOT_APPLICABLE).all(): - raise ValueError("ok/unknown predictions must use severity=nie_dotyczy.") - - def run_pipeline( val: pd.DataFrame, test: pd.DataFrame, diff --git a/requirements-lock.txt b/requirements-lock.txt index 21cefe7..6f0271f 100644 --- a/requirements-lock.txt +++ b/requirements-lock.txt @@ -1,6 +1,50 @@ +altair==6.2.2 +anyio==4.14.2 +attrs==26.1.0 +blinker==1.9.0 +certifi==2026.7.22 +charset-normalizer==3.5.1 +click==8.4.2 +contourpy==1.3.3 +cycler==0.12.1 +fonttools==4.63.0 +h11==0.16.0 +httptools==0.8.0 +idna==3.19 +itsdangerous==2.2.0 +Jinja2==3.1.6 +joblib==1.5.3 +jsonschema==4.26.0 +jsonschema-specifications==2025.9.1 +kiwisolver==1.5.0 +MarkupSafe==3.0.3 matplotlib==3.10.8 +narwhals==2.25.0 numpy==2.3.5 +packaging==26.3 pandas==2.2.3 +pillow==12.3.0 plotly==6.9.0 +protobuf==7.36.0 +pyarrow==25.0.1 +pydeck==0.9.3 +pyparsing==3.3.2 +python-dateutil==2.9.0.post0 +python-multipart==0.0.32 +pytz==2026.3.post1 +referencing==0.37.0 +requests==2.34.2 +rpds-py==2026.6.3 scikit-learn==1.8.0 +scipy==1.18.1 +six==1.17.0 +starlette==1.6.0 streamlit==1.62.0 +threadpoolctl==3.6.0 +toml==0.10.2 +typing_extensions==4.16.0 +tzdata==2026.3 +urllib3==2.7.0 +uvicorn==0.52.4 +watchdog==6.0.0 +websockets==16.1.1 diff --git a/scripts/__init__.py b/scripts/__init__.py new file mode 100644 index 0000000..e88d6f3 --- /dev/null +++ b/scripts/__init__.py @@ -0,0 +1 @@ +"""Operational scripts for building and verifying ENGIN releases.""" diff --git a/scripts/build_model_artifact.py b/scripts/build_model_artifact.py new file mode 100644 index 0000000..b95ba03 --- /dev/null +++ b/scripts/build_model_artifact.py @@ -0,0 +1,130 @@ +"""Build or verify the versioned ENGIN production model artifact.""" + +from __future__ import annotations + +import argparse +from pathlib import Path + +import pandas as pd +from pandas.testing import assert_frame_equal + +from engin.artifact import ( + LoadedModelArtifact, + load_model_artifact, + save_model_artifact, + sha256_file, +) +from engin.inference import predict_test, validate_submission +from final_pipeline import ( + LABEL_MODEL_NAME, + SEVERITY_CANDIDATE_ID, + train_models, +) + +DEFAULT_MODEL_VERSION = "engin-2026.08.25-1" + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--val", type=Path, default=Path("val.csv")) + parser.add_argument("--test", type=Path, default=Path("test.csv")) + parser.add_argument( + "--expected-submission", type=Path, default=Path("predictions.csv") + ) + parser.add_argument( + "--expected-diagnostics", + type=Path, + default=Path("prediction_diagnostics.csv"), + ) + parser.add_argument("--model-version", default=DEFAULT_MODEL_VERSION) + parser.add_argument("--source-revision", default="unknown") + parser.add_argument("--random-state", type=int, default=42) + parser.add_argument("--n-jobs", type=int, default=-1) + parser.add_argument("--artifact-dir", type=Path) + parser.add_argument( + "--verify-only", + action="store_true", + help="Load and verify the existing artifact without retraining.", + ) + return parser.parse_args() + + +def _artifact_dir(args: argparse.Namespace) -> Path: + return args.artifact_dir or Path("artifacts") / args.model_version + + +def build_artifact(args: argparse.Namespace) -> LoadedModelArtifact: + artifact_dir = _artifact_dir(args) + if not args.verify_only: + val = pd.read_csv(args.val).reset_index(drop=True) + models = train_models( + val, + random_state=args.random_state, + n_jobs=args.n_jobs, + ) + save_model_artifact( + models, + artifact_dir, + model_version=args.model_version, + source_revision=args.source_revision, + training_data_sha256=sha256_file(args.val), + random_state=args.random_state, + n_jobs=args.n_jobs, + label_model=LABEL_MODEL_NAME, + severity_model=SEVERITY_CANDIDATE_ID, + ) + return load_model_artifact( + artifact_dir, + expected_model_version=args.model_version, + ) + + +def verify_predictions( + artifact: LoadedModelArtifact, + *, + test_path: Path, + expected_submission_path: Path, + expected_diagnostics_path: Path, +) -> None: + test = pd.read_csv(test_path).reset_index(drop=True) + submission, diagnostics = predict_test(artifact.models, test) + validate_submission(submission, test) + + expected_submission = pd.read_csv(expected_submission_path).reset_index(drop=True) + expected_diagnostics = pd.read_csv(expected_diagnostics_path).reset_index(drop=True) + assert_frame_equal( + submission, + expected_submission, + check_dtype=False, + check_exact=True, + ) + assert_frame_equal( + diagnostics, + expected_diagnostics, + check_dtype=False, + check_exact=False, + rtol=1e-12, + atol=1e-12, + ) + + +def main() -> None: + args = parse_args() + artifact = build_artifact(args) + verify_predictions( + artifact, + test_path=args.test, + expected_submission_path=args.expected_submission, + expected_diagnostics_path=args.expected_diagnostics, + ) + metadata = artifact.metadata + print(f"Artifact verified: {_artifact_dir(args).resolve()}") + print(f"Model version: {metadata.model_version}") + print(f"Model SHA-256: {metadata.model_sha256}") + print(f"Training data SHA-256: {metadata.training_data_sha256}") + print(f"Python: {metadata.python_version} | inference n_jobs: {metadata.n_jobs}") + print("Predictions: exact 600-row regression match") + + +if __name__ == "__main__": + main() diff --git a/scripts/verify.sh b/scripts/verify.sh index 387aa23..d4648bb 100755 --- a/scripts/verify.sh +++ b/scripts/verify.sh @@ -1,5 +1,7 @@ #!/usr/bin/env sh set -eu -python -m py_compile app.py engin/*.py tests/*.py +python -m compileall -q app.py engin tests final_pipeline.py scripts +python -m pip check python -m unittest discover -v +python -m scripts.build_model_artifact --verify-only diff --git a/severity_benchmark.py b/severity_benchmark.py index 330e526..e6ba140 100644 --- a/severity_benchmark.py +++ b/severity_benchmark.py @@ -32,6 +32,7 @@ from sklearn.pipeline import make_pipeline as sklearn_pipeline from sklearn.preprocessing import StandardScaler from sklearn.svm import SVC +from engin.features import FEATURE_SETS, SeverityFeatures from benchmark_grouped import ( FAULT_LABELS, FREQ_COLS, @@ -50,7 +51,6 @@ from robustness_grouped import mask_spectral_cells, missing_scenario_name DEFAULT_SEEDS = [7, 21, 42, 77, 123] SEVERITY_TO_INT = {"male": 0, "srednie": 1, "duze": 2} INT_TO_SEVERITY = np.asarray(SEVERITIES, dtype=object) -FEATURE_SETS = ("raw", "deviation", "all") @dataclass(frozen=True) @@ -91,146 +91,6 @@ CANDIDATES = [ CANDIDATE_BY_ID = {candidate.candidate_id: candidate for candidate in CANDIDATES} -class SeverityFeatures: - """Create severity-oriented features without using validation labels. - - Engine-relative values use a leave-one-cylinder-out median reference. Input - to ``transform`` must therefore contain every cylinder of each engine. - ``diagnostic_label`` is optional and may contain only the model-predicted - label at validation or test time. - """ - - def __init__(self, feature_set: str, include_fault_type: bool = False) -> None: - if feature_set not in FEATURE_SETS: - raise ValueError(f"Unknown feature_set={feature_set!r}") - self.feature_set = feature_set - self.include_fault_type = include_fault_type - - def fit( - self, X: pd.DataFrame, y: Iterable[str] | None = None - ) -> "SeverityFeatures": - self._validate(X) - spectra = X[FREQ_COLS].apply(pd.to_numeric, errors="coerce") - fallback = spectra.median(axis=0).to_numpy(dtype=float) - if np.isnan(fallback).any(): - raise ValueError("A frequency column is entirely NaN in training.") - self.fallback_medians_ = fallback - return self - - def transform(self, X: pd.DataFrame) -> np.ndarray: - self._validate(X) - if not hasattr(self, "fallback_medians_"): - raise RuntimeError("SeverityFeatures must be fitted before transform().") - raw = self._clean_spectra(X) - reference = self._leave_one_out_engine_median(raw, X["engine_id"].to_numpy()) - relative = raw - reference - absolute_relative = np.abs(relative) - ratio_delta = raw / np.maximum(np.abs(reference), 1e-6) - 1.0 - - if self.feature_set == "raw": - features = raw - elif self.feature_set == "deviation": - features = np.hstack( - [relative, absolute_relative, ratio_delta, self._summary_features(raw, relative, ratio_delta)] - ) - else: - features = np.hstack( - [raw, relative, absolute_relative, ratio_delta, self._summary_features(raw, relative, ratio_delta)] - ) - - if self.include_fault_type: - if "diagnostic_label" not in X.columns: - raise ValueError("diagnostic_label is required when include_fault_type=True") - labels = X["diagnostic_label"].to_numpy(dtype=object) - one_hot = np.column_stack([labels == label for label in FAULT_LABELS]).astype(float) - features = np.hstack([features, one_hot]) - if np.isnan(features).any() or np.isinf(features).any(): - raise ValueError("Severity features contain NaN or infinity.") - return features - - def fit_transform( - self, X: pd.DataFrame, y: Iterable[str] | None = None - ) -> np.ndarray: - return self.fit(X).transform(X) - - def _clean_spectra(self, X: pd.DataFrame) -> np.ndarray: - spectra = X[FREQ_COLS].apply(pd.to_numeric, errors="coerce") - spectra = spectra.interpolate(axis=1, limit_direction="both") - spectra = spectra.fillna(pd.Series(self.fallback_medians_, index=FREQ_COLS)) - return spectra.to_numpy(dtype=float) - - @staticmethod - def _leave_one_out_engine_median(raw: np.ndarray, groups: np.ndarray) -> np.ndarray: - reference = np.empty_like(raw) - for engine_id in pd.unique(groups): - positions = np.flatnonzero(groups == engine_id) - if len(positions) < 2: - raise ValueError( - f"Engine {engine_id!r} has fewer than two cylinders in transform()." - ) - engine_values = raw[positions] - for local_position, global_position in enumerate(positions): - keep = np.arange(len(positions)) != local_position - reference[global_position] = np.median(engine_values[keep], axis=0) - return reference - - @staticmethod - def _summary_features( - raw: np.ndarray, relative: np.ndarray, ratio_delta: np.ndarray - ) -> np.ndarray: - absolute_relative = np.abs(relative) - gradients = np.diff(raw, axis=1) - frequency = np.arange(raw.shape[1], dtype=float) - centered_frequency = frequency - frequency.mean() - slope = raw @ centered_frequency / np.sum(centered_frequency**2) - # ``np.trapezoid`` was added after NumPy 1.26, which is still allowed by - # requirements.txt. Keep the fresh-install path compatible with both - # NumPy 1.x and 2.x. - if hasattr(np, "trapezoid"): - spectral_auc = np.trapezoid(raw, axis=1) - else: # pragma: no cover - exercised only with NumPy 1.x - spectral_auc = np.trapz(raw, axis=1) - - columns = [ - raw.mean(axis=1), - raw.std(axis=1), - raw.min(axis=1), - raw.max(axis=1), - np.ptp(raw, axis=1), - spectral_auc, - raw.argmax(axis=1) / 20.0, - raw.argmin(axis=1) / 20.0, - slope, - np.abs(gradients).mean(axis=1), - np.abs(gradients).max(axis=1), - relative.mean(axis=1), - relative.std(axis=1), - relative.min(axis=1), - relative.max(axis=1), - absolute_relative.mean(axis=1), - absolute_relative.max(axis=1), - np.sqrt(np.mean(relative**2, axis=1)), - ratio_delta.mean(axis=1), - ratio_delta.std(axis=1), - ] - for start, stop in ((0, 5), (5, 10), (10, 15), (15, 21)): - columns.extend( - [ - raw[:, start:stop].mean(axis=1), - relative[:, start:stop].mean(axis=1), - absolute_relative[:, start:stop].mean(axis=1), - ] - ) - return np.column_stack(columns) - - @staticmethod - def _validate(X: pd.DataFrame) -> None: - required = {"engine_id", *FREQ_COLS} - missing = sorted(required.difference(X.columns)) - if missing: - raise ValueError(f"Missing columns: {missing}") - - class OrdinalRidge(BaseEstimator): """Weighted ridge regression rounded to the three ordered severities.""" diff --git a/tests/test_model_artifact.py b/tests/test_model_artifact.py new file mode 100644 index 0000000..db31ebb --- /dev/null +++ b/tests/test_model_artifact.py @@ -0,0 +1,114 @@ +from __future__ import annotations + +import json +import tempfile +import unittest +from pathlib import Path +from unittest.mock import patch + +import pandas as pd +from pandas.testing import assert_frame_equal + +import app +from engin.artifact import ( + MANIFEST_FILENAME, + MODEL_FILENAME, + ModelArtifactError, + load_model_artifact, + save_model_artifact, +) +from engin.inference import DiagnosticModels +from engin.model import SklearnPredictionModel + +ROOT = Path(__file__).resolve().parents[1] +MODEL_VERSION = "engin-2026.08.25-1" +ARTIFACT_DIR = ROOT / "artifacts" / MODEL_VERSION + + +def _save_fake_artifact(directory: Path) -> None: + save_model_artifact( + DiagnosticModels("label", "transformer", "estimator"), + directory, + model_version="test-v1", + source_revision="unit-test", + training_data_sha256="a" * 64, + random_state=42, + n_jobs=1, + label_model="fake-label", + severity_model="fake-severity", + ) + + +class ModelArtifactContractTests(unittest.TestCase): + def test_round_trip_preserves_metadata_and_models(self) -> None: + with tempfile.TemporaryDirectory() as temporary_directory: + directory = Path(temporary_directory) + _save_fake_artifact(directory) + loaded = load_model_artifact( + directory, expected_model_version="test-v1" + ) + self.assertEqual(loaded.metadata.model_version, "test-v1") + self.assertEqual(loaded.models.label_pipeline, "label") + + def test_corrupt_model_is_rejected_before_deserialization(self) -> None: + with tempfile.TemporaryDirectory() as temporary_directory: + directory = Path(temporary_directory) + _save_fake_artifact(directory) + with (directory / MODEL_FILENAME).open("ab") as handle: + handle.write(b"corruption") + with self.assertRaisesRegex(ModelArtifactError, "checksum"): + load_model_artifact(directory) + + def test_runtime_version_mismatch_is_rejected(self) -> None: + with tempfile.TemporaryDirectory() as temporary_directory: + directory = Path(temporary_directory) + _save_fake_artifact(directory) + manifest_path = directory / MANIFEST_FILENAME + manifest = json.loads(manifest_path.read_text(encoding="utf-8")) + manifest["runtime_versions"]["scikit-learn"] = "0.0-invalid" + manifest_path.write_text( + json.dumps(manifest), encoding="utf-8" + ) + with self.assertRaisesRegex(ModelArtifactError, "runtime versions"): + load_model_artifact(directory) + + +class ShippedArtifactTests(unittest.TestCase): + def test_artifact_reproduces_frozen_submission(self) -> None: + model = SklearnPredictionModel.from_artifact( + ARTIFACT_DIR, + expected_model_version=MODEL_VERSION, + ) + test = pd.read_csv(ROOT / "test.csv").reset_index(drop=True) + expected = pd.read_csv(ROOT / "predictions.csv").reset_index(drop=True) + actual = model.predict(test).submission + assert_frame_equal(actual, expected, check_dtype=False, check_exact=True) + + def test_diagnostics_identify_model_release(self) -> None: + model = SklearnPredictionModel.from_artifact( + ARTIFACT_DIR, + expected_model_version=MODEL_VERSION, + ) + test = pd.read_csv(ROOT / "test.csv").reset_index(drop=True) + diagnostics = model.predict(test).diagnostics + self.assertEqual(set(diagnostics["model_version"]), {MODEL_VERSION}) + self.assertEqual( + set(diagnostics["model_artifact_sha256"]), + {model.metadata.model_sha256}, + ) + + def test_application_startup_does_not_read_training_data(self) -> None: + app.build_dependencies.clear() + with patch.object( + app.PandasCsvReader, + "read_path", + side_effect=AssertionError("runtime attempted to read training data"), + ): + dependencies = app.build_dependencies() + app.build_dependencies.clear() + self.assertTrue(dependencies.live_inference) + self.assertEqual(dependencies.model_version, MODEL_VERSION) + + +if __name__ == "__main__": + unittest.main(verbosity=2)