File size: 12,303 Bytes
bf8df4f
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
import numpy as np


#%% Cell 1 — Imports

import pandas as pd
import numpy as np
import torch
import warnings
import joblib
import os

from sklearn.preprocessing import MinMaxScaler, FunctionTransformer
from sklearn.pipeline import Pipeline
from sklearn.model_selection import train_test_split
from torch.utils.data import TensorDataset, DataLoader


warnings.filterwarnings("ignore")


def log1p_base10(x):
    return np.log10(1 + x)
#%% Cell 2 — IDSDataPipeline class

class IDSDataPipeline:
    """End-to-end preprocessing + DataLoader builder for one attack class.

    Loads benign and attack CSVs, applies the shared preprocessing artifacts
    (column drops, flag binning, one-hot encoding), filters invalid rows,
    splits into train/val/test, fits the log+minmax pipeline on train only,
    and exposes ready-to-use PyTorch DataLoaders.

    Attributes
    ----------
    train_loader, val_loader, test_loader : torch.utils.data.DataLoader
        Shuffled-train, non-shuffled val/test loaders over (X, y) tensors.
    input_dim : int
        Number of features fed to the model.
    numeric_pipeline : sklearn.pipeline.Pipeline
        The fitted log1p_base10 + MinMaxScaler pipeline (saved to disk).
    feature_names : list[str]
        Final ordered column names of the feature matrix.
    """

    def __init__(
        self,
        attack: str,
        benign_data_path: str = 'Data/benign_only/all_days_benign.csv',
        attack_data_dir:  str = 'Data/attacks_only',
        preprocessing_dir: str = '_prepcosessing_artefacts/',
        checkpoints_dir:   str = 'checkpoints_MLP',
        batch_size: int = 256,
        test_size:  float = 0.10,
        val_size:   float = 0.10,
        random_state: int = 42,
        verbose: bool = True,
    ):
        self.attack             = attack
        self.benign_data_path   = benign_data_path
        self.attack_data_path   = os.path.join(attack_data_dir, f"{attack}.csv")
        self.preprocessing_dir  = preprocessing_dir
        self.checkpoint_dir     = os.path.join(checkpoints_dir, attack)
        self.batch_size         = batch_size
        self.test_size          = test_size
        self.val_size           = val_size
        self.random_state       = random_state
        self.verbose            = verbose

        os.makedirs(self.checkpoint_dir, exist_ok=True)

        # Filled in by .build()
        self.train_loader = None
        self.val_loader   = None
        self.test_loader  = None
        self.input_dim    = None
        self.numeric_pipeline = None
        self.feature_names = None

        self.build()

    # ---------- internal steps ----------

    def _log(self, msg):
        if self.verbose:
            print(msg)

    def _load_artifacts(self):
        self.ohe = joblib.load(os.path.join(self.preprocessing_dir, "onehot_encoder.pkl"))
        schema  = joblib.load(os.path.join(self.preprocessing_dir, "column_schema.pkl"))

        self.col_to_drop         = schema["col_to_drop"]
        self.onehot_cols         = schema["onehot_cols"]
        self.flag_cols           = schema["flag_cols"]
        self.flag_bin_config     = schema["flag_bin_config"]
        self.numerical_cols      = schema["numerical_cols"]
        self.ohe_feature_names   = schema["ohe_feature_names"]

        self._log(
            f"Loaded schema: {len(self.numerical_cols)} numerical, "
            f"{len(self.flag_cols)} flag, {len(self.onehot_cols)} OHE source "
            f"-> {len(self.ohe_feature_names)} OHE features"
        )

    def _load_data(self):
        df_benign = pd.read_csv(self.benign_data_path)
        df_attack = pd.read_csv(self.attack_data_path)

        df_benign.drop(columns=self.col_to_drop, inplace=True)
        df_attack.drop(columns=self.col_to_drop, inplace=True)
        if 'AttackFamily' in df_attack.columns:
            df_attack.drop(columns=['AttackFamily'], inplace=True)

        assert list(df_benign.columns) == list(df_attack.columns), \
            "Benign and attack dataframes have different columns after dropping."

        return df_benign, df_attack

    @staticmethod
    def _bin_flag_columns(df, bin_config):
        df = df.copy()
        for col, edges in bin_config.items():
            df[col] = pd.cut(df[col], bins=edges, labels=False, right=True).astype(np.int8)
        return df

    def _encode(self, df):
        """Apply flag binning and OHE; return (numerical, flag, ohe) concatenated df."""
        df = self._bin_flag_columns(df, self.flag_bin_config)
        ohe_arr = self.ohe.transform(df[self.onehot_cols])
        ohe_df  = pd.DataFrame(ohe_arr, columns=self.ohe_feature_names, index=df.index)

        return pd.concat(
            [df[self.numerical_cols].reset_index(drop=True),
             df[self.flag_cols].reset_index(drop=True),
             ohe_df.reset_index(drop=True)],
            axis=1,
        )

    def _filter_invalid(self, df):
        """Drop rows with negative or non-finite numerical values."""
        mask = (df[self.numerical_cols] >= 0).all(axis=1) & \
               np.isfinite(df[self.numerical_cols]).all(axis=1)
        return df[mask]

    def _split(self, X, y):
        X_trainval, X_test, y_trainval, y_test = train_test_split(
            X, y,
            test_size=self.test_size,
            stratify=y,
            random_state=self.random_state,
        )
        # val_size is fraction of the original; convert to fraction of remaining.
        val_relative = self.val_size / (1.0 - self.test_size)
        X_train, X_val, y_train, y_val = train_test_split(
            X_trainval, y_trainval,
            test_size=val_relative,
            stratify=y_trainval,
            random_state=self.random_state,
        )
        return X_train, X_val, X_test, y_train, y_val, y_test

    def _fit_numeric(self, X_train):
        log_transformer = FunctionTransformer(
            func=log1p_base10, validate=False, feature_names_out="one-to-one",
        )
        self.numeric_pipeline = Pipeline([
            ("log_transform", log_transformer),
            ("minmax_scaler", MinMaxScaler()),
        ])
        self.numeric_pipeline.fit(X_train[self.numerical_cols])

    def _transform_numeric(self, X):
        return pd.DataFrame(
            self.numeric_pipeline.transform(X[self.numerical_cols]),
            columns=self.numerical_cols, index=X.index,
        )

    def _assemble(self, X, X_num):
        """Reassemble: scaled numerics + flag bins + OHE columns."""
        return pd.concat(
            [X_num, X[self.flag_cols], X[self.ohe_feature_names]],
            axis=1,
        )

    def _make_loader(self, X_df, y_series, shuffle):
        X_t = torch.tensor(X_df.values, dtype=torch.float32)
        y_t = torch.tensor(y_series.values, dtype=torch.float32).unsqueeze(1)
        ds  = TensorDataset(X_t, y_t)
        return DataLoader(ds, batch_size=self.batch_size, shuffle=shuffle, drop_last=False)

    # ---------- orchestration ----------

    def build(self):
        self._load_artifacts()
        df_benign, df_attack = self._load_data()

        # Encode + filter (processed versions)
        df_benign_proc = self._encode(df_benign)
        df_attack_proc = self._encode(df_attack)
        df_benign_proc = self._filter_invalid(df_benign_proc)
        df_attack_proc = self._filter_invalid(df_attack_proc)

        # Keep raw versions aligned with the filtered processed versions.
        # _encode reset the index of df_*_proc to 0..N-1 (via reset_index(drop=True)
        # inside the concat), so we align by position rather than by original index.
        df_benign_raw = df_benign.iloc[: len(df_benign_proc)].reset_index(drop=True)
        df_attack_raw = df_attack.iloc[: len(df_attack_proc)].reset_index(drop=True)
        # The above assumes _filter_invalid drops rows from the END only — which
        # is not generally true. Safer: align by the actual surviving positions.

        # Actually: track surviving positions explicitly.
        # Re-do the filtering with explicit index preservation.
        df_benign_proc_unfiltered = self._encode(df_benign)
        df_attack_proc_unfiltered = self._encode(df_attack)
        mask_benign = (df_benign_proc_unfiltered[self.numerical_cols] >= 0).all(axis=1) & \
                    np.isfinite(df_benign_proc_unfiltered[self.numerical_cols]).all(axis=1)
        mask_attack = (df_attack_proc_unfiltered[self.numerical_cols] >= 0).all(axis=1) & \
                    np.isfinite(df_attack_proc_unfiltered[self.numerical_cols]).all(axis=1)

        df_benign_proc = df_benign_proc_unfiltered[mask_benign].reset_index(drop=True)
        df_attack_proc = df_attack_proc_unfiltered[mask_attack].reset_index(drop=True)
        df_benign_raw  = df_benign[mask_benign.values].reset_index(drop=True)
        df_attack_raw  = df_attack[mask_attack.values].reset_index(drop=True)

        # Add labels
        df_benign_proc["label"] = 0
        df_attack_proc["label"] = 1
        df_benign_raw["label"]  = 0
        df_attack_raw["label"]  = 1

        # Combine
        df_full     = pd.concat([df_benign_proc, df_attack_proc], axis=0, ignore_index=True)
        df_raw_full = pd.concat([df_benign_raw,  df_attack_raw],  axis=0, ignore_index=True)

        # Shuffle with same random_state — apply same permutation to both
        perm        = df_full.sample(frac=1, random_state=self.random_state).index
        df_full     = df_full.loc[perm].reset_index(drop=True)
        df_raw_full = df_raw_full.loc[perm].reset_index(drop=True)

        self._log(f"\nFull dataset shape: {df_full.shape}")
        self._log(f"Class balance:\n{df_full['label'].value_counts(normalize=True)}")

        X = df_full.drop(columns=["label"])
        y = df_full["label"]

        # Split (returns indices we can apply to df_raw_full too)
        indices = np.arange(len(df_full))
        train_idx, test_idx = train_test_split(
            indices, test_size=self.test_size, stratify=y, random_state=self.random_state
        )
        val_relative = self.val_size / (1.0 - self.test_size)
        train_idx, val_idx = train_test_split(
            train_idx, test_size=val_relative,
            stratify=y.iloc[train_idx], random_state=self.random_state
        )

        X_train, X_val, X_test = X.iloc[train_idx], X.iloc[val_idx], X.iloc[test_idx]
        y_train, y_val, y_test = y.iloc[train_idx], y.iloc[val_idx], y.iloc[test_idx]

        # The raw splits, aligned to processed:
        raw_train = df_raw_full.iloc[train_idx].reset_index(drop=True)
        raw_val   = df_raw_full.iloc[val_idx].reset_index(drop=True)
        raw_test  = df_raw_full.iloc[test_idx].reset_index(drop=True)

        self._log(f"\nTrain: {X_train.shape}, Val: {X_val.shape}, Test: {X_test.shape}")

        # Fit numeric pipeline on train only, transform all splits
        self._fit_numeric(X_train)
        X_train_num = self._transform_numeric(X_train)
        X_val_num   = self._transform_numeric(X_val)
        X_test_num  = self._transform_numeric(X_test)

        X_train_final = self._assemble(X_train, X_train_num)
        X_val_final   = self._assemble(X_val,   X_val_num)
        X_test_final  = self._assemble(X_test,  X_test_num)

        self.X_train = X_train_final.values
        self.X_val   = X_val_final.values
        self.X_test  = X_test_final.values

        self.y_train = y_train.values
        self.y_val   = y_val.values
        self.y_test  = y_test.values

        # NEW: raw dataframes per split, same row ordering as processed splits.
        self.raw_train = raw_train
        self.raw_val   = raw_val
        self.raw_test  = raw_test

        self.feature_names = X_train_final.columns.tolist()
        self.input_dim     = X_train_final.shape[1]

        self._log(f"\nFinal feature count: {self.input_dim}")

        # Persist fitted numeric pipeline
        joblib.dump(
            self.numeric_pipeline,
            os.path.join(self.checkpoint_dir, "numeric_pipeline.pkl"),
        )

        # Loaders (unchanged)
        self.train_loader = self._make_loader(X_train_final, y_train, shuffle=True)
        self.val_loader   = self._make_loader(X_val_final,   y_val,   shuffle=False)
        self.test_loader  = self._make_loader(X_test_final,  y_test,  shuffle=False)

#%% Cell 3 — Usage