源码精读:AuditService — 跨进程审计的 Kafka Consumer
AuditService 是 coomia-dip 元数据治理层(Metadata & Governance Layer,合并至 Control Layer)的审计引擎,负责跨进程的操作审计、决策追踪和合规性记录。它通过 gRPC 暴露 10 个 RPC 方法,覆盖事件记录(RecordEvent/BatchRecordEvents)、日志查询(QueryAuditLogs/GetEntityHistory)、决策追踪(GetDecisionTrace/GetDecisionHistory)和数据导出(ExportAuditLogs)四大能力域。本文将深入剖析 Proto 合约中 13 种审计事件类型的设计、AuditEmitter 的双通道发射策略(gRPC 优先 + JSONL 降级)、ComputeAuditLogger 的 Kafka 三级降级机制、以及 DecisionTrace 的输入快照与推理步骤记录模型。
源码精读:AuditService — 跨进程审计的 Kafka Consumer
“系列:S9 源码精读 · 第 17 篇 | 难度:高级 | 阅读时间:25 分钟
#TL;DR
AuditService 是 coomia-dip 元数据治理层(Metadata & Governance Layer,合并至 Control Layer)的审计引擎,负责跨进程的操作审计、决策追踪和合规性记录。它通过 gRPC 暴露 10 个 RPC 方法,覆盖事件记录(RecordEvent/BatchRecordEvents)、日志查询(QueryAuditLogs/GetEntityHistory)、决策追踪(GetDecisionTrace/GetDecisionHistory)和数据导出(ExportAuditLogs)四大能力域。本文将深入剖析 Proto 合约中 13 种审计事件类型的设计、AuditEmitter 的双通道发射策略(gRPC 优先 + JSONL 降级)、ComputeAuditLogger 的 Kafka 三级降级机制、以及 DecisionTrace 的输入快照与推理步骤记录模型。
#目录
- 整体架构:多源审计事件的汇聚点
- Proto 合约:13 种审计事件类型
- AuditEvent 数据模型:16 字段的合规设计
- AuditEmitter:双通道发射策略
- ComputeAuditLogger:Kafka 三级降级
- DecisionTrace:决策追踪与可解释性
- QueryAuditLogs:多维度过滤查询
- 敏感数据脱敏:_sanitize 函数
- 数据生命周期:CleanupOldEvents
- 导出与合规:ExportAuditLogs
- Key Takeaways
#1. 整体架构:多源审计事件的汇聚点
AuditService 是所有 Layer 的审计事件的汇聚点,接收来自三个方向的事件:
Reasoning & Decision Layer (推理决策) ──→ ComputeAuditLogger ──→ Kafka topic / gRPC
Agent Runtime Layer (行动执行) ──→ AuditEmitter ──→ gRPC / JSONL fallback
Data Layer (数据操作) ──→ 直接 gRPC 调用 ──→ AuditService
intelligence-Layer/src/
├── reasoning_decision_plane/security/
│ └── compute_audit.py # Reasoning & Decision Layer 审计日志器
├── agent_runtime_plane/action/
│ └── audit_emitter.py # Agent Runtime Layer 审计发射器
proto/plane_g/
└── audit_service.proto # 298 行审计合约
设计哲学:审计事件的记录必须是"fire-and-forget"——永远不能因为审计系统的故障而阻塞或失败主业务流程。
#2. Proto 合约:13 种审计事件类型
AuditService 定义了 13 种覆盖全生命周期的事件类型:
enum AuditEventType {
AUDIT_EVENT_TYPE_CREATE = 1; // 资源创建
AUDIT_EVENT_TYPE_UPDATE = 2; // 资源更新
AUDIT_EVENT_TYPE_DELETE = 3; // 资源删除
AUDIT_EVENT_TYPE_QUERY = 4; // 查询操作
AUDIT_EVENT_TYPE_REASONING = 5; // 推理执行
AUDIT_EVENT_TYPE_DECISION = 6; // 决策生成
AUDIT_EVENT_TYPE_ACTION = 7; // 行动执行
AUDIT_EVENT_TYPE_DECISION_APPROVED = 8; // 决策批准
AUDIT_EVENT_TYPE_DECISION_REJECTED = 9; // 决策拒绝
AUDIT_EVENT_TYPE_DECISION_EXECUTED = 10; // 决策执行
AUDIT_EVENT_TYPE_ACCESS = 11; // 数据访问
AUDIT_EVENT_TYPE_EXPORT = 12; // 数据导出
AUDIT_EVENT_TYPE_MASKING = 13; // 数据脱敏
}
这 13 种类型可以分为四个逻辑组:
- CRUD 组(1-4):基础资源操作
- 决策组(5-10):从推理到执行的完整决策链
- 合规组(11-13):数据访问、导出和脱敏
决策组的 6 种类型完整地记录了决策的生命周期:REASONING → DECISION → APPROVED/REJECTED → EXECUTED/ACTION。
#3. AuditEvent 数据模型:16 字段的合规设计
AuditEvent 消息包含 16 个字段,兼顾业务审计和合规性要求:
message AuditEvent {
string log_id = 1; // 事件唯一标识
AuditEventType event_type = 2;
string entity_id = 3; // 关联实体 ID
string user_id = 4; // 操作用户 ID
string world_id = 5; // World 上下文
google.protobuf.Timestamp timestamp = 6;
string trace_id = 7; // 分布式追踪 ID
AuditEventStatus status = 8;
string summary = 9; // 操作摘要
google.protobuf.Struct details = 10; // 详细信息
string decision_id = 11; // 决策 ID
string tenant_id = 12; // 租户 ID
// 合规性字段
string data_classification = 13; // 数据分类
string access_purpose = 14; // 访问目的
string client_ip = 15; // 客户端 IP
string user_agent = 16; // User Agent
}
合规性字段(13-16)是为满足 GDPR/SOX 等合规要求而设计的:
data_classification:标记访问的数据级别(如"机密"/"内部")access_purpose:记录访问目的(如"客户服务"/"数据分析"),GDPR 要求的"目的限制"原则client_ip+user_agent:源头追溯和异常检测
#4. AuditEmitter:双通道发射策略
AuditEmitter(audit_emitter.py)是 Agent Runtime Layer 的审计事件发射器,实现了优雅的双通道降级:
class AuditEmitter:
def __init__(
self,
plane_g_endpoint: str | None = None,
local_log_path: str = "logs/audit_fallback.jsonl",
) -> None:
self._local_log_path = Path(local_log_path)
self._channel: Any = None
self._stub: Any = None
if plane_g_endpoint:
self._init_grpc(plane_g_endpoint)
懒加载 Proto Stub:gRPC 初始化使用双重 try/except,因为 Proto stubs 可能还没编译:
def _init_grpc(self, endpoint: str) -> None:
try:
import grpc
self._channel = grpc.insecure_channel(endpoint)
except ImportError:
logger.warning("grpc package not available")
return
try:
try:
from reasoning_decision_plane.generated.plane_g import audit_service_pb2_grpc
except ImportError:
from agent_runtime_plane.api.generated.plane_g import audit_service_pb2_grpc
self._stub = audit_service_pb2_grpc.AuditServiceStub(self._channel)
except ImportError:
self._stub = None
双通道发射:_emit 方法实现了 gRPC → JSONL 的降级链:
async def _emit(self, event: AuditEvent) -> None:
"""NEVER raises; all errors are logged as warnings."""
try:
if self._stub is not None:
self._send_to_plane_g(event)
return
except Exception:
logger.warning("Failed to send to Metadata & Governance Layer, falling back to local log")
try:
self._write_local_log(event)
except Exception:
logger.warning("Failed to write to local log")
核心设计原则:NEVER raises。即使两个通道都失败,也只记录警告,绝不传播异常到调用方。
#5. ComputeAuditLogger:Kafka 三级降级
ComputeAuditLogger(compute_audit.py)是 Reasoning & Decision Layer 的审计日志器,实现了三级降级:Kafka → gRPC → Log Warning:
def flush(self) -> int:
events = self.drain_events()
if not events:
return 0
# 第一级:尝试 Kafka
sent = self._send_kafka(events)
if sent > 0:
return sent
# 第二级:回退到 gRPC
sent = self._send_grpc(events)
if sent > 0:
return sent
# 第三级:最后手段——仅日志
logger.warning("No external audit sink available, %d events dropped", len(events))
return 0
内存缓冲设计:事件先存入内存 buffer(self._events),通过 drain_events() 批量获取并清空:
def drain_events(self) -> list[DerivedPropertyComputeAudit]:
events = list(self._events)
self._events.clear()
return events
Kafka 发送使用 kafka-python 库,发送到 onto.audit.derived-property topic:
def _send_kafka(self, events):
self._kafka_producer = KafkaProducer(
bootstrap_servers=bootstrap,
value_serializer=lambda v: json.dumps(v, default=str).encode('utf-8'),
acks=1, retries=1, request_timeout_ms=5000,
)
for event in events:
self._kafka_producer.send('onto.audit.derived-property', value=payload)
self._kafka_producer.flush(timeout=5)
acks=1 和 retries=1 的配置体现了"尽力而为"的设计——审计日志宁可丢失也不能阻塞计算路径。
#6. DecisionTrace:决策追踪与可解释性
DecisionTrace 是 AuditService 中最复杂的数据模型,提供决策的完整可解释性:
message DecisionTrace {
string decision_id = 1;
string world_id = 2;
string commit_id = 3; // Nessie commit ID
google.protobuf.Timestamp timestamp = 4;
string outcome = 5;
repeated InputSnapshot inputs = 6; // 输入快照
repeated ReasoningStep reasoning_steps = 7; // 推理步骤
repeated FiredRule fired_rules = 8; // 触发的规则
google.protobuf.Struct metadata = 9;
}
InputSnapshot 冻结了决策时刻的实体状态:
message InputSnapshot {
string entity_id = 1;
string entity_type = 2;
google.protobuf.Struct attributes = 3;
int64 version = 4;
}
ReasoningStep 记录了每一步推理的输入/输出和置信度:
message ReasoningStep {
int32 step_number = 1;
string step_name = 2;
string engine_type = 3;
google.protobuf.Struct input = 4;
google.protobuf.Struct output = 5;
double confidence = 6;
int64 duration_ms = 7;
}
FiredRule 记录了每条规则的贡献权重:
message FiredRule {
string rule_id = 1;
string rule_name = 2;
string rule_version = 3;
double contribution_weight = 4;
google.protobuf.Struct context = 5;
}
这三层结构使得事后审计时可以完整重现:"看到了什么数据" → "经过了什么推理" → "哪些规则参与了" → "最终决策是什么"。
#7. QueryAuditLogs:多维度过滤查询
审计日志查询支持 9 个过滤维度:
message QueryAuditLogsRequest {
string entity_id = 2;
string user_id = 3;
string world_id = 4;
string decision_id = 5;
google.protobuf.Timestamp start_time = 6;
google.protobuf.Timestamp end_time = 7;
repeated AuditEventType event_types = 8;
repeated AuditEventStatus statuses = 9;
int32 page_size = 10;
string page_token = 11;
string sort_by = 12;
bool sort_asc = 13;
}
Cursor 分页:使用 page_token(而非 offset)实现游标分页,避免了大偏移量查询的性能问题。next_page_token 在响应中返回:
message QueryAuditLogsResponse {
repeated AuditEvent events = 1;
int64 total_count = 2;
string next_page_token = 3;
bool has_more = 4;
}
#8. 敏感数据脱敏:_sanitize 函数
AuditEmitter 在记录审计事件前会对敏感字段进行脱敏:
_SENSITIVE_FIELDS = frozenset({"password", "secret", "token", "api_key"})
def _sanitize(data: dict[str, Any] | None) -> dict[str, Any]:
if not data:
return {}
result: dict[str, Any] = {}
for key, value in data.items():
if key.lower() in _SENSITIVE_FIELDS:
result[key] = "***"
elif isinstance(value, dict):
result[key] = _sanitize(value)
else:
result[key] = value
return result
递归脱敏:嵌套的字典也会被递归处理。frozenset 的使用确保了 O(1) 的字段名查找性能。
_make_json_safe 函数处理 Proto Struct 不支持的类型(如 datetime):
def _make_json_safe(data: dict[str, Any] | None) -> dict[str, Any]:
for key, value in data.items():
if isinstance(value, datetime):
result[key] = value.isoformat()
elif isinstance(value, (str, int, float, bool, type(None))):
result[key] = value
else:
result[key] = str(value)
#9. 数据生命周期:CleanupOldEvents
审计日志的存储不是无限的,CleanupOldEvents 提供了基于保留天数的清理机制:
message CleanupOldEventsRequest {
int32 retention_days = 2; // 保留天数
bool dry_run = 3; // 是否仅模拟
}
dry_run 模式允许在实际删除前预览影响范围,返回将被删除的事件数量。
#10. 导出与合规:ExportAuditLogs
审计日志支持三种格式的导出:
enum ExportFormat {
EXPORT_FORMAT_JSON = 1;
EXPORT_FORMAT_CSV = 2;
EXPORT_FORMAT_PARQUET = 3;
}
Parquet 格式的支持是为大规模数据分析场景设计的——合规团队可以将审计日志导出为 Parquet 文件,在 Spark/Trino 中进行交互式分析。
#11. Key Takeaways
- Fire-and-forget 原则:审计记录永远不能阻塞或失败主业务流程,所有异常都被捕获并降级处理。
- 三级降级:Kafka → gRPC → JSONL,确保审计事件在任何基础设施故障下都不会完全丢失。
- 13 种事件类型完整覆盖了从 CRUD 到决策到合规的全生命周期。
- DecisionTrace 三层结构:InputSnapshot → ReasoningStep → FiredRule 提供了完整的决策可解释性。
- 敏感字段递归脱敏:
_sanitize函数确保审计日志本身不会泄露敏感信息。 - 游标分页:
page_token替代 offset 避免了大偏移量查询的性能问题。
#下一篇
S9-18:LineageService — 实体级+字段级血缘,我们将深入 Metadata & Governance Layer 的血缘追踪服务,了解双视角血缘图和影响分析的实现。
Tags: #coomia-dip #source-code-reading #audit-service #kafka-consumer #compliance #grpc #Layer-g