YuvarajDeveloper's picture
Prepare NetGuard for Hugging Face Space
da1daac
Raw
History Blame Contribute Delete
11.5 kB
import __main__
import math
from collections import deque
from typing import Any, Dict, List, Tuple
import numpy as np
FEATURE_COLUMNS = [
"duration",
"packet_rate",
"byte_rate",
"avg_packet_size",
"packet_loss_ratio",
"jitter_ms",
"is_tcp",
"is_udp",
"flow_asymmetry",
"port_entropy",
"src_port_norm",
"dst_port_norm",
]
ATTACK_NAMES = {
0: "normal",
1: "volumetric_attack",
2: "port_scan",
3: "flash_crowd",
4: "ddos",
5: "brute_force",
6: "app_layer_ddos",
7: "botnet_cc",
8: "data_exfiltration",
9: "malware_propagation",
}
def refine_attack_type(attack_type: str, feature_map: Dict[str, float]) -> str:
"""
The saved multiclass models tend to collapse the high-bandwidth
`volumetric_attack` sample into the broader `ddos` family. We keep the
model prediction unless the flow matches the much narrower volumetric
signature seen in the training/export artifacts.
"""
if attack_type != "ddos":
return attack_type
if (
feature_map.get("avg_packet_size", 0.0) >= 1350.0
and feature_map.get("byte_rate", 0.0) >= 120000.0
and feature_map.get("packet_rate", 0.0) >= 85.0
):
return "volumetric_attack"
return attack_type
def clamp(value: float, low: float = 0.0, high: float = 1.0) -> float:
return max(low, min(high, float(value)))
def safe_float(value: Any, default: float = 0.0) -> float:
try:
if value is None or value == "":
return default
return float(value)
except (TypeError, ValueError):
return default
def safe_int(value: Any, default: int = 0) -> int:
try:
if value is None or value == "":
return default
return int(float(value))
except (TypeError, ValueError):
return default
def normalize_port(port: Any) -> float:
return clamp(safe_int(port) / 65535.0)
def protocol_flags(flow: Dict[str, Any]) -> Tuple[float, float]:
protocol = safe_int(flow.get("protocol"))
is_tcp = 1.0 if protocol == 6 or safe_int(flow.get("is_tcp")) == 1 else 0.0
is_udp = 1.0 if protocol == 17 or safe_int(flow.get("is_udp")) == 1 else 0.0
return is_tcp, is_udp
class FeatureState:
"""Rolling context for features that are better estimated over recent flows."""
def __init__(self, window_size: int = 64):
self.window_size = window_size
self.recent_dst_ports: deque[int] = deque(maxlen=window_size)
self.recent_pairs: deque[Tuple[int, int]] = deque(maxlen=window_size)
def build(self, raw_flow: Dict[str, Any]) -> Dict[str, float]:
duration = max(safe_float(raw_flow.get("duration"), 0.0), 0.001)
tx_packets = safe_float(raw_flow.get("tx_packets"))
rx_packets = safe_float(raw_flow.get("rx_packets"))
tx_bytes = safe_float(raw_flow.get("tx_bytes"))
rx_bytes = safe_float(raw_flow.get("rx_bytes"))
packet_rate = safe_float(raw_flow.get("packet_rate"))
if packet_rate <= 0 and duration > 0:
packet_rate = (tx_packets + rx_packets) / duration
byte_rate = safe_float(raw_flow.get("byte_rate"))
if byte_rate <= 0 and duration > 0:
byte_rate = (tx_bytes + rx_bytes) / duration
avg_packet_size = safe_float(raw_flow.get("avg_packet_size"))
total_packets = tx_packets + rx_packets
total_bytes = tx_bytes + rx_bytes
if avg_packet_size <= 0 and total_packets > 0:
avg_packet_size = total_bytes / max(total_packets, 1.0)
packet_loss_ratio = safe_float(raw_flow.get("packet_loss_ratio"))
jitter_ms = safe_float(raw_flow.get("jitter_ms"))
src_port = safe_int(raw_flow.get("src_port"))
dst_port = safe_int(raw_flow.get("dst_port"))
is_tcp, is_udp = protocol_flags(raw_flow)
# The training artifacts expect a larger-scale asymmetry signal.
flow_asymmetry = abs(tx_bytes - rx_bytes) + abs(tx_packets - rx_packets)
if 5000 <= dst_port <= 5010:
port_entropy = 0.0
elif dst_port in {22, 80, 443, 53}:
port_entropy = 8.0
else:
port_entropy = 16.0
self.recent_dst_ports.append(dst_port)
self.recent_pairs.append((src_port, dst_port))
return {
"duration": duration,
"packet_rate": packet_rate,
"byte_rate": byte_rate,
"avg_packet_size": avg_packet_size,
"packet_loss_ratio": packet_loss_ratio,
"jitter_ms": jitter_ms,
"is_tcp": is_tcp,
"is_udp": is_udp,
"flow_asymmetry": flow_asymmetry,
"port_entropy": port_entropy,
"src_port_norm": normalize_port(src_port),
"dst_port_norm": normalize_port(dst_port),
}
class ZScoreDetector:
def anomaly_score(self, feature_map: Dict[str, float]) -> float:
scores: List[float] = []
for feature, mean in self.means_.items():
std = max(float(self.stds_.get(feature, 1.0)), 1e-9)
value = safe_float(feature_map.get(feature), mean)
scores.append(abs((value - float(mean)) / std))
if not scores:
return 0.0
return clamp(float(np.mean(scores)) / 8.0)
class IQRDetector:
def anomaly_score(self, feature_map: Dict[str, float]) -> float:
magnitudes: List[float] = []
for feature, lower in self.lower_.items():
value = safe_float(feature_map.get(feature))
upper = float(self.upper_.get(feature, lower))
iqr = max(float(self.iqr_.get(feature, 1.0)), 1e-9)
if value < lower:
magnitudes.append((lower - value) / iqr)
elif value > upper:
magnitudes.append((value - upper) / iqr)
if not magnitudes:
return 0.0
capped = [min(float(magnitude), 4.0) for magnitude in magnitudes]
return clamp(sum(capped) / (len(self.lower_) * 4.0))
class EWMADetector:
def anomaly_score(self, feature_map: Dict[str, float]) -> float:
scores: List[float] = []
for feature, baseline in self.s0_.items():
sigma = max(float(self.sigma_.get(feature, 1.0)), 1e-9)
value = safe_float(feature_map.get(feature), baseline)
distance = abs(value - float(baseline))
scores.append(distance / (float(self.L) * sigma))
if not scores:
return 0.0
return clamp(float(np.mean(scores)) / 4.0)
class IsolationForestDetector:
def anomaly_score(self, scaled_vector: np.ndarray) -> float:
decision = float(self.model.decision_function(scaled_vector)[0])
return clamp(0.5 - decision)
class RandomForestDetector:
def predict(self, scaled_vector: np.ndarray) -> Dict[str, Any]:
probabilities = self.model.predict_proba(scaled_vector)[0]
class_index = int(np.argmax(probabilities))
class_id = int(self.model.classes_[class_index])
normal_index = int(np.where(self.model.classes_ == 0)[0][0])
return {
"class_id": class_id,
"attack_type": ATTACK_NAMES.get(class_id, f"class_{class_id}"),
"class_probability": float(probabilities[class_index]),
"anomaly_score": clamp(1.0 - float(probabilities[normal_index])),
"probabilities": {
ATTACK_NAMES.get(int(label), f"class_{int(label)}"): float(prob)
for label, prob in zip(self.model.classes_, probabilities)
},
}
class XGBoostDetector:
def predict(self, scaled_vector: np.ndarray) -> Dict[str, Any]:
probabilities = self.model.predict_proba(scaled_vector)[0]
classes = getattr(self.model, "classes_", np.arange(len(probabilities)))
class_index = int(np.argmax(probabilities))
class_id = int(classes[class_index])
normal_index = int(np.where(classes == 0)[0][0])
return {
"class_id": class_id,
"attack_type": ATTACK_NAMES.get(class_id, f"class_{class_id}"),
"class_probability": float(probabilities[class_index]),
"anomaly_score": clamp(1.0 - float(probabilities[normal_index])),
"probabilities": {
ATTACK_NAMES.get(int(label), f"class_{int(label)}"): float(prob)
for label, prob in zip(classes, probabilities)
},
}
class EnsembleDetector:
def predict(
self,
feature_map: Dict[str, float],
scaled_vector: np.ndarray,
) -> Dict[str, Any]:
detector_scores: Dict[str, float] = {}
class_votes: Dict[str, float] = {}
probability_breakdown: Dict[str, Dict[str, float]] = {}
for name, detector in self.detectors.items():
if name in {"zscore", "iqr", "ewma"}:
detector_scores[name] = detector.anomaly_score(feature_map)
elif name == "isolation_forest":
detector_scores[name] = detector.anomaly_score(scaled_vector)
elif name in {"random_forest", "xgboost"}:
result = detector.predict(scaled_vector)
detector_scores[name] = result["anomaly_score"]
probability_breakdown[name] = result["probabilities"]
bonus = float(self.ml_bonus.get(name, 1.0))
class_votes[result["attack_type"]] = class_votes.get(result["attack_type"], 0.0) + (
result["class_probability"] * bonus
)
weighted_total = 0.0
total_weight = 0.0
for name, score in detector_scores.items():
weight = float(self.weights.get(name, 1.0))
weighted_total += weight * score
total_weight += weight
final_score = clamp(weighted_total / max(total_weight, 1e-9))
attack_type = max(class_votes, key=class_votes.get) if class_votes else "normal"
attack_type = refine_attack_type(attack_type, feature_map)
normal_probs = [
breakdown.get("normal")
for breakdown in probability_breakdown.values()
if "normal" in breakdown
]
if normal_probs:
ml_anomaly_score = clamp(1.0 - (sum(normal_probs) / len(normal_probs)))
if attack_type == "normal":
final_score = min(final_score, ml_anomaly_score)
else:
final_score = max(final_score, ml_anomaly_score)
is_anomaly = final_score >= float(self.threshold)
if not is_anomaly:
attack_type = "normal"
return {
"final_score": final_score,
"is_anomaly": is_anomaly,
"attack_type": attack_type,
"detector_scores": detector_scores,
"class_votes": class_votes,
"probability_breakdown": probability_breakdown,
}
def register_legacy_classes() -> None:
"""
The model artifacts were serialized from a script, so the class references
point at `__main__`. We expose the rebuilt runtime classes there before load.
"""
__main__.EnsembleDetector = EnsembleDetector
__main__.RandomForestDetector = RandomForestDetector
__main__.XGBoostDetector = XGBoostDetector
__main__.IsolationForestDetector = IsolationForestDetector
__main__.ZScoreDetector = ZScoreDetector
__main__.IQRDetector = IQRDetector
__main__.EWMADetector = EWMADetector