"""Train the selected ENGIN models and create a validated submission. Selected, leakage-safe architecture: * label: signed/absolute leave-one-cylinder-out deviations, ratio deltas and summary statistics with balanced Logistic Regression, * severity: leave-one-cylinder-out deviation features with Extra Trees. The script trains only on labeled ``val.csv``. The unlabeled archive is not 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 severity_benchmark import ( CANDIDATE_BY_ID, SeverityFeatures, fit_candidate, predict_candidate, prepare_labeled_frame, ) 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: """Build the multi-seed winner used for final label predictions.""" return Pipeline( [ ("features", SeverityFeatures(LABEL_FEATURE_SET)), ("scale", StandardScaler()), ( "model", LogisticRegression( C=10.0, class_weight="balanced", max_iter=5_000, random_state=random_state, ), ), ] ) 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")) parser.add_argument("--test", type=Path, default=Path("test.csv")) parser.add_argument( "--sample-submit", type=Path, default=Path("sample_submit.csv") ) parser.add_argument("--output", type=Path, default=Path("predictions.csv")) parser.add_argument( "--diagnostics", type=Path, default=Path("prediction_diagnostics.csv"), help="Auxiliary model scores and anomaly data for the application.", ) parser.add_argument("--random-state", type=int, default=42) parser.add_argument("--n-jobs", type=int, default=-1) 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: raise ValueError(f"Sample submission is missing key columns: {missing}") if sample.duplicated(KEY_COLUMNS).any(): raise ValueError("Sample submission contains duplicate keys.") sample_keys = set(map(tuple, sample[KEY_COLUMNS].to_numpy())) test_keys = set(map(tuple, test[KEY_COLUMNS].to_numpy())) if sample_keys != test_keys or len(sample) != len(test): raise ValueError("sample_submit.csv keys do not match test.csv.") def train_models( val: pd.DataFrame, random_state: int = 42, n_jobs: int = -1, ) -> DiagnosticModels: validate_data(val, n_splits=2) y_label = val["label"].reset_index(drop=True) y_severity = val["severity"].reset_index(drop=True) label_pipeline = make_final_label_pipeline(random_state=random_state).fit( val, y_label ) fault_mask = y_label.isin(FAULT_LABELS).to_numpy() observed_severities = set(y_severity.to_numpy()[fault_mask]) if observed_severities != set(SEVERITIES): raise ValueError( "Training faults must contain every allowed severity; " f"observed={sorted(observed_severities)}" ) candidate = CANDIDATE_BY_ID[SEVERITY_CANDIDATE_ID] val_with_labels = prepare_labeled_frame( val, y_label.to_numpy(dtype=object) ) transformer, estimator = fit_candidate( candidate=candidate, X_train_full=val_with_labels, train_fault_mask=fault_mask, y_severity_fault=y_severity.to_numpy(dtype=object)[fault_mask], y_fault_label=y_label.to_numpy(dtype=object)[fault_mask], random_state=random_state + 1_002, n_jobs=n_jobs, ) 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, sample_submit: pd.DataFrame | None = None, random_state: int = 42, n_jobs: int = -1, ) -> tuple[pd.DataFrame, pd.DataFrame]: if sample_submit is not None: validate_sample_keys(sample_submit, test) models = train_models(val, random_state=random_state, n_jobs=n_jobs) submission, diagnostics = predict_test(models, test) validate_submission(submission, test) return submission, diagnostics def main() -> None: args = parse_args() val = pd.read_csv(args.val).reset_index(drop=True) test = pd.read_csv(args.test).reset_index(drop=True) sample_submit = pd.read_csv(args.sample_submit).reset_index(drop=True) submission, diagnostics = run_pipeline( val=val, test=test, sample_submit=sample_submit, random_state=args.random_state, n_jobs=args.n_jobs, ) args.output.parent.mkdir(parents=True, exist_ok=True) args.diagnostics.parent.mkdir(parents=True, exist_ok=True) submission.to_csv(args.output, index=False) diagnostics.to_csv(args.diagnostics, index=False) print(f"Saved validated submission: {args.output.resolve()}") print(f"Saved application diagnostics: {args.diagnostics.resolve()}") print(f"Rows: {len(submission)} | engines: {submission['engine_id'].nunique()}") print("\nLabel distribution:") print(submission["label"].value_counts().reindex(LABELS, fill_value=0).to_string()) print("\nSeverity distribution:") print( submission["severity"] .value_counts() .reindex([NOT_APPLICABLE, *SEVERITIES], fill_value=0) .to_string() ) if __name__ == "__main__": main()