diff --git a/.github/workflows/python-app.yml b/.github/workflows/python-app.yml index e9727e9a..a0f4fbe4 100644 --- a/.github/workflows/python-app.yml +++ b/.github/workflows/python-app.yml @@ -30,7 +30,7 @@ jobs: run: uv run --dev prek run -a - name: Test with pytest - run: uv run --dev pytest -s + run: uv run pytest -s --run-ignored # integration tests for DetectMateService - name: Checkout DetectMateService diff --git a/docs/detectors.md b/docs/detectors.md index 183659a3..2e4056a0 100644 --- a/docs/detectors.md +++ b/docs/detectors.md @@ -92,6 +92,8 @@ List of detectors: * [Rule Based](detectors/rule_based.md): Detect anomalies based in a set of rules. * [Bigram Frequency](detectors/bigram_frequency.md): Detect bigram-frequency-based anomalies in the logs. * [Charset](detectors/charset.md): Detect new characters in the variables in the logs. +* [Deeplog](detectors/deeplog.md): Detect anomalies of a sequence of evend IDs with a LSTM. +* [LogBert](detectors/logbert.md): Detect anomalies of a sequence of evend IDs with a Transformer. * [SCVS Detector](detectors/scvs_detector.md): Detect anomalies by looking at different sequence count vectors. * [ECVC Detector](detectors/ecvc_detector.md): Detect anomalies by calculating the distance between different sequence count vectors. @@ -104,6 +106,8 @@ detectors: NewValueDetector: method_type: new_value_detector auto_config: False + data_use_configure: None # Data used for configuration + data_use_training: 199 # Data used for training params: {} # global parameters events: # event-specific configuration 1: # event_id diff --git a/docs/detectors/deeplog.md b/docs/detectors/deeplog.md new file mode 100644 index 00000000..f9e7e3f9 --- /dev/null +++ b/docs/detectors/deeplog.md @@ -0,0 +1,63 @@ +# Deeplog Detector + +The Deeplog Detector is inspired from [Deeplog paper](https://dl.acm.org/doi/10.1145/3133956.3134015). + +| | Schema | Description | +|------------|----------------------------|--------------------| +| **Input** | [ParserSchema](../schemas.md) | Structured log | +| **Output** | [DetectorSchema](../schemas.md) | Combined alert / finding | + +## Description +Deep learning method that looks at the event ID sequence + +## Configuration + +```yaml +detectors: + DeeplogDetector: + method_type: deeplog_detector + auto_config: False + data_use_training: 10 + window_size: 3 + hyperparameters: + Model: + hidden_dim: 64 + n_layers: 2 + Train: + seed: 0 + batch_size: 2048 + learning_rate: 0.01 + epochs: 10 + patience: 3 + Finetune: + - ["Model", "hidden_dim", [128, 256, 512]] + - ["Model", "n_layers", [1, 2, 3]] + - ["Train", "learning_rate", [0.01, 0.02, 0.03]] +``` + +## Example usage + +```python +from detectmatelibrary.detectors.deeplog_detector import DeeplogDetector + +import detectmatelibrary.schemas as schemas + +detector = DeeplogDetector(name="DeeplogDetector", config=cfg) + +test_data = schemas.ParserSchema({ + "parserType": "test", + "EventID": 12, + "template": "test template", + "variables": ["adsasd", "asdasd"], + "logID": "2", + "parsedLogID": "2", + "parserID": "test_parser", + "log": "test log message", + "logFormatVariables": {"level": "CRITICAL"} +}) +output = schemas.DetectorSchema() + +result = detector.detect(test_data, output) + +``` +Go back [Index](../index.md) diff --git a/docs/detectors/logbert.md b/docs/detectors/logbert.md new file mode 100644 index 00000000..6c379014 --- /dev/null +++ b/docs/detectors/logbert.md @@ -0,0 +1,68 @@ +# LogBert Detector + +The LogBert Detector is inspired from [LogBert paper](https://ieeexplore.ieee.org/stamp/stamp.jsp?arnumber=9534113). + +| | Schema | Description | +|------------|----------------------------|--------------------| +| **Input** | [ParserSchema](../schemas.md) | Structured log | +| **Output** | [DetectorSchema](../schemas.md) | Combined alert / finding | + +## Description +Deep learning method that looks at the event ID sequence + +## Configuration + +```yaml +detectors: + LogBertDetector: + method_type: logbert_detector + auto_config: False + data_use_training: 10 + window_size: 4 + hyperparameters: + Model: + hidden: 32 + num_heads: 2 + n_layers: 1 + dropout: 0.0 + max_len: 1000 + Train: + seed: 0 + batch_size: 256 + learning_rate: 0.01 + epochs: 10 + mask_per: 0.4 + alpha: 0.0 + patience: 3 + Finetune: + - ["Model", "hidden", [64, 128, 256]] + - ["Model", "n_layers", [1, 2, 3]] + - ["Train", "learning_rate", [0.002, 0.001, 0.005]] +``` + +## Example usage + +```python +from detectmatelibrary.detectors.logbert_detector import LogBertDetector + +import detectmatelibrary.schemas as schemas + +detector = LogBertDetector(name="LogBertDetector", config=cfg) + +test_data = schemas.ParserSchema({ + "parserType": "test", + "EventID": 12, + "template": "test template", + "variables": ["adsasd", "asdasd"], + "logID": "2", + "parsedLogID": "2", + "parserID": "test_parser", + "log": "test log message", + "logFormatVariables": {"level": "CRITICAL"} +}) +output = schemas.DetectorSchema() + +result = detector.detect(test_data, output) + +``` +Go back [Index](../index.md) diff --git a/mkdocs.yml b/mkdocs.yml index ff502825..9e18d370 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -31,6 +31,8 @@ nav: - BiGram Frequency: detectors/bigram_frequency.md - CharSet: detectors/charset.md - Value Range: detectors/value_range.md + - Deeplog Detector: detectors/deeplog.md + - LogBert Detector: detectors/logbert.md - SCVS Detector: detectors/scvs_detector.md - ECVC Detector: detectors/ecvc_detector.md - Alert Aggregation Methods: diff --git a/pyproject.toml b/pyproject.toml index 904edf5f..35339252 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -17,6 +17,9 @@ dependencies = [ "msgpack>=1.0.0", "fsspec>=2024.1.0", "pyarrow>=24.0.0", + "jax>=0.11.0", + "flax>=0.12.8", + "optax>=0.2.8", ] [dependency-groups] diff --git a/pytest.ini b/pytest.ini new file mode 100644 index 00000000..ed39a206 --- /dev/null +++ b/pytest.ini @@ -0,0 +1,5 @@ +[pytest] +# You can use addopts to change the default terminal behavior globally +addopts = --strict-markers +markers = + ignored: tests that should only run on demand diff --git a/src/detectmatelibrary/common/_core_op/_fit_logic.py b/src/detectmatelibrary/common/_core_op/_fit_logic.py index d0699191..b99a352a 100644 --- a/src/detectmatelibrary/common/_core_op/_fit_logic.py +++ b/src/detectmatelibrary/common/_core_op/_fit_logic.py @@ -67,7 +67,7 @@ class FitLogicState(Enum): def describe(self) -> str: descriptions = [ "Configuring", - "Training.", + "Training", "Default" ] return descriptions[self.value] diff --git a/src/detectmatelibrary/common/deeplearning_detector.py b/src/detectmatelibrary/common/deeplearning_detector.py new file mode 100644 index 00000000..1189483e --- /dev/null +++ b/src/detectmatelibrary/common/deeplearning_detector.py @@ -0,0 +1,90 @@ + +from detectmatelibrary.common.detector import CoreDetector, CoreDetectorConfig + +from detectmatelibrary.utils.deep_learning.imodel import DeepModel +from detectmatelibrary.utils.data_buffer import BufferMode + +from detectmatelibrary import schemas + +from typing import Any + + +class DeepLearningDetectorConfig(CoreDetectorConfig): + window_size: int = 10 + validation_per: float = 0.2 + finetune_epochs: int = 2 + + hyperparameters: dict[str, Any] = { + "Model": { + + }, + "Train": { + + }, + "Finetune": [], + } + + +def build_seq(input_: list[schemas.ParserSchema]) -> tuple[int]: + return tuple([in_["EventID"] for in_ in input_]) + + +class DeepLearningDetector(CoreDetector): + def __init__( + self, + model_cls: DeepModel, + name: str = "CoreDetector", + config: DeepLearningDetectorConfig = DeepLearningDetectorConfig() + ) -> None: + + if isinstance(config, dict): + config = DeepLearningDetectorConfig.from_dict(config, name) + + super().__init__( + name=name, + buffer_mode=BufferMode.WINDOW, + buffer_size=config.window_size, + config=config + ) + self.config: DeepLearningDetectorConfig + self.model: DeepModel = model_cls(config=self.config.hyperparameters) # type: ignore + + self.train_seqs: list[tuple[int]] = [] + self.config_seqs: list[tuple[int]] = [] + self.stats: dict[str, float | int] = {} + self.top_k: int = 0 + + def train(self, input_: list[schemas.ParserSchema]) -> None: # type: ignore + self.train_seqs.append(build_seq(input_)) + + def configure(self, input_: list[schemas.ParserSchema]) -> None: # type: ignore + self.config_seqs.append(build_seq(input_)) + + def set_configuration(self) -> None: + self.model.finetune( + self.config_seqs, var_per=self.config.validation_per, epochs=self.config.finetune_epochs + ) + self.config_seqs = [] + + def post_train(self) -> None: + self.stats = self.model.train(self.train_seqs, var_per=self.config.validation_per) + self.train_seqs = [] + + if "top_k" in self.stats: + self.top_k = int(self.stats["top_k"]) + print(self.model) + print("Top k assigned", self.top_k) + + def detect( + self, + input_: list[schemas.ParserSchema], # type: ignore + output_: schemas.DetectorSchema, + ) -> bool: + + alert = self.model.check_anomaly(build_seq(input_), top_k=self.top_k) + if alert: + output_["score"] = 1.0 + output_["description"] = f"{self.name} found an anomaly in the sequence" + return True + + return False diff --git a/src/detectmatelibrary/detectors/deeplog_detector.py b/src/detectmatelibrary/detectors/deeplog_detector.py new file mode 100644 index 00000000..f575887f --- /dev/null +++ b/src/detectmatelibrary/detectors/deeplog_detector.py @@ -0,0 +1,46 @@ +from detectmatelibrary.common.deeplearning_detector import ( + DeepLearningDetectorConfig, DeepLearningDetector +) + +from detectmatelibrary.utils.deep_learning.deeplog import DeepLog + + +from typing import Any + + +class DeeplogDetectorConfig(DeepLearningDetectorConfig): + method_type: str = "deeplog_detector" + + hyperparameters: dict[str, Any] = { # type: ignore + "Model": { + "hidden_dim": 64, + "n_layers": 2, + }, + "Train": { + "seed": 0, + "batch_size": 2048, + "learning_rate": 0.01, + "epochs": 10, + "patience": 3, + }, + "Finetune": [ + ["Model", "hidden_dim", [128, 256, 512]], + ["Model", "n_layers", [1, 2, 3]], + ["Train", "learning_rate", [0.01, 0.02, 0.03]], + ], + } + + +class DeeplogDetector(DeepLearningDetector): + def __init__( + self, + name: str = "DeeplogDetector", + config: DeeplogDetectorConfig | dict[str, Any] = DeeplogDetectorConfig(), + ) -> None: + + if isinstance(config, dict): + config = DeeplogDetectorConfig.from_dict(config, name) + + super().__init__( + name=name, model_cls=DeepLog, config=config + ) \ No newline at end of file diff --git a/src/detectmatelibrary/detectors/logbert_detector.py b/src/detectmatelibrary/detectors/logbert_detector.py new file mode 100644 index 00000000..d9a0e5e5 --- /dev/null +++ b/src/detectmatelibrary/detectors/logbert_detector.py @@ -0,0 +1,52 @@ +from detectmatelibrary.common.deeplearning_detector import ( + DeepLearningDetectorConfig, DeepLearningDetector +) + +from detectmatelibrary.utils.deep_learning.logbert import LogBert + + +from typing import Any + + +class LogBertDetectorConfig(DeepLearningDetectorConfig): + method_type: str = "logbert_detector" + + hyperparameters: dict[str, Any] = { # type: ignore + "Model": { + "n_embed": 10, + "hidden": 32, + "num_heads": 2, + "n_layers": 1, + "dropout": 0.0, + "max_len": 1000, + }, + "Train": { + "seed": 0, + "batch_size": 256, + "learning_rate": 0.01, + "epochs": 10, + "mask_per": 0.4, + "alpha": 0.0, + "patience": 3, + }, + "Finetune": [ + ["Model", "hidden", [64, 128, 256]], + ["Model", "n_layers", [1, 2, 3]], + ["Train", "learning_rate", [0.002, 0.001, 0.005]], + ] + } + + +class LogBertDetector(DeepLearningDetector): + def __init__( + self, + name: str = "LogBertDetector", + config: LogBertDetectorConfig | dict[str, Any] = LogBertDetectorConfig(), + ) -> None: + + if isinstance(config, dict): + config = LogBertDetectorConfig.from_dict(config, name) + + super().__init__( + name=name, model_cls=LogBert, config=config + ) \ No newline at end of file diff --git a/src/detectmatelibrary/utils/deep_learning/_op.py b/src/detectmatelibrary/utils/deep_learning/_op.py new file mode 100644 index 00000000..dbdfec92 --- /dev/null +++ b/src/detectmatelibrary/utils/deep_learning/_op.py @@ -0,0 +1,46 @@ +import jax +import jax.numpy as jnp + +from typing import Any + + +class CheckPoint: + def __init__(self, patience: int) -> None: + self.last_loss = jnp.inf + self.epoch = -1 + self.param: None | dict[str, Any] = None + self.patience, self.i = patience, 0 + + def __call__(self, loss: float, epoch: int, param: dict[str, Any]) -> bool: + if self.last_loss > loss or self.param is None: + self.last_loss = loss + self.epoch = epoch + self.param = param + self.i = 0 + else: + self.i += 1 + + return self.i >= self.patience + + def load_checkpoint(self) -> tuple[int, dict[str, Any] | None]: + return self.epoch, self.param + + +class Mask: + def __init__(self, seq_size: int, mask_per: float) -> None: + self.i = 0 + self.seq_size = seq_size + self.mask_per = mask_per + self.s0 = int(seq_size * mask_per) + self.s1 = self.seq_size - self.s0 + + def __call__(self, batch_size: int) -> jnp.ndarray: + mask = jnp.concat([ + jnp.ones((batch_size, self.s1)), jnp.zeros((batch_size, self.s0)) + ], axis=1).astype(jnp.int32) + + seed = jax.random.key(self.i) + self.i += 1 + + mask = jax.random.permutation(seed, mask, independent=True, axis=1) + return mask diff --git a/src/detectmatelibrary/utils/deep_learning/deeplog.py b/src/detectmatelibrary/utils/deep_learning/deeplog.py new file mode 100644 index 00000000..40bac7dc --- /dev/null +++ b/src/detectmatelibrary/utils/deep_learning/deeplog.py @@ -0,0 +1,216 @@ +import jax.numpy as jnp +import jax + +import flax.linen as nn +import optax + +from functools import lru_cache + +from dataclasses import dataclass +from typing import Any +from tqdm import tqdm +from math import ceil + +from detectmatelibrary.utils.deep_learning.imodel import DeepModel +from detectmatelibrary.utils.deep_learning._op import CheckPoint +from detectmatelibrary.utils.finetune import Combinations + +import logging + + +## Model Deeplog +class DeepLogModel(nn.Module): + hidden_dim: int + n_layers: int + output_size: int = 1 + + @nn.compact + def __call__(self, x: jnp.ndarray) -> jnp.ndarray: + for _ in range(self.n_layers): + lstm_cell = nn.OptimizedLSTMCell(features=self.hidden_dim) + x = nn.RNN(lstm_cell)(x) + + last_step = x[:, -1, :] + return nn.Dense(features=self.output_size)(last_step) + + +## Train script +@dataclass +class TrainConfig: + seed: int = 0 + epochs: int = 40 + learning_rate: float = 0.05 + batch_size: int = 2 + patience: int = 3 + + +def loss_f(model: nn.Module, params: dict[str, Any], x: jnp.ndarray, y: jnp.ndarray) -> jnp.ndarray: + return optax.softmax_cross_entropy_with_integer_labels( + logits=model.apply({'params': params}, x), labels=y + ).mean() + + +def train( + model: nn.Module, + x: jnp.ndarray, + y: jnp.ndarray, + x_val: jnp.ndarray, + y_val: jnp.ndarray, + trainConfig: TrainConfig = TrainConfig() +) -> tuple[dict[str, Any], dict[str, float]]: + @jax.jit + def train_step( + params: dict[str, Any], opt_state: optax.OptState, x: jnp.ndarray, y: jnp.ndarray + ) -> jnp.ndarray: + def loss_fn(params: dict[str, Any]) -> jnp.ndarray: + return loss_f(params=params, x=x, y=y, model=model) + + loss, grads = jax.value_and_grad(loss_fn)(params) + updates, opt_state = optimizer.update(grads, opt_state) + params = optax.apply_updates(params, updates) + return params, opt_state, loss + + key = jax.random.PRNGKey(trainConfig.seed) + n_steps = ceil(x.shape[0] / trainConfig.batch_size) + + variables = model.init(key, x[:1]) + params = variables['params'] + optimizer = optax.adam(learning_rate=trainConfig.learning_rate) + opt_state = optimizer.init(params) + + idx = jnp.arange(x.shape[0]) + idx = jax.random.permutation(jax.random.key(trainConfig.seed), idx) + + losses_epoch, losses_step, loss_val = [], [], [] + checkpoint = CheckPoint(trainConfig.patience) + for epoch in tqdm(range(trainConfig.epochs), desc="training..."): + step_loss = 0 + for step_idx in jnp.array_split(idx, n_steps): + params, opt_state, loss = train_step( + params, opt_state, x[step_idx], y[step_idx] + ) + step_loss += loss + losses_step.append(loss) + losses_epoch.append(step_loss / n_steps) + loss_val.append(loss_f(model=model, params=params, x=x_val, y=y_val)) + idx = jax.random.permutation(jax.random.key(epoch), idx) + if checkpoint(loss=loss_val[-1], epoch=epoch, param=params): + logging.info("Early stop") + break + + best_e, params = checkpoint.load_checkpoint() + logging.info(f"Best epoch {best_e} -> Train {losses_epoch[best_e]} Val {loss_val[best_e]}") + return params, { + "Loss Epoch": losses_epoch, + "Loss Step": losses_step, + "Loss Val": loss_val, + "Best val": loss_val[best_e] + } + + +def do_train( + model: nn.Module, train_seqs: jnp.ndarray, val_seqs: jnp.ndarray, config: TrainConfig +) -> tuple[dict[str, Any], dict[str, float]]: + return train( + model=model, + x=train_seqs[:, :-1, :], + y=train_seqs[:, -1, :].reshape((train_seqs.shape[0])), + x_val=val_seqs[:, :-1, :], + y_val=val_seqs[:, -1, :].reshape((val_seqs.shape[0])), + trainConfig=config + ) + + +## Final model +default_config = { + "Model": { + "hidden_dim": 64, + "n_layers": 2, + }, + "Train": { + "seed": 0, + "batch_size": 2048, + "learning_rate": 0.01, + "epochs": 10, + "patience": 3, + }, +} + + +class DeepLog(DeepModel): + def __init__(self, config: dict = default_config) -> None: + self.config = config + self.params = {} + self.config_train = TrainConfig(**config["Train"]) + self.model_trained = False + self.model: DeepLogModel | None = None + + def __str__(self) -> str: + return str(self.model) + "\n" + str(self.config_train) + + def top_pred(self, x: jnp.ndarray) -> jnp.ndarray: + return jnp.argsort( + self.model.apply({"params": self.params}, x[None, ..., None]).flatten(), + descending=True + ) + + def get_best_k(self, seq: jnp.ndarray) -> int: + if seq.shape[0] == 0: + return 0 + + x_s, y_s = seq[:, :-1, :], seq[:, -1] + x_s = jnp.argsort( + self.model.apply({"params": self.params}, x_s), descending=True + ) + + return int(jax.scipy.stats.mode( + (jnp.arange(x_s.shape[1]) * (x_s == y_s)).sum(1) + ).mode + 2) # give a little space for variation + + @lru_cache + def check_anomaly(self, seq: tuple[int], top_k: int) -> bool: + if not self.model_trained: + return False + + seq = jnp.array(seq) + x, y = seq[:-1], seq[-1] + return not jnp.isin(y, self.top_pred(x)[:top_k]) + + def _prepare_data(self, seqs: list[tuple[int]], var_per: float) -> tuple[jnp.ndarray]: + seed = jax.random.key(self.config_train.seed) + idx = jax.random.permutation(seed, len(seqs)) + seqs = jnp.array(seqs)[..., None][idx] + train_seqs = seqs[:ceil(len(seqs) * (1 - var_per))] + val_seqs = seqs[ceil(len(seqs) * (1 - var_per)):] + return train_seqs, val_seqs + + def train(self, seqs: list[tuple[int]], var_per: float) -> dict[str, int | float]: + train_seqs, val_seqs = self._prepare_data(seqs=seqs, var_per=var_per) + + self.config_train = TrainConfig(**self.config["Train"]) + self.config["Model"]["output_size"] = train_seqs.max() + 1 + logging.info(f"Output shape: {self.config["Model"]["output_size"]}") + + self.model = DeepLogModel(**self.config["Model"]) + self.params, stats = do_train( + self.model, train_seqs=train_seqs, val_seqs=val_seqs, config=self.config_train + ) + self.model_trained = True + stats["top_k"] = self.get_best_k(val_seqs) + + return stats + + def finetune(self, seqs: list[tuple[int]], var_per: float, epochs: int = 2) -> None: + train_seqs, val_seqs = self._prepare_data(seqs=seqs, var_per=var_per) + combos = Combinations(config=self.config) + + for comb in combos(): + comb["Model"]["output_size"] = train_seqs.max() + 1 + comb["Train"]["epochs"] = epochs + model = DeepLogModel(**comb["Model"]) + config_train = TrainConfig(**comb["Train"]) + + _, stats = do_train(model, train_seqs=train_seqs, val_seqs=val_seqs, config=config_train) + combos.add_value(stats["Best val"]) + self.config = combos.get_best() + logging.info(self.config) \ No newline at end of file diff --git a/src/detectmatelibrary/utils/deep_learning/imodel.py b/src/detectmatelibrary/utils/deep_learning/imodel.py new file mode 100644 index 00000000..51446d50 --- /dev/null +++ b/src/detectmatelibrary/utils/deep_learning/imodel.py @@ -0,0 +1,19 @@ + +from abc import ABC, abstractmethod + + +class DeepModel(ABC): + def __init__(self) -> None: + pass + + @abstractmethod + def check_anomaly(self, seq: tuple[int], top_k: int) -> bool: + pass + + @abstractmethod + def train(self, seqs: list[tuple[int]], var_per: float) -> dict[str, int | float]: + pass + + @abstractmethod + def finetune(self, seqs: list[tuple[int]], var_per: float, epochs: int = 2) -> None: + pass diff --git a/src/detectmatelibrary/utils/deep_learning/logbert.py b/src/detectmatelibrary/utils/deep_learning/logbert.py new file mode 100644 index 00000000..cb7f8793 --- /dev/null +++ b/src/detectmatelibrary/utils/deep_learning/logbert.py @@ -0,0 +1,290 @@ + +import jax.numpy as jnp +import jax + +import flax.linen as nn +import optax + +from functools import lru_cache + +from dataclasses import dataclass +from typing import Any +from tqdm import tqdm +from math import ceil + +from detectmatelibrary.utils.deep_learning._op import CheckPoint, Mask +from detectmatelibrary.utils.deep_learning.imodel import DeepModel +from detectmatelibrary.utils.finetune import Combinations + +import logging + + +class PositionEmbedding(nn.Module): + """ + Equation: + pi,2j = sin(i / 10000^(2j/d)) + pi,2j+1 = cos(i / 10000^(2j/d)) + """ + hidden: int + max_len: int = 1000 + + def setup(self) -> None: + x = jnp.arange(self.max_len, dtype=jnp.float32).reshape(-1, 1) + jd = jnp.power(10000, jnp.arange(0, self.hidden, 2, dtype=jnp.float32) / self.hidden) + + position = jnp.zeros((1, self.max_len, self.hidden)) + position = position.at[:, :, 0::2].set(jnp.sin(x / jd)) + position = position.at[:, :, 1::2].set(jnp.cos(x / jd)) + self.position = position + + @nn.compact + def __call__(self, X: jnp.ndarray) -> jnp.ndarray: + X = X + self.position[:, :X.shape[1] ,:] + return X + + +class Embedding(nn.Module): + n_embed: int + hidden: int + dropout: float + max_len: int = 1000 + + def setup(self) -> None: + self.position = PositionEmbedding(hidden=self.hidden, max_len=self.max_len) + self.embed = nn.Embed(num_embeddings=self.n_embed, features=self.hidden) + + @nn.compact + def __call__(self, x: jnp.ndarray, training: bool = False) -> jnp.ndarray: + h = self.embed(x) + h = h + self.position(h) + return nn.Dropout(self.dropout)(h, deterministic=not training) + + +class TransformerBlock(nn.Module): + hidden: int + num_heads: int + dropout: float + + @nn.compact + def __call__(self, x: jnp.ndarray, training: bool = False) -> jnp.ndarray: + h = nn.MultiHeadAttention( + num_heads=self.num_heads, + dropout_rate=self.dropout, + out_features=self.hidden, + deterministic=not training + )(inputs_q=x) + + h = nn.LayerNorm()(h) + + h = nn.gelu(nn.Dense(self.hidden * 2)(h)) + h = nn.Dropout(self.dropout)(h, deterministic=not training) + h = nn.Dense(self.hidden)(h) + + return h + x + + +class LogBertModel(nn.Module): + n_embed: int + hidden: int + num_heads: int + n_layers: int + dropout: float + max_len: int = 1000 + + def setup(self) -> None: + self.special_tokens = 2 + + @nn.compact + def __call__(self, x: jnp.ndarray, training: bool = False) -> tuple[jnp.ndarray, jnp.ndarray]: + x = jnp.concat([x, jnp.ones((x.shape[0], 1), dtype=jnp.int32)], axis=1) + h = Embedding( + n_embed=self.n_embed + self.special_tokens, hidden=self.hidden, dropout=self.dropout + )(x + self.special_tokens) + + for _ in range(self.n_layers): + h = TransformerBlock( + hidden=self.hidden, num_heads=self.num_heads, dropout=self.dropout + )(h, training=training) + h, dist = h[:, :-1, :], h[:, -1, :] + + return nn.Dense(self.n_embed)(h), dist + + +@dataclass +class TrainConfig: + seed: int = 0 + epochs: int = 40 + learning_rate: float = 0.01 + batch_size: int = 2 + mask_per: float = 0.4 + alpha: float = 0.0 + patience: int = 3 + + +def loss_( + model: nn.Module, params: dict[str, Any], x: jnp.ndarray, m: jnp.ndarray, alpha: float +) -> jnp.ndarray: + + logist, h_dist = model.apply({"params": params}, x * m, training=True) + loss_mlkp = ((1 - m) * optax.softmax_cross_entropy_with_integer_labels( + logits=logist, labels=x + )) + loss_vhm = jnp.linalg.norm(h_dist - h_dist.mean(axis=1)[..., None], axis=1) ** 2 + + return loss_mlkp.mean() + alpha * loss_vhm.mean() + + +def train( + model: nn.Module, x: jnp.ndarray, x_val: jnp.ndarray, mask: Mask, trainConfig=TrainConfig() +) -> tuple[dict[str, Any], dict[str, float]]: + @jax.jit + def train_step(params, opt_state, x, m, alpha): + def loss_f(params): + return loss_(model=model, params=params, x=x, m=m, alpha=alpha) + + loss, grads = jax.value_and_grad(loss_f)(params) + updates, opt_state = optimizer.update(grads, opt_state) + params = optax.apply_updates(params, updates) + return params, opt_state, loss + + seed = jax.random.key(trainConfig.seed) + n_steps = ceil(x.shape[0] / trainConfig.batch_size) + + params = model.init(seed, jnp.zeros((1, x.shape[-1]), dtype=jnp.int32))["params"] + optimizer = optax.adam(learning_rate=trainConfig.learning_rate) + opt_state = optimizer.init(params) + + idx = jnp.arange(x.shape[0]) + idx = jax.random.permutation(jax.random.key(trainConfig.seed), idx) + + losses_epoch, losses_step, loss_val = [], [], [] + checkpoint, alpha = CheckPoint(TrainConfig.patience), trainConfig.alpha + for epoch in tqdm(range(trainConfig.epochs), desc="training..."): + step_loss = 0 + for step_idx in jnp.array_split(idx, n_steps): + x_step = x[step_idx] + m = mask(x_step.shape[0]) + params, opt_state, loss = train_step(params, opt_state, x_step, m, alpha=alpha) + step_loss += loss + losses_step.append(loss) + + m = mask(x_val.shape[0]) + loss_val.append(loss_(model=model, params=params, x=x_val, m=m, alpha=alpha)) + losses_epoch.append(step_loss / n_steps) + idx = jax.random.permutation(jax.random.key(epoch), idx) + if checkpoint(loss=loss_val[-1], epoch=epoch, param=params): + logging.info("Early stop") + break + + best_e, params = checkpoint.load_checkpoint() + logging.info(f"Best epoch {best_e} -> Train {losses_epoch[best_e]} Val {loss_val[best_e]}") + return params, { + "Loss Epoch": losses_epoch, "Loss Step": losses_step, "Loss Val": loss_val, "Best val": loss_val[best_e] + } + + +## Final model +default_config = { + "Model": { + "hidden": 256, + "num_heads": 2, + "n_layers": 4, + "dropout": 0.0, + "max_len": 1000, + }, + "Train": { + "seed": 0, + "batch_size": 256, + "learning_rate": 0.01, + "epochs": 10, + "mask_per": 0.4, + "alpha": 0.0, + "patience": 3, + }, +} + + +class LogBert(DeepModel): + def __init__(self, config: dict = default_config) -> None: + self.config = config + self.params = {} + self.model: LogBertModel | None = None + self.mask: Mask | None = None + self.config_train = TrainConfig(**self.config["Train"]) + + def __str__(self) -> str: + return str(self.model) + "\n" + str(self.config_train) + + def top_pred(self, x: jnp.ndarray) -> tuple[jnp.ndarray]: + m = self.mask(x.shape[0]) + y, _ = self.model.apply({"params": self.params}, x * m, training=False) + + idx = jnp.nonzero(m == 0) + pred = jnp.argsort(y[idx], axis=1, descending=True) + return pred, x[idx] + + def get_best_k(self, seq: jnp.ndarray) -> int: + if seq.shape[0] == 0: + return 0 + + pred, y = self.top_pred(seq) + return int(jax.scipy.stats.mode( + ((y[..., None] == pred) * jnp.arange(pred.shape[1])[None, ...]).sum(1) + ).mode + 2) + + @lru_cache + def check_anomaly(self, seq: tuple[int], top_k: int) -> int: + if self.model is None: + return False + + pred, y = self.top_pred(jnp.array([seq])) + pred = pred[:, :top_k] + score = 0 + for i in range(pred.shape[0]): + score += not bool(jnp.isin(y[i], pred[i])) + return score + + def _prepare_data(self, seqs: list[tuple[int]], var_per: float) -> tuple[jnp.ndarray]: + seed = jax.random.key(self.config_train.seed) + idx = jax.random.permutation(seed, len(seqs)) + seqs = jnp.array(seqs)[idx] + train_seqs = seqs[:ceil(len(seqs) * (1 - var_per))] + val_seqs = seqs[ceil(len(seqs) * (1 - var_per)):] + return train_seqs, val_seqs + + def train(self, seqs: list[tuple[int]], var_per: float) -> dict[str, int | float]: + train_seqs, val_seqs = self._prepare_data(seqs=seqs, var_per=var_per) + + self.config_train = TrainConfig(**self.config["Train"]) + self.config["Model"]["n_embed"] = int(train_seqs.max() + 1) + self.model = LogBertModel(**self.config["Model"]) + self.mask = Mask( + seq_size=train_seqs.shape[-1], mask_per=self.config_train.mask_per + ) + self.params, stats = train( + model=self.model, mask=self.mask, x=train_seqs, x_val=val_seqs, trainConfig=self.config_train + ) + stats["top_k"] = self.get_best_k(val_seqs) + + return stats + + def finetune(self, seqs: list[tuple[int]], var_per: float, epochs: int = 2) -> None: + train_seqs, val_seqs = self._prepare_data(seqs=seqs, var_per=var_per) + combos = Combinations(config=self.config) + + for comb in combos(): + comb["Model"]["n_embed"] = int(train_seqs.max() + 1) + comb["Train"]["epochs"] = epochs + model = LogBertModel(**comb["Model"]) + config_train = TrainConfig(**comb["Train"]) + mask = Mask( + seq_size=train_seqs.shape[-1], mask_per=config_train.mask_per + ) + + _, stats = train( + model=model, mask=mask, x=train_seqs, x_val=val_seqs, trainConfig=config_train + ) + combos.add_value(stats["Best val"]) + self.config = combos.get_best() + logging.info(self.config) + \ No newline at end of file diff --git a/src/detectmatelibrary/utils/finetune.py b/src/detectmatelibrary/utils/finetune.py new file mode 100644 index 00000000..8015531a --- /dev/null +++ b/src/detectmatelibrary/utils/finetune.py @@ -0,0 +1,86 @@ +from detectmatelibrary.common.core import CoreConfig + +from typing import Any +import numpy as np + +import itertools +import warnings +import typing +import copy + + +class CombOp: + @staticmethod + def is_finetune_in_there(config: CoreConfig | dict[str, Any]) -> bool: + variables = dir(config) if isinstance(config, CoreConfig) else config + return "Finetune" in variables + + @staticmethod + def get_value(config: CoreConfig | dict[str, Any], path: str) -> Any: + return getattr(config, path) if isinstance(config, CoreConfig) else config[path] + + @staticmethod + def set_value(config: CoreConfig | dict[str, Any], path: list[str], value: Any) -> None: + if isinstance(config, CoreConfig): + setattr(config, path[0], value) + else: + path = [path] if isinstance(path, str) else path + dict_ = config + for v in path[:-1]: + dict_ = dict_[v] + + dict_[path[-1]] = value + + @staticmethod + def value_exist(config: CoreConfig | dict[str, Any], path: list[str]) -> bool: + if isinstance(config, CoreConfig): + return path[0] in dir(config) + + path = [path] if isinstance(path, str) else path + dict_ = config + for v in path[:-1]: + dict_ = dict_[v] + + return path[-1] in dict_ + + +class Combinations: + def __init__(self, config: CoreConfig | dict[str, Any]) -> None: + self.config = config + + self.paths: list[list[str]] = [] + self.combs: list[tuple[Any, ...]] = [] + if not CombOp.is_finetune_in_there(config): + warnings.warn("No finetune options found") + else: + self.paths = [path[:-1] for path in CombOp.get_value(config, "Finetune")] + self.combs = list(itertools.product( + *[path[-1] for path in CombOp.get_value(config, "Finetune")] + )) + self.values: list[float] = [] + + def add_value(self, value: float) -> None: + self.values.append(value) + + def get_best(self) -> CoreConfig | dict[str, Any]: + if self.values == []: + return self.config + + idx = np.argmin(self.values) + for i, combo in enumerate(self.combs): + if i == idx: + config = copy.deepcopy(self.config) + for path, value in zip(self.paths, combo): + if CombOp.value_exist(config, path): + CombOp.set_value(config, path, value) + return config + return self.config + + def __call__(self) -> typing.Iterable[CoreConfig | dict[str, Any]]: + for combo in self.combs: + config = copy.deepcopy(self.config) + for path, value in zip(self.paths, combo): + if CombOp.value_exist(config, path): + CombOp.set_value(config, path, value) + + yield config diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 00000000..e9ae1503 --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,23 @@ +# conftest.py +import pytest + + +def pytest_addoption(parser): + # This creates the command line flag + parser.addoption( + "--run-ignored", + action="store_true", + default=False + ) + + +def pytest_collection_modifyitems(config, items): + """Automatically skips tests tagged as 'ignored' unless flag is passed.""" + if config.getoption("--run-ignored"): + # If terminal flag is present, don't skip anything + return + + skip_marker = pytest.mark.skip(reason="Skipped: requires --run-ignored flag") + for item in items: + if "ignored" in item.keywords: + item.add_marker(skip_marker) diff --git a/tests/test_common/test_deeplearning_detectors.py b/tests/test_common/test_deeplearning_detectors.py new file mode 100644 index 00000000..c2260fe6 --- /dev/null +++ b/tests/test_common/test_deeplearning_detectors.py @@ -0,0 +1,69 @@ + +from detectmatelibrary.common.deeplearning_detector import ( + DeepLearningDetector, DeepLearningDetectorConfig +) + +from detectmatelibrary.utils.deep_learning.imodel import DeepModel +from detectmatelibrary import schemas + +import pytest + + +class Flag(Exception): + pass + + +class DummyDeepModel(DeepModel): + def __init__(self, *args, **kargs): + super().__init__() + + def check_anomaly(self, seq, top_k): + raise Flag() + + def train(self, seqs, var_per): + return {} + + def finetune(self, seqs, var_per, epochs=2): + return None + + +class TestDeepLearning: + def test_normal_run_configure(self): + deep_learning_detector = DeepLearningDetector( + model_cls=DummyDeepModel, config=DeepLearningDetectorConfig(window_size=1) + ) + + for i in range(10): + deep_learning_detector.configure( + [schemas.ParserSchema({"EventID": i})] + ) + + for i in range(10): + assert deep_learning_detector.config_seqs[i] == (i,) + + deep_learning_detector.set_configuration() + assert len(deep_learning_detector.config_seqs) == 0 + + def test_normal_run_train(self): + deep_learning_detector = DeepLearningDetector( + model_cls=DummyDeepModel, config=DeepLearningDetectorConfig(window_size=1) + ) + + for i in range(10): + deep_learning_detector.train( + [schemas.ParserSchema({"EventID": i})] + ) + + for i in range(10): + assert deep_learning_detector.train_seqs[i] == (i,) + + deep_learning_detector.post_train() + assert len(deep_learning_detector.train_seqs) == 0 + + def test_check_that_check_anomaly_is_call(self): + deep_learning_detector = DeepLearningDetector( + model_cls=DummyDeepModel, config=DeepLearningDetectorConfig(window_size=1) + ) + + with pytest.raises(Flag): + deep_learning_detector.process(schemas.ParserSchema({"EventID": 0})) diff --git a/tests/test_detectors/test_bigram_frequency_detector.py b/tests/test_detectors/test_bigram_frequency_detector.py index 399b1bd8..878630e8 100644 --- a/tests/test_detectors/test_bigram_frequency_detector.py +++ b/tests/test_detectors/test_bigram_frequency_detector.py @@ -22,6 +22,9 @@ from detectmatelibrary.utils.aux import time_test_mode from tests.test_data import AUDIT_LOG, AUDIT_TEMPLATES, TRAIN_UNTIL # Set time test mode for consistent timestamps +import pytest + + time_test_mode() @@ -238,6 +241,7 @@ def test_detect_known_value_alert(self): class TestBigramFrequencyDetectorEndToEnd: """Regression test: full configure/train/detect pipeline on audit.log.""" + @pytest.mark.ignored def test_audit_log_anomalies(self): parser = MatcherParser(config=_PARSER_CONFIG) detector = BigramFrequencyDetector( @@ -268,6 +272,7 @@ class TestBigramFrequencyDetectorAutoConfig: """Test that process() drives configure/set_configuration/train/detect automatically.""" + @pytest.mark.ignored def test_audit_log_anomalies_via_process(self): parser = MatcherParser(config=_PARSER_CONFIG) detector = BigramFrequencyDetector(config=_SKIP_REPETITIONS_CONFIG, name="MultipleDetector") diff --git a/tests/test_detectors/test_charset_detector.py b/tests/test_detectors/test_charset_detector.py index 7d113f20..80c1c1db 100644 --- a/tests/test_detectors/test_charset_detector.py +++ b/tests/test_detectors/test_charset_detector.py @@ -19,6 +19,8 @@ from detectmatelibrary.utils.aux import time_test_mode from tests.test_data import AUDIT_LOG, AUDIT_TEMPLATES, TRAIN_UNTIL +import pytest + # Set time test mode for consistent timestamps time_test_mode() @@ -289,6 +291,7 @@ def test_detect_unknown_chars_reported_per_variable(self): class TestCharsetDetectorEndToEnd: """Regression test: full configure/train/detect pipeline on audit.log.""" + @pytest.mark.ignored def test_audit_log_anomalies(self): parser = MatcherParser(config=_PARSER_CONFIG) detector = CharsetDetector() @@ -315,6 +318,7 @@ class TestCharsetDetectorAutoConfig: """Test that process() drives configure/set_configuration/train/detect automatically.""" + @pytest.mark.ignored def test_audit_log_anomalies_via_process(self): parser = MatcherParser(config=_PARSER_CONFIG) detector = CharsetDetector() diff --git a/tests/test_detectors/test_deeplog_detector.py b/tests/test_detectors/test_deeplog_detector.py new file mode 100644 index 00000000..c4b9d8ad --- /dev/null +++ b/tests/test_detectors/test_deeplog_detector.py @@ -0,0 +1,148 @@ +""" +Main tests are done at low level in the deeplearning module as DeeplogDetector is an interface. + +Tests at this level are only end to end pipeline. +""" + +from detectmatelibrary.detectors.deeplog_detector import DeeplogDetector +from detectmatelibrary.parsers.template_matcher import MatcherParser +from detectmatelibrary.helper.from_to import From + +from tests.test_data import AUDIT_LOG, AUDIT_TEMPLATES, TRAIN_UNTIL + +import detectmatelibrary.schemas as schemas + +import pytest + + +class TestDeeplog: + @pytest.mark.ignored + def test_end2end_autoconfig(self) -> None: + config = { + "detectors": { + "DeeplogDetector": { + "method_type": "deeplog_detector", + "auto_config": True, + "data_use_configure": 10, + "data_use_training": 10, + "window_size": 3, + } + } + } + + deeplog = DeeplogDetector(config=config) + assert deeplog.get_state() == "Default" + + for j in range(2): + for i in [1, 2, 3, 4, 5, 1, 2]: + deeplog.process(schemas.ParserSchema({"EventID": i})) + + if j == 0: + assert deeplog.get_state() == "Configuring" + + assert deeplog.get_state() == "Training" + + for _ in range(2): + for i in [1, 2, 3, 4, 5, 1, 2]: + deeplog.process(schemas.ParserSchema({"EventID": i})) + assert deeplog.get_state() == "Default" + + @pytest.mark.ignored + def test_end2end_no_autoconfig(self) -> None: + config = { + "detectors": { + "DeeplogDetector": { + "method_type": "deeplog_detector", + "auto_config": False, + "data_use_training": 10, + "window_size": 3, + "hyperparameters": { + "Model": { + "hidden_dim": 64, + "n_layers": 2, + }, + "Train": { + "seed": 0, + "batch_size": 2048, + "learning_rate": 0.01, + "epochs": 10, + "patience": 3, + }, + "Finetune": [], + } + } + } + } + + deeplog = DeeplogDetector(config=config) + assert deeplog.get_state() == "Default" + + for j in range(2): + for i in [1, 2, 3, 4, 5, 1, 2]: + deeplog.process(schemas.ParserSchema({"EventID": i})) + if j == 0: + assert deeplog.get_state() == "Training" + + assert deeplog.get_state() == "Default" + + +PIPELINE_CONFIG = { + "parsers": { + "MatcherParser": { + "method_type": "matcher_parser", + "auto_config": False, + "log_format": "type= msg=audit(