aps-agent/server/agent_core/audit_mirror.py

105 lines
4.3 KiB
Python
Raw Normal View History

# ============================================================
# 审计事件独立介质镜像 v1(moduleId: core-audit-mirror, 可重生 ✅)
# plan.md §3.6 / 矩阵「审计 append-only 且可检测篡改」:
# APS_AUDIT_MIRROR=1 时,write_audit 将每条事件同步镜像到独立
# append-only JSONL(按租户/项目隔离);业务进程改 world 无法抹掉
# 独立介质证据。镜像尽力而为:写入失败只记日志,不影响业务写。
# ============================================================
from __future__ import annotations # 前向类型引用
import json # JSONL 序列化
import logging # 失败日志
import os # 环境变量
import threading # 并发写锁
from datetime import datetime # 时间戳
from pathlib import Path # 路径处理
from typing import Any # 类型标注
logger = logging.getLogger(__name__)
def mirror_enabled() -> bool:
"""镜像开关:APS_AUDIT_MIRROR=1(或 true/yes)时启用。"""
raw = (os.environ.get("APS_AUDIT_MIRROR") or "").strip().lower()
return raw in ("1", "true", "yes", "on")
def _default_mirror_dir() -> str:
env = os.environ.get("APS_AUDIT_MIRROR_DIR")
if env:
return env
try:
from server.aps_home import path_under_data
return str(path_under_data("audit-mirror"))
except Exception:
return os.path.join("server", "data", "audit-mirror")
def _safe_scope(value: str) -> str:
cleaned = "".join(c for c in (value or "default") if c.isalnum() or c in "-_")
return cleaned[:64] or "default"
class AuditMirror:
"""独立 append-only 审计事件镜像(按租户/项目隔离,JSONL)。"""
def __init__(self, tenant_uuid: str = "platform", world_key: str = "default",
mirror_dir: str | None = None) -> None:
self.tenant_uuid = _safe_scope(tenant_uuid)
self.world_key = _safe_scope(world_key)
base = Path(mirror_dir or _default_mirror_dir())
self.path = base / self.tenant_uuid / f"{self.world_key}.jsonl"
self._lock = threading.Lock()
def _ensure_dir(self) -> None:
self.path.parent.mkdir(parents=True, exist_ok=True)
def append_event(self, event: dict[str, Any]) -> dict[str, Any] | None:
"""追加一条审计事件镜像(原子:单行写入 + flush + fsync)。"""
record = {
"event": event,
"mirroredAt": _now_iso(),
}
with self._lock:
try:
self._ensure_dir()
# 跨进程文件锁防并发 append 行交错(矩阵 116 行)
from server.agent_core.audit_filelock import locked_append
with locked_append(self.path) as f:
f.write(json.dumps(record, ensure_ascii=False, sort_keys=True) + "\n")
return record
except OSError as exc:
logger.warning("audit mirror append failed (%s): %s", self.path, exc)
return None
def read_events(self) -> list[dict[str, Any]]:
"""读取镜像中的全部事件(损坏行跳过并记日志)。"""
if not self.path.exists():
return []
with self._lock:
try:
text = self.path.read_text(encoding="utf-8")
except OSError as exc:
logger.warning("audit mirror read failed (%s): %s", self.path, exc)
return []
events: list[dict[str, Any]] = []
for line in text.splitlines():
line = line.strip()
if not line:
continue
try:
record = json.loads(line)
if isinstance(record, dict) and isinstance(record.get("event"), dict):
events.append(record["event"])
except json.JSONDecodeError:
logger.warning("audit mirror corrupt line in %s", self.path)
return events
def count(self) -> int:
return len(self.read_events())
def _now_iso() -> str:
from server.timeutil import fmt_dt
return fmt_dt(datetime.now())