Dennislee:)
Deploy autoencoder-unsw (auto)
b41fe25
Raw
History Blame Contribute Delete
9.02 kB
"""
Axiom Autoencoder UNSW-NB15 - 高維 log、行為壓縮與偵測(49 維)
使用 Autoencoder 分析 UNSW-NB15 高維特徵。
支援 POST /train {"csv_url": "..."} 從 Supabase 匯出資料訓練(5 維 pad 至 49)。
"""
import csv
import io
import os
import pickle
import threading
import time
import urllib.request
from typing import List, Optional
import numpy as np
from fastapi import Body, FastAPI
from fastapi.middleware.cors import CORSMiddleware
from pydantic import BaseModel
app = FastAPI(title="Axiom Autoencoder UNSW", version="0.1.0")
app.add_middleware(CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"])
ANOMALY_THRESHOLD = float(os.environ.get("ANOMALY_THRESHOLD", "0.7"))
_raw_repo = os.environ.get("HF_REPO") or os.environ.get("hf_repo", "")
HF_REPO = _raw_repo.strip() if "/" in _raw_repo else ""
if _raw_repo and not HF_REPO:
print(f"HF_REPO 格式無效(需 username/repo_name): {_raw_repo!r},跳過模型載入")
HF_TOKEN = (os.environ.get("HF_TOKEN") or os.environ.get("hf_token", "")).strip()
SKLEARN_N_SAMPLES = int(os.environ.get("SKLEARN_N_SAMPLES", "2000"))
FEATURE_DIM = 49
_model = None
_scaler = None
_threshold = 0.1
_training = False
_training_lock = threading.Lock()
class ScoreRequest(BaseModel):
tenant_id: str
device_id: str
features: List[float]
sequence: Optional[List[List[float]]] = None
class ScoreResponse(BaseModel):
anomaly_score: float
is_anomaly: bool
details: Optional[dict] = None
def _load_model():
global _model, _scaler, _threshold
if not HF_REPO or not HF_TOKEN:
return
try:
from huggingface_hub import hf_hub_download
path = hf_hub_download(repo_id=HF_REPO, filename="model.pt", token=HF_TOKEN)
with open(path, "rb") as f:
b = pickle.load(f)
_model = b.get("model")
_scaler = b.get("scaler")
_threshold = b.get("threshold", 0.1)
print(f"Loaded Autoencoder from {HF_REPO}")
except Exception as e:
print(f"Load failed: {e}")
_model = _scaler = None
def _heuristic_score(features: List[float]) -> tuple[float, dict]:
if not features:
return 0.0, {"reason": "empty_features"}
arr = np.array(features, dtype=np.float64)
std = float(np.std(arr))
score = min(1.0, std / 2.0) if std > 0 else 0.0
return score, {"feature_count": len(arr)}
def compute_score(features: List[float]) -> tuple[float, dict]:
if _model is not None and _scaler is not None:
try:
import torch
x = np.array([features], dtype=np.float64)
x_scaled = _scaler.transform(x)
x_t = torch.tensor(x_scaled, dtype=torch.float32)
with torch.no_grad():
recon = _model(x_t)
mse = float(((x_t - recon) ** 2).mean().item())
score = min(1.0, mse / (_threshold + 1e-8))
return min(1.0, max(0.0, score)), {"source": "autoencoder", "mse": mse, "feature_count": len(features)}
except Exception:
return _heuristic_score(features)
return _heuristic_score(features)
def _load_csv_from_url(csv_url: str, feature_dim: int = 49) -> np.ndarray:
"""下載 CSV,解析為 5 維,pad 至 feature_dim。"""
with urllib.request.urlopen(csv_url, timeout=60) as resp:
raw = resp.read().decode("utf-8")
rows = []
for r in csv.DictReader(io.StringIO(raw)):
risk = min(1.0, max(0.0, float(r.get("risk_score", 0) or 0) / 100.0))
sev = float(r.get("severity", 0.25) or 0.25)
layer = min(1.0, float(r.get("rule_layer", 1) or 1) / 5.0)
dev = (hash(r.get("device_id", "") or "") % 10000) / 10000.0
ag = (hash(r.get("agent_id", "") or "") % 10000) / 10000.0
rows.append([risk, sev, layer, dev, ag])
if not rows:
return np.empty((0, feature_dim))
X5 = np.array(rows, dtype=np.float64)
if X5.shape[1] < feature_dim:
pad = np.zeros((X5.shape[0], feature_dim - X5.shape[1]))
X5 = np.hstack([X5, pad])
return X5
def _fetch_unsw(n_samples: int, rs: int) -> np.ndarray:
try:
from datasets import load_dataset
try:
ds = load_dataset("Mouwiya/UNSW-NB15", split="train", trust_remote_code=True)
except Exception:
ds = load_dataset("wwydmanski/UNSW-NB15", split="train", trust_remote_code=True)
df = ds.to_pandas()
if "label" in df.columns:
df = df.drop(columns=["label"], errors="ignore")
X = df.select_dtypes(include=[np.number]).values.astype(float)
if X.shape[1] > FEATURE_DIM:
X = X[:, :FEATURE_DIM]
elif X.shape[1] < FEATURE_DIM:
pad = np.zeros((X.shape[0], FEATURE_DIM - X.shape[1]))
X = np.hstack([X, pad])
except Exception:
X = np.random.uniform(0, 1, (n_samples, FEATURE_DIM))
rng = np.random.default_rng(rs)
n = min(n_samples, len(X))
idx = rng.choice(len(X), size=n, replace=False)
return X[idx]
def _run_training(csv_url: Optional[str] = None):
global _model, _training
with _training_lock:
if _training:
return
_training = True
try:
import torch
import torch.nn as nn
from sklearn.preprocessing import StandardScaler
from huggingface_hub import HfApi, login
if csv_url:
X = _load_csv_from_url(csv_url)
if len(X) < 10:
rs = int(time.time()) % 10000
X = _fetch_unsw(SKLEARN_N_SAMPLES, rs)
else:
rs = int(time.time()) % 10000
X = _fetch_unsw(SKLEARN_N_SAMPLES, rs)
scaler = StandardScaler()
X_scaled = scaler.fit_transform(X)
class Autoencoder(nn.Module):
def __init__(self):
super().__init__()
self.enc = nn.Sequential(nn.Linear(FEATURE_DIM, 16), nn.ReLU(), nn.Linear(16, 8))
self.dec = nn.Sequential(nn.Linear(8, 16), nn.ReLU(), nn.Linear(16, FEATURE_DIM))
def forward(self, x):
return self.dec(self.enc(x))
model = Autoencoder()
opt = torch.optim.Adam(model.parameters(), lr=0.01)
x_t = torch.tensor(X_scaled, dtype=torch.float32)
for _ in range(50):
recon = model(x_t)
loss = nn.functional.mse_loss(recon, x_t)
opt.zero_grad()
loss.backward()
opt.step()
model.eval()
with torch.no_grad():
recon = model(x_t)
threshold = float(((x_t - recon) ** 2).mean().item()) * 2.0
bundle = {"model": model, "scaler": scaler, "threshold": threshold, "feature_dim": FEATURE_DIM, "dataset": "unsw_nb15"}
os.makedirs("/tmp/models", exist_ok=True)
path = "/tmp/models/model.pt"
with open(path, "wb") as f:
pickle.dump(bundle, f)
if HF_TOKEN and HF_REPO:
login(token=HF_TOKEN)
api = HfApi()
api.create_repo(repo_id=HF_REPO, repo_type="model", exist_ok=True, token=HF_TOKEN)
api.upload_file(path_or_fileobj=path, path_in_repo="model.pt", repo_id=HF_REPO, repo_type="model", token=HF_TOKEN)
print(f"Uploaded to {HF_REPO}")
_load_model()
except Exception as e:
print(f"Training failed: {e}")
finally:
with _training_lock:
_training = False
@app.on_event("startup")
def startup():
_load_model()
@app.get("/")
def root():
return {"service": "Axiom Autoencoder UNSW", "version": "0.1.0", "endpoints": ["/health", "/score", "/reload", "/train", "/train/status"]}
@app.get("/model-url")
def model_url():
"""回傳此 Space 的 HF Model 下載 URL(依 HF_REPO 環境變數)。"""
if not HF_REPO:
return {"model_url": None, "repo_id": None, "filename": "model.pt", "error": "HF_REPO not configured"}
url = f"https://huggingface.co/{HF_REPO}/resolve/main/model.pt"
return {"model_url": url, "repo_id": HF_REPO, "filename": "model.pt"}
@app.get("/health")
def health():
return {"status": "ok", "model_loaded": _model is not None}
@app.post("/reload")
def reload():
_load_model()
return {"status": "ok", "model_loaded": _model is not None}
class TrainRequest(BaseModel):
csv_url: Optional[str] = None
@app.post("/train")
def train(req: TrainRequest | None = Body(None)):
csv_url = req.csv_url if req else None
t = threading.Thread(target=_run_training, args=(csv_url,), daemon=True)
t.start()
return {"status": "training_started", "dataset": "csv" if csv_url else "UNSW-NB15"}
@app.get("/train/status")
def train_status():
return {"training": _training}
@app.post("/score", response_model=ScoreResponse)
def score(req: ScoreRequest):
anomaly_score, details = compute_score(req.features)
return ScoreResponse(anomaly_score=round(anomaly_score, 4), is_anomaly=anomaly_score >= ANOMALY_THRESHOLD, details=details)