返回博客

源码精读: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 的输入快照与推理步骤记录模型。

Coomia发布于 2025年12月13日11 分钟阅读
分享本文Twitter / X

源码精读: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 的输入快照与推理步骤记录模型。

#目录

  1. 整体架构:多源审计事件的汇聚点
  2. Proto 合约:13 种审计事件类型
  3. AuditEvent 数据模型:16 字段的合规设计
  4. AuditEmitter:双通道发射策略
  5. ComputeAuditLogger:Kafka 三级降级
  6. DecisionTrace:决策追踪与可解释性
  7. QueryAuditLogs:多维度过滤查询
  8. 敏感数据脱敏:_sanitize 函数
  9. 数据生命周期:CleanupOldEvents
  10. 导出与合规:ExportAuditLogs
  11. Key Takeaways

#1. 整体架构:多源审计事件的汇聚点

AuditService 是所有 Layer 的审计事件的汇聚点,接收来自三个方向的事件:

Code
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 种覆盖全生命周期的事件类型:

PROTOBUF
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 个字段,兼顾业务审计和合规性要求:

PROTOBUF
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:双通道发射策略

AuditEmitteraudit_emitter.py)是 Agent Runtime Layer 的审计事件发射器,实现了优雅的双通道降级:

Python
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 可能还没编译:

Python
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 的降级链:

Python
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 三级降级

ComputeAuditLoggercompute_audit.py)是 Reasoning & Decision Layer 的审计日志器,实现了三级降级:Kafka → gRPC → Log Warning:

Python
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() 批量获取并清空:

Python
def drain_events(self) -> list[DerivedPropertyComputeAudit]:
    events = list(self._events)
    self._events.clear()
    return events

Kafka 发送使用 kafka-python 库,发送到 onto.audit.derived-property topic:

Python
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=1retries=1 的配置体现了"尽力而为"的设计——审计日志宁可丢失也不能阻塞计算路径。

#6. DecisionTrace:决策追踪与可解释性

DecisionTrace 是 AuditService 中最复杂的数据模型,提供决策的完整可解释性:

PROTOBUF
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 冻结了决策时刻的实体状态:

PROTOBUF
message InputSnapshot {
  string entity_id = 1;
  string entity_type = 2;
  google.protobuf.Struct attributes = 3;
  int64 version = 4;
}

ReasoningStep 记录了每一步推理的输入/输出和置信度:

PROTOBUF
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 记录了每条规则的贡献权重:

PROTOBUF
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 个过滤维度:

PROTOBUF
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 在响应中返回:

PROTOBUF
message QueryAuditLogsResponse {
  repeated AuditEvent events = 1;
  int64 total_count = 2;
  string next_page_token = 3;
  bool has_more = 4;
}

#8. 敏感数据脱敏:_sanitize 函数

AuditEmitter 在记录审计事件前会对敏感字段进行脱敏:

Python
_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):

Python
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 提供了基于保留天数的清理机制:

PROTOBUF
message CleanupOldEventsRequest {
  int32 retention_days = 2;               // 保留天数
  bool dry_run = 3;                       // 是否仅模拟
}

dry_run 模式允许在实际删除前预览影响范围,返回将被删除的事件数量。

#10. 导出与合规:ExportAuditLogs

审计日志支持三种格式的导出:

PROTOBUF
enum ExportFormat {
  EXPORT_FORMAT_JSON = 1;
  EXPORT_FORMAT_CSV = 2;
  EXPORT_FORMAT_PARQUET = 3;
}

Parquet 格式的支持是为大规模数据分析场景设计的——合规团队可以将审计日志导出为 Parquet 文件,在 Spark/Trino 中进行交互式分析。

#11. Key Takeaways

  1. Fire-and-forget 原则:审计记录永远不能阻塞或失败主业务流程,所有异常都被捕获并降级处理。
  2. 三级降级:Kafka → gRPC → JSONL,确保审计事件在任何基础设施故障下都不会完全丢失。
  3. 13 种事件类型完整覆盖了从 CRUD 到决策到合规的全生命周期。
  4. DecisionTrace 三层结构:InputSnapshot → ReasoningStep → FiredRule 提供了完整的决策可解释性。
  5. 敏感字段递归脱敏_sanitize 函数确保审计日志本身不会泄露敏感信息。
  6. 游标分页page_token 替代 offset 避免了大偏移量查询的性能问题。

#下一篇

S9-18:LineageService — 实体级+字段级血缘,我们将深入 Metadata & Governance Layer 的血缘追踪服务,了解双视角血缘图和影响分析的实现。

Tags: #coomia-dip #source-code-reading #audit-service #kafka-consumer #compliance #grpc #Layer-g