{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"name":"python","pygments_lexer":"ipython3"},"kaggle":{"accelerator":"none","dataSources":[{"sourceType":"competition","sourceId":119083,"databundleVersionId":15997195,"isSourceIdPinned":false}],"dockerImageVersionId":31328,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"id":"e40fadae-4ef0-4dba-9711-fd2bf03d1234","cell_type":"code","source":"from pathlib import Path\nimport sys\n\nPROJECT_ROOT = Path.cwd()\nSOURCE_DIRS = [\n    \"src\",\n]\n\nfor relative_dir in SOURCE_DIRS:\n    (PROJECT_ROOT / relative_dir).mkdir(parents=True, exist_ok=True)\n\nif str(PROJECT_ROOT) not in sys.path:\n    sys.path.insert(0, str(PROJECT_ROOT))\n\nprint(f\"project_root={PROJECT_ROOT}\")\nprint(f\"competition_data_exists={Path('/kaggle/input/competitions/cyber-physical-anomaly-detection-for-der-systems').exists()}\")\nprint(f\"kaggle_working_exists={Path('/kaggle/working').exists()}\")","metadata":{},"outputs":[],"execution_count":null},{"id":"96567a75-e3ae-4328-a912-b468377d580e","cell_type":"code","source":"%%writefile src/__init__.py\n\"\"\"DER anomaly baseline modules.\"\"\"","metadata":{},"outputs":[],"execution_count":null},{"id":"efd93884-40b6-4297-8dcd-8fe12d0b7340","cell_type":"code","source":"%%writefile src/contracts.py\n\"\"\"Shared runtime constants, dataclasses, and run configuration.\"\"\"\n\nimport math\nfrom dataclasses import dataclass, field\nfrom pathlib import Path\nfrom typing import Any, Dict, List, Optional, Tuple\n\nKAGGLE_WORKING_DIR = Path(\"/kaggle/working\")\nTRAIN_CSV_PATH = Path(\"/kaggle/input/competitions/cyber-physical-anomaly-detection-for-der-systems/train.csv\")\nTEST_CSV_PATH = Path(\"/kaggle/input/competitions/cyber-physical-anomaly-detection-for-der-systems/test.csv\")\nDEFAULT_SEED = 42\nSQRT3 = math.sqrt(3.0)\nCANON1 = \"DERSec|DER Simulator|10 kW DER|1.2.3|SN-Three-Phase\"\nCANON2 = \"DERSec|DER Simulator 100 kW|1.2.3.1|1.0.0|1100058974\"\nDEVICE_FAMILY_MAP = {\"canon10\": 0, \"canon100\": 1}\nRESIDUAL_TAIL_LEVELS = {\"tail\": 0.95, \"extreme\": 0.99, \"ultra\": 0.999}\nRESIDUAL_TAIL_FALLBACKS = {\"tail\": 0.05, \"extreme\": 0.10, \"ultra\": 0.20}\nFAMILY_THRESHOLD_FLOOR = 0.02\nMAX_THRESHOLD = 0.60\nCANON100_NEGATIVE_WEIGHT = 1.50\nHARD_OVERRIDE_TRAIN_WEIGHT = 0.35\nSCENARIO_SMOOTHING = 50.0\nAUDIT_TOLERANCE = 0.003\nMIN_OVERRIDE_PRECISION = 0.995\nCANON100_INTERACTION_FEATURES = [\n    \"hard_rule_score\",\n    \"scenario_rate\",\n    \"scenario_output_rate\",\n    \"resid_quantile_score\",\n    \"mode_dispatch_w_resid\",\n]\nSURROGATE_TARGETS = {\n    \"w\": (\"DERMeasureAC_0_W\", \"DERCapacity_0_WMaxRtg\"),\n    \"va\": (\"DERMeasureAC_0_VA\", \"DERCapacity_0_VAMaxRtg\"),\n    \"var\": (\"DERMeasureAC_0_Var\", \"DERCapacity_0_VarMaxInjRtg\"),\n    \"pf\": (\"DERMeasureAC_0_PF\", None),\n    \"a\": (\"DERMeasureAC_0_A\", \"DERCapacity_0_AMaxRtg\"),\n}\nSURROGATE_LEAKY_FEATURES = {\n    *(\n        f\"DERMeasureAC_0_{field}\"\n        for field in \"\"\"\n    W VA Var PF A WL1 WL2 WL3 VAL1 VAL2 VAL3 VarL1 VarL2 VarL3 PFL1 PFL2 PFL3\n    AL1 AL2 AL3\n    \"\"\".split()\n    ),\n    *\"\"\"\n    w_over_wmaxrtg w_over_wmax va_over_vamax va_over_vamaxrtg var_over_injmax\n    var_over_absmax a_over_amax w_minus_wmax w_minus_wmaxrtg va_minus_vamax\n    var_minus_injmax var_plus_absmax w_eq_wmaxrtg w_eq_wmax var_eq_varmaxinj\n    var_eq_neg_varmaxabs pf_sign_mismatch w_gt_wmax_tol w_gt_wmaxrtg_tol\n    va_gt_vamax_tol var_gt_injmax_tol var_lt_absmax_tol va_minus_pqmag\n    va_over_pqmag pf_from_w_va pf_error w_phase_sum_error va_phase_sum_error\n    var_phase_sum_error phase_w_spread phase_var_spread wset_abs_error\n    wsetpct_target wsetpct_abs_error wmaxlim_target wmaxlim_excess\n    varset_abs_error varsetpct_target varsetpct_abs_error wset_enabled_far\n    wsetpct_enabled_far wmaxlim_enabled_far varsetpct_enabled_far w_pct_of_rtg\n    var_pct_of_limit enter_service_blocked_power enter_service_blocked_va\n    enter_service_blocked_current pf_inj_target_error pf_inj_reversion_error\n    pf_reactive_near_limit trip_lv_power_when_outside trip_hv_power_when_outside\n    trip_lf_power_when_outside trip_hf_power_when_outside\n    trip_any_power_when_outside voltvar_curve_error voltwatt_curve_error\n    wattvar_curve_expected wattvar_curve_error freqdroop_w_over_pmin_pct\n    dcw_over_w dcw_over_abs_w ac_zero_dc_positive ac_positive_dc_zero\n    ac_dc_same_sign\n    \"\"\".split(),\n}\n\n\n@dataclass\nclass MetricSummary:\n    f2: float\n    precision: float\n    recall: float\n    positive_rate: float\n    rows: int\n\n\n@dataclass\nclass ScenarioStats:\n    \"\"\"Keep the four scenario lookup maps together as one learned state block.\"\"\"\n\n    sum_map: Dict[int, float] = field(default_factory=dict)\n    count_map: Dict[int, int] = field(default_factory=dict)\n    output_sum_map: Dict[int, float] = field(default_factory=dict)\n    output_count_map: Dict[int, int] = field(default_factory=dict)\n\n\n@dataclass\nclass FamilySemanticContext:\n    \"\"\"Learned semantic calibration state for one device family.\"\"\"\n\n    surrogate_feature_cols: List[str] = field(default_factory=list)\n    surrogate_models: Dict[Tuple[str, str], Any] = field(default_factory=dict)\n    residual_quantiles: Dict[str, Dict[str, Dict[str, float]]] = field(default_factory=dict)\n    family_base_rates: Dict[str, float] = field(default_factory=dict)\n    scenario_stats: ScenarioStats = field(default_factory=ScenarioStats)\n\n\n@dataclass\nclass FamilyModelBundle:\n    \"\"\"Store every learned component needed to score one family.\"\"\"\n\n    semantic_model: Optional[Any]\n    cat_model: Optional[Any]\n    semantic_context: FamilySemanticContext\n    semantic_feature_cols: List[str]\n    cat_feature_cols: List[str]\n    threshold: float\n    blend_weight: float\n\n\n@dataclass(frozen=True)\nclass BaselineConfig:\n    \"\"\"Immutable training knobs for the reusable baseline model.\"\"\"\n\n    chunksize: int = 5000\n    cv_folds: int = 5\n    n_estimators: int = 180\n    max_depth: int = 8\n    learning_rate: float = 0.05\n    subsample: float = 0.8\n    colsample_bytree: float = 0.8\n    cat_iterations: int = 400\n    cat_depth: int = 8\n    cat_learning_rate: float = 0.05\n    n_jobs: int = 4\n    seed: int = DEFAULT_SEED\n\n\n@dataclass(frozen=True)\nclass RunConfig:\n    train_path: Path = TRAIN_CSV_PATH\n    test_path: Path = TEST_CSV_PATH\n    submission_path: Path = KAGGLE_WORKING_DIR / \"submission.csv\"\n    chunksize: int = 5000\n    cv_folds: int = 5\n    xgb_n_estimators: int = 180\n    xgb_max_depth: int = 8\n    xgb_learning_rate: float = 0.05\n    xgb_subsample: float = 0.8\n    xgb_colsample_bytree: float = 0.8\n    cat_iterations: int = 400\n    cat_depth: int = 8\n    cat_learning_rate: float = 0.05\n    n_jobs: int = 4\n    seed: int = DEFAULT_SEED\n\n    def baseline_config(self) -> BaselineConfig:\n        \"\"\"Project a run config down to the model's immutable training knobs.\"\"\"\n\n        return BaselineConfig(\n            chunksize=self.chunksize,\n            cv_folds=self.cv_folds,\n            n_estimators=self.xgb_n_estimators,\n            max_depth=self.xgb_max_depth,\n            learning_rate=self.xgb_learning_rate,\n            subsample=self.xgb_subsample,\n            colsample_bytree=self.xgb_colsample_bytree,\n            cat_iterations=self.cat_iterations,\n            cat_depth=self.cat_depth,\n            cat_learning_rate=self.cat_learning_rate,\n            n_jobs=self.n_jobs,\n            seed=self.seed,\n        )\n","metadata":{},"outputs":[],"execution_count":null},{"id":"9762cfe7-c96e-4a54-95d6-1c80601fa747","cell_type":"code","source":"%%writefile src/feature_engineering.py\n\"\"\"Feature construction for the DER anomaly baseline.\"\"\"\n\nfrom typing import Dict, Optional, Sequence, Tuple\n\nimport numpy as np\nimport pandas as pd\n\nfrom src.contracts import (\n    CANON1,\n    CANON2,\n    SQRT3,\n)\nfrom src.rules import (\n    DEFAULT_HARD_OVERRIDE_NAMES,\n    compute_hard_rule_outputs,\n)\nfrom src.schema import (\n    BLOCK_SOURCE_COLUMNS,\n    COMMON_STR,\n    CURVE_FEATURE_SPECS,\n    ENTER_SERVICE_COLUMNS,\n    EXPECTED_MODEL_META,\n    NUMERIC_SOURCE_COLUMNS,\n    RAW_NUMERIC,\n    RAW_STRING_COLUMNS,\n    SAFE_RAW,\n    SAFE_STR,\n    TRIP_SPECS,\n)\n\n# These are the engineered columns that CatBoost consumes alongside raw strings.\nCAT_ENGINEERED_COLUMNS = [\n    \"device_fingerprint\",\n    \"common_missing_pattern\",\n    \"enter_service_missing_pattern\",\n    \"missing_selected_total\",\n    \"missing_selected_blocks\",\n    \"common_missing_any\",\n    \"common_missing_count\",\n    \"common_sn_has_decimal_suffix\",\n]\n\n\ndef _nanextreme_rows(arr: np.ndarray, *, fill_value: float, use_max: bool) -> np.ndarray:\n    arr = np.asarray(arr, dtype=np.float32)\n    mask = np.isfinite(arr)\n    out = np.full(arr.shape[0], np.nan, dtype=np.float32)\n    if arr.shape[1] == 0:\n        return out\n    reduced = np.where(mask, arr, fill_value)\n    reduced = reduced.max(axis=1) if use_max else reduced.min(axis=1)\n    valid_rows = mask.any(axis=1)\n    out[valid_rows] = reduced[valid_rows]\n    return out\n\n\ndef nanmin_rows(arr: np.ndarray) -> np.ndarray:\n    return _nanextreme_rows(arr, fill_value=np.inf, use_max=False)\n\n\ndef nanmax_rows(arr: np.ndarray) -> np.ndarray:\n    return _nanextreme_rows(arr, fill_value=-np.inf, use_max=True)\n\n\ndef nanmean_rows(arr: np.ndarray) -> np.ndarray:\n    arr = np.asarray(arr, dtype=np.float32)\n    mask = np.isfinite(arr)\n    out = np.full(arr.shape[0], np.nan, dtype=np.float32)\n    counts = mask.sum(axis=1)\n    valid_rows = counts > 0\n    if valid_rows.any():\n        totals = np.where(mask, arr, 0.0).sum(axis=1)\n        out[valid_rows] = totals[valid_rows] / counts[valid_rows]\n    return out\n\n\ndef curve_index(raw_idx: np.ndarray, num_options: int) -> np.ndarray:\n    idx = np.nan_to_num(np.asarray(raw_idx, dtype=np.float32), nan=1.0)\n    idx = idx.astype(np.int16) - 1\n    idx[(idx < 0) | (idx >= num_options)] = 0\n    return idx.astype(np.int8)\n\n\ndef select_curve_scalar(curves: Sequence[np.ndarray], idx: np.ndarray) -> np.ndarray:\n    stacked = np.stack(curves, axis=1)\n    return np.take_along_axis(stacked, idx[:, None], axis=1)[:, 0]\n\n\ndef select_curve_points(curves: Sequence[np.ndarray], idx: np.ndarray) -> np.ndarray:\n    stacked = np.stack(curves, axis=1)\n    return np.take_along_axis(stacked, idx[:, None, None], axis=1)[:, 0, :]\n\n\ndef pair_point_count(x_points: np.ndarray, y_points: np.ndarray) -> np.ndarray:\n    return (np.isfinite(np.asarray(x_points, dtype=np.float32)) & np.isfinite(np.asarray(y_points, dtype=np.float32))).sum(axis=1).astype(np.int16)\n\n\ndef curve_reverse_steps(x_points: np.ndarray) -> np.ndarray:\n    x_points = np.asarray(x_points, dtype=np.float32)\n    finite_pair = np.isfinite(x_points[:, :-1]) & np.isfinite(x_points[:, 1:])\n    return (((np.diff(x_points, axis=1) < -1e-6) & finite_pair).sum(axis=1)).astype(np.int8)\n\n\ndef curve_slope_stats(x_points: np.ndarray, y_points: np.ndarray) -> Tuple[np.ndarray, np.ndarray]:\n    x_points = np.asarray(x_points, dtype=np.float32)\n    y_points = np.asarray(y_points, dtype=np.float32)\n    dx = np.diff(x_points, axis=1)\n    dy = np.diff(y_points, axis=1)\n    valid = np.isfinite(x_points[:, :-1]) & np.isfinite(x_points[:, 1:]) & np.isfinite(y_points[:, :-1]) & np.isfinite(y_points[:, 1:]) & (np.abs(dx) > 1e-6)\n    slopes = np.full(dx.shape, np.nan, dtype=np.float32)\n    slopes[valid] = dy[valid] / dx[valid]\n    return nanmean_rows(slopes), nanmax_rows(np.abs(slopes))\n\n\ndef piecewise_interp(x: np.ndarray, x_points: np.ndarray, y_points: np.ndarray) -> np.ndarray:\n    x = np.asarray(x, dtype=np.float32)\n    x_points = np.asarray(x_points, dtype=np.float32)\n    y_points = np.asarray(y_points, dtype=np.float32)\n    n_rows, n_points = x_points.shape\n    result = np.full(n_rows, np.nan, dtype=np.float32)\n    valid_points = np.isfinite(x_points) & np.isfinite(y_points)\n    has_valid = valid_points.any(axis=1)\n    if n_points == 0:\n        return result\n\n    row_idx = np.arange(n_rows)\n    first_valid = np.argmax(valid_points, axis=1)\n    last_valid = n_points - 1 - np.argmax(valid_points[:, ::-1], axis=1)\n\n    first_x = np.full(n_rows, np.nan, dtype=np.float32)\n    first_y = np.full(n_rows, np.nan, dtype=np.float32)\n    last_x = np.full(n_rows, np.nan, dtype=np.float32)\n    last_y = np.full(n_rows, np.nan, dtype=np.float32)\n    first_x[has_valid] = x_points[row_idx[has_valid], first_valid[has_valid]]\n    first_y[has_valid] = y_points[row_idx[has_valid], first_valid[has_valid]]\n    last_x[has_valid] = x_points[row_idx[has_valid], last_valid[has_valid]]\n    last_y[has_valid] = y_points[row_idx[has_valid], last_valid[has_valid]]\n\n    for seg in range(n_points - 1):\n        x0 = x_points[:, seg]\n        x1 = x_points[:, seg + 1]\n        y0 = y_points[:, seg]\n        y1 = y_points[:, seg + 1]\n        valid_seg = np.isfinite(x0) & np.isfinite(x1) & np.isfinite(y0) & np.isfinite(y1) & (np.abs(x1 - x0) > 1e-6)\n        lo = np.minimum(x0, x1)\n        hi = np.maximum(x0, x1)\n        mask = valid_seg & np.isfinite(x) & np.isnan(result) & (x >= lo) & (x <= hi)\n        if mask.any():\n            frac = (x[mask] - x0[mask]) / (x1[mask] - x0[mask])\n            result[mask] = y0[mask] + frac * (y1[mask] - y0[mask])\n\n    low_mask = has_valid & np.isfinite(x) & np.isnan(result) & (x <= np.minimum(first_x, last_x))\n    result[low_mask] = first_y[low_mask]\n    high_mask = has_valid & np.isfinite(x) & np.isnan(result) & (x >= np.maximum(first_x, last_x))\n    result[high_mask] = last_y[high_mask]\n    return result\n\n\ndef var_pct(var: np.ndarray, varmaxinj: np.ndarray, varmaxabs: np.ndarray) -> np.ndarray:\n    var = np.asarray(var, dtype=np.float32)\n    denom = np.where(\n        var >= 0,\n        np.asarray(varmaxinj, dtype=np.float32),\n        np.asarray(varmaxabs, dtype=np.float32),\n    )\n    return 100.0 * np.divide(\n        var,\n        denom,\n        out=np.full_like(var, np.nan, dtype=np.float32),\n        where=np.isfinite(var) & np.isfinite(denom) & (np.abs(denom) > 1e-6),\n    )\n\n\ndef coerce_numeric(df: pd.DataFrame) -> None:\n    for col in NUMERIC_SOURCE_COLUMNS:\n        if df[col].dtype == object:\n            df[col] = pd.to_numeric(df[col], errors=\"coerce\")\n\n\ndef _initialize_feature_data(df: pd.DataFrame) -> Dict[str, np.ndarray]:\n    \"\"\"Capture the raw identity fields before the heavier numeric feature work.\"\"\"\n\n    fingerprint = df[COMMON_STR].fillna(\"<NA>\").agg(\"|\".join, axis=1)\n    device_family = np.where(\n        fingerprint == CANON1,\n        \"canon10\",\n        np.where(fingerprint == CANON2, \"canon100\", \"other\"),\n    )\n    return {\n        \"Id\": df[\"Id\"].to_numpy(),\n        \"device_fingerprint\": fingerprint.to_numpy(dtype=object),\n        \"device_family\": device_family,\n        \"common_missing_any\": df[COMMON_STR].isna().any(axis=1).astype(np.int8).to_numpy(),\n        \"common_missing_count\": df[COMMON_STR].isna().sum(axis=1).astype(np.int16).to_numpy(),\n        \"common_sn_has_decimal_suffix\": df[\"common[0].SN\"].fillna(\"\").astype(str).str.endswith(\".0\").astype(np.int8).to_numpy(),\n        \"noncanonical\": (device_family == \"other\").astype(np.int8),\n    }\n\n\ndef _copy_raw_source_columns(data: Dict[str, np.ndarray], df: pd.DataFrame) -> None:\n    \"\"\"Mirror the selected raw columns using the sanitized downstream names.\"\"\"\n\n    for col in RAW_NUMERIC:\n        arr = df[col].to_numpy()\n        if np.issubdtype(arr.dtype, np.floating):\n            arr = arr.astype(np.float32, copy=False)\n        data[SAFE_RAW[col]] = arr\n    for col in RAW_STRING_COLUMNS:\n        data[SAFE_STR[col]] = df[col].fillna(\"<NA>\").astype(str).to_numpy(dtype=object)\n\n\ndef add_block_missingness(data: Dict[str, np.ndarray], df: pd.DataFrame) -> None:\n    block_missing_total = np.zeros(len(df), dtype=np.int16)\n    block_missing_any = np.zeros(len(df), dtype=np.int16)\n    for block_name, cols in BLOCK_SOURCE_COLUMNS.items():\n        missing = df[cols].isna()\n        missing_count = missing.sum(axis=1).astype(np.int16).to_numpy()\n        data[f\"missing_{block_name}_count\"] = missing_count\n        data[f\"missing_{block_name}_any\"] = (missing_count > 0).astype(np.int8)\n        block_missing_total += missing_count\n        block_missing_any += (missing_count > 0).astype(np.int16)\n    data[\"missing_selected_total\"] = block_missing_total\n    data[\"missing_selected_blocks\"] = block_missing_any.astype(np.int8)\n    common_missing = df[[*COMMON_STR, \"common[0].ID\", \"common[0].L\"]].isna().to_numpy(dtype=np.uint16)\n    common_weights = (1 << np.arange(common_missing.shape[1], dtype=np.uint16)).reshape(1, -1)\n    data[\"common_missing_pattern\"] = (common_missing * common_weights).sum(axis=1).astype(np.int16)\n    enter_missing = df[ENTER_SERVICE_COLUMNS].isna().to_numpy(dtype=np.uint16)\n    enter_weights = (1 << np.arange(enter_missing.shape[1], dtype=np.uint16)).reshape(1, -1)\n    data[\"enter_service_missing_pattern\"] = (enter_missing * enter_weights).sum(axis=1).astype(np.int16)\n\n\ndef add_model_integrity_features(data: Dict[str, np.ndarray], df: pd.DataFrame) -> None:\n    anomaly_sum = np.zeros(len(df), dtype=np.int16)\n    missing_sum = np.zeros(len(df), dtype=np.int16)\n    for block_name, (id_col, len_col, expected_id, expected_len) in EXPECTED_MODEL_META.items():\n        raw_id = df[id_col].to_numpy(float)\n        raw_len = df[len_col].to_numpy(float)\n        id_missing = ~np.isfinite(raw_id)\n        len_missing = ~np.isfinite(raw_len)\n        id_match = np.isclose(raw_id, expected_id, equal_nan=False)\n        len_match = np.isclose(raw_len, expected_len, equal_nan=False)\n        data[f\"{block_name}_model_id_missing\"] = id_missing.astype(np.int8)\n        data[f\"{block_name}_model_len_missing\"] = len_missing.astype(np.int8)\n        data[f\"{block_name}_model_id_match\"] = id_match.astype(np.int8)\n        data[f\"{block_name}_model_len_match\"] = len_match.astype(np.int8)\n        data[f\"{block_name}_model_integrity_ok\"] = (id_match & len_match).astype(np.int8)\n        mismatch = (~id_missing & ~id_match) | (~len_missing & ~len_match)\n        data[f\"{block_name}_model_structure_anomaly\"] = mismatch.astype(np.int8)\n        anomaly_sum += mismatch.astype(np.int16)\n        missing_sum += (id_missing | len_missing).astype(np.int16)\n    data[\"model_structure_anomaly_count\"] = anomaly_sum.astype(np.int8)\n    data[\"model_structure_missing_count\"] = missing_sum.astype(np.int8)\n    data[\"model_structure_anomaly_any\"] = (anomaly_sum > 0).astype(np.int8)\n\n\ndef add_capacity_extension_features(\n    data: Dict[str, np.ndarray],\n    *,\n    wmaxrtg: np.ndarray,\n    wmax: np.ndarray,\n    vamaxrtg: np.ndarray,\n    vamax: np.ndarray,\n    varmaxinjrtg: np.ndarray,\n    varmaxinj: np.ndarray,\n    varmaxabsrtg: np.ndarray,\n    varmaxabs: np.ndarray,\n    vnomrtg: np.ndarray,\n    vnom: np.ndarray,\n    vmaxrtg: np.ndarray,\n    vmax: np.ndarray,\n    vminrtg: np.ndarray,\n    vmin: np.ndarray,\n    amaxrtg: np.ndarray,\n    amax: np.ndarray,\n    wcha_rtg: np.ndarray,\n    wdis_rtg: np.ndarray,\n    vacha_rtg: np.ndarray,\n    vadis_rtg: np.ndarray,\n    wcha: np.ndarray,\n    wdis: np.ndarray,\n    vacha: np.ndarray,\n    vadis: np.ndarray,\n    pfover_rtg: np.ndarray,\n    pfover: np.ndarray,\n    pfunder_rtg: np.ndarray,\n    pfunder: np.ndarray,\n) -> None:\n    data[\"vnom_setting_delta\"] = (vnom - vnomrtg).astype(np.float32)\n    data[\"vmax_setting_delta\"] = (vmax - vmaxrtg).astype(np.float32)\n    data[\"vmin_setting_delta\"] = (vmin - vminrtg).astype(np.float32)\n    data[\"amax_setting_delta\"] = (amax - amaxrtg).astype(np.float32)\n    data[\"pfover_setting_delta\"] = (pfover - pfover_rtg).astype(np.float32)\n    data[\"pfunder_setting_delta\"] = (pfunder - pfunder_rtg).astype(np.float32)\n    for name, numerator, denominator in [\n        (\"charge_rate_share_rtg\", wcha_rtg, wmaxrtg),\n        (\"discharge_rate_share_rtg\", wdis_rtg, wmaxrtg),\n        (\"charge_va_share_rtg\", vacha_rtg, vamaxrtg),\n        (\"discharge_va_share_rtg\", vadis_rtg, vamaxrtg),\n        (\"charge_rate_share_setting\", wcha, wmax),\n        (\"discharge_rate_share_setting\", wdis, wmax),\n        (\"charge_va_share_setting\", vacha, vamax),\n        (\"discharge_va_share_setting\", vadis, vamax),\n    ]:\n        numerator = np.asarray(numerator, dtype=np.float32)\n        denominator = np.asarray(denominator, dtype=np.float32)\n        data[name] = np.divide(\n            numerator,\n            denominator,\n            out=np.full_like(numerator, np.nan, dtype=np.float32),\n            where=np.isfinite(numerator) & np.isfinite(denominator) & (np.abs(denominator) > 1e-6),\n        )\n    rating_pairs = [\n        (wmaxrtg, wmax),\n        (vamaxrtg, vamax),\n        (varmaxinjrtg, varmaxinj),\n        (varmaxabsrtg, varmaxabs),\n        (vnomrtg, vnom),\n        (vmaxrtg, vmax),\n        (vminrtg, vmin),\n        (amaxrtg, amax),\n    ]\n    gap_count = np.zeros(len(wmaxrtg), dtype=np.int16)\n    for rating, setting in rating_pairs:\n        tol = np.maximum(1.0, 0.01 * np.nan_to_num(np.abs(rating), nan=0.0)).astype(np.float32)\n        gap = np.isfinite(rating) & np.isfinite(setting) & (np.abs(setting - rating) > tol)\n        gap_count += gap.astype(np.int16)\n    data[\"rating_setting_gap_count\"] = gap_count.astype(np.int8)\n\n\ndef add_temperature_features(data: Dict[str, np.ndarray], df: pd.DataFrame) -> None:\n    temp_cols = [\n        \"DERMeasureAC[0].TmpAmb\",\n        \"DERMeasureAC[0].TmpCab\",\n        \"DERMeasureAC[0].TmpSnk\",\n        \"DERMeasureAC[0].TmpTrns\",\n        \"DERMeasureAC[0].TmpSw\",\n        \"DERMeasureAC[0].TmpOt\",\n    ]\n    temps = df[temp_cols].to_numpy(float)\n    temp_min = nanmin_rows(temps)\n    temp_max = nanmax_rows(temps)\n    temp_mean = nanmean_rows(temps)\n    amb = df[\"DERMeasureAC[0].TmpAmb\"].to_numpy(float)\n    data[\"temp_min\"] = temp_min\n    data[\"temp_max\"] = temp_max\n    data[\"temp_mean\"] = temp_mean\n    data[\"temp_spread\"] = (temp_max - temp_min).astype(np.float32)\n    data[\"temp_max_over_ambient\"] = (temp_max - amb).astype(np.float32)\n\n\ndef add_enter_service_features(\n    data: Dict[str, np.ndarray],\n    df: pd.DataFrame,\n    *,\n    voltage_pct: np.ndarray,\n    hz: np.ndarray,\n    abs_w: np.ndarray,\n    va: np.ndarray,\n    a: np.ndarray,\n    tolw: np.ndarray,\n    tolva: np.ndarray,\n    amax: np.ndarray,\n) -> None:\n    es = df[\"DEREnterService[0].ES\"].to_numpy(float)\n    es_v_hi = df[\"DEREnterService[0].ESVHi\"].to_numpy(float)\n    es_v_lo = df[\"DEREnterService[0].ESVLo\"].to_numpy(float)\n    es_hz_hi = df[\"DEREnterService[0].ESHzHi\"].to_numpy(float)\n    es_hz_lo = df[\"DEREnterService[0].ESHzLo\"].to_numpy(float)\n    es_delay = df[\"DEREnterService[0].ESDlyTms\"].to_numpy(float)\n    es_random = df[\"DEREnterService[0].ESRndTms\"].to_numpy(float)\n    es_ramp = df[\"DEREnterService[0].ESRmpTms\"].to_numpy(float)\n    es_delay_rem = df[\"DEREnterService[0].ESDlyRemTms\"].to_numpy(float)\n\n    inside_v = np.isfinite(voltage_pct) & np.isfinite(es_v_hi) & np.isfinite(es_v_lo) & (voltage_pct >= es_v_lo) & (voltage_pct <= es_v_hi)\n    inside_hz = np.isfinite(hz) & np.isfinite(es_hz_hi) & np.isfinite(es_hz_lo) & (hz >= es_hz_lo) & (hz <= es_hz_hi)\n    inside_window = inside_v & inside_hz\n    enabled = np.isfinite(es) & (es == 1.0)\n    state_anomaly = np.isfinite(es) & (es >= 1.5)\n    should_idle = (~enabled) | (~inside_window)\n    current_tol = np.maximum(1.0, 0.02 * np.nan_to_num(amax, nan=0.0))\n\n    data[\"enter_service_enabled\"] = enabled.astype(np.int8)\n    data[\"enter_service_state_anomaly\"] = state_anomaly.astype(np.int8)\n    data[\"enter_service_inside_window\"] = inside_window.astype(np.int8)\n    data[\"enter_service_outside_window\"] = (~inside_window).astype(np.int8)\n    data[\"enter_service_should_idle\"] = should_idle.astype(np.int8)\n    data[\"enter_service_v_window_width\"] = (es_v_hi - es_v_lo).astype(np.float32)\n    data[\"enter_service_hz_window_width\"] = (es_hz_hi - es_hz_lo).astype(np.float32)\n    data[\"enter_service_v_margin_low\"] = (voltage_pct - es_v_lo).astype(np.float32)\n    data[\"enter_service_v_margin_high\"] = (es_v_hi - voltage_pct).astype(np.float32)\n    data[\"enter_service_hz_margin_low\"] = (hz - es_hz_lo).astype(np.float32)\n    data[\"enter_service_hz_margin_high\"] = (es_hz_hi - hz).astype(np.float32)\n    data[\"enter_service_total_delay\"] = (es_delay + es_random).astype(np.float32)\n    data[\"enter_service_delay_remaining\"] = es_delay_rem.astype(np.float32)\n    data[\"enter_service_ramp_time\"] = es_ramp.astype(np.float32)\n    data[\"enter_service_delay_active\"] = (np.nan_to_num(es_delay_rem, nan=0.0) > 0).astype(np.int8)\n\n    blocked_power = should_idle & (abs_w > tolw)\n    blocked_va = should_idle & (va > tolva)\n    blocked_current = should_idle & (a > current_tol)\n    data[\"enter_service_blocked_power\"] = blocked_power.astype(np.int8)\n    data[\"enter_service_blocked_va\"] = blocked_va.astype(np.int8)\n    data[\"enter_service_blocked_current\"] = blocked_current.astype(np.int8)\n\n\ndef add_pf_control_features(\n    data: Dict[str, np.ndarray],\n    df: pd.DataFrame,\n    *,\n    pf: np.ndarray,\n    var: np.ndarray,\n    varmaxinj: np.ndarray,\n    varmaxabs: np.ndarray,\n) -> None:\n    pfinj_ena = np.nan_to_num(df[\"DERCtlAC[0].PFWInjEna\"].to_numpy(float), nan=0.0)\n    pfinj_ena_rvrt = np.nan_to_num(df[\"DERCtlAC[0].PFWInjEnaRvrt\"].to_numpy(float), nan=0.0)\n    pfabs_ena = np.nan_to_num(df[\"DERCtlAC[0].PFWAbsEna\"].to_numpy(float), nan=0.0)\n    pfabs_ena_rvrt = np.nan_to_num(df[\"DERCtlAC[0].PFWAbsEnaRvrt\"].to_numpy(float), nan=0.0)\n    pfinj_target = df[\"DERCtlAC[0].PFWInj.PF\"].to_numpy(float)\n    pfinj_rvrt_target = df[\"DERCtlAC[0].PFWInjRvrt.PF\"].to_numpy(float)\n    pfinj_ext = df[\"DERCtlAC[0].PFWInj.Ext\"].to_numpy(float)\n    pfinj_rvrt_ext = df[\"DERCtlAC[0].PFWInjRvrt.Ext\"].to_numpy(float)\n    pfabs_ext = df[\"DERCtlAC[0].PFWAbs.Ext\"].to_numpy(float)\n    pfabs_rvrt_ext = df[\"DERCtlAC[0].PFWAbsRvrt.Ext\"].to_numpy(float)\n\n    observed_var_pct = var_pct(var, varmaxinj, varmaxabs)\n    inj_target_error = np.where(\n        (pfinj_ena > 0) & np.isfinite(pfinj_target),\n        np.abs(np.abs(pf) - pfinj_target),\n        np.nan,\n    )\n    inj_rvrt_error = np.where(\n        (pfinj_ena_rvrt > 0) & np.isfinite(pfinj_rvrt_target),\n        np.abs(np.abs(pf) - pfinj_rvrt_target),\n        np.nan,\n    )\n    data[\"pf_control_any_enabled\"] = ((pfinj_ena > 0) | (pfabs_ena > 0)).astype(np.int8)\n    data[\"pf_control_any_reversion\"] = ((pfinj_ena_rvrt > 0) | (pfabs_ena_rvrt > 0)).astype(np.int8)\n    data[\"pf_inj_target_error\"] = inj_target_error.astype(np.float32)\n    data[\"pf_inj_reversion_error\"] = inj_rvrt_error.astype(np.float32)\n    data[\"pf_inj_ext_present\"] = np.isfinite(pfinj_ext).astype(np.int8)\n    data[\"pf_inj_rvrt_ext_present\"] = np.isfinite(pfinj_rvrt_ext).astype(np.int8)\n    data[\"pf_abs_ext_present\"] = np.isfinite(pfabs_ext).astype(np.int8)\n    data[\"pf_abs_rvrt_ext_present\"] = np.isfinite(pfabs_rvrt_ext).astype(np.int8)\n    data[\"pf_inj_enabled_missing_target\"] = ((pfinj_ena > 0) & ~np.isfinite(pfinj_target)).astype(np.int8)\n    data[\"pf_reactive_near_limit\"] = (np.abs(observed_var_pct) >= 95.0).astype(np.int8)\n\n\ndef add_trip_block_features(\n    data: Dict[str, np.ndarray],\n    df: pd.DataFrame,\n    *,\n    short_name: str,\n    prefix: str,\n    axis_name: str,\n    mode: str,\n    measure_value: np.ndarray,\n    abs_w: np.ndarray,\n    tolw: np.ndarray,\n) -> Tuple[np.ndarray, np.ndarray]:\n    adpt_idx = curve_index(df[f\"{prefix}.AdptCrvRslt\"].to_numpy(float), 2)\n    must_actpt = select_curve_scalar([df[f\"{prefix}.Crv[{curve}].MustTrip.ActPt\"].to_numpy(float) for curve in range(2)], adpt_idx)\n    mom_actpt = select_curve_scalar([df[f\"{prefix}.Crv[{curve}].MomCess.ActPt\"].to_numpy(float) for curve in range(2)], adpt_idx)\n    must_x = select_curve_points(\n        [np.column_stack([df[f\"{prefix}.Crv[{curve}].MustTrip.Pt[{point}].{axis_name}\"].to_numpy(float) for point in range(5)]) for curve in range(2)],\n        adpt_idx,\n    )\n    must_t = select_curve_points(\n        [np.column_stack([df[f\"{prefix}.Crv[{curve}].MustTrip.Pt[{point}].Tms\"].to_numpy(float) for point in range(5)]) for curve in range(2)],\n        adpt_idx,\n    )\n    mom_x = select_curve_points(\n        [np.column_stack([df[f\"{prefix}.Crv[{curve}].MomCess.Pt[{point}].{axis_name}\"].to_numpy(float) for point in range(5)]) for curve in range(2)],\n        adpt_idx,\n    )\n    mom_t = select_curve_points(\n        [np.column_stack([df[f\"{prefix}.Crv[{curve}].MomCess.Pt[{point}].Tms\"].to_numpy(float) for point in range(5)]) for curve in range(2)],\n        adpt_idx,\n    )\n    may_present = np.column_stack([df[f\"{prefix}.Crv[{curve}].MayTrip.Pt[{point}].{axis_name}\"].to_numpy(float) for curve in range(2) for point in range(5)])\n\n    enabled = np.nan_to_num(df[f\"{prefix}.Ena\"].to_numpy(float), nan=0.0) > 0\n    must_count = pair_point_count(must_x, must_t)\n    mom_count = pair_point_count(mom_x, mom_t)\n    must_x_min = nanmin_rows(must_x)\n    must_x_max = nanmax_rows(must_x)\n    must_t_min = nanmin_rows(must_t)\n    must_t_max = nanmax_rows(must_t)\n    mom_x_min = nanmin_rows(mom_x)\n    mom_x_max = nanmax_rows(mom_x)\n    mom_t_min = nanmin_rows(mom_t)\n    mom_t_max = nanmax_rows(mom_t)\n\n    if mode == \"low\":\n        margin = measure_value - must_x_max\n    else:\n        margin = must_x_min - measure_value\n    outside = enabled & np.isfinite(margin) & (margin < 0)\n    power_when_outside = outside & (abs_w > tolw)\n    envelope_gap = np.where(\n        np.isfinite(mom_x_min) & np.isfinite(must_x_max),\n        np.abs(mom_x_min - must_x_max),\n        np.nan,\n    )\n\n    data[f\"trip_{short_name}_curve_idx\"] = adpt_idx.astype(np.int8)\n    data[f\"trip_{short_name}_enabled\"] = enabled.astype(np.int8)\n    data[f\"trip_{short_name}_curve_req_gap\"] = (df[f\"{prefix}.AdptCrvReq\"].to_numpy(float) - df[f\"{prefix}.AdptCrvRslt\"].to_numpy(float)).astype(np.float32)\n    data[f\"trip_{short_name}_musttrip_count\"] = must_count\n    data[f\"trip_{short_name}_musttrip_actpt_gap\"] = (must_actpt - must_count).astype(np.float32)\n    data[f\"trip_{short_name}_musttrip_axis_min\"] = must_x_min\n    data[f\"trip_{short_name}_musttrip_axis_max\"] = must_x_max\n    data[f\"trip_{short_name}_musttrip_axis_span\"] = (must_x_max - must_x_min).astype(np.float32)\n    data[f\"trip_{short_name}_musttrip_tms_span\"] = (must_t_max - must_t_min).astype(np.float32)\n    data[f\"trip_{short_name}_musttrip_reverse_steps\"] = curve_reverse_steps(must_x)\n    data[f\"trip_{short_name}_momcess_count\"] = mom_count\n    data[f\"trip_{short_name}_momcess_actpt_gap\"] = (mom_actpt - mom_count).astype(np.float32)\n    data[f\"trip_{short_name}_momcess_axis_span\"] = (mom_x_max - mom_x_min).astype(np.float32)\n    data[f\"trip_{short_name}_momcess_tms_span\"] = (mom_t_max - mom_t_min).astype(np.float32)\n    data[f\"trip_{short_name}_momcess_reverse_steps\"] = curve_reverse_steps(mom_x)\n    data[f\"trip_{short_name}_maytrip_present_any\"] = np.isfinite(may_present).any(axis=1).astype(np.int8)\n    data[f\"trip_{short_name}_musttrip_margin\"] = margin.astype(np.float32)\n    data[f\"trip_{short_name}_outside_musttrip\"] = outside.astype(np.int8)\n    data[f\"trip_{short_name}_power_when_outside\"] = power_when_outside.astype(np.int8)\n    data[f\"trip_{short_name}_momcess_musttrip_gap\"] = envelope_gap.astype(np.float32)\n    return outside.astype(np.int8), power_when_outside.astype(np.int8)\n\n\ndef add_curve_block_features(\n    data: Dict[str, np.ndarray],\n    *,\n    name: str,\n    raw_idx: np.ndarray,\n    curve_x: Sequence[np.ndarray],\n    curve_y: Sequence[np.ndarray],\n    curve_actpt: Sequence[np.ndarray],\n    curve_meta: Dict[str, Sequence[np.ndarray]],\n    measure_value: np.ndarray,\n    observed_value: Optional[np.ndarray] = None,\n) -> None:\n    adpt_idx = curve_index(raw_idx, len(curve_x))\n    selected_x = select_curve_points(curve_x, adpt_idx)\n    selected_y = select_curve_points(curve_y, adpt_idx)\n    selected_actpt = select_curve_scalar(curve_actpt, adpt_idx)\n    data[f\"{name}_curve_idx\"] = adpt_idx.astype(np.int8)\n    point_count = pair_point_count(selected_x, selected_y)\n    data[f\"{name}_curve_point_count\"] = point_count\n    data[f\"{name}_curve_actpt_gap\"] = (selected_actpt - point_count).astype(np.float32)\n    x_min = nanmin_rows(selected_x)\n    x_max = nanmax_rows(selected_x)\n    y_min = nanmin_rows(selected_y)\n    y_max = nanmax_rows(selected_y)\n    mean_slope, max_abs_slope = curve_slope_stats(selected_x, selected_y)\n    data[f\"{name}_curve_x_span\"] = (x_max - x_min).astype(np.float32)\n    data[f\"{name}_curve_y_span\"] = (y_max - y_min).astype(np.float32)\n    data[f\"{name}_curve_reverse_steps\"] = curve_reverse_steps(selected_x)\n    data[f\"{name}_curve_mean_slope\"] = mean_slope\n    data[f\"{name}_curve_max_abs_slope\"] = max_abs_slope\n    data[f\"{name}_curve_measure_margin_low\"] = (measure_value - x_min).astype(np.float32)\n    data[f\"{name}_curve_measure_margin_high\"] = (x_max - measure_value).astype(np.float32)\n    if observed_value is not None:\n        expected_value = piecewise_interp(measure_value, selected_x, selected_y)\n        data[f\"{name}_curve_expected\"] = expected_value.astype(np.float32)\n        data[f\"{name}_curve_error\"] = (observed_value - expected_value).astype(np.float32)\n    for meta_name, curves in curve_meta.items():\n        data[f\"{name}_curve_{meta_name}\"] = select_curve_scalar(curves, adpt_idx).astype(np.float32)\n\n\ndef add_freq_droop_features(\n    data: Dict[str, np.ndarray],\n    df: pd.DataFrame,\n    *,\n    hz: np.ndarray,\n    w_pct: np.ndarray,\n) -> None:\n    raw_idx = df[\"DERFreqDroop[0].AdptCtlRslt\"].to_numpy(float)\n    ctl_idx = curve_index(raw_idx, 3)\n    dbof_curves = [df[f\"DERFreqDroop[0].Ctl[{i}].DbOf\"].to_numpy(float) for i in range(3)]\n    dbuf_curves = [df[f\"DERFreqDroop[0].Ctl[{i}].DbUf\"].to_numpy(float) for i in range(3)]\n    kof_curves = [df[f\"DERFreqDroop[0].Ctl[{i}].KOf\"].to_numpy(float) for i in range(3)]\n    kuf_curves = [df[f\"DERFreqDroop[0].Ctl[{i}].KUf\"].to_numpy(float) for i in range(3)]\n    rsp_curves = [df[f\"DERFreqDroop[0].Ctl[{i}].RspTms\"].to_numpy(float) for i in range(3)]\n    pmin_curves = [df[f\"DERFreqDroop[0].Ctl[{i}].PMin\"].to_numpy(float) for i in range(3)]\n    ro_curves = [df[f\"DERFreqDroop[0].Ctl[{i}].ReadOnly\"].to_numpy(float) for i in range(3)]\n    dbof = select_curve_scalar(dbof_curves, ctl_idx)\n    dbuf = select_curve_scalar(dbuf_curves, ctl_idx)\n    kof = select_curve_scalar(kof_curves, ctl_idx)\n    kuf = select_curve_scalar(kuf_curves, ctl_idx)\n    rsp = select_curve_scalar(rsp_curves, ctl_idx)\n    pmin = select_curve_scalar(pmin_curves, ctl_idx)\n    readonly = select_curve_scalar(ro_curves, ctl_idx)\n\n    over_activation = np.maximum(hz - (60.0 + dbof), 0.0)\n    under_activation = np.maximum((60.0 - dbuf) - hz, 0.0)\n    over_activation = np.asarray(over_activation, dtype=np.float32)\n    under_activation = np.asarray(under_activation, dtype=np.float32)\n    kof = np.asarray(kof, dtype=np.float32)\n    kuf = np.asarray(kuf, dtype=np.float32)\n    expected_delta_pct = 100.0 * np.divide(\n        over_activation,\n        kof,\n        out=np.full_like(over_activation, np.nan, dtype=np.float32),\n        where=np.isfinite(over_activation) & np.isfinite(kof) & (np.abs(kof) > 1e-6),\n    ) - 100.0 * np.divide(\n        under_activation,\n        kuf,\n        out=np.full_like(under_activation, np.nan, dtype=np.float32),\n        where=np.isfinite(under_activation) & np.isfinite(kuf) & (np.abs(kuf) > 1e-6),\n    )\n    dbof_stack = np.column_stack(dbof_curves)\n    dbuf_stack = np.column_stack(dbuf_curves)\n    k_stack = np.column_stack(kof_curves + kuf_curves)\n    pmin_stack = np.column_stack(pmin_curves)\n\n    data[\"freqdroop_ctl_idx\"] = ctl_idx.astype(np.int8)\n    data[\"freqdroop_dbof\"] = dbof.astype(np.float32)\n    data[\"freqdroop_dbuf\"] = dbuf.astype(np.float32)\n    data[\"freqdroop_kof\"] = kof.astype(np.float32)\n    data[\"freqdroop_kuf\"] = kuf.astype(np.float32)\n    data[\"freqdroop_rsp\"] = rsp.astype(np.float32)\n    data[\"freqdroop_pmin\"] = pmin.astype(np.float32)\n    data[\"freqdroop_readonly\"] = readonly.astype(np.float32)\n    data[\"freqdroop_deadband_width\"] = (dbof + dbuf).astype(np.float32)\n    data[\"freqdroop_over_activation\"] = over_activation.astype(np.float32)\n    data[\"freqdroop_under_activation\"] = under_activation.astype(np.float32)\n    data[\"freqdroop_expected_delta_pct\"] = expected_delta_pct.astype(np.float32)\n    data[\"freqdroop_outside_deadband\"] = ((over_activation > 0) | (under_activation > 0)).astype(np.int8)\n    data[\"freqdroop_w_over_pmin_pct\"] = (w_pct - pmin).astype(np.float32)\n    data[\"freqdroop_db_span\"] = (nanmax_rows(np.column_stack([dbof_stack, dbuf_stack])) - nanmin_rows(np.column_stack([dbof_stack, dbuf_stack]))).astype(np.float32)\n    data[\"freqdroop_k_span\"] = (nanmax_rows(k_stack) - nanmin_rows(k_stack)).astype(np.float32)\n    data[\"freqdroop_pmin_span\"] = (nanmax_rows(pmin_stack) - nanmin_rows(pmin_stack)).astype(np.float32)\n\n\ndef add_dc_features(\n    data: Dict[str, np.ndarray],\n    df: pd.DataFrame,\n    *,\n    w: np.ndarray,\n    abs_w: np.ndarray,\n) -> None:\n    dcw = df[\"DERMeasureDC[0].DCW\"].to_numpy(float)\n    dca = df[\"DERMeasureDC[0].DCA\"].to_numpy(float)\n    prt0 = df[\"DERMeasureDC[0].Prt[0].DCW\"].to_numpy(float)\n    prt1 = df[\"DERMeasureDC[0].Prt[1].DCW\"].to_numpy(float)\n    prt0_v = df[\"DERMeasureDC[0].Prt[0].DCV\"].to_numpy(float)\n    prt1_v = df[\"DERMeasureDC[0].Prt[1].DCV\"].to_numpy(float)\n    prt0_a = df[\"DERMeasureDC[0].Prt[0].DCA\"].to_numpy(float)\n    prt1_a = df[\"DERMeasureDC[0].Prt[1].DCA\"].to_numpy(float)\n    prt0_t = df[\"DERMeasureDC[0].Prt[0].PrtTyp\"].to_numpy(float)\n    prt1_t = df[\"DERMeasureDC[0].Prt[1].PrtTyp\"].to_numpy(float)\n\n    dcw32 = np.asarray(dcw, dtype=np.float32)\n    w32 = np.asarray(w, dtype=np.float32)\n    abs_w32 = np.asarray(abs_w, dtype=np.float32)\n    data[\"dcw_over_w\"] = np.divide(\n        dcw32,\n        w32,\n        out=np.full_like(dcw32, np.nan, dtype=np.float32),\n        where=np.isfinite(dcw32) & np.isfinite(w32) & (np.abs(w32) > 1e-6),\n    )\n    data[\"dcw_over_abs_w\"] = np.divide(\n        dcw32,\n        abs_w32,\n        out=np.full_like(dcw32, np.nan, dtype=np.float32),\n        where=np.isfinite(dcw32) & np.isfinite(abs_w32) & (np.abs(abs_w32) > 1e-6),\n    )\n    data[\"dcw_minus_port_sum\"] = (dcw - (prt0 + prt1)).astype(np.float32)\n    data[\"dcv_spread\"] = np.abs(prt0_v - prt1_v).astype(np.float32)\n    data[\"dca_spread\"] = np.abs(prt0_a - prt1_a).astype(np.float32)\n    prt032 = np.asarray(prt0, dtype=np.float32)\n    prt0_plus_prt1 = np.asarray(prt0 + prt1, dtype=np.float32)\n    data[\"dc_port0_share\"] = np.divide(\n        prt032,\n        prt0_plus_prt1,\n        out=np.full_like(prt032, np.nan, dtype=np.float32),\n        where=np.isfinite(prt032) & np.isfinite(prt0_plus_prt1) & (np.abs(prt0_plus_prt1) > 1e-6),\n    )\n    data[\"dc_port_type_mismatch\"] = (np.isfinite(prt0_t) & np.isfinite(prt1_t) & (prt0_t != prt1_t)).astype(np.int8)\n    rare_type = (prt0_t == 7) | (prt1_t == 7)\n    data[\"dc_port_type_rare_any\"] = rare_type.astype(np.int8)\n    data[\"ac_zero_dc_positive\"] = ((np.abs(w) <= 1e-6) & (dcw > 0)).astype(np.int8)\n    data[\"ac_positive_dc_zero\"] = ((w > 0) & (np.abs(dcw) <= 1e-6)).astype(np.int8)\n    data[\"ac_dc_same_sign\"] = (np.sign(np.nan_to_num(w, nan=0.0)) == np.sign(np.nan_to_num(dcw, nan=0.0))).astype(np.int8)\n    dca32 = np.asarray(dca, dtype=np.float32)\n    prt_total_a = np.asarray(prt0_a + prt1_a, dtype=np.float32)\n    data[\"dca_over_total\"] = np.divide(\n        dca32,\n        prt_total_a,\n        out=np.full_like(dca32, np.nan, dtype=np.float32),\n        where=np.isfinite(dca32) & np.isfinite(prt_total_a) & (np.abs(prt_total_a) > 1e-6),\n    )\n\n\ndef _add_adaptive_curve_features(\n    data: Dict[str, np.ndarray],\n    df: pd.DataFrame,\n    *,\n    voltage_pct: np.ndarray,\n    reactive_pct: np.ndarray,\n    w_pct: np.ndarray,\n) -> None:\n    \"\"\"Bind live measurements to the schema-owned adaptive curve layouts.\"\"\"\n\n    runtime_measurements = {\n        \"voltvar\": (\n            voltage_pct - 100.0 + df[\"DERVoltVar[0].Crv[0].VRef\"].fillna(100.0).to_numpy(float),\n            reactive_pct,\n        ),\n        \"voltwatt\": (voltage_pct, w_pct),\n        \"wattvar\": (w_pct, reactive_pct),\n    }\n    for name, spec in CURVE_FEATURE_SPECS.items():\n        measure_value, observed_value = runtime_measurements[name]\n        add_curve_block_features(\n            data,\n            name=name,\n            raw_idx=df[f\"{spec.prefix}.AdptCrvRslt\"].to_numpy(float),\n            curve_x=[np.column_stack([df[f\"{spec.prefix}.Crv[{curve}].Pt[{point}].{spec.x_field}\"].to_numpy(float) for point in range(spec.point_count)]) for curve in range(3)],\n            curve_y=[np.column_stack([df[f\"{spec.prefix}.Crv[{curve}].Pt[{point}].{spec.y_field}\"].to_numpy(float) for point in range(spec.point_count)]) for curve in range(3)],\n            curve_actpt=[df[f\"{spec.prefix}.Crv[{curve}].ActPt\"].to_numpy(float) for curve in range(3)],\n            curve_meta={meta_name: [df[f\"{spec.prefix}.Crv[{curve}].{field}\"].to_numpy(float) for curve in range(3)] for meta_name, field in spec.meta_fields.items()},\n            measure_value=measure_value,\n            observed_value=observed_value,\n        )\n\n\ndef _apply_hard_rule_features(\n    data: Dict[str, np.ndarray],\n    *,\n    hard_override_names: Sequence[str],\n) -> None:\n    \"\"\"Derive aggregate hard-rule signals from the centralized rule metadata.\"\"\"\n\n    outputs = compute_hard_rule_outputs(data, hard_override_names)\n    data[\"hard_override_anomaly\"] = outputs.hard_override_anomaly\n    data[\"hard_rule_count\"] = outputs.hard_rule_count\n    data[\"hard_rule_score\"] = outputs.hard_rule_score\n    data[\"hard_rule_anomaly\"] = outputs.hard_rule_anomaly\n\n\ndef build_features(df: pd.DataFrame, hard_override_names: Optional[Sequence[str]] = None) -> pd.DataFrame:\n    if hard_override_names is None:\n        hard_override_names = DEFAULT_HARD_OVERRIDE_NAMES\n\n    def block_values(prefix: str, fields: Sequence[str]) -> Dict[str, np.ndarray]:\n        return {field: df[f\"{prefix}.{field}\"].to_numpy(float) for field in fields}\n\n    # The feature frame is built in ordered phases so downstream model code sees\n    # the same columns and dtypes as the original monolithic implementation.\n    coerce_numeric(df)\n\n    data = _initialize_feature_data(df)\n    _copy_raw_source_columns(data, df)\n\n    # Missingness and structural integrity are cheap signals, so populate them\n    # before the heavier physical-control feature set below.\n    add_block_missingness(data, df)\n    add_model_integrity_features(data, df)\n    add_temperature_features(data, df)\n\n    measure_ac = block_values(\n        \"DERMeasureAC[0]\",\n        [\"W\", \"VA\", \"Var\", \"PF\", \"A\", \"LLV\", \"LNV\", \"Hz\", \"ACType\"],\n    )\n    capacity = block_values(\n        \"DERCapacity[0]\",\n        [\n            \"WMaxRtg\",\n            \"VAMaxRtg\",\n            \"VarMaxInjRtg\",\n            \"VarMaxAbsRtg\",\n            \"WMax\",\n            \"VAMax\",\n            \"VarMaxInj\",\n            \"VarMaxAbs\",\n            \"AMax\",\n            \"VNom\",\n            \"VMax\",\n            \"VMin\",\n        ],\n    )\n    ctl_ac = block_values(\n        \"DERCtlAC[0]\",\n        [\"WSetEna\", \"WSet\", \"WSetPct\", \"WMaxLimPctEna\", \"WMaxLimPct\", \"VarSetEna\", \"VarSet\", \"VarSetPct\"],\n    )\n    phase_ac = block_values(\n        \"DERMeasureAC[0]\",\n        [\"WL1\", \"WL2\", \"WL3\", \"VAL1\", \"VAL2\", \"VAL3\", \"VarL1\", \"VarL2\", \"VarL3\", \"VL1L2\", \"VL2L3\", \"VL3L1\", \"VL1\", \"VL2\", \"VL3\"],\n    )\n    capacity_ext = block_values(\n        \"DERCapacity[0]\",\n        [\n            \"VNomRtg\",\n            \"VMaxRtg\",\n            \"VMinRtg\",\n            \"AMaxRtg\",\n            \"WChaRteMaxRtg\",\n            \"WDisChaRteMaxRtg\",\n            \"VAChaRteMaxRtg\",\n            \"VADisChaRteMaxRtg\",\n            \"WChaRteMax\",\n            \"WDisChaRteMax\",\n            \"VAChaRteMax\",\n            \"VADisChaRteMax\",\n            \"PFOvrExtRtg\",\n            \"PFOvrExt\",\n            \"PFUndExtRtg\",\n            \"PFUndExt\",\n        ],\n    )\n\n    w = measure_ac[\"W\"]\n    abs_w = np.abs(w)\n    va = measure_ac[\"VA\"]\n    var = measure_ac[\"Var\"]\n    pf = measure_ac[\"PF\"]\n    a = measure_ac[\"A\"]\n    llv = measure_ac[\"LLV\"]\n    lnv = measure_ac[\"LNV\"]\n    lnv_scaled = np.asarray(lnv * SQRT3, dtype=np.float32)\n    hz = measure_ac[\"Hz\"]\n\n    wmaxrtg = capacity[\"WMaxRtg\"]\n    vamaxrtg = capacity[\"VAMaxRtg\"]\n    varmaxinjrtg = capacity[\"VarMaxInjRtg\"]\n    varmaxabsrtg = capacity[\"VarMaxAbsRtg\"]\n    wmax = capacity[\"WMax\"]\n    vamax = capacity[\"VAMax\"]\n    varmaxinj = capacity[\"VarMaxInj\"]\n    varmaxabs = capacity[\"VarMaxAbs\"]\n    amax = capacity[\"AMax\"]\n    vnom = capacity[\"VNom\"]\n    vmax = capacity[\"VMax\"]\n    vmin = capacity[\"VMin\"]\n\n    for name, numerator, denominator in [\n        (\"w_over_wmaxrtg\", w, wmaxrtg),\n        (\"w_over_wmax\", w, wmax),\n        (\"va_over_vamax\", va, vamax),\n        (\"va_over_vamaxrtg\", va, vamaxrtg),\n        (\"var_over_injmax\", var, varmaxinj),\n        (\"var_over_absmax\", var, varmaxabs),\n        (\"a_over_amax\", a, amax),\n        (\"llv_over_vnom\", llv, vnom),\n        (\"lnv_over_vnom\", lnv * SQRT3, vnom),\n    ]:\n        numerator = np.asarray(numerator, dtype=np.float32)\n        denominator = np.asarray(denominator, dtype=np.float32)\n        data[name] = np.divide(\n            numerator,\n            denominator,\n            out=np.full_like(numerator, np.nan, dtype=np.float32),\n            where=np.isfinite(numerator) & np.isfinite(denominator) & (np.abs(denominator) > 1e-6),\n        )\n\n    for name, value in [\n        (\"w_minus_wmax\", w - wmax),\n        (\"w_minus_wmaxrtg\", w - wmaxrtg),\n        (\"va_minus_vamax\", va - vamax),\n        (\"var_minus_injmax\", var - varmaxinj),\n        (\"var_plus_absmax\", var + varmaxabs),\n        (\"llv_minus_lnv_sqrt3\", llv - lnv * SQRT3),\n        (\"hz_delta_60\", hz - 60.0),\n    ]:\n        data[name] = value.astype(np.float32)\n\n    for name, left, right in [\n        (\"w_eq_wmaxrtg\", w, wmaxrtg),\n        (\"w_eq_wmax\", w, wmax),\n        (\"var_eq_varmaxinj\", var, varmaxinj),\n        (\"var_eq_neg_varmaxabs\", var, -varmaxabs),\n    ]:\n        data[name] = np.isclose(left, right, equal_nan=False).astype(np.int8)\n    data[\"pf_sign_mismatch\"] = ((np.sign(np.nan_to_num(pf)) != np.sign(np.nan_to_num(w))) & (np.nan_to_num(pf) != 0) & (np.nan_to_num(w) != 0)).astype(np.int8)\n\n    tolw = np.maximum(50.0, 0.02 * np.nan_to_num(wmaxrtg, nan=0.0)).astype(np.float32)\n    tolva = np.maximum(50.0, 0.02 * np.nan_to_num(vamax, nan=0.0)).astype(np.float32)\n    tolvi = np.maximum(20.0, 0.02 * np.nan_to_num(varmaxinj, nan=0.0)).astype(np.float32)\n    tolva2 = np.maximum(20.0, 0.02 * np.nan_to_num(varmaxabs, nan=0.0)).astype(np.float32)\n    for name, value, upper_bound in [\n        (\"w_gt_wmax_tol\", w, wmax + tolw),\n        (\"w_gt_wmaxrtg_tol\", w, wmaxrtg + tolw),\n        (\"va_gt_vamax_tol\", va, vamax + tolva),\n        (\"var_gt_injmax_tol\", var, varmaxinj + tolvi),\n    ]:\n        data[name] = (value > upper_bound).astype(np.int8)\n    data[\"var_lt_absmax_tol\"] = (var < (-varmaxabs - tolva2)).astype(np.int8)\n\n    pq = np.sqrt(np.square(w.astype(np.float32)) + np.square(var.astype(np.float32)))\n    data[\"va_minus_pqmag\"] = (va - pq).astype(np.float32)\n    va32 = np.asarray(va, dtype=np.float32)\n    pq = np.asarray(pq, dtype=np.float32)\n    data[\"va_over_pqmag\"] = np.divide(\n        va32,\n        pq,\n        out=np.full_like(va32, np.nan, dtype=np.float32),\n        where=np.isfinite(va32) & np.isfinite(pq) & (np.abs(pq) > 1e-6),\n    )\n    w32 = np.asarray(w, dtype=np.float32)\n    pf_from_w_va = np.divide(\n        w32,\n        va32,\n        out=np.full_like(w32, np.nan, dtype=np.float32),\n        where=np.isfinite(w32) & np.isfinite(va32) & (np.abs(va32) > 1e-6),\n    )\n    data[\"pf_from_w_va\"] = pf_from_w_va\n    data[\"pf_error\"] = (pf - pf_from_w_va).astype(np.float32)\n\n    for name, total, suffixes in [\n        (\"w_phase_sum_error\", w, [\"WL1\", \"WL2\", \"WL3\"]),\n        (\"va_phase_sum_error\", va, [\"VAL1\", \"VAL2\", \"VAL3\"]),\n        (\"var_phase_sum_error\", var, [\"VarL1\", \"VarL2\", \"VarL3\"]),\n    ]:\n        phase_sum = sum(phase_ac[suffix] for suffix in suffixes)\n        data[name] = (total - phase_sum).astype(np.float32)\n    for name, suffixes in [\n        (\"phase_ll_spread\", [\"VL1L2\", \"VL2L3\", \"VL3L1\"]),\n        (\"phase_ln_spread\", [\"VL1\", \"VL2\", \"VL3\"]),\n        (\"phase_w_spread\", [\"WL1\", \"WL2\", \"WL3\"]),\n        (\"phase_var_spread\", [\"VarL1\", \"VarL2\", \"VarL3\"]),\n    ]:\n        phase_values = np.column_stack([phase_ac[suffix] for suffix in suffixes])\n        data[name] = (nanmax_rows(phase_values) - nanmin_rows(phase_values)).astype(np.float32)\n\n    for name, numerator, denominator in [\n        (\"wmax_over_wmaxrtg\", wmax, wmaxrtg),\n        (\"vamax_over_vamaxrtg\", vamax, vamaxrtg),\n        (\"vmax_over_vnom\", vmax, vnom),\n        (\"vmin_over_vnom\", vmin, vnom),\n    ]:\n        numerator = np.asarray(numerator, dtype=np.float32)\n        denominator = np.asarray(denominator, dtype=np.float32)\n        data[name] = np.divide(\n            numerator,\n            denominator,\n            out=np.full_like(numerator, np.nan, dtype=np.float32),\n            where=np.isfinite(numerator) & np.isfinite(denominator) & (np.abs(denominator) > 1e-6),\n        )\n\n    wsetena = np.nan_to_num(ctl_ac[\"WSetEna\"], nan=0.0)\n    wset = ctl_ac[\"WSet\"]\n    wsetpct = ctl_ac[\"WSetPct\"]\n    wmaxlimena = np.nan_to_num(ctl_ac[\"WMaxLimPctEna\"], nan=0.0)\n    wmaxlimpct = ctl_ac[\"WMaxLimPct\"]\n    varsetena = np.nan_to_num(ctl_ac[\"VarSetEna\"], nan=0.0)\n    varset = ctl_ac[\"VarSet\"]\n    varsetpct = ctl_ac[\"VarSetPct\"]\n    wset_abs_error = np.where(wsetena > 0, np.abs(w - wset), np.nan)\n    wsetpct_target = wmaxrtg * (wsetpct / 100.0)\n    wsetpct_abs_error = np.where(wsetena > 0, np.abs(w - wsetpct_target), np.nan)\n    wmaxlim_target = wmaxrtg * (wmaxlimpct / 100.0)\n    wmaxlim_excess = np.where(wmaxlimena > 0, w - wmaxlim_target, np.nan)\n    varset_abs_error = np.where(varsetena > 0, np.abs(var - varset), np.nan)\n    varsetpct_target = varmaxinj * (varsetpct / 100.0)\n    varsetpct_abs_error = np.where(varsetena > 0, np.abs(var - varsetpct_target), np.nan)\n    for name, value in [\n        (\"wset_abs_error\", wset_abs_error),\n        (\"wsetpct_target\", wsetpct_target),\n        (\"wsetpct_abs_error\", wsetpct_abs_error),\n        (\"wmaxlim_target\", wmaxlim_target),\n        (\"wmaxlim_excess\", wmaxlim_excess),\n        (\"varset_abs_error\", varset_abs_error),\n        (\"varsetpct_target\", varsetpct_target),\n        (\"varsetpct_abs_error\", varsetpct_abs_error),\n    ]:\n        data[name] = value.astype(np.float32)\n    for name, value in [\n        (\"wset_enabled_far\", (wsetena > 0) & (wset_abs_error > np.maximum(50.0, 0.05 * np.nan_to_num(wmaxrtg, nan=0.0)))),\n        (\"wsetpct_enabled_far\", (wsetena > 0) & (wsetpct_abs_error > np.maximum(50.0, 0.05 * np.nan_to_num(wmaxrtg, nan=0.0)))),\n        (\"wmaxlim_enabled_far\", (wmaxlimena > 0) & (wmaxlim_excess > np.maximum(50.0, 0.05 * np.nan_to_num(wmaxrtg, nan=0.0)))),\n        (\"varsetpct_enabled_far\", (varsetena > 0) & (varsetpct_abs_error > np.maximum(20.0, 0.05 * np.nan_to_num(varmaxinj, nan=0.0)))),\n    ]:\n        data[name] = value.astype(np.int8)\n\n    add_capacity_extension_features(\n        data,\n        wmaxrtg=wmaxrtg,\n        wmax=wmax,\n        vamaxrtg=vamaxrtg,\n        vamax=vamax,\n        varmaxinjrtg=varmaxinjrtg,\n        varmaxinj=varmaxinj,\n        varmaxabsrtg=varmaxabsrtg,\n        varmaxabs=varmaxabs,\n        vnomrtg=capacity_ext[\"VNomRtg\"],\n        vnom=vnom,\n        vmaxrtg=capacity_ext[\"VMaxRtg\"],\n        vmax=vmax,\n        vminrtg=capacity_ext[\"VMinRtg\"],\n        vmin=vmin,\n        amaxrtg=capacity_ext[\"AMaxRtg\"],\n        amax=amax,\n        wcha_rtg=capacity_ext[\"WChaRteMaxRtg\"],\n        wdis_rtg=capacity_ext[\"WDisChaRteMaxRtg\"],\n        vacha_rtg=capacity_ext[\"VAChaRteMaxRtg\"],\n        vadis_rtg=capacity_ext[\"VADisChaRteMaxRtg\"],\n        wcha=capacity_ext[\"WChaRteMax\"],\n        wdis=capacity_ext[\"WDisChaRteMax\"],\n        vacha=capacity_ext[\"VAChaRteMax\"],\n        vadis=capacity_ext[\"VADisChaRteMax\"],\n        pfover_rtg=capacity_ext[\"PFOvrExtRtg\"],\n        pfover=capacity_ext[\"PFOvrExt\"],\n        pfunder_rtg=capacity_ext[\"PFUndExtRtg\"],\n        pfunder=capacity_ext[\"PFUndExt\"],\n    )\n\n    llv = np.asarray(llv, dtype=np.float32)\n    lnv = np.asarray(lnv, dtype=np.float32)\n    vnom = np.asarray(vnom, dtype=np.float32)\n    wmaxrtg = np.asarray(wmaxrtg, dtype=np.float32)\n    voltage_pct = 100.0 * np.divide(\n        llv,\n        vnom,\n        out=np.full_like(llv, np.nan, dtype=np.float32),\n        where=np.isfinite(llv) & np.isfinite(vnom) & (np.abs(vnom) > 1e-6),\n    )\n    line_neutral_voltage_pct = 100.0 * np.divide(\n        lnv_scaled,\n        vnom,\n        out=np.full_like(lnv_scaled, np.nan, dtype=np.float32),\n        where=np.isfinite(lnv_scaled) & np.isfinite(vnom) & (np.abs(vnom) > 1e-6),\n    )\n    w_pct = 100.0 * np.divide(\n        w32,\n        wmaxrtg,\n        out=np.full_like(w32, np.nan, dtype=np.float32),\n        where=np.isfinite(w32) & np.isfinite(wmaxrtg) & (np.abs(wmaxrtg) > 1e-6),\n    )\n    reactive_pct = var_pct(var, varmaxinj, varmaxabs)\n\n    for name, value in [\n        (\"voltage_pct\", voltage_pct),\n        (\"line_neutral_voltage_pct\", line_neutral_voltage_pct),\n        (\"w_pct_of_rtg\", w_pct),\n        (\"var_pct_of_limit\", reactive_pct),\n    ]:\n        data[name] = value.astype(np.float32)\n\n    add_enter_service_features(\n        data,\n        df,\n        voltage_pct=voltage_pct,\n        hz=hz,\n        abs_w=abs_w,\n        va=va,\n        a=a,\n        tolw=tolw,\n        tolva=tolva,\n        amax=amax,\n    )\n    add_pf_control_features(\n        data,\n        df,\n        pf=pf,\n        var=var,\n        varmaxinj=varmaxinj,\n        varmaxabs=varmaxabs,\n    )\n\n    trip_outside_flags = []\n    trip_power_flags = []\n    for short_name, (prefix, axis_name, mode) in TRIP_SPECS.items():\n        measure_value = voltage_pct if axis_name == \"V\" else hz\n        outside, power_when_outside = add_trip_block_features(\n            data,\n            df,\n            short_name=short_name,\n            prefix=prefix,\n            axis_name=axis_name,\n            mode=mode,\n            measure_value=measure_value,\n            abs_w=abs_w,\n            tolw=tolw,\n        )\n        trip_outside_flags.append(outside)\n        trip_power_flags.append(power_when_outside)\n    if trip_outside_flags:\n        trip_any_outside = np.column_stack(trip_outside_flags).any(axis=1).astype(np.int8)\n        trip_any_power_when_outside = np.column_stack(trip_power_flags).any(axis=1).astype(np.int8)\n    else:\n        trip_any_outside = np.zeros(len(df), dtype=np.int8)\n        trip_any_power_when_outside = np.zeros(len(df), dtype=np.int8)\n    data[\"trip_any_outside_musttrip\"] = trip_any_outside\n    data[\"trip_any_power_when_outside\"] = trip_any_power_when_outside\n\n    _add_adaptive_curve_features(\n        data,\n        df,\n        voltage_pct=voltage_pct,\n        reactive_pct=reactive_pct,\n        w_pct=w_pct,\n    )\n\n    add_freq_droop_features(data, df, hz=hz, w_pct=w_pct)\n    add_dc_features(data, df, w=w, abs_w=abs_w)\n\n    ac_type = measure_ac[\"ACType\"]\n    data[\"ac_type_is_rare\"] = (np.isfinite(ac_type) & (ac_type == 3.0)).astype(np.int8)\n\n    _apply_hard_rule_features(data, hard_override_names=hard_override_names)\n    return pd.DataFrame(data)\n","metadata":{},"outputs":[],"execution_count":null},{"id":"a001ad62-93bb-4212-be94-01d5be451e92","cell_type":"code","source":"%%writefile src/modeling.py\n\"\"\"Model training and inference for the DER anomaly baseline.\"\"\"\n\nimport gc\nimport logging\nfrom copy import deepcopy\nfrom pathlib import Path\nfrom typing import Dict, Iterator, List, Optional, Sequence, Tuple\n\nimport numpy as np\nimport pandas as pd\nfrom catboost import CatBoostClassifier\nfrom sklearn.metrics import fbeta_score, precision_score, recall_score\nfrom tqdm.auto import tqdm\nfrom xgboost import XGBClassifier, XGBRegressor\n\nfrom src.contracts import (\n    AUDIT_TOLERANCE,\n    BaselineConfig,\n    CANON100_INTERACTION_FEATURES,\n    CANON100_NEGATIVE_WEIGHT,\n    DEVICE_FAMILY_MAP,\n    FAMILY_THRESHOLD_FLOOR,\n    FamilyModelBundle,\n    FamilySemanticContext,\n    HARD_OVERRIDE_TRAIN_WEIGHT,\n    MAX_THRESHOLD,\n    MetricSummary,\n    MIN_OVERRIDE_PRECISION,\n    RESIDUAL_TAIL_FALLBACKS,\n    RESIDUAL_TAIL_LEVELS,\n    SCENARIO_SMOOTHING,\n    ScenarioStats,\n    SURROGATE_LEAKY_FEATURES,\n    SURROGATE_TARGETS,\n)\nfrom src.feature_engineering import CAT_ENGINEERED_COLUMNS, build_features\nfrom src.rules import (\n    DEFAULT_HARD_OVERRIDE_NAMES,\n    build_named_rule_flags,\n    compute_active_override_anomaly,\n)\nfrom src.schema import (\n    RAW_NUMERIC,\n    RAW_STRING_COLUMNS,\n    SAFE_RAW,\n    SAFE_STR,\n    USECOLS_TEST,\n    USECOLS_TRAIN,\n    dedupe,\n)\n\nLOGGER = logging.getLogger(__name__)\n\n\nclass ResearchBaseline:\n    \"\"\"Own the learned per-family model bundles and inference stack.\"\"\"\n\n    def __init__(self, config: BaselineConfig | None = None) -> None:\n        # Keep fixed model knobs grouped together so fit-time state is easier to\n        # distinguish from learned model state.\n        self.config = config or BaselineConfig()\n        self._reset_fit_state()\n\n    def _reset_fit_state(self) -> None:\n        \"\"\"Reset learned state while keeping immutable config untouched.\"\"\"\n\n        self.active_hard_override_names = list(DEFAULT_HARD_OVERRIDE_NAMES)\n        self.family_models: Dict[str, FamilyModelBundle] = {}\n        self.semantic_context = FamilySemanticContext()\n\n    def iter_raw_chunks(\n        self,\n        csv_path: Path,\n        usecols: Sequence[str],\n        limit_rows: int = 0,\n    ) -> Iterator[pd.DataFrame]:\n        yielded = 0\n        for chunk in pd.read_csv(\n            csv_path,\n            usecols=list(usecols),\n            chunksize=self.config.chunksize,\n            low_memory=False,\n        ):\n            if limit_rows > 0:\n                remaining = limit_rows - yielded\n                if remaining <= 0:\n                    break\n                if len(chunk) > remaining:\n                    chunk = chunk.iloc[:remaining].copy()\n            yielded += len(chunk)\n            yield chunk\n            if limit_rows > 0 and yielded >= limit_rows:\n                break\n\n    def _encode_device_family(self, df: pd.DataFrame) -> pd.DataFrame:\n        out = df.copy()\n        out[\"device_family\"] = out[\"device_family\"].map(DEVICE_FAMILY_MAP).fillna(-1).astype(np.int8)\n        return out\n\n    def _get_surrogate_feature_cols(self, columns: Sequence[str]) -> List[str]:\n        excluded = {\n            \"Id\",\n            \"Label\",\n            \"fold_id\",\n            \"audit_fold_id\",\n            \"device_fingerprint\",\n            \"hard_rule_anomaly\",\n            \"hard_rule_count\",\n            \"hard_rule_score\",\n            \"hard_override_anomaly\",\n        }\n        excluded.update(SAFE_STR.values())\n        return [col for col in columns if col not in excluded and col not in SURROGATE_LEAKY_FEATURES]\n\n    def _build_sample_weights(self, x_df: pd.DataFrame, y: np.ndarray) -> np.ndarray:\n        weights = np.ones(len(x_df), dtype=np.float32)\n        family = x_df[\"device_family\"].to_numpy()\n        hard_override = x_df[\"hard_override_anomaly\"].to_numpy() == 1\n        weights[(family == \"canon100\") & (y == 0)] *= CANON100_NEGATIVE_WEIGHT\n        weights[hard_override] *= HARD_OVERRIDE_TRAIN_WEIGHT\n        return weights\n\n    @staticmethod\n    def _bucketize(\n        values: pd.Series,\n        *,\n        fill_value: int | float,\n        dtype: np.dtype,\n        scale: float = 1.0,\n        round_values: bool = True,\n    ) -> pd.Series:\n        out = pd.to_numeric(values, errors=\"coerce\")\n        if scale != 1.0:\n            out = out / scale\n        if round_values:\n            out = out.round()\n        return out.fillna(fill_value).astype(dtype)\n\n    def _build_scenario_frame(self, x_df: pd.DataFrame, *, include_output_bins: bool) -> pd.DataFrame:\n        frame: Dict[str, pd.Series] = {\n            \"family\": x_df[\"device_family\"].astype(str),\n            \"throt_src\": self._bucketize(x_df[\"DERMeasureAC_0_ThrotSrc\"], fill_value=-1, dtype=np.int16),\n            \"throt_pct\": self._bucketize(\n                x_df[\"DERMeasureAC_0_ThrotPct\"],\n                scale=5.0,\n                fill_value=-1,\n                dtype=np.int16,\n            ),\n            \"wmaxlim_pct\": self._bucketize(x_df[\"DERCtlAC_0_WMaxLimPct\"], scale=5.0, fill_value=-1, dtype=np.int16),\n            \"wset_pct\": self._bucketize(x_df[\"DERCtlAC_0_WSetPct\"], scale=5.0, fill_value=-1, dtype=np.int16),\n            \"varset_pct\": self._bucketize(x_df[\"DERCtlAC_0_VarSetPct\"], scale=5.0, fill_value=-1, dtype=np.int16),\n            \"pf_set\": self._bucketize(x_df[\"DERCtlAC_0_PFWInj_PF\"], scale=0.02, fill_value=-1, dtype=np.int16),\n            \"fd_idx\": self._bucketize(x_df[\"DERFreqDroop_0_AdptCtlRslt\"], fill_value=-1, dtype=np.int16),\n            \"vv_idx\": self._bucketize(x_df[\"DERVoltVar_0_AdptCrvRslt\"], fill_value=-1, dtype=np.int16),\n            \"vw_idx\": self._bucketize(x_df[\"DERVoltWatt_0_AdptCrvRslt\"], fill_value=-1, dtype=np.int16),\n            \"wv_idx\": self._bucketize(x_df[\"DERWattVar_0_AdptCrvRslt\"], fill_value=-1, dtype=np.int16),\n            \"volt_bin\": self._bucketize(x_df[\"voltage_pct\"], fill_value=-999, dtype=np.int16),\n            \"hz_bin\": self._bucketize(x_df[\"DERMeasureAC_0_Hz\"], scale=0.1, fill_value=-999, dtype=np.int16),\n            \"enter_idle\": self._bucketize(\n                x_df[\"enter_service_should_idle\"],\n                fill_value=0,\n                dtype=np.int8,\n                round_values=False,\n            ),\n            \"droop_active\": self._bucketize(\n                x_df[\"freqdroop_outside_deadband\"],\n                fill_value=0,\n                dtype=np.int8,\n                round_values=False,\n            ),\n        }\n        if include_output_bins:\n            frame[\"w_bin\"] = self._bucketize(x_df[\"w_pct_of_rtg\"], scale=5.0, fill_value=-999, dtype=np.int16)\n            frame[\"var_bin\"] = self._bucketize(\n                x_df[\"var_pct_of_limit\"],\n                scale=5.0,\n                fill_value=-999,\n                dtype=np.int16,\n            )\n            frame[\"pf_mode\"] = self._bucketize(\n                x_df[\"pf_control_any_enabled\"],\n                fill_value=0,\n                dtype=np.int8,\n                round_values=False,\n            )\n        return pd.DataFrame(frame)\n\n    def _build_scenario_keys(self, x_df: pd.DataFrame, *, include_output_bins: bool = False) -> np.ndarray:\n        return pd.util.hash_pandas_object(self._build_scenario_frame(x_df, include_output_bins=include_output_bins), index=False).to_numpy(np.uint64)\n\n    @staticmethod\n    def _lookup_scenario_stats(\n        keys: np.ndarray,\n        sum_map: Dict[int, float],\n        count_map: Dict[int, int],\n    ) -> Tuple[np.ndarray, np.ndarray]:\n        key_series = pd.Series(keys)\n        sum_values = key_series.map(sum_map).fillna(0.0).to_numpy(np.float32)\n        count_values = key_series.map(count_map).fillna(0).to_numpy(np.int32)\n        return sum_values, count_values\n\n    @staticmethod\n    def _assign_scenario_features(\n        out: pd.DataFrame,\n        *,\n        family_prior: np.ndarray,\n        scenario_rate: np.ndarray,\n        scenario_count: np.ndarray,\n        scenario_output_rate: np.ndarray,\n        scenario_output_count: np.ndarray,\n    ) -> pd.DataFrame:\n        out[\"scenario_rate\"] = scenario_rate.astype(np.float32)\n        out[\"scenario_rate_delta\"] = (scenario_rate - family_prior).astype(np.float32)\n        out[\"scenario_count\"] = scenario_count.astype(np.int32)\n        out[\"scenario_log_count\"] = np.log1p(scenario_count).astype(np.float32)\n        out[\"scenario_low_support\"] = (scenario_count < 20).astype(np.int8)\n        out[\"scenario_output_rate\"] = scenario_output_rate.astype(np.float32)\n        out[\"scenario_output_rate_delta\"] = (scenario_output_rate - family_prior).astype(np.float32)\n        out[\"scenario_output_count\"] = scenario_output_count.astype(np.int32)\n        out[\"scenario_output_log_count\"] = np.log1p(scenario_output_count).astype(np.float32)\n        out[\"scenario_output_low_support\"] = (scenario_output_count < 20).astype(np.int8)\n        return out\n\n    @staticmethod\n    def _group_key_stats(keys: np.ndarray, y_values: np.ndarray) -> pd.DataFrame:\n        return pd.DataFrame({\"key\": keys, \"y\": y_values}).groupby(\"key\")[\"y\"].agg([\"sum\", \"count\"])\n\n    def _family_prior(self, family_series: pd.Series, global_rate: float) -> np.ndarray:\n        return family_series.map(self.semantic_context.family_base_rates).fillna(global_rate).to_numpy(np.float32)\n\n    @staticmethod\n    def _smoothed_rate(\n        sum_values: np.ndarray,\n        count_values: np.ndarray,\n        prior: np.ndarray,\n    ) -> np.ndarray:\n        return (sum_values + SCENARIO_SMOOTHING * prior) / (count_values + SCENARIO_SMOOTHING)\n\n    def _fit_transform_scenario_features(self, x_train: pd.DataFrame, y_train: pd.Series) -> pd.DataFrame:\n        out = x_train.copy()\n        y_arr = y_train.to_numpy(np.float32)\n        family_series = out[\"device_family\"].astype(str)\n        self.semantic_context.family_base_rates = pd.DataFrame({\"family\": family_series, \"y\": y_arr}).groupby(\"family\")[\"y\"].mean().to_dict()\n        keys = self._build_scenario_keys(out)\n        output_keys = self._build_scenario_keys(out, include_output_bins=True)\n        fold_ids = (out[\"Id\"].to_numpy(np.int64) % self.config.cv_folds).astype(np.int8)\n        scenario_rate = np.zeros(len(out), dtype=np.float32)\n        scenario_count = np.zeros(len(out), dtype=np.int32)\n        scenario_output_rate = np.zeros(len(out), dtype=np.float32)\n        scenario_output_count = np.zeros(len(out), dtype=np.int32)\n        global_rate = float(np.mean(y_arr))\n        family_base_rates = self.semantic_context.family_base_rates\n\n        for fold in range(self.config.cv_folds):\n            train_mask = fold_ids != fold\n            valid_mask = fold_ids == fold\n            if not valid_mask.any():\n                continue\n            stats = self._group_key_stats(keys[train_mask], y_arr[train_mask])\n            output_stats = self._group_key_stats(output_keys[train_mask], y_arr[train_mask])\n            valid_keys = pd.Series(keys[valid_mask])\n            valid_sum = valid_keys.map(stats[\"sum\"]).fillna(0.0).to_numpy(np.float32)\n            valid_count = valid_keys.map(stats[\"count\"]).fillna(0).to_numpy(np.int32)\n            valid_output_keys = pd.Series(output_keys[valid_mask])\n            valid_output_sum = valid_output_keys.map(output_stats[\"sum\"]).fillna(0.0).to_numpy(np.float32)\n            valid_output_count = valid_output_keys.map(output_stats[\"count\"]).fillna(0).to_numpy(np.int32)\n            valid_family = family_series.loc[valid_mask].tolist()\n            prior = np.array(\n                [family_base_rates.get(name, global_rate) for name in valid_family],\n                dtype=np.float32,\n            )\n            scenario_rate[valid_mask] = self._smoothed_rate(valid_sum, valid_count, prior)\n            scenario_count[valid_mask] = valid_count\n            scenario_output_rate[valid_mask] = self._smoothed_rate(\n                valid_output_sum,\n                valid_output_count,\n                prior,\n            )\n            scenario_output_count[valid_mask] = valid_output_count\n\n        full_stats = self._group_key_stats(keys, y_arr)\n        full_output_stats = self._group_key_stats(output_keys, y_arr)\n        self.semantic_context.scenario_stats = ScenarioStats(\n            sum_map={int(idx): float(val) for idx, val in full_stats[\"sum\"].items()},\n            count_map={int(idx): int(val) for idx, val in full_stats[\"count\"].items()},\n            output_sum_map={int(idx): float(val) for idx, val in full_output_stats[\"sum\"].items()},\n            output_count_map={int(idx): int(val) for idx, val in full_output_stats[\"count\"].items()},\n        )\n\n        family_prior = self._family_prior(family_series, global_rate)\n        return self._assign_scenario_features(\n            out,\n            family_prior=family_prior,\n            scenario_rate=scenario_rate,\n            scenario_count=scenario_count,\n            scenario_output_rate=scenario_output_rate,\n            scenario_output_count=scenario_output_count,\n        )\n\n    def _apply_scenario_features(self, x_df: pd.DataFrame) -> pd.DataFrame:\n        scenario_stats = self.semantic_context.scenario_stats\n        if not scenario_stats.count_map:\n            return x_df\n\n        out = x_df.copy()\n        keys = self._build_scenario_keys(out)\n        output_keys = self._build_scenario_keys(out, include_output_bins=True)\n        sum_values, count_values = self._lookup_scenario_stats(\n            keys,\n            scenario_stats.sum_map,\n            scenario_stats.count_map,\n        )\n        output_sum_values, output_count_values = self._lookup_scenario_stats(\n            output_keys,\n            scenario_stats.output_sum_map,\n            scenario_stats.output_count_map,\n        )\n        family_base_rates = self.semantic_context.family_base_rates\n        global_rate = float(np.mean(list(family_base_rates.values()))) if family_base_rates else 0.5\n        family_prior = self._family_prior(out[\"device_family\"].astype(str), global_rate)\n        scenario_rate = self._smoothed_rate(sum_values, count_values, family_prior)\n        scenario_output_rate = self._smoothed_rate(\n            output_sum_values,\n            output_count_values,\n            family_prior,\n        )\n        return self._assign_scenario_features(\n            out,\n            family_prior=family_prior,\n            scenario_rate=scenario_rate,\n            scenario_count=count_values,\n            scenario_output_rate=scenario_output_rate,\n            scenario_output_count=output_count_values,\n        )\n\n    def _add_family_interaction_features(self, x_df: pd.DataFrame) -> pd.DataFrame:\n        out = x_df.copy()\n        canon100_mask = out[\"device_family\"].astype(str) == \"canon100\"\n        for feature_name in CANON100_INTERACTION_FEATURES:\n            if feature_name not in out.columns:\n                continue\n            values = pd.to_numeric(out[feature_name], errors=\"coerce\").to_numpy(np.float32)\n            out[f\"canon100_{feature_name}\"] = np.where(canon100_mask.to_numpy(), values, 0.0).astype(np.float32)\n        return out\n\n    def _surrogate_partition_mask(self, ids: Sequence[int], *, fit_partition: bool) -> np.ndarray:\n        ids_arr = np.asarray(ids, dtype=np.int64)\n        fit_mask = (ids_arr % 2) == 0\n        return fit_mask if fit_partition else ~fit_mask\n\n    @staticmethod\n    def _normal_training_mask(\n        x_df: pd.DataFrame,\n        y_train: pd.Series,\n        valid_mask: pd.Series,\n    ) -> np.ndarray:\n        return (y_train.to_numpy(np.int8) == 0) & (x_df[\"hard_override_anomaly\"].to_numpy(np.int8) == 0) & (x_df[\"device_family\"].to_numpy(dtype=object) != \"other\") & (~valid_mask.to_numpy())\n\n    @staticmethod\n    def _initialize_residual_columns(out: pd.DataFrame) -> None:\n        for target_name in SURROGATE_TARGETS:\n            out[f\"pred_{target_name}\"] = np.nan\n            out[f\"resid_{target_name}\"] = np.nan\n            out[f\"abs_resid_{target_name}\"] = np.nan\n            out[f\"norm_resid_{target_name}\"] = np.nan\n            out[f\"abs_norm_resid_{target_name}\"] = np.nan\n            out[f\"tail_resid_{target_name}\"] = 0\n            out[f\"extreme_resid_{target_name}\"] = 0\n            out[f\"ultra_resid_{target_name}\"] = 0\n            out[f\"q99_ratio_resid_{target_name}\"] = np.nan\n\n    @staticmethod\n    def _curve_mode(\n        out: pd.DataFrame,\n        *,\n        enabled_col: str,\n        expected_col: str,\n    ) -> np.ndarray:\n        enabled = np.nan_to_num(out[enabled_col].to_numpy(np.float32), nan=0.0) > 0\n        expected = out[expected_col].to_numpy(np.float32)\n        return enabled & np.isfinite(expected)\n\n    @staticmethod\n    def _cat_categorical_cols(feature_cols: Sequence[str]) -> List[str]:\n        categorical_candidates = [*SAFE_STR.values(), \"device_fingerprint\"]\n        return [col for col in categorical_candidates if col in feature_cols]\n\n    def _fit_surrogate_models(self, x_train: pd.DataFrame, y_train: pd.Series, valid_mask: pd.Series) -> None:\n        self.semantic_context.surrogate_feature_cols = self._get_surrogate_feature_cols(x_train.columns)\n        fit_partition = self._surrogate_partition_mask(x_train[\"Id\"], fit_partition=True)\n        normal_mask = self._normal_training_mask(x_train, y_train, valid_mask) & fit_partition\n        surrogate_df = x_train.loc[normal_mask].copy()\n        if surrogate_df.empty:\n            raise RuntimeError(\"No rows available to train surrogate models.\")\n\n        self.semantic_context.surrogate_models = {}\n        model_params = {\n            \"n_estimators\": max(80, self.config.n_estimators // 2),\n            \"max_depth\": max(4, self.config.max_depth - 2),\n            \"learning_rate\": min(0.08, self.config.learning_rate * 1.2),\n            \"objective\": \"reg:squarederror\",\n            \"subsample\": self.config.subsample,\n            \"colsample_bytree\": self.config.colsample_bytree,\n            \"eval_metric\": \"rmse\",\n            \"tree_method\": \"hist\",\n            \"n_jobs\": self.config.n_jobs,\n            \"random_state\": self.config.seed,\n            \"seed\": self.config.seed,\n            \"verbosity\": 0,\n        }\n        for family in DEVICE_FAMILY_MAP:\n            family_df = surrogate_df.loc[surrogate_df[\"device_family\"] == family].copy()\n            if family_df.empty:\n                continue\n            x_surrogate = self._encode_device_family(family_df[self.semantic_context.surrogate_feature_cols])\n            for target_name, (target_col, _) in SURROGATE_TARGETS.items():\n                model = XGBRegressor(**model_params)\n                y_target = family_df[target_col].to_numpy(np.float32)\n                LOGGER.info(f\"[surrogate] training {family}/{target_name} on {len(family_df):,} normal rows\")\n                model.fit(x_surrogate, y_target)\n                self.semantic_context.surrogate_models[(family, target_name)] = model\n\n    def _augment_with_surrogates(self, x_df: pd.DataFrame) -> pd.DataFrame:\n        if not self.semantic_context.surrogate_feature_cols or not self.semantic_context.surrogate_models:\n            return x_df\n\n        out = x_df.copy()\n        self._initialize_residual_columns(out)\n        x_surrogate = self._encode_device_family(out[self.semantic_context.surrogate_feature_cols])\n        for family in DEVICE_FAMILY_MAP:\n            family_mask = out[\"device_family\"] == family\n            if not family_mask.any():\n                continue\n            x_family = x_surrogate.loc[family_mask]\n            for target_name, (target_col, scale_col) in SURROGATE_TARGETS.items():\n                model = self.semantic_context.surrogate_models.get((family, target_name))\n                if model is None:\n                    continue\n                pred = model.predict(x_family).astype(np.float32)\n                actual = out.loc[family_mask, target_col].to_numpy(np.float32)\n                resid = actual - pred\n                out.loc[family_mask, f\"pred_{target_name}\"] = pred\n                out.loc[family_mask, f\"resid_{target_name}\"] = resid\n                out.loc[family_mask, f\"abs_resid_{target_name}\"] = np.abs(resid).astype(np.float32)\n                if scale_col is not None:\n                    scale = out.loc[family_mask, scale_col].to_numpy(np.float32)\n                    norm_resid = np.divide(\n                        resid,\n                        scale,\n                        out=np.full_like(resid, np.nan, dtype=np.float32),\n                        where=np.isfinite(resid) & np.isfinite(scale) & (np.abs(scale) > 1e-6),\n                    )\n                else:\n                    scale = np.maximum(0.05, np.abs(actual))\n                    norm_resid = (resid / scale).astype(np.float32)\n                out.loc[family_mask, f\"norm_resid_{target_name}\"] = norm_resid.astype(np.float32)\n                out.loc[family_mask, f\"abs_norm_resid_{target_name}\"] = np.abs(norm_resid).astype(np.float32)\n\n        abs_resid_cols = [\n            \"abs_resid_w\",\n            \"abs_resid_va\",\n            \"abs_resid_var\",\n            \"abs_resid_pf\",\n            \"abs_resid_a\",\n        ]\n        out[\"resid_energy_total\"] = out[abs_resid_cols].sum(axis=1).astype(np.float32)\n        out[\"resid_va_minus_pq\"] = (out[\"pred_va\"] - np.sqrt(np.square(out[\"pred_w\"]) + np.square(out[\"pred_var\"]))).astype(np.float32)\n        abs_resid_w = out[\"abs_resid_w\"].to_numpy(np.float32)\n        abs_resid_var = out[\"abs_resid_var\"].to_numpy(np.float32) + 1e-3\n        out[\"resid_w_var_ratio\"] = np.divide(\n            abs_resid_w,\n            abs_resid_var,\n            out=np.full_like(abs_resid_w, np.nan, dtype=np.float32),\n            where=np.isfinite(abs_resid_w) & np.isfinite(abs_resid_var) & (np.abs(abs_resid_var) > 1e-6),\n        )\n        return out\n\n    def _compute_residual_quantiles(\n        self,\n        x_train: pd.DataFrame,\n        y_train: pd.Series,\n        valid_mask: pd.Series,\n    ) -> None:\n        calibration_partition = self._surrogate_partition_mask(x_train[\"Id\"], fit_partition=False)\n        base_mask = self._normal_training_mask(x_train, y_train, valid_mask)\n        self.semantic_context.residual_quantiles = {}\n        for family in DEVICE_FAMILY_MAP:\n            family_mask = base_mask & (x_train[\"device_family\"] == family)\n            family_calibration = family_mask & calibration_partition\n            if not family_calibration.any():\n                family_calibration = family_mask\n            family_quantiles: Dict[str, Dict[str, float]] = {}\n            for target_name in SURROGATE_TARGETS:\n                series = x_train.loc[family_calibration, f\"abs_norm_resid_{target_name}\"]\n                values = pd.to_numeric(series, errors=\"coerce\").to_numpy(np.float32)\n                values = values[np.isfinite(values)]\n                quantiles = RESIDUAL_TAIL_FALLBACKS.copy()\n                if values.size > 0:\n                    for level_name, q in RESIDUAL_TAIL_LEVELS.items():\n                        quantiles[level_name] = float(np.quantile(values, q))\n                family_quantiles[target_name] = {key: max(1e-6, value) for key, value in quantiles.items()}\n            self.semantic_context.residual_quantiles[family] = family_quantiles\n\n    def _apply_residual_calibration_features(self, x_df: pd.DataFrame) -> pd.DataFrame:\n        residual_quantiles = self.semantic_context.residual_quantiles\n        if not residual_quantiles:\n            return x_df\n\n        out = x_df.copy()\n        calibration_levels = (\n            (\"tail\", 0),\n            (\"extreme\", 0),\n            (\"ultra\", 0),\n            (\"q99_ratio\", np.nan),\n        )\n        for target_name in SURROGATE_TARGETS:\n            for prefix, default_value in calibration_levels:\n                out[f\"{prefix}_resid_{target_name}\"] = default_value\n\n        for family in DEVICE_FAMILY_MAP:\n            family_mask = out[\"device_family\"] == family\n            if not family_mask.any():\n                continue\n            family_quantiles = residual_quantiles.get(family, {})\n            for target_name in SURROGATE_TARGETS:\n                abs_norm = out.loc[family_mask, f\"abs_norm_resid_{target_name}\"].to_numpy(np.float32)\n                q = family_quantiles.get(target_name, RESIDUAL_TAIL_FALLBACKS)\n                calibrated_values = (\n                    (f\"tail_resid_{target_name}\", abs_norm >= q[\"tail\"], np.int8),\n                    (f\"extreme_resid_{target_name}\", abs_norm >= q[\"extreme\"], np.int8),\n                    (f\"ultra_resid_{target_name}\", abs_norm >= q[\"ultra\"], np.int8),\n                    (\n                        f\"q99_ratio_resid_{target_name}\",\n                        np.divide(\n                            abs_norm,\n                            np.full_like(abs_norm, q[\"extreme\"], dtype=np.float32),\n                            out=np.full_like(abs_norm, np.nan, dtype=np.float32),\n                            where=np.isfinite(abs_norm) & np.isfinite(q[\"extreme\"]) & (abs(q[\"extreme\"]) > 1e-6),\n                        ),\n                        np.float32,\n                    ),\n                )\n                for column_name, values, dtype in calibrated_values:\n                    out.loc[family_mask, column_name] = values.astype(dtype)\n\n        abs_norm_values = {\n            \"w\": np.nan_to_num(out[\"abs_norm_resid_w\"].to_numpy(np.float32), nan=0.0),\n            \"var\": np.nan_to_num(out[\"abs_norm_resid_var\"].to_numpy(np.float32), nan=0.0),\n            \"pf\": np.nan_to_num(out[\"abs_norm_resid_pf\"].to_numpy(np.float32), nan=0.0),\n            \"a\": np.nan_to_num(out[\"abs_norm_resid_a\"].to_numpy(np.float32), nan=0.0),\n        }\n        mode_flags = {\n            \"pf\": np.nan_to_num(out[\"pf_control_any_enabled\"].to_numpy(np.float32), nan=0.0) > 0,\n            \"voltvar\": self._curve_mode(\n                out,\n                enabled_col=\"DERVoltVar_0_Ena\",\n                expected_col=\"voltvar_curve_expected\",\n            ),\n            \"voltwatt\": self._curve_mode(\n                out,\n                enabled_col=\"DERVoltWatt_0_Ena\",\n                expected_col=\"voltwatt_curve_expected\",\n            ),\n            \"wattvar\": self._curve_mode(\n                out,\n                enabled_col=\"DERWattVar_0_Ena\",\n                expected_col=\"wattvar_curve_expected\",\n            ),\n            \"droop\": np.nan_to_num(out[\"freqdroop_outside_deadband\"].to_numpy(np.float32), nan=0.0) > 0,\n            \"enter_idle\": np.nan_to_num(out[\"enter_service_should_idle\"].to_numpy(np.float32), nan=0.0) > 0,\n        }\n        mode_flags[\"curve_var\"] = mode_flags[\"voltvar\"] | mode_flags[\"wattvar\"] | mode_flags[\"pf\"]\n        mode_flags[\"dispatch\"] = mode_flags[\"voltwatt\"] | mode_flags[\"droop\"] | mode_flags[\"enter_idle\"]\n        extreme_residuals = {\n            \"var\": np.nan_to_num(out[\"extreme_resid_var\"].to_numpy(np.float32), nan=0.0),\n            \"w\": np.nan_to_num(out[\"extreme_resid_w\"].to_numpy(np.float32), nan=0.0),\n        }\n\n        mode_residual_specs = (\n            (\"mode_resid_pf_pf\", \"pf\", \"pf\"),\n            (\"mode_resid_var_pf\", \"var\", \"pf\"),\n            (\"mode_resid_var_voltvar\", \"var\", \"voltvar\"),\n            (\"mode_resid_w_voltwatt\", \"w\", \"voltwatt\"),\n            (\"mode_resid_var_wattvar\", \"var\", \"wattvar\"),\n            (\"mode_resid_w_droop\", \"w\", \"droop\"),\n            (\"mode_resid_w_enter_idle\", \"w\", \"enter_idle\"),\n            (\"mode_resid_a_enter_idle\", \"a\", \"enter_idle\"),\n            (\"mode_curve_var_resid\", \"var\", \"curve_var\"),\n            (\"mode_dispatch_w_resid\", \"w\", \"dispatch\"),\n        )\n        for column_name, residual_name, mode_name in mode_residual_specs:\n            out[column_name] = (abs_norm_values[residual_name] * mode_flags[mode_name]).astype(np.float32)\n\n        extreme_mode_specs = (\n            (\"mode_extreme_var_curve\", \"var\", \"curve_var\"),\n            (\"mode_extreme_w_dispatch\", \"w\", \"dispatch\"),\n        )\n        for column_name, residual_name, mode_name in extreme_mode_specs:\n            out[column_name] = (extreme_residuals[residual_name] * mode_flags[mode_name]).astype(np.int8)\n\n        out[\"mode_tail_count\"] = out[[\"mode_extreme_var_curve\", \"mode_extreme_w_dispatch\"]].sum(axis=1).astype(np.int8)\n        count_specs = (\n            (\"resid_tail_count\", \"tail\"),\n            (\"resid_extreme_count\", \"extreme\"),\n            (\"resid_ultra_count\", \"ultra\"),\n        )\n        for column_name, prefix in count_specs:\n            out[column_name] = out[[f\"{prefix}_resid_{name}\" for name in SURROGATE_TARGETS]].sum(axis=1).astype(np.int8)\n        out[\"resid_quantile_score\"] = out[[f\"q99_ratio_resid_{name}\" for name in SURROGATE_TARGETS]].sum(axis=1).astype(np.float32)\n        return out\n\n    @staticmethod\n    def tune_threshold(\n        y_true: np.ndarray,\n        prob: np.ndarray,\n        *,\n        low: float = FAMILY_THRESHOLD_FLOOR,\n        high: float = MAX_THRESHOLD,\n        step: float = 0.01,\n    ) -> Tuple[float, float]:\n        if len(y_true) == 0:\n            return 0.5, 0.0\n        best_thr, best_f2 = 0.5, -1.0\n        thresholds = np.arange(low, high + 1e-9, step, dtype=np.float32)\n        for thr in thresholds:\n            pred = (prob >= thr).astype(np.int8)\n            score = fbeta_score(y_true, pred, beta=2)\n            if score > best_f2:\n                best_thr, best_f2 = float(thr), float(score)\n        return best_thr, best_f2\n\n    @staticmethod\n    def _metric_summary_from_pred(y_true: np.ndarray, pred: np.ndarray) -> MetricSummary:\n        if len(y_true) == 0:\n            return MetricSummary(f2=0.0, precision=0.0, recall=0.0, positive_rate=0.0, rows=0)\n        return MetricSummary(\n            f2=float(fbeta_score(y_true, pred, beta=2)),\n            precision=float(precision_score(y_true, pred, zero_division=0)),\n            recall=float(recall_score(y_true, pred, zero_division=0)),\n            positive_rate=float(np.mean(pred)),\n            rows=int(len(y_true)),\n        )\n\n    @staticmethod\n    def _blend_probs(primary: np.ndarray, secondary: Optional[np.ndarray], weight: float) -> np.ndarray:\n        if secondary is None:\n            return primary.astype(np.float32)\n        return (weight * primary + (1.0 - weight) * secondary).astype(np.float32)\n\n    @staticmethod\n    def _select_nonconstant_columns(df: pd.DataFrame, candidates: Sequence[str]) -> List[str]:\n        keep: List[str] = []\n        for col in candidates:\n            if col not in df.columns:\n                continue\n            series = df[col]\n            if series.notna().sum() == 0:\n                continue\n            if series.nunique(dropna=True) <= 1:\n                continue\n            keep.append(col)\n        return keep\n\n    def _build_train_chunk_features(self, chunk: pd.DataFrame) -> Tuple[pd.DataFrame, np.ndarray]:\n        labels = chunk[\"Label\"].astype(np.int8).to_numpy()\n        features = build_features(\n            chunk.drop(columns=[\"Label\"]),\n            hard_override_names=DEFAULT_HARD_OVERRIDE_NAMES,\n        )\n        return features, labels\n\n    def _audit_hard_override_rules(self, train_path: Path) -> Tuple[Dict[str, int], List[str], pd.DataFrame]:\n        row_counts = {\"canon10\": 0, \"canon100\": 0, \"other\": 0}\n        total_positive = 0\n        total_rows = 0\n        per_rule_counts = {name: {\"count\": 0, \"positives\": 0} for name in DEFAULT_HARD_OVERRIDE_NAMES}\n        other_frames: List[pd.DataFrame] = []\n\n        LOGGER.info(\"[fit] auditing hard override rules from %s\", train_path)\n        with tqdm(\n            self.iter_raw_chunks(train_path, USECOLS_TRAIN, 0),\n            desc=\"audit chunks\",\n            unit=\"chunk\",\n            dynamic_ncols=True,\n        ) as progress:\n            for chunk in progress:\n                feats, labels = self._build_train_chunk_features(chunk)\n                rule_flags = build_named_rule_flags(feats)\n                total_rows += len(feats)\n                total_positive += int(labels.sum())\n                family_values = feats[\"device_family\"]\n                for family in row_counts:\n                    family_mask = family_values == family\n                    family_count = int(family_mask.sum())\n                    if family_count == 0:\n                        continue\n                    row_counts[family] += family_count\n                    if family == \"other\":\n                        other_idx = family_mask.to_numpy()\n                        other_frames.append(\n                            pd.DataFrame(\n                                {\n                                    \"Id\": feats.loc[family_mask, \"Id\"].to_numpy(np.int64),\n                                    \"Label\": labels[other_idx],\n                                }\n                            )\n                        )\n                for rule_name in DEFAULT_HARD_OVERRIDE_NAMES:\n                    mask = rule_flags[rule_name]\n                    if not mask.any():\n                        continue\n                    per_rule_counts[rule_name][\"count\"] += int(mask.sum())\n                    per_rule_counts[rule_name][\"positives\"] += int(labels[mask].sum())\n                progress.set_postfix(rows=f\"{total_rows:,}\", positives=f\"{total_positive:,}\")\n                del feats\n                gc.collect()\n\n        demoted: List[str] = []\n        for rule_name, counts in per_rule_counts.items():\n            count = counts[\"count\"]\n            positives = counts[\"positives\"]\n            precision = float(positives / count) if count else 1.0\n            recall = float(positives / total_positive) if total_positive else 0.0\n            LOGGER.info(\n                \"[fit] hard override rule=%s count=%s positives=%s precision=%.4f recall=%.4f\",\n                rule_name,\n                f\"{count:,}\",\n                f\"{positives:,}\",\n                precision,\n                recall,\n            )\n            if count > 0 and precision < MIN_OVERRIDE_PRECISION:\n                demoted.append(rule_name)\n        self.active_hard_override_names = [name for name in DEFAULT_HARD_OVERRIDE_NAMES if name not in demoted]\n        LOGGER.info(\"[fit] hard override audit active=%s demoted=%s\", self.active_hard_override_names, demoted)\n        other_df = pd.concat(other_frames, ignore_index=True) if other_frames else pd.DataFrame(columns=[\"Id\", \"Label\"])\n        return row_counts, demoted, other_df\n\n    def _load_family_training_frame(self, train_path: Path, family: str) -> pd.DataFrame:\n        frames: List[pd.DataFrame] = []\n        total_rows = 0\n        LOGGER.info(\"[fit] materializing %s rows from %s\", family, train_path)\n        with tqdm(\n            self.iter_raw_chunks(train_path, USECOLS_TRAIN, 0),\n            desc=f\"{family} chunks\",\n            unit=\"chunk\",\n            dynamic_ncols=True,\n        ) as progress:\n            for chunk in progress:\n                feats, labels = self._build_train_chunk_features(chunk)\n                family_mask = feats[\"device_family\"] == family\n                if not family_mask.any():\n                    continue\n                family_idx = family_mask.to_numpy()\n                family_df = feats.loc[family_mask].copy()\n                family_df[\"Label\"] = labels[family_idx]\n                family_df[\"fold_id\"] = (family_df[\"Id\"].to_numpy(np.int64) % self.config.cv_folds).astype(np.int8)\n                family_df[\"audit_fold_id\"] = (self._build_scenario_keys(family_df) % self.config.cv_folds).astype(np.int8)\n                frames.append(family_df)\n                total_rows += len(family_df)\n                progress.set_postfix(rows=f\"{total_rows:,}\")\n                del feats\n                gc.collect()\n        if not frames:\n            return pd.DataFrame()\n        return pd.concat(frames, ignore_index=True)\n\n    def _refresh_override_columns(self, df: pd.DataFrame) -> pd.DataFrame:\n        out = df.copy()\n        rule_flags = build_named_rule_flags(out)\n        row_count = len(out)\n        out[\"hard_override_anomaly\"] = compute_active_override_anomaly(\n            out,\n            self.active_hard_override_names,\n            rule_flags=rule_flags,\n            row_count=row_count,\n        )\n        out[\"hard_rule_anomaly\"] = compute_active_override_anomaly(\n            out,\n            DEFAULT_HARD_OVERRIDE_NAMES,\n            rule_flags=rule_flags,\n            row_count=row_count,\n        )\n        return out\n\n    def _capture_semantic_context(self) -> FamilySemanticContext:\n        return FamilySemanticContext(\n            surrogate_feature_cols=list(self.semantic_context.surrogate_feature_cols),\n            surrogate_models=dict(self.semantic_context.surrogate_models),\n            residual_quantiles=deepcopy(self.semantic_context.residual_quantiles),\n            family_base_rates=dict(self.semantic_context.family_base_rates),\n            scenario_stats=ScenarioStats(\n                sum_map=dict(self.semantic_context.scenario_stats.sum_map),\n                count_map=dict(self.semantic_context.scenario_stats.count_map),\n                output_sum_map=dict(self.semantic_context.scenario_stats.output_sum_map),\n                output_count_map=dict(self.semantic_context.scenario_stats.output_count_map),\n            ),\n        )\n\n    def _activate_semantic_context(self, context: FamilySemanticContext) -> None:\n        self.semantic_context = deepcopy(context)\n\n    def _apply_semantic_transforms(\n        self,\n        augmented_df: pd.DataFrame,\n        *,\n        y: Optional[pd.Series] = None,\n    ) -> pd.DataFrame:\n        out = self._apply_residual_calibration_features(augmented_df)\n        if y is None:\n            out = self._apply_scenario_features(out)\n        else:\n            out = self._fit_transform_scenario_features(out, y)\n        return self._add_family_interaction_features(out)\n\n    def _prepare_semantic_frame(self, base_df: pd.DataFrame) -> pd.DataFrame:\n        work = self._refresh_override_columns(base_df)\n        work = self._augment_with_surrogates(work)\n        return self._apply_semantic_transforms(work)\n\n    def _prepare_family_semantic_frame(\n        self,\n        base_df: pd.DataFrame,\n        y: pd.Series,\n    ) -> Tuple[pd.DataFrame, FamilySemanticContext]:\n        work = self._refresh_override_columns(base_df)\n        no_valid = pd.Series(np.zeros(len(work), dtype=bool), index=work.index)\n        self._fit_surrogate_models(work, y, no_valid)\n        work = self._augment_with_surrogates(work)\n        self._compute_residual_quantiles(work, y, no_valid)\n        work = self._apply_semantic_transforms(work, y=y)\n        return work, self._capture_semantic_context()\n\n    def _semantic_feature_candidates(self, semantic_df: pd.DataFrame) -> List[str]:\n        excluded = {\n            \"Id\",\n            \"Label\",\n            \"fold_id\",\n            \"audit_fold_id\",\n            \"hard_override_anomaly\",\n            \"device_fingerprint\",\n        }\n        excluded.update(SAFE_STR.values())\n        return [col for col in semantic_df.columns if col not in excluded and pd.api.types.is_numeric_dtype(semantic_df[col])]\n\n    def _prepare_cat_frame(self, base_df: pd.DataFrame) -> pd.DataFrame:\n        out = self._refresh_override_columns(base_df)\n        for col in [*SAFE_STR.values(), \"device_fingerprint\"]:\n            if col in out.columns:\n                out[col] = out[col].fillna(\"<NA>\").astype(str)\n        return out\n\n    def _cat_feature_candidates(self, cat_df: pd.DataFrame) -> List[str]:\n        raw_numeric_cols = [SAFE_RAW[col] for col in RAW_NUMERIC if SAFE_RAW[col] in cat_df.columns]\n        missing_cols = [col for col in cat_df.columns if col.startswith(\"missing_\")]\n        categorical_cols = [SAFE_STR[col] for col in RAW_STRING_COLUMNS if SAFE_STR[col] in cat_df.columns]\n        candidates = dedupe(\n            [\n                *raw_numeric_cols,\n                *categorical_cols,\n                \"device_fingerprint\",\n                *CAT_ENGINEERED_COLUMNS,\n                *missing_cols,\n            ]\n        )\n        excluded = {\n            \"Id\",\n            \"Label\",\n            \"fold_id\",\n            \"audit_fold_id\",\n            \"hard_override_anomaly\",\n            \"hard_rule_anomaly\",\n        }\n        return [col for col in candidates if col in cat_df.columns and col not in excluded]\n\n    def _train_semantic_oof(\n        self,\n        semantic_df: pd.DataFrame,\n        y: np.ndarray,\n        feature_cols: Sequence[str],\n        *,\n        fold_col: str,\n        fit_final: bool,\n    ) -> Tuple[np.ndarray, Optional[XGBClassifier]]:\n        probs = np.ones(len(semantic_df), dtype=np.float32)\n        model_mask = semantic_df[\"hard_override_anomaly\"].to_numpy(np.int8) == 0\n        fold_ids = semantic_df[fold_col].to_numpy(np.int8)\n        final_model: Optional[XGBClassifier] = None\n        model_params = {\n            \"n_estimators\": self.config.n_estimators,\n            \"max_depth\": self.config.max_depth,\n            \"learning_rate\": self.config.learning_rate,\n            \"objective\": \"binary:logistic\",\n            \"subsample\": self.config.subsample,\n            \"colsample_bytree\": self.config.colsample_bytree,\n            \"eval_metric\": \"logloss\",\n            \"tree_method\": \"hist\",\n            \"n_jobs\": self.config.n_jobs,\n            \"random_state\": self.config.seed,\n            \"seed\": self.config.seed,\n            \"verbosity\": 1,\n        }\n        for fold in range(self.config.cv_folds):\n            train_mask = model_mask & (fold_ids != fold)\n            valid_mask = model_mask & (fold_ids == fold)\n            if not valid_mask.any():\n                continue\n            model = XGBClassifier(**model_params)\n            x_train = semantic_df.loc[train_mask, feature_cols]\n            y_train = y[train_mask]\n            weights = self._build_sample_weights(semantic_df.loc[train_mask], y_train)\n            model.fit(x_train, y_train, sample_weight=weights)\n            probs[valid_mask] = model.predict_proba(semantic_df.loc[valid_mask, feature_cols])[:, 1].astype(np.float32)\n        if fit_final and model_mask.any():\n            final_model = XGBClassifier(**model_params)\n            weights = self._build_sample_weights(semantic_df.loc[model_mask], y[model_mask])\n            final_model.fit(\n                semantic_df.loc[model_mask, feature_cols],\n                y[model_mask],\n                sample_weight=weights,\n            )\n        return probs, final_model\n\n    def _train_cat_oof(\n        self,\n        cat_df: pd.DataFrame,\n        y: np.ndarray,\n        feature_cols: Sequence[str],\n        categorical_cols: Sequence[str],\n        *,\n        fold_col: str,\n        fit_final: bool,\n    ) -> Tuple[np.ndarray, Optional[CatBoostClassifier]]:\n        probs = np.ones(len(cat_df), dtype=np.float32)\n        model_mask = cat_df[\"hard_override_anomaly\"].to_numpy(np.int8) == 0\n        fold_ids = cat_df[fold_col].to_numpy(np.int8)\n        final_model: Optional[CatBoostClassifier] = None\n        cat_features = list(categorical_cols)\n        model_params = {\n            \"iterations\": self.config.cat_iterations,\n            \"depth\": self.config.cat_depth,\n            \"learning_rate\": self.config.cat_learning_rate,\n            \"loss_function\": \"Logloss\",\n            \"eval_metric\": \"Logloss\",\n            \"random_seed\": self.config.seed,\n            \"thread_count\": self.config.n_jobs,\n            \"allow_writing_files\": False,\n            \"verbose\": False,\n        }\n        for fold in range(self.config.cv_folds):\n            train_mask = model_mask & (fold_ids != fold)\n            valid_mask = model_mask & (fold_ids == fold)\n            if not valid_mask.any():\n                continue\n            model = CatBoostClassifier(**model_params)\n            weights = self._build_sample_weights(cat_df.loc[train_mask], y[train_mask])\n            model.fit(\n                cat_df.loc[train_mask, feature_cols],\n                y[train_mask],\n                cat_features=cat_features,\n                sample_weight=weights,\n                verbose=False,\n            )\n            probs[valid_mask] = model.predict_proba(cat_df.loc[valid_mask, feature_cols])[:, 1].astype(np.float32)\n        if fit_final and model_mask.any():\n            final_model = CatBoostClassifier(**model_params)\n            weights = self._build_sample_weights(cat_df.loc[model_mask], y[model_mask])\n            final_model.fit(\n                cat_df.loc[model_mask, feature_cols],\n                y[model_mask],\n                cat_features=cat_features,\n                sample_weight=weights,\n                verbose=False,\n            )\n        return probs, final_model\n\n    def _select_family_blend(\n        self,\n        y: np.ndarray,\n        hard_override: np.ndarray,\n        semantic_primary: np.ndarray,\n        semantic_audit: np.ndarray,\n        cat_primary: Optional[np.ndarray],\n        cat_audit: Optional[np.ndarray],\n    ) -> Tuple[float, float, np.ndarray, np.ndarray]:\n        baseline_primary_prob = semantic_primary.copy()\n        baseline_primary_prob[hard_override] = 1.0\n        baseline_thr, _ = self.tune_threshold(y, baseline_primary_prob)\n        baseline_pred_primary = (baseline_primary_prob >= baseline_thr).astype(np.int8)\n        baseline_audit_prob = semantic_audit.copy()\n        baseline_audit_prob[hard_override] = 1.0\n        baseline_pred_audit = (baseline_audit_prob >= baseline_thr).astype(np.int8)\n        baseline_primary_score = float(fbeta_score(y, baseline_pred_primary, beta=2))\n        baseline_audit_score = float(fbeta_score(y, baseline_pred_audit, beta=2))\n\n        best_weight = 1.0\n        best_thr = baseline_thr\n        best_primary_pred = baseline_pred_primary\n        best_audit_pred = baseline_pred_audit\n        best_primary_score = baseline_primary_score\n        best_audit_for_best = baseline_audit_score\n\n        weight_grid = [round(step / 20.0, 2) for step in range(21)] if cat_primary is not None else [1.0]\n        for weight in weight_grid:\n            blended_primary = self._blend_probs(semantic_primary, cat_primary, weight)\n            blended_primary[hard_override] = 1.0\n            thr, _ = self.tune_threshold(y, blended_primary)\n            pred_primary = (blended_primary >= thr).astype(np.int8)\n            blended_audit = self._blend_probs(semantic_audit, cat_audit, weight)\n            blended_audit[hard_override] = 1.0\n            pred_audit = (blended_audit >= thr).astype(np.int8)\n            primary_score = float(fbeta_score(y, pred_primary, beta=2))\n            audit_score = float(fbeta_score(y, pred_audit, beta=2))\n            if audit_score < baseline_audit_score - AUDIT_TOLERANCE:\n                continue\n            improved_primary = primary_score > best_primary_score + 1e-12\n            tied_primary = abs(primary_score - best_primary_score) <= 1e-12\n            improved_audit = audit_score > best_audit_for_best\n            if improved_primary or (tied_primary and improved_audit):\n                best_weight = weight\n                best_thr = thr\n                best_primary_pred = pred_primary\n                best_audit_pred = pred_audit\n                best_primary_score = primary_score\n                best_audit_for_best = audit_score\n        return (\n            best_weight,\n            best_thr,\n            best_primary_pred.astype(np.int8),\n            best_audit_pred.astype(np.int8),\n        )\n\n    def fit(self, train_path: Path) -> \"ResearchBaseline\":\n        self._reset_fit_state()\n        row_counts, demoted_rules, other_df = self._audit_hard_override_rules(train_path)\n        LOGGER.info(\"[fit] training row counts=%s\", row_counts)\n        primary_metrics: Dict[str, MetricSummary] = {}\n        audit_metrics: Dict[str, MetricSummary] = {}\n        prediction_rows: List[pd.DataFrame] = []\n\n        LOGGER.info(\"[fit] starting family training\")\n        with tqdm(\n            list(DEVICE_FAMILY_MAP),\n            desc=\"fit families\",\n            unit=\"family\",\n            dynamic_ncols=True,\n        ) as family_progress:\n            for family in family_progress:\n                family_progress.set_postfix(family=family)\n                if row_counts[family] == 0:\n                    LOGGER.info(\"[fit] skipping %s; no training rows\", family)\n                    continue\n                base_df = self._load_family_training_frame(train_path, family)\n                if base_df.empty:\n                    LOGGER.info(\"[fit] skipping %s; no training rows\", family)\n                    continue\n                LOGGER.info(\"[fit] training %s on %s rows\", family, f\"{len(base_df):,}\")\n                y_series = base_df[\"Label\"].astype(np.int8)\n                semantic_df, context = self._prepare_family_semantic_frame(base_df.copy(), y_series)\n                semantic_feature_cols = self._select_nonconstant_columns(\n                    semantic_df,\n                    self._semantic_feature_candidates(semantic_df),\n                )\n\n                cat_df = self._prepare_cat_frame(base_df.copy())\n                cat_feature_cols = self._select_nonconstant_columns(cat_df, self._cat_feature_candidates(cat_df))\n                cat_categorical_cols = self._cat_categorical_cols(cat_feature_cols)\n\n                y = y_series.to_numpy(np.int8)\n                hard_override = semantic_df[\"hard_override_anomaly\"].to_numpy(np.int8) == 1\n\n                semantic_primary_prob, semantic_model = self._train_semantic_oof(\n                    semantic_df,\n                    y,\n                    semantic_feature_cols,\n                    fold_col=\"fold_id\",\n                    fit_final=True,\n                )\n                semantic_audit_prob, _ = self._train_semantic_oof(\n                    semantic_df,\n                    y,\n                    semantic_feature_cols,\n                    fold_col=\"audit_fold_id\",\n                    fit_final=False,\n                )\n\n                cat_primary_prob: Optional[np.ndarray] = None\n                cat_audit_prob: Optional[np.ndarray] = None\n                cat_model: Optional[CatBoostClassifier] = None\n                if cat_feature_cols:\n                    cat_primary_prob, cat_model = self._train_cat_oof(\n                        cat_df,\n                        y,\n                        cat_feature_cols,\n                        cat_categorical_cols,\n                        fold_col=\"fold_id\",\n                        fit_final=True,\n                    )\n                    cat_audit_prob, _ = self._train_cat_oof(\n                        cat_df,\n                        y,\n                        cat_feature_cols,\n                        cat_categorical_cols,\n                        fold_col=\"audit_fold_id\",\n                        fit_final=False,\n                    )\n\n                weight, threshold, primary_pred, audit_pred = self._select_family_blend(\n                    y,\n                    hard_override,\n                    semantic_primary_prob,\n                    semantic_audit_prob,\n                    cat_primary_prob,\n                    cat_audit_prob,\n                )\n                self.family_models[family] = FamilyModelBundle(\n                    semantic_model=semantic_model,\n                    cat_model=cat_model,\n                    semantic_context=context,\n                    semantic_feature_cols=semantic_feature_cols,\n                    cat_feature_cols=cat_feature_cols,\n                    threshold=threshold,\n                    blend_weight=weight,\n                )\n\n                family_rows = pd.DataFrame(\n                    {\n                        \"Id\": semantic_df[\"Id\"].to_numpy(np.int64),\n                        \"Label\": y,\n                        \"family\": family,\n                        \"pred_primary\": primary_pred,\n                        \"pred_audit\": audit_pred,\n                    }\n                )\n                prediction_rows.append(family_rows)\n                primary_metrics[family] = self._metric_summary_from_pred(family_rows[\"Label\"].to_numpy(np.int8), family_rows[\"pred_primary\"].to_numpy(np.int8))\n                audit_metrics[family] = self._metric_summary_from_pred(family_rows[\"Label\"].to_numpy(np.int8), family_rows[\"pred_audit\"].to_numpy(np.int8))\n                LOGGER.info(\n                    \"[fit] %s primary F2=%.6f, audit F2=%.6f, threshold=%.3f, blend_weight=%.2f\",\n                    family,\n                    primary_metrics[family].f2,\n                    audit_metrics[family].f2,\n                    threshold,\n                    weight,\n                )\n                del base_df, semantic_df, cat_df, family_rows\n                gc.collect()\n\n        if not other_df.empty:\n            other_rows = other_df.copy()\n            other_rows[\"family\"] = \"other\"\n            other_rows[\"pred_primary\"] = 1\n            other_rows[\"pred_audit\"] = 1\n            prediction_rows.append(other_rows[[\"Id\", \"Label\", \"family\", \"pred_primary\", \"pred_audit\"]])\n            primary_metrics[\"other\"] = self._metric_summary_from_pred(other_rows[\"Label\"].to_numpy(np.int8), other_rows[\"pred_primary\"].to_numpy(np.int8))\n            audit_metrics[\"other\"] = self._metric_summary_from_pred(other_rows[\"Label\"].to_numpy(np.int8), other_rows[\"pred_audit\"].to_numpy(np.int8))\n\n        if not prediction_rows:\n            raise RuntimeError(\"No training rows were available for model fitting.\")\n        all_predictions = pd.concat(prediction_rows, ignore_index=True)\n        primary_metrics[\"overall\"] = self._metric_summary_from_pred(all_predictions[\"Label\"].to_numpy(np.int8), all_predictions[\"pred_primary\"].to_numpy(np.int8))\n        audit_metrics[\"overall\"] = self._metric_summary_from_pred(all_predictions[\"Label\"].to_numpy(np.int8), all_predictions[\"pred_audit\"].to_numpy(np.int8))\n        LOGGER.info(\n            \"[fit] overall primary F2=%.6f, precision=%.4f, recall=%.4f\",\n            primary_metrics[\"overall\"].f2,\n            primary_metrics[\"overall\"].precision,\n            primary_metrics[\"overall\"].recall,\n        )\n        LOGGER.info(\n            \"[fit] overall audit F2=%.6f, precision=%.4f, recall=%.4f\",\n            audit_metrics[\"overall\"].f2,\n            audit_metrics[\"overall\"].precision,\n            audit_metrics[\"overall\"].recall,\n        )\n        if demoted_rules:\n            LOGGER.info(\"[fit] demoted hard override rules=%s\", demoted_rules)\n        return self\n\n    def _predict_family_chunk(self, family: str, base_df: pd.DataFrame) -> np.ndarray:\n        bundle = self.family_models.get(family)\n        if bundle is None or bundle.semantic_model is None:\n            raise RuntimeError(f\"Missing fitted semantic model bundle for family {family}.\")\n        self._activate_semantic_context(bundle.semantic_context)\n        semantic_df = self._prepare_semantic_frame(base_df.copy())\n        hard_override = semantic_df[\"hard_override_anomaly\"].to_numpy(np.int8) == 1\n        semantic_prob = np.ones(len(semantic_df), dtype=np.float32)\n        if (~hard_override).any():\n            semantic_features = semantic_df.loc[~hard_override, bundle.semantic_feature_cols]\n            semantic_prob[~hard_override] = bundle.semantic_model.predict_proba(semantic_features)[:, 1].astype(np.float32)\n\n        cat_prob: Optional[np.ndarray] = None\n        if bundle.cat_model is not None and bundle.cat_feature_cols:\n            cat_df = self._prepare_cat_frame(base_df.copy())\n            cat_prob = np.ones(len(cat_df), dtype=np.float32)\n            if (~hard_override).any():\n                cat_features = cat_df.loc[~hard_override, bundle.cat_feature_cols]\n                cat_prob[~hard_override] = bundle.cat_model.predict_proba(cat_features)[:, 1].astype(np.float32)\n        blend_prob = self._blend_probs(semantic_prob, cat_prob, bundle.blend_weight)\n        blend_prob[hard_override] = 1.0\n        pred = (blend_prob >= bundle.threshold).astype(np.int8)\n        pred[hard_override] = 1\n        return pred\n\n    def predict_test(self, test_path: Path, out_csv: Path) -> None:\n        if not self.family_models:\n            raise RuntimeError(\"Model is not fitted.\")\n        out_csv.parent.mkdir(parents=True, exist_ok=True)\n        total_rows = 0\n        positive_rows = 0\n        LOGGER.info(\"[test] generating predictions from %s\", test_path)\n        with out_csv.open(\"w\", encoding=\"utf-8\") as fh:\n            fh.write(\"Id,Label\\n\")\n            with tqdm(\n                self.iter_raw_chunks(test_path, USECOLS_TEST, 0),\n                desc=\"test chunks\",\n                unit=\"chunk\",\n                dynamic_ncols=True,\n            ) as progress:\n                for chunk in progress:\n                    feats = build_features(chunk, hard_override_names=self.active_hard_override_names)\n                    pred = feats[\"hard_override_anomaly\"].astype(np.int8).to_numpy()\n                    for family in DEVICE_FAMILY_MAP:\n                        family_mask = feats[\"device_family\"] == family\n                        if not family_mask.any():\n                            continue\n                        family_df = feats.loc[family_mask].copy()\n                        family_pred = self._predict_family_chunk(family, family_df)\n                        pred[np.flatnonzero(family_mask.to_numpy())] = family_pred\n                    out = pd.DataFrame(\n                        {\n                            \"Id\": feats[\"Id\"].astype(np.int64),\n                            \"Label\": pred.astype(np.int8),\n                        }\n                    )\n                    out.to_csv(fh, index=False, header=False)\n                    total_rows += len(out)\n                    positive_rows += int(out[\"Label\"].sum())\n                    progress.set_postfix(\n                        rows=f\"{total_rows:,}\",\n                        positives=f\"{positive_rows:,}\",\n                    )\n        LOGGER.info(\n            \"[test] done; total_rows=%s, positive_rows=%s, positive_rate=%.6f\",\n            f\"{total_rows:,}\",\n            f\"{positive_rows:,}\",\n            positive_rows / max(total_rows, 1),\n        )\n","metadata":{},"outputs":[],"execution_count":null},{"id":"069bf829-4285-4b89-a514-69487d53fa6e","cell_type":"code","source":"%%writefile src/pipeline.py\n\"\"\"Top-level orchestration for the reproducible DER baseline run.\"\"\"\n\nimport hashlib\nimport logging\nimport random\nfrom pathlib import Path\n\nimport numpy as np\n\nfrom src.contracts import RunConfig\nfrom src.modeling import ResearchBaseline\n\nDEFAULT_RUN_CONFIG = RunConfig()\nLOGGER = logging.getLogger(__name__)\n\n\ndef seed_everything(seed: int) -> None:\n    \"\"\"Seed the local RNGs used by the training pipeline.\"\"\"\n\n    random.seed(seed)\n    np.random.seed(seed)\n\n\ndef file_sha256(path: Path) -> str:\n    \"\"\"Hash a file in chunks so large submissions do not spike memory.\"\"\"\n\n    digest = hashlib.sha256()\n    with path.open(\"rb\") as fh:\n        for chunk in iter(lambda: fh.read(1024 * 1024), b\"\"):\n            digest.update(chunk)\n    return digest.hexdigest()\n\n\ndef run_pipeline(config: RunConfig = DEFAULT_RUN_CONFIG) -> ResearchBaseline:\n    \"\"\"Train the baseline and emit the submission file for the pinned config.\"\"\"\n\n    seed_everything(config.seed)\n    LOGGER.info(\n        \"[run] starting pipeline with train=%s test=%s submission=%s\",\n        config.train_path,\n        config.test_path,\n        config.submission_path,\n    )\n    baseline = ResearchBaseline(config.baseline_config()).fit(config.train_path)\n    baseline.predict_test(config.test_path, config.submission_path)\n    LOGGER.info(\n        \"[solution] submission_sha256=%s path=%s\",\n        file_sha256(config.submission_path),\n        config.submission_path,\n    )\n    return baseline\n","metadata":{},"outputs":[],"execution_count":null},{"id":"7befaacb-26bb-47c0-b6fd-ca579051ddaf","cell_type":"code","source":"%%writefile src/rules.py\n\"\"\"Hard-rule metadata shared by feature engineering and modeling.\"\"\"\n\nfrom dataclasses import dataclass\nfrom typing import Mapping, Sequence\n\nimport numpy as np\nimport pandas as pd\n\n\n@dataclass(frozen=True)\nclass HardRuleContribution:\n    \"\"\"One boolean source column that contributes to hard-rule aggregates.\"\"\"\n\n    name: str | None\n    column: str\n    score_weight: float\n    count_weight: int = 1\n    default_override: bool = False\n\n\n@dataclass(frozen=True)\nclass RuleGroup:\n    \"\"\"Semantic grouping for related hard-rule contributions.\"\"\"\n\n    name: str\n    rules: tuple[HardRuleContribution, ...]\n\n\n@dataclass(frozen=True)\nclass HardRuleOutputs:\n    \"\"\"Aggregate hard-rule feature values derived from raw rule columns.\"\"\"\n\n    hard_override_anomaly: np.ndarray\n    hard_rule_count: np.ndarray\n    hard_rule_score: np.ndarray\n    hard_rule_anomaly: np.ndarray\n\n\ndef _rule(name: str, column: str, score_weight: float, count_weight: int = 1) -> HardRuleContribution:\n    return HardRuleContribution(name, column, score_weight, count_weight, True)\n\n\ndef _flatten_contributions(groups: Sequence[RuleGroup]) -> tuple[HardRuleContribution, ...]:\n    return tuple(rule for group in groups for rule in group.rules)\n\n\n# Order matters because feature generation, override handling, and reporting all\n# consume a stable rule ordering, so the grouped source of truth must preserve\n# the existing flattened order exactly.\nHARD_RULE_GROUPS = (\n    RuleGroup(\n        \"identity\",\n        (\n            _rule(\"noncanonical\", \"noncanonical\", 3.0),\n        ),\n    ),\n    RuleGroup(\n        \"power_envelope\",\n        (\n            _rule(\"w_gt_wmax\", \"w_gt_wmax_tol\", 2.0),\n            _rule(\"w_gt_wmaxrtg\", \"w_gt_wmaxrtg_tol\", 2.0),\n            _rule(\"va_gt_vamax\", \"va_gt_vamax_tol\", 2.0),\n            _rule(\"var_gt_injmax\", \"var_gt_injmax_tol\", 2.0),\n            _rule(\"var_lt_absmax\", \"var_lt_absmax_tol\", 2.0),\n        ),\n    ),\n    RuleGroup(\n        \"control_state\",\n        (\n            # `wsetpct_far` absorbs the legacy `wset_far` aggregate contribution\n            # because the two source columns are identical on the fixed\n            # train/test data.\n            _rule(\"wsetpct_far\", \"wsetpct_enabled_far\", 3.0, 2),\n            _rule(\"ac_type_rare\", \"ac_type_is_rare\", 1.5),\n            _rule(\"dc_type_rare\", \"dc_port_type_rare_any\", 1.5),\n            _rule(\"enter_state\", \"enter_service_state_anomaly\", 2.0),\n        ),\n    ),\n    RuleGroup(\n        \"protection\",\n        (\n            _rule(\"pf_abs\", \"pf_abs_ext_present\", 1.5),\n            _rule(\"pf_abs_rvrt\", \"pf_abs_rvrt_ext_present\", 1.5),\n            _rule(\"trip_power\", \"trip_any_power_when_outside\", 2.0),\n        ),\n    ),\n    RuleGroup(\n        \"aggregate_only\",\n        (\n            HardRuleContribution(None, \"common_missing_any\", 2.5),\n            HardRuleContribution(None, \"varsetpct_enabled_far\", 1.0),\n            HardRuleContribution(None, \"enter_service_blocked_power\", 0.35),\n            HardRuleContribution(None, \"enter_service_blocked_current\", 0.35),\n        ),\n    ),\n)\nHARD_RULE_CONTRIBUTIONS = _flatten_contributions(HARD_RULE_GROUPS)\n\nHARD_RULE_SPECS = tuple(spec for spec in HARD_RULE_CONTRIBUTIONS if spec.name is not None)\nHARD_RULE_COLUMNS = [spec.column for spec in HARD_RULE_CONTRIBUTIONS]\nHARD_RULE_NAMES = [spec.name for spec in HARD_RULE_SPECS]\nDEFAULT_HARD_OVERRIDE_NAMES = [spec.name for spec in HARD_RULE_SPECS if spec.default_override]\nRULE_COLUMN_MAP = {spec.name: spec.column for spec in HARD_RULE_SPECS}\n\n\ndef _infer_row_count(data: Mapping[str, Sequence[object]]) -> int:\n    for values in data.values():\n        return len(values)\n    return 0\n\n\ndef _coerce_rule_flag(values: Sequence[object] | pd.Series) -> np.ndarray:\n    if isinstance(values, pd.Series):\n        arr = pd.to_numeric(values, errors=\"coerce\").fillna(0).to_numpy()\n    else:\n        arr = np.asarray(values)\n        if np.issubdtype(arr.dtype, np.number) or np.issubdtype(arr.dtype, np.bool_):\n            arr = np.nan_to_num(arr, nan=0.0)\n        else:\n            arr = pd.to_numeric(pd.Series(arr), errors=\"coerce\").fillna(0).to_numpy()\n    return np.asarray(arr, dtype=np.int8) == 1\n\n\ndef build_hard_rule_column_flags(\n    data: Mapping[str, Sequence[object]] | pd.DataFrame,\n    *,\n    columns: Sequence[str] = HARD_RULE_COLUMNS,\n) -> dict[str, np.ndarray]:\n    return {column: _coerce_rule_flag(data[column]) for column in columns}\n\n\ndef build_named_rule_flags(\n    data: Mapping[str, Sequence[object]] | pd.DataFrame,\n    *,\n    column_flags: Mapping[str, np.ndarray] | None = None,\n) -> dict[str, np.ndarray]:\n    column_flags = column_flags or build_hard_rule_column_flags(data)\n    return {spec.name: column_flags[spec.column] for spec in HARD_RULE_SPECS}\n\n\ndef compute_active_override_anomaly(\n    data: Mapping[str, Sequence[object]] | pd.DataFrame,\n    active_override_names: Sequence[str],\n    *,\n    rule_flags: Mapping[str, np.ndarray] | None = None,\n    row_count: int | None = None,\n) -> np.ndarray:\n    if rule_flags is None:\n        rule_flags = build_named_rule_flags(data)\n    if row_count is None:\n        row_count = _infer_row_count(rule_flags)\n    if not active_override_names:\n        return np.zeros(row_count, dtype=np.int8)\n    return np.column_stack([rule_flags[name] for name in active_override_names]).any(axis=1).astype(np.int8)\n\n\ndef compute_hard_rule_outputs(\n    data: Mapping[str, Sequence[object]] | pd.DataFrame,\n    active_override_names: Sequence[str],\n    *,\n    row_count: int | None = None,\n) -> HardRuleOutputs:\n    column_flags = build_hard_rule_column_flags(data)\n    if row_count is None:\n        row_count = _infer_row_count(column_flags)\n    rule_flags = build_named_rule_flags(data, column_flags=column_flags)\n    hard_override_anomaly = compute_active_override_anomaly(\n        data,\n        active_override_names,\n        rule_flags=rule_flags,\n        row_count=row_count,\n    )\n    hard_rule_count = np.zeros(row_count, dtype=np.int16)\n    hard_rule_score = np.zeros(row_count, dtype=np.float32)\n    hard_rule_anomaly = np.zeros(row_count, dtype=bool)\n    for spec in HARD_RULE_CONTRIBUTIONS:\n        flags = column_flags[spec.column]\n        hard_rule_count += spec.count_weight * flags.astype(np.int16)\n        if spec.score_weight != 0.0:\n            hard_rule_score += spec.score_weight * flags.astype(np.float32)\n        hard_rule_anomaly |= flags\n    return HardRuleOutputs(\n        hard_override_anomaly=hard_override_anomaly,\n        hard_rule_count=hard_rule_count.astype(np.int8),\n        hard_rule_score=hard_rule_score,\n        hard_rule_anomaly=hard_rule_anomaly.astype(np.int8),\n    )\n","metadata":{},"outputs":[],"execution_count":null},{"id":"6e791364-cd87-4c9e-9aeb-5cd881c2c780","cell_type":"code","source":"%%writefile src/schema.py\n\"\"\"Schema-owned metadata for raw DER columns and structured block layouts.\"\"\"\n\nimport re\nfrom dataclasses import dataclass\nfrom typing import Dict, List, Sequence, Tuple\n\n\ndef dedupe(columns: Sequence[str]) -> List[str]:\n    return list(dict.fromkeys(columns))\n\n\ndef prefixed(prefix: str, fields: Sequence[str]) -> List[str]:\n    return [f\"{prefix}.{field}\" for field in fields]\n\n\nADAPTIVE_CURVE_HEADER_FIELDS = (\n    \"ID\",\n    \"L\",\n    \"Ena\",\n    \"AdptCrvReq\",\n    \"AdptCrvRslt\",\n    \"NPt\",\n    \"NCrv\",\n    \"RvrtTms\",\n    \"RvrtRem\",\n    \"RvrtCrv\",\n)\n\n\ndef build_adaptive_curve_columns(\n    prefix: str,\n    *,\n    curve_scalar_fields: Sequence[str],\n    point_fields: Sequence[str],\n    point_count: int,\n    curve_count: int = 3,\n) -> List[str]:\n    cols = prefixed(prefix, ADAPTIVE_CURVE_HEADER_FIELDS)\n    for curve in range(curve_count):\n        curve_prefix = f\"{prefix}.Crv[{curve}]\"\n        cols.extend(f\"{curve_prefix}.{field}\" for field in curve_scalar_fields)\n        for point in range(point_count):\n            cols.extend(f\"{curve_prefix}.Pt[{point}].{field}\" for field in point_fields)\n    return cols\n\n\ndef build_repeated_child_columns(\n    prefix: str,\n    *,\n    header_fields: Sequence[str],\n    child_label: str,\n    child_fields: Sequence[str],\n    child_count: int,\n) -> List[str]:\n    cols = prefixed(prefix, header_fields)\n    for child in range(child_count):\n        child_prefix = f\"{prefix}.{child_label}[{child}]\"\n        cols.extend(f\"{child_prefix}.{field}\" for field in child_fields)\n    return cols\n\n\ndef build_trip_columns(prefix: str, axis_name: str) -> List[str]:\n    cols = [\n        f\"{prefix}.ID\",\n        f\"{prefix}.L\",\n        f\"{prefix}.Ena\",\n        f\"{prefix}.AdptCrvReq\",\n        f\"{prefix}.AdptCrvRslt\",\n        f\"{prefix}.NPt\",\n        f\"{prefix}.NCrvSet\",\n    ]\n    for curve in range(2):\n        curve_prefix = f\"{prefix}.Crv[{curve}]\"\n        cols.append(f\"{curve_prefix}.ReadOnly\")\n        for group in [\"MustTrip\", \"MayTrip\", \"MomCess\"]:\n            group_prefix = f\"{curve_prefix}.{group}\"\n            cols.append(f\"{group_prefix}.ActPt\")\n            for point in range(5):\n                cols.extend(\n                    [\n                        f\"{group_prefix}.Pt[{point}].{axis_name}\",\n                        f\"{group_prefix}.Pt[{point}].Tms\",\n                    ]\n                )\n    return cols\n\n\n@dataclass(frozen=True)\nclass CurveFeatureSpec:\n    \"\"\"Static layout for one adaptive curve block.\"\"\"\n\n    prefix: str\n    point_count: int\n    x_field: str\n    y_field: str\n    meta_fields: Dict[str, str]\n\n\nCOMMON_FIELDS = \"Mn Md Opt Vr SN\".split()\nCOMMON_STR = prefixed(\"common[0]\", COMMON_FIELDS)\nCOMMON_COLUMNS = prefixed(\"common[0]\", [\"ID\", \"L\", *COMMON_FIELDS, \"DA\"])\n\nMEASURE_AC_FIELDS = \"\"\"\nID L ACType W VA Var PF A LLV LNV Hz TmpAmb TmpCab TmpSnk TmpTrns TmpSw TmpOt\nThrotPct ThrotSrc WL1 WL2 WL3 VAL1 VAL2 VAL3 VarL1 VarL2 VarL3 PFL1 PFL2 PFL3\nAL1 AL2 AL3 VL1L2 VL2L3 VL3L1 VL1 VL2 VL3\n\"\"\".split()\nMEASURE_AC_COLUMNS = prefixed(\"DERMeasureAC[0]\", MEASURE_AC_FIELDS)\n\nCAPACITY_FIELDS = \"\"\"\nID L WMaxRtg VAMaxRtg VarMaxInjRtg VarMaxAbsRtg WChaRteMaxRtg WDisChaRteMaxRtg\nVAChaRteMaxRtg VADisChaRteMaxRtg VNomRtg VMaxRtg VMinRtg AMaxRtg PFOvrExtRtg\nPFUndExtRtg NorOpCatRtg AbnOpCatRtg IntIslandCatRtg WMax WMaxOvrExt WOvrExtPF\nWMaxUndExt WUndExtPF VAMax VarMaxInj VarMaxAbs WChaRteMax WDisChaRteMax\nVAChaRteMax VADisChaRteMax VNom VMax VMin AMax PFOvrExt PFUndExt CtrlModes\nIntIslandCat\n\"\"\".split()\nCAPACITY_COLUMNS = prefixed(\"DERCapacity[0]\", CAPACITY_FIELDS)\n\nENTER_SERVICE_FIELDS = \"ID L ES ESVHi ESVLo ESHzHi ESHzLo ESDlyTms ESRndTms ESRmpTms ESDlyRemTms\".split()\nENTER_SERVICE_COLUMNS = prefixed(\"DEREnterService[0]\", ENTER_SERVICE_FIELDS)\n\nCTL_AC_FIELDS = \"\"\"\nID L PFWInjEna PFWInjEnaRvrt PFWInjRvrtTms PFWInjRvrtRem PFWAbsEna PFWAbsEnaRvrt\nPFWAbsRvrtTms PFWAbsRvrtRem WMaxLimPctEna WMaxLimPct WMaxLimPctRvrt\nWMaxLimPctEnaRvrt WMaxLimPctRvrtTms WMaxLimPctRvrtRem WSetEna WSetMod WSet\nWSetRvrt WSetPct WSetPctRvrt WSetEnaRvrt WSetRvrtTms WSetRvrtRem VarSetEna\nVarSetMod VarSetPri VarSet VarSetRvrt VarSetPct VarSetPctRvrt VarSetEnaRvrt\nVarSetRvrtTms VarSetRvrtRem WRmp WRmpRef VarRmp AntiIslEna PFWInj.PF PFWInj.Ext\nPFWInjRvrt.PF PFWInjRvrt.Ext PFWAbs.Ext PFWAbsRvrt.Ext\n\"\"\".split()\nCTL_AC_COLUMNS = prefixed(\"DERCtlAC[0]\", CTL_AC_FIELDS)\n\nVOLT_VAR_COLUMNS = build_adaptive_curve_columns(\n    \"DERVoltVar[0]\",\n    curve_scalar_fields=(\"ActPt\", \"DeptRef\", \"Pri\", \"VRef\", \"VRefAuto\", \"VRefAutoEna\", \"VRefAutoTms\", \"RspTms\", \"ReadOnly\"),\n    point_fields=(\"V\", \"Var\"),\n    point_count=4,\n)\nVOLT_WATT_COLUMNS = build_adaptive_curve_columns(\n    \"DERVoltWatt[0]\",\n    curve_scalar_fields=(\"ActPt\", \"DeptRef\", \"RspTms\", \"ReadOnly\"),\n    point_fields=(\"V\", \"W\"),\n    point_count=2,\n)\nFREQ_DROOP_COLUMNS = build_repeated_child_columns(\n    \"DERFreqDroop[0]\",\n    header_fields=(\"ID\", \"L\", \"Ena\", \"AdptCtlReq\", \"AdptCtlRslt\", \"NCtl\", \"RvrtTms\", \"RvrtRem\", \"RvrtCtl\"),\n    child_label=\"Ctl\",\n    child_fields=(\"DbOf\", \"DbUf\", \"KOf\", \"KUf\", \"RspTms\", \"PMin\", \"ReadOnly\"),\n    child_count=3,\n)\nWATT_VAR_COLUMNS = build_adaptive_curve_columns(\n    \"DERWattVar[0]\",\n    curve_scalar_fields=(\"ActPt\", \"DeptRef\", \"Pri\", \"ReadOnly\"),\n    point_fields=(\"W\", \"Var\"),\n    point_count=6,\n)\n\nTRIP_SPECS: Dict[str, Tuple[str, str, str]] = {\n    \"lv\": (\"DERTripLV[0]\", \"V\", \"low\"),\n    \"hv\": (\"DERTripHV[0]\", \"V\", \"high\"),\n    \"lf\": (\"DERTripLF[0]\", \"Hz\", \"low\"),\n    \"hf\": (\"DERTripHF[0]\", \"Hz\", \"high\"),\n}\nTRIP_COLUMNS = {short_name: build_trip_columns(prefix, axis_name) for short_name, (prefix, axis_name, _) in TRIP_SPECS.items()}\n\n# These expected IDs and lengths belong with the raw schema definition rather\n# than with model-training configuration.\nEXPECTED_MODEL_META = {\n    \"common\": (\"common[0].ID\", \"common[0].L\", 1.0, 66.0),\n    \"measure_ac\": (\"DERMeasureAC[0].ID\", \"DERMeasureAC[0].L\", 701.0, 153.0),\n    \"capacity\": (\"DERCapacity[0].ID\", \"DERCapacity[0].L\", 702.0, 50.0),\n    \"enter_service\": (\"DEREnterService[0].ID\", \"DEREnterService[0].L\", 703.0, 17.0),\n    \"measure_dc\": (\"DERMeasureDC[0].ID\", \"DERMeasureDC[0].L\", 714.0, 68.0),\n}\n\nCURVE_FEATURE_SPECS = {\n    \"voltvar\": CurveFeatureSpec(\n        prefix=\"DERVoltVar[0]\",\n        point_count=4,\n        x_field=\"V\",\n        y_field=\"Var\",\n        meta_fields={\n            \"deptref\": \"DeptRef\",\n            \"pri\": \"Pri\",\n            \"vref\": \"VRef\",\n            \"vref_auto\": \"VRefAuto\",\n            \"vref_auto_ena\": \"VRefAutoEna\",\n            \"vref_auto_tms\": \"VRefAutoTms\",\n            \"rsp\": \"RspTms\",\n            \"readonly\": \"ReadOnly\",\n        },\n    ),\n    \"voltwatt\": CurveFeatureSpec(\n        prefix=\"DERVoltWatt[0]\",\n        point_count=2,\n        x_field=\"V\",\n        y_field=\"W\",\n        meta_fields={\n            \"deptref\": \"DeptRef\",\n            \"rsp\": \"RspTms\",\n            \"readonly\": \"ReadOnly\",\n        },\n    ),\n    \"wattvar\": CurveFeatureSpec(\n        prefix=\"DERWattVar[0]\",\n        point_count=6,\n        x_field=\"W\",\n        y_field=\"Var\",\n        meta_fields={\n            \"deptref\": \"DeptRef\",\n            \"pri\": \"Pri\",\n            \"readonly\": \"ReadOnly\",\n        },\n    ),\n}\n\nMEASURE_DC_FIELDS = \"\"\"\nID L NPrt DCA DCW Prt[0].PrtTyp Prt[0].ID Prt[0].DCA Prt[0].DCV Prt[0].DCW\nPrt[0].Tmp Prt[1].PrtTyp Prt[1].ID Prt[1].DCA Prt[1].DCV Prt[1].DCW Prt[1].Tmp\n\"\"\".split()\nMEASURE_DC_COLUMNS = prefixed(\"DERMeasureDC[0]\", MEASURE_DC_FIELDS)\n\nBLOCK_SOURCE_COLUMNS: Dict[str, List[str]] = {\n    \"common\": COMMON_COLUMNS,\n    \"measure_ac\": MEASURE_AC_COLUMNS,\n    \"capacity\": CAPACITY_COLUMNS,\n    \"enter_service\": ENTER_SERVICE_COLUMNS,\n    \"ctl_ac\": CTL_AC_COLUMNS,\n    \"volt_var\": VOLT_VAR_COLUMNS,\n    \"volt_watt\": VOLT_WATT_COLUMNS,\n    \"freq_droop\": FREQ_DROOP_COLUMNS,\n    \"watt_var\": WATT_VAR_COLUMNS,\n    \"measure_dc\": MEASURE_DC_COLUMNS,\n}\nfor short_name, cols in TRIP_COLUMNS.items():\n    BLOCK_SOURCE_COLUMNS[f\"trip_{short_name}\"] = cols\n\nCURVE_BLOCK_META_FIELDS = \"Ena AdptCrvReq AdptCrvRslt NPt NCrv RvrtTms RvrtRem RvrtCrv\".split()\nFREQ_DROOP_META_FIELDS = \"Ena AdptCtlReq AdptCtlRslt NCtl RvrtTms RvrtRem RvrtCtl\".split()\nTRIP_META_FIELDS = \"Ena AdptCrvReq AdptCrvRslt NPt NCrvSet\".split()\n\nRAW_NUMERIC = dedupe(\n    [\n        \"common[0].DA\",\n        *prefixed(\"DERMeasureAC[0]\", MEASURE_AC_FIELDS[2:]),\n        *prefixed(\"DERCapacity[0]\", CAPACITY_FIELDS[2:]),\n        *prefixed(\"DEREnterService[0]\", ENTER_SERVICE_FIELDS[2:]),\n        *prefixed(\"DERCtlAC[0]\", CTL_AC_FIELDS[2:]),\n        *prefixed(\"DERVoltVar[0]\", CURVE_BLOCK_META_FIELDS),\n        *prefixed(\"DERVoltWatt[0]\", CURVE_BLOCK_META_FIELDS),\n        *prefixed(\"DERFreqDroop[0]\", FREQ_DROOP_META_FIELDS),\n        *prefixed(\"DERWattVar[0]\", CURVE_BLOCK_META_FIELDS),\n        *prefixed(\"DERMeasureDC[0]\", MEASURE_DC_FIELDS[2:]),\n    ]\n)\n\nTRIP_META_COLUMNS = [f\"{prefix}.{field}\" for prefix, _, _ in TRIP_SPECS.values() for field in TRIP_META_FIELDS]\nRAW_EXTRA_NUMERIC_COLUMNS = [\n    \"DERMeasureAC[0].A_SF\",\n    \"DERMeasureAC[0].V_SF\",\n    \"DERMeasureAC[0].Hz_SF\",\n    \"DERMeasureAC[0].W_SF\",\n    \"DERMeasureAC[0].PF_SF\",\n    \"DERMeasureAC[0].VA_SF\",\n    \"DERMeasureAC[0].Var_SF\",\n    \"DERCapacity[0].WOvrExtRtg\",\n    \"DERCapacity[0].WOvrExtRtgPF\",\n    \"DERCapacity[0].WUndExtRtg\",\n    \"DERCapacity[0].WUndExtRtgPF\",\n    \"DERCapacity[0].W_SF\",\n    \"DERCapacity[0].PF_SF\",\n    \"DERCapacity[0].VA_SF\",\n    \"DERCapacity[0].Var_SF\",\n    \"DERCapacity[0].V_SF\",\n    \"DERCapacity[0].A_SF\",\n    \"DERCtlAC[0].WSet_SF\",\n    \"DERMeasureDC[0].DCA_SF\",\n    \"DERMeasureDC[0].DCW_SF\",\n]\nRAW_EXTRA_STRING_COLUMNS = [\n    \"DERMeasureDC[0].Prt[0].IDStr\",\n    \"DERMeasureDC[0].Prt[1].IDStr\",\n]\nRAW_NUMERIC = dedupe([*RAW_NUMERIC, *TRIP_META_COLUMNS, *RAW_EXTRA_NUMERIC_COLUMNS])\nRAW_STRING_COLUMNS = dedupe([*COMMON_STR, *RAW_EXTRA_STRING_COLUMNS])\n\nTRIP_SOURCE_COLUMNS = [col for cols in TRIP_COLUMNS.values() for col in cols]\nALL_SOURCE_COLUMNS = dedupe(\n    [\n        *COMMON_COLUMNS,\n        *MEASURE_AC_COLUMNS,\n        *CAPACITY_COLUMNS,\n        *ENTER_SERVICE_COLUMNS,\n        *CTL_AC_COLUMNS,\n        *VOLT_VAR_COLUMNS,\n        *VOLT_WATT_COLUMNS,\n        *FREQ_DROOP_COLUMNS,\n        *WATT_VAR_COLUMNS,\n        *MEASURE_DC_COLUMNS,\n        *TRIP_SOURCE_COLUMNS,\n        *RAW_EXTRA_NUMERIC_COLUMNS,\n        *RAW_EXTRA_STRING_COLUMNS,\n    ]\n)\nNUMERIC_SOURCE_COLUMNS = [c for c in ALL_SOURCE_COLUMNS if c not in RAW_STRING_COLUMNS]\n\nUSECOLS_TRAIN = dedupe([\"Id\", \"Label\", *ALL_SOURCE_COLUMNS])\nUSECOLS_TEST = dedupe([\"Id\", *ALL_SOURCE_COLUMNS])\n\nSAFE_RAW = {c: re.sub(r\"[^0-9A-Za-z_]+\", \"_\", c) for c in RAW_NUMERIC}\nSAFE_STR = {c: re.sub(r\"[^0-9A-Za-z_]+\", \"_\", c) for c in RAW_STRING_COLUMNS}\n","metadata":{},"outputs":[],"execution_count":null},{"id":"6376c36c-065d-462c-812f-5d7a12b73b71","cell_type":"code","source":"%%writefile main.py\n#!/usr/bin/env python3\n\nimport logging\n\nfrom src.pipeline import run_pipeline\n\n\ndef configure_logging() -> None:\n    \"\"\"Configure the CLI logging format once at process startup.\"\"\"\n\n    logging.basicConfig(\n        level=logging.INFO,\n        format=\"%(asctime)s | %(levelname)s | %(message)s\",\n        force=True,\n    )\n\n\ndef main() -> None:\n    \"\"\"Run the pinned training-and-inference pipeline.\"\"\"\n\n    configure_logging()\n    run_pipeline()\n\n\nif __name__ == \"__main__\":\n    main()","metadata":{},"outputs":[],"execution_count":null},{"id":"e88d3a4b-7241-4bf0-b468-a1ac980f1012","cell_type":"code","source":"import sys\n!{sys.executable} main.py","metadata":{},"outputs":[],"execution_count":null}]}