cypher-v12-finalized / modules /audit_logger.py
jescy525's picture
Upload folder using huggingface_hub
076a67c verified
Raw
History Blame Contribute Delete
7.62 kB
"""CYPHER V12 M14 — Audit logger (history complet inference).
Log every /chat request to JSONL: timestamp, prompt, category, banks_used,
response, latency, scores. Used by post-mortem analysis, regression detection,
QS continuous eval (M8).
Append-only, atomic writes, rotation on size.
"""
from __future__ import annotations
import json
import logging
import os
import threading
import time
from pathlib import Path
from typing import Any
logger = logging.getLogger(__name__)
class AuditLogger:
"""Thread-safe JSONL audit logger with size rotation."""
def __init__(
self,
log_path: str = "/workspace/CYPHER_V12/logs/audit.jsonl",
max_size_mb: int = 100,
rotation_keep: int = 5,
):
self.log_path = Path(log_path)
self.log_path.parent.mkdir(parents=True, exist_ok=True)
self.max_size_bytes = max_size_mb * 1024 * 1024
self.rotation_keep = rotation_keep
self._lock = threading.Lock()
def _rotate_if_needed(self) -> None:
if not self.log_path.exists():
return
try:
if self.log_path.stat().st_size < self.max_size_bytes:
return
except OSError:
return
# Rotate: audit.jsonl → audit.jsonl.1, audit.jsonl.1 → audit.jsonl.2, ...
for i in range(self.rotation_keep, 0, -1):
src = self.log_path.with_suffix(f".jsonl.{i-1}") if i > 1 else self.log_path
dst = self.log_path.with_suffix(f".jsonl.{i}")
try:
if src.exists():
if dst.exists():
dst.unlink()
src.rename(dst)
except OSError as e:
logger.warning(f"rotate fail at {i}: {e}")
def log_chat(
self,
prompt: str,
response: str,
category: str,
banks_used: list[str] | None = None,
latency_ms: int = 0,
generation_ms: int = 0,
max_tokens: int = 200,
temperature: float = 0.35,
extra: dict | None = None,
) -> bool:
record = {
"ts": int(time.time()),
"type": "chat",
"prompt": prompt[:1000],
"prompt_len": len(prompt),
"response": response[:2000],
"response_len": len(response),
"category": category,
"banks_used": banks_used or [],
"latency_ms": latency_ms,
"generation_ms": generation_ms,
"max_tokens": max_tokens,
"temperature": temperature,
}
if extra:
for k, v in extra.items():
if isinstance(v, (str, int, float, bool, list, dict)) and k not in record:
record[k] = v
return self._append(record)
def log_event(self, event_type: str, **kwargs: Any) -> bool:
record = {"ts": int(time.time()), "type": event_type}
for k, v in kwargs.items():
if isinstance(v, (str, int, float, bool, list, dict)):
record[k] = v
return self._append(record)
def _append(self, record: dict) -> bool:
line = json.dumps(record, ensure_ascii=False) + "\n"
with self._lock:
try:
self._rotate_if_needed()
with self.log_path.open("a", encoding="utf-8") as f:
f.write(line)
return True
except OSError as e:
logger.error(f"audit append failed: {e}")
return False
def tail(self, n: int = 100) -> list[dict]:
if not self.log_path.exists():
return []
try:
lines = self.log_path.read_text(encoding="utf-8").splitlines()
except OSError:
return []
out: list[dict] = []
for line in lines[-n:]:
line = line.strip()
if not line:
continue
try:
out.append(json.loads(line))
except json.JSONDecodeError:
continue
return out
def filter(
self,
category: str | None = None,
since_ts: int | None = None,
until_ts: int | None = None,
event_type: str | None = None,
limit: int = 1000,
) -> list[dict]:
records = self.tail(limit * 5)
out: list[dict] = []
for r in records:
if category and r.get("category") != category:
continue
if event_type and r.get("type") != event_type:
continue
ts = r.get("ts", 0)
if since_ts is not None and ts < since_ts:
continue
if until_ts is not None and ts > until_ts:
continue
out.append(r)
if len(out) >= limit:
break
return out
def stats(self) -> dict:
if not self.log_path.exists():
return {"total": 0, "size_bytes": 0, "by_category": {}, "by_type": {}}
records = self.tail(100000)
by_cat: dict[str, int] = {}
by_type: dict[str, int] = {}
latencies: list[int] = []
for r in records:
by_cat[r.get("category", "?")] = by_cat.get(r.get("category", "?"), 0) + 1
by_type[r.get("type", "?")] = by_type.get(r.get("type", "?"), 0) + 1
if r.get("type") == "chat":
latencies.append(r.get("latency_ms", 0))
return {
"total": len(records),
"size_bytes": self.log_path.stat().st_size,
"by_category": by_cat,
"by_type": by_type,
"avg_latency_ms": sum(latencies) / max(1, len(latencies)),
"p50_latency_ms": sorted(latencies)[len(latencies) // 2] if latencies else 0,
"p95_latency_ms": sorted(latencies)[int(len(latencies) * 0.95)] if latencies else 0,
}
_singleton_lock = threading.Lock()
_singleton: AuditLogger | None = None
def get_audit_logger(log_path: str | None = None) -> AuditLogger:
global _singleton
if _singleton is None:
with _singleton_lock:
if _singleton is None:
_singleton = AuditLogger(log_path=log_path or "/workspace/CYPHER_V12/logs/audit.jsonl")
return _singleton
__all__ = ["AuditLogger", "get_audit_logger"]
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO)
print("=== M14 audit_logger SMOKE ===")
test_path = "/tmp/smoke_audit.jsonl"
if Path(test_path).exists():
Path(test_path).unlink()
log = AuditLogger(log_path=test_path, max_size_mb=1, rotation_keep=3)
log.log_chat("What is SQL injection?", "SQL injection is...", "CYBERSEC",
banks_used=["mitre"], latency_ms=450, generation_ms=420)
log.log_chat("Bonjour", "Bonjour, comment puis-je t'aider?", "CONV",
latency_ms=120, generation_ms=110)
log.log_chat("What is Order Block?", "An OB is...", "TRADING",
banks_used=["predator", "patterns"], latency_ms=380, generation_ms=350)
log.log_event("startup", model_path="/workspace/CYPHER_V12/ckpts/v12.pt", version="V12.2")
tail = log.tail(5)
print(f"Tail count: {len(tail)}")
for r in tail:
print(f" type={r['type']} cat={r.get('category', '?')} latency={r.get('latency_ms', '-')}")
stats = log.stats()
print(f"Stats: by_category={stats['by_category']} by_type={stats['by_type']}")
print(f" avg_latency_ms={stats['avg_latency_ms']:.1f}")
filtered = log.filter(category="CYBERSEC")
print(f"Filter CYBERSEC: {len(filtered)} records")
print("=== SMOKE PASS ===")