diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml new file mode 100644 index 0000000..218c691 --- /dev/null +++ b/.github/workflows/tests.yml @@ -0,0 +1,16 @@ +name: Tests +on: + push: + pull_request: + +jobs: + test: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: '3.12' + cache: pip + - run: python -m pip install -r requirements-dev.txt + - run: python -m pytest -q diff --git a/README.md b/README.md index ae40f0f..22f0f5f 100644 --- a/README.md +++ b/README.md @@ -367,8 +367,8 @@ sentinel-ml/ Clone the repository: ```bash -git clone -cd sentinel-ml +git clone https://github.com/Ayushdevo/SentinelML.git +cd SentinelML ``` Create a virtual environment: @@ -483,9 +483,17 @@ models/model_health_report.json Instead of executing each monitoring component manually: ```bash -python src/run_pipeline.py +python -m src.run_pipeline ``` +Run from the repository root. The pipeline generates both sample datasets when +they are absent and trains the baseline model when its artifact is absent. +Existing datasets and model artifacts are reused. To retrain after changing +the reference data, run `python -m src.train` first. + +Run regression checks with `python -m pip install -r requirements-dev.txt` +and `python -m pytest -q`. + Pipeline: ```text @@ -632,4 +640,4 @@ Portfolio: **www.ayushtiwari.tech** ## ⭐ SentinelML -> Detect drift. Identify risky predictions. Understand why. Know when to retrain. \ No newline at end of file +> Detect drift. Identify risky predictions. Understand why. Know when to retrain. diff --git a/api/main.py b/api/main.py index ad44dcc..8c738ef 100644 --- a/api/main.py +++ b/api/main.py @@ -1,10 +1,11 @@ from pathlib import Path import json +from functools import lru_cache import joblib import pandas as pd -from fastapi import FastAPI -from pydantic import BaseModel +from fastapi import FastAPI, HTTPException +from pydantic import BaseModel, ConfigDict MODEL_PATH = Path("models/baseline_model.joblib") @@ -16,10 +17,15 @@ version="1.0.0", ) -model = joblib.load(MODEL_PATH) +@lru_cache(maxsize=1) +def load_model(): + if not MODEL_PATH.exists(): + raise HTTPException(status_code=503, detail="Baseline model is not available") + return joblib.load(MODEL_PATH) class PredictionRequest(BaseModel): + model_config = ConfigDict(extra="forbid", allow_inf_nan=False) tenure: float monthly_charges: float usage_hours: float @@ -40,6 +46,8 @@ def root(): @app.get("/health") def model_health(): + if not HEALTH_PATH.exists(): + raise HTTPException(status_code=503, detail="Model health report is not available") with open( HEALTH_PATH, "r", @@ -50,8 +58,17 @@ def model_health(): return health +@app.get("/ready") +def readiness(): + missing = [str(path) for path in (MODEL_PATH, HEALTH_PATH) if not path.exists()] + if missing: + raise HTTPException(status_code=503, detail={"missing_artifacts": missing}) + return {"status": "ready"} + + @app.post("/predict") def predict(data: PredictionRequest): + model = load_model() input_df = pd.DataFrame( [ @@ -67,6 +84,8 @@ def predict(data: PredictionRequest): } ] ) + if hasattr(model, "feature_names_in_"): + input_df = input_df[list(model.feature_names_in_)] probability = float( model.predict_proba(input_df)[0, 1] @@ -96,4 +115,4 @@ def predict(data: PredictionRequest): 4, ), "prediction_risk": risk, - } \ No newline at end of file + } diff --git a/dashboard/app.py b/dashboard/app.py index 01076d3..a224429 100644 --- a/dashboard/app.py +++ b/dashboard/app.py @@ -24,6 +24,13 @@ def load_json(path): @st.cache_data def load_data(): + required = (HEALTH_PATH, DRIFT_PATH, ROOT_CAUSE_PATH, PRODUCTION_PATH) + missing = [str(path) for path in required if not path.exists()] + if missing: + raise FileNotFoundError( + "Generate monitoring reports with `python -m src.run_pipeline` first. " + f"Missing: {', '.join(missing)}" + ) health = load_json(HEALTH_PATH) drift = load_json(DRIFT_PATH) root_cause = load_json(ROOT_CAUSE_PATH) @@ -32,7 +39,11 @@ def load_data(): return health, drift, root_cause, production -health, drift, root_cause, production = load_data() +try: + health, drift, root_cause, production = load_data() +except (FileNotFoundError, ValueError) as exc: + st.error(str(exc)) + st.stop() st.title("🛡️ SentinelML") @@ -116,6 +127,8 @@ def load_data(): "PSI": values["psi"], "KS Statistic": values["ks_statistic"], "p-value": values["p_value"], + "Reference missing": values.get("reference_missing_rate", 0), + "Production missing": values.get("production_missing_rate", 0), "Status": values["status"], } ) @@ -298,4 +311,4 @@ def load_data(): These signals are combined into a transparent model-health score and a rules-based retraining recommendation. """ - ) \ No newline at end of file + ) diff --git a/requirements-dev.txt b/requirements-dev.txt new file mode 100644 index 0000000..7af6bb0 --- /dev/null +++ b/requirements-dev.txt @@ -0,0 +1,3 @@ +-r requirements.txt +pytest +httpx diff --git a/requirements.txt b/requirements.txt index 1b930bd..8e02c5a 100644 --- a/requirements.txt +++ b/requirements.txt @@ -6,4 +6,6 @@ scipy shap joblib streamlit -matplotlib \ No newline at end of file +matplotlib +fastapi +uvicorn diff --git a/src/anomaly_detector.py b/src/anomaly_detector.py index 8c145c4..f57f775 100644 --- a/src/anomaly_detector.py +++ b/src/anomaly_detector.py @@ -18,14 +18,25 @@ def analyze_production(): reference = pd.read_csv(REFERENCE_PATH) production = pd.read_csv(PRODUCTION_PATH) + model = joblib.load(MODEL_PATH) - features = [ - column for column in reference.columns - if column != TARGET + features = list(model.feature_names_in_) if hasattr(model, "feature_names_in_") else [ + column for column in reference.columns if column != TARGET ] + for name, frame in (("Reference", reference), ("Production", production)): + missing = sorted(set(features) - set(frame.columns)) + if missing: + raise ValueError(f"{name} data is missing model features: {', '.join(missing)}") + if frame.empty: + raise ValueError(f"{name} data has no rows") X_reference = reference[features] X_production = production[features] + for name, frame in (("Reference", X_reference), ("Production", X_production)): + if not all(pd.api.types.is_numeric_dtype(dtype) for dtype in frame.dtypes): + raise ValueError(f"{name} features must be numeric") + if not np.isfinite(frame.to_numpy(dtype=float)).all(): + raise ValueError(f"{name} features contain missing or infinite values") # ---------------------------- # 1. Scale features @@ -65,8 +76,6 @@ def analyze_production(): # ---------------------------- # 3. Model predictions # ---------------------------- - model = joblib.load(MODEL_PATH) - probabilities = model.predict_proba( X_production )[:, 1] @@ -187,4 +196,4 @@ def analyze_production(): if __name__ == "__main__": - analyze_production() \ No newline at end of file + analyze_production() diff --git a/src/drift_detector.py b/src/drift_detector.py index 382899b..e7c86ef 100644 --- a/src/drift_detector.py +++ b/src/drift_detector.py @@ -22,8 +22,16 @@ def calculate_psi(expected, actual, bins=10): PSI >= 0.25 -> Significant drift """ - expected = np.asarray(expected) - actual = np.asarray(actual) + expected = np.asarray(expected, dtype=float) + actual = np.asarray(actual, dtype=float) + if not len(expected) or not len(actual): + raise ValueError("PSI requires nonempty reference and production samples") + if not np.isfinite(expected).all() or not np.isfinite(actual).all(): + raise ValueError("PSI requires finite numeric values") + if bins < 2: + raise ValueError("PSI requires at least two bins") + if np.min(expected) == np.max(expected) == np.min(actual) == np.max(actual): + return 0.0 # Quantile-based bins from reference distribution breakpoints = np.unique( @@ -40,6 +48,8 @@ def calculate_psi(expected, actual, bins=10): max(expected.max(), actual.max()), bins + 1, ) + if breakpoints[0] == breakpoints[-1]: + breakpoints = np.array([breakpoints[0] - 0.5, breakpoints[0] + 0.5]) # Make sure all production values are captured breakpoints[0] = -np.inf @@ -61,6 +71,8 @@ def calculate_psi(expected, actual, bins=10): # Prevent division by zero expected_pct = np.clip(expected_pct, 1e-6, None) actual_pct = np.clip(actual_pct, 1e-6, None) + expected_pct /= expected_pct.sum() + actual_pct /= actual_pct.sum() psi = np.sum( (actual_pct - expected_pct) @@ -90,12 +102,19 @@ def classify_drift(psi, p_value): def detect_drift(): reference = pd.read_csv(REFERENCE_PATH) production = pd.read_csv(PRODUCTION_PATH) + if TARGET not in reference.columns: + raise ValueError(f"Reference data is missing target column: {TARGET}") features = [ column for column in reference.columns if column != TARGET ] + if not features: + raise ValueError("Reference data has no feature columns") + missing = sorted(set(features) - set(production.columns)) + if missing: + raise ValueError(f"Production data is missing features: {', '.join(missing)}") report = {} @@ -116,6 +135,8 @@ def detect_drift(): ref_values = reference[feature].dropna() prod_values = production[feature].dropna() + if ref_values.empty or prod_values.empty: + raise ValueError(f"Feature {feature!r} has no usable samples") psi = calculate_psi( ref_values, @@ -131,8 +152,18 @@ def detect_drift(): psi, p_value, ) + missing_shift = abs( + reference[feature].isna().mean() - production[feature].isna().mean() + ) + if missing_shift >= 0.10 and status in {"STABLE", "LOW"}: + status = "MODERATE" report[feature] = { + "reference_count": int(len(ref_values)), + "production_count": int(len(prod_values)), + "reference_missing_rate": round(float(reference[feature].isna().mean()), 6), + "production_missing_rate": round(float(production[feature].isna().mean()), 6), + "missing_rate_shift": round(float(missing_shift), 6), "psi": round(float(psi), 6), "ks_statistic": round( float(ks_statistic), 6 @@ -197,4 +228,4 @@ def detect_drift(): if __name__ == "__main__": - detect_drift() \ No newline at end of file + detect_drift() diff --git a/src/generate_data.py b/src/generate_data.py index c715cd6..bfa4917 100644 --- a/src/generate_data.py +++ b/src/generate_data.py @@ -9,6 +9,8 @@ def build_dataset(n_samples=10000): + if not isinstance(n_samples, int) or n_samples < 10: + raise ValueError("n_samples must be an integer of at least 10") X, y = make_classification( n_samples=n_samples, n_features=8, @@ -77,6 +79,8 @@ def create_production_data(reference): behaviour has changed. """ + if reference.empty: + raise ValueError("Cannot simulate production data from an empty reference frame") production = reference.sample( n=2500, replace=True, @@ -126,4 +130,4 @@ def create_production_data(reference): print("SentinelML datasets generated.") print(f"Reference samples: {len(reference)}") print(f"Production samples: {len(production)}") - print(f"Churn rate: {reference['churn'].mean():.2%}") \ No newline at end of file + print(f"Churn rate: {reference['churn'].mean():.2%}") diff --git a/src/retraining_engine.py b/src/retraining_engine.py index 847504e..ca209fe 100644 --- a/src/retraining_engine.py +++ b/src/retraining_engine.py @@ -22,6 +22,12 @@ def calculate_health(): production = pd.read_csv(PRODUCTION_PATH) total = len(production) + if total == 0: + raise ValueError("Cannot calculate model health without production observations") + required = {"is_anomaly", "high_uncertainty", "prediction_risk"} + missing = required - set(production.columns) + if missing: + raise ValueError(f"Production analysis is missing columns: {sorted(missing)}") anomaly_rate = ( production["is_anomaly"].mean() @@ -163,6 +169,11 @@ def calculate_health(): # --------------------------------- report = { + "production_observations": total, + "scoring_note": ( + "Heuristic monitoring score; retraining requires review and labeled validation, " + "not automatic deployment." + ), "health_score": round( health_score, 2, @@ -316,4 +327,4 @@ def calculate_health(): if __name__ == "__main__": - calculate_health() \ No newline at end of file + calculate_health() diff --git a/src/root_cause.py b/src/root_cause.py index 58fc02a..102871e 100644 --- a/src/root_cause.py +++ b/src/root_cause.py @@ -38,6 +38,11 @@ def analyze_root_causes(): for column in production.columns if column not in NON_FEATURE_COLUMNS ] + if hasattr(model, "feature_names_in_"): + features = list(model.feature_names_in_) + missing = sorted(set(features) - set(production.columns)) + if missing or production.empty: + raise ValueError(f"Production analysis is empty or missing features: {missing}") X = production[features] @@ -54,6 +59,13 @@ def analyze_root_causes(): shap_values = shap_values[-1] shap_values = np.asarray(shap_values) + if shap_values.ndim == 3: + # Multiclass explainers return (rows, features, classes). + shap_values = shap_values[:, :, 1] + if shap_values.shape != (len(X), len(features)): + raise ValueError( + f"Unexpected SHAP shape {shap_values.shape}; expected {X.shape}" + ) # -------------------------------- # 2. Global feature importance @@ -220,4 +232,4 @@ def analyze_root_causes(): if __name__ == "__main__": - analyze_root_causes() \ No newline at end of file + analyze_root_causes() diff --git a/src/run_pipeline.py b/src/run_pipeline.py index 7851124..701335f 100644 --- a/src/run_pipeline.py +++ b/src/run_pipeline.py @@ -1,8 +1,9 @@ -from generate_data import build_dataset, create_production_data -from drift_detector import detect_drift -from anomaly_detector import analyze_production -from root_cause import analyze_root_causes -from retraining_engine import calculate_health +from src.generate_data import build_dataset, create_production_data +from src.drift_detector import detect_drift +from src.anomaly_detector import analyze_production +from src.root_cause import analyze_root_causes +from src.retraining_engine import calculate_health +from src.train import train import pandas as pd from pathlib import Path @@ -17,8 +18,15 @@ def main(): # Only regenerate data if files do not exist reference_path = Path("data/reference.csv") production_path = Path("data/production.csv") + model_path = Path("models/baseline_model.joblib") - if not reference_path.exists() or not production_path.exists(): + if reference_path.exists() != production_path.exists(): + raise FileNotFoundError( + "Only one dataset exists; restore its matching dataset or remove both to regenerate" + ) + + if not reference_path.exists(): + reference_path.parent.mkdir(parents=True, exist_ok=True) dataset = build_dataset() @@ -40,6 +48,10 @@ def main(): print("Datasets generated.") + if not model_path.exists(): + print("\n[setup] Training baseline model...") + train() + print("\n[1/4] Detecting drift...") detect_drift() @@ -58,4 +70,4 @@ def main(): if __name__ == "__main__": - main() \ No newline at end of file + main() diff --git a/src/train.py b/src/train.py index 88e2550..cdd9230 100644 --- a/src/train.py +++ b/src/train.py @@ -23,6 +23,14 @@ def train(): df = pd.read_csv(DATA_PATH) + if TARGET not in df: + raise ValueError(f"Training data is missing target column: {TARGET}") + if df[TARGET].isna().any() or set(df[TARGET].unique()) != {0, 1}: + raise ValueError("Training target must contain both binary classes without missing values") + if df.empty or df.isna().any().any(): + raise ValueError("Training features must contain nonempty, complete data") + if df[TARGET].value_counts().min() < 2 or len(df) < 10: + raise ValueError("Training requires at least 10 rows and two examples per class") X = df.drop(columns=[TARGET]) y = df[TARGET] @@ -88,4 +96,4 @@ def train(): if __name__ == "__main__": - train() \ No newline at end of file + train() diff --git a/tests/test_api.py b/tests/test_api.py new file mode 100644 index 0000000..47786eb --- /dev/null +++ b/tests/test_api.py @@ -0,0 +1,27 @@ +from fastapi.testclient import TestClient + +from api import main + + +def test_missing_model_returns_503_without_preventing_startup(tmp_path, monkeypatch): + monkeypatch.setattr(main, "MODEL_PATH", tmp_path / "missing.joblib") + main.load_model.cache_clear() + client = TestClient(main.app) + assert client.get("/").status_code == 200 + payload = dict.fromkeys(main.PredictionRequest.model_fields, 1.0) + response = client.post("/predict", json=payload) + assert response.status_code == 503 + + +def test_invalid_payload_is_rejected_before_model_load(tmp_path, monkeypatch): + monkeypatch.setattr(main, "MODEL_PATH", tmp_path / "missing.joblib") + main.load_model.cache_clear() + client = TestClient(main.app) + payload = dict.fromkeys(main.PredictionRequest.model_fields, 1.0) + payload["extra_feature"] = 2 + assert client.post("/predict", json=payload).status_code == 422 + + +def test_missing_health_report_returns_503(tmp_path, monkeypatch): + monkeypatch.setattr(main, "HEALTH_PATH", tmp_path / "missing.json") + assert TestClient(main.app).get("/health").status_code == 503 diff --git a/tests/test_drift.py b/tests/test_drift.py new file mode 100644 index 0000000..d418c2b --- /dev/null +++ b/tests/test_drift.py @@ -0,0 +1,38 @@ +import numpy as np +import pytest +import pandas as pd + +from src import drift_detector +from src.drift_detector import calculate_psi + + +def test_identical_constant_populations_have_zero_drift(): + assert calculate_psi([7] * 30, [7] * 20) == 0 + + +def test_constant_population_shift_is_detected(): + assert calculate_psi([7] * 30, [9] * 20) > 0.25 + + +@pytest.mark.parametrize("reference, production", [([], [1]), ([1], []), ([1, np.inf], [1])]) +def test_invalid_psi_samples_fail_clearly(reference, production): + with pytest.raises(ValueError): + calculate_psi(reference, production) + + +def test_missingness_shift_is_reported_even_with_stable_values(tmp_path, monkeypatch): + reference = pd.DataFrame({"signal": [1.0] * 30, "churn": [0, 1] * 15}) + production = pd.DataFrame({"signal": [1.0] * 20 + [np.nan] * 10}) + reference.to_csv(tmp_path / "reference.csv", index=False) + production.to_csv(tmp_path / "production.csv", index=False) + monkeypatch.setattr(drift_detector, "REFERENCE_PATH", tmp_path / "reference.csv") + monkeypatch.setattr(drift_detector, "PRODUCTION_PATH", tmp_path / "production.csv") + monkeypatch.setattr(drift_detector, "OUTPUT_PATH", tmp_path / "drift.json") + + drift_detector.detect_drift() + + import json + report = json.loads((tmp_path / "drift.json").read_text()) + assert report["overall_drift_detected"] is True + assert report["features"]["signal"]["status"] == "MODERATE" + assert report["features"]["signal"]["production_count"] == 20