Wazuh Inventory Sync FlatBuffer 协议解析:消息结构、会话流程与 Manager 侧实现
【免费下载链接】wazuhWazuh - The Open Source Security Platform. Unified XDR and SIEM protection for endpoints and cloud workloads.项目地址: https://gitcode.com/GitHub_Trending/wa/wazuh
本文以 Wazuh 仓库中 Inventory Sync 的 FlatBuffer 协议文档 为主体,完整讲解 Agent 与 Manager 之间的在线协议(on-the-wire protocol):根消息封装、全部消息表与枚举、会话生命周期、重传机制,并结合 schema 源文件 与 Manager 侧 C++ 实现(Facade、AgentSession、GapSet、ResponseDispatcher)说明每条消息在接收端如何被解析、校验和落地。读完后你将能够独立读懂(或生成/校验)Inventory Sync 流量,并为协议工具开发、集成测试提供可验证的字段级依据。
协议总览:schema 文件是唯一事实源
Inventory Sync 是 Manager 侧的 Agent 状态同步服务:它从 Router 主题inventory-states接收 Agent 发来的 FlatBuffer 消息,把会话载荷暂存到 RocksDB,再把状态文档索引到wazuh-states-*索引族,最后向 Agent 确认完成(详见 模块总览)。整个在线协议由一份 FlatBuffer schema 定义:
- schema 文件:inventorySync.fbs,命名空间为
Wazuh.SyncSchema; - 构建时由
flatc -c生成 C 头文件inventorySync_generated.h,规则见 schemas/CMakeLists.txt; - Manager 侧解析入口在 inventorySyncFacade.hpp,QA 集成测试与
inventory_sync_testtool也复用同一份 live schema(见 test-tools 文档)。
从源码结构看,该模块对协议的消费分为两层:InventorySyncFacadeImpl::run()按MessageType分发(inventorySyncFacade.hpp#L125-L431),每个会话的状态机由AgentSessionImpl承担(agentSession.hpp)。此外,Facade 在工作线程中会先用flatbuffers::Verifier+VerifyMessageBuffer做缓冲区完整性校验,校验失败的消息直接丢弃并记录日志(inventorySyncFacade.hpp#L672-L683)。
根消息:所有协议报文都封装在 Message 中
所有协议消息都包装在一个Message表内,content字段是一个 union,用类型判别位区分实际载荷:
table Message { content: MessageType; } union MessageType { DataValue, DataClean, ChecksumModule, Start, StartAck, End, EndAck, ReqRet, DataContext, DataBatch } root_type Message;这与 inventorySync.fbs 中的定义逐字一致。10 种消息可以按方向分为三类:
- Agent → Manager:
Start、DataValue、DataBatch、DataContext、DataClean、ChecksumModule、End; - Manager → Agent:
StartAck、EndAck、ReqRet; - 双向语义由
Status枚举表达。
Manager 侧的分发逻辑印证了这个划分:run()中按content_type()依次匹配DataValue、DataClean、DataContext、DataBatch、Start、ChecksumModule、End,未知类型抛出InventorySyncException("Invalid message type")(inventorySyncFacade.hpp#L127-L431)。而StartAck/EndAck/ReqRet只出现在 Manager 的发送路径中,由 responseDispatcher.hpp 构造并经由 Agent 回传队列发回。
四个核心枚举:Mode、Operation、Status、Option
Mode:会话工作模式
enum Mode: byte { ModuleFull, ModuleDelta, ModuleCheck, MetadataDelta, MetadataCheck, GroupDelta, GroupCheck }七种模式对应模块级全量同步、增量同步、完整性校验,以及元数据/分组的增量更新与灾难恢复式核对。AgentSessionImpl的构造函数中可以看到模式与size的强关联:MetadataDelta、MetadataCheck、GroupDelta、GroupCheck、ModuleCheck这五种模式不携带数据消息,允许size为 0;其余模式下size == 0会直接回StartAck(Error)并抛异常拒绝会话(agentSession.hpp#L174-L183)。
Operation:文档级操作
enum Operation: byte { Upsert, Delete }DataValue用它声明该条载荷对目标索引文档是插入/更新还是删除。
Status:应答状态
enum Status: byte { Ok, Error, Offline, ChecksumMismatch, Processing }Manager 侧的实际使用方式(源码印证):
Ok:会话受理成功,或索引/更新操作完成后在EndAck中返回(inventorySyncFacade.hpp#L795-L796);Error:Agent 被锁、或Start.size非法等参数级错误;Offline:索引器不可用、会话数达到上限、DataValue 配额耗尽时,Manager 用StartAck(Offline, session=-1)拒绝新会话,Agent 稍后重试(inventorySyncFacade.hpp#L296-L344);Processing:收到End且序号齐全时,Manager 先回EndAck(Processing)表示“会话已提交处理”,索引完成后另发最终EndAck(Ok)(agentSession.hpp#L504-L524);ChecksumMismatch:ModuleCheck模式校验不一致时返回。
Option:漏洞扫描联动
enum Option: byte { Sync, VDFirst, VDSync }VDFirst/VDSync表示本次会话持久化完成后需要触发 Vulnerability Scanner 编排。Facade 会识别名为syscollector_vd的模块并为其设置单独的加锁策略(inventorySyncFacade.hpp#L282-L295);模块总览也确认漏洞扫描是在库存数据落库之后用同一会话上下文触发的(README)。
Start 消息:打开会话并携带 Manager 侧上下文
table Start { module: string; mode: Mode; size: ulong; index: [string]; option: Option; architecture: string; hostname: string; osname: string; osplatform: string; ostype: string; osversion: string; agentversion: string; agentname: string; agentid: string; groups: [string]; global_version: ulong; cluster_name: string; cluster_node: string; }Start 打开会话,并把后续索引与下游处理所需的 Manager 侧上下文一次性带上。各重要字段及其在实现中的语义:
module:当前为syscollector、fim或sca(漏洞扫描会话为syscollector_vd);mode:full、delta、integrity-check、metadata 或 group 模式;size:本次会话预期收到的、受序号跟踪的消息数。Manager 用它初始化GapSet,并作为 DataValue 配额的预留量(inventorySyncFacade.hpp#L322-L344);index:本次会话的目标索引列表。Manager 会过滤只保留属于 Agent 作用域的wazuh-states-*家族索引,不匹配的索引会被告警并忽略,白名单函数isAgentScopedStateIndex接受wazuh-states-inventory-*、wazuh-states-vulnerabilities、wazuh-states-fim-*、wazuh-states-sca、wazuh-states-sca-*(agentSession.hpp#L44-L48);option:漏洞扫描联动行为;global_version:metadata / group 更新流程使用的版本;cluster_name、cluster_node:Manager 侧会话上下文传播的集群元数据;agentid:Manager 侧会自动做零填充——不足 3 位时左侧补 0(agentSession.hpp#L128-L131),与响应帧头中 3 位 Agent ID 的约定一致;- OS 系列字段(
osname、osversion等)与agentname/agentversion/groups:在MetadataDelta/GroupDelta会话中被用来构造updateByQuery文档更新(inventorySyncFacade.hpp#L807-L856)。
会话 ID 不是 Agent 指定的,而是 Manager 收到合法 Start 后用std::random_device生成 64 位随机值,并通过StartAck回传(inventorySyncFacade.hpp#L347-L377)。同一 Agent + 模块组合的陈旧会话会在开新会话前被清理,以覆盖 Agent 重启 / modulesd 重启场景(inventorySyncFacade.hpp#L306-L308)。
数据消息:DataValue、DataBatch、DataContext、DataClean、ChecksumModule
DataValue:主索引载荷
table DataValue { seq: ulong; session: ulong; operation: Operation; id: string; index: string; version: ulong; data: [byte]; }这是协议中主要的可索引载荷类型:
seq:序号,被GapSet用于缺口跟踪与重传请求;session:StartAck中分配的会话 ID;operation:Upsert或Delete;id:逻辑文档 ID 片段;index:目标状态索引;version:可选的文档版本,会传播到索引器;data:JSON 载荷字节。
Manager 侧AgentSessionImpl::handleData()的处理要点:先拒绝seq >= size的越界消息(防止留下孤儿键),然后把整条 FlatBuffer 原始字节以{session}_{seq}为 key 写入 RocksDB,再调GapSet::observe(seq)标记该序号已收到(agentSession.hpp#L258-L279)。也就是说 Manager 先按序暂存原始报文,会话结束后才由索引器线程统一回放,这是“会话级暂存、结束时索引”的实现基础。
DataBatch:批量载荷
table DataBatch { values: [DataValue]; }DataBatch允许在一条消息里携带多个DataValue。DataBatch属于 live 协议,生成或校验 Inventory Sync 流量的工具应当支持它。Manager 侧的处理是在内部展开批次:Facade 遍历values,把每个DataValue重新序列化为独立的Message{DataValue}FlatBuffer,再逐条走handleData()——因为 RocksDB 下游消费者期望“一个 key 对应一条独立 DataValue 消息”(inventorySyncFacade.hpp#L199-L269)。因此每个批次内项都作为独立的会话记录存储并参与GapSet计数。
DataContext:会话级上下文数据
table DataContext { seq: ulong; session: ulong; id: string; index: string; data: [byte]; }DataContext的当前行为在实现中有明确印证:
- 以
_context后缀存 RocksDB,key 为{session}_{seq}_context,与DataValue区分(agentSession.hpp#L395-L397); - 同样参与
GapSet跟踪,纳入重传与“会话结束完整性”判定; - 不直接索引——源码注释明确其用途是留给下游(如漏洞扫描)读取的会话上下文数据(agentSession.hpp#L355-L364)。
DataClean:索引清理请求
table DataClean { seq: ulong; session: ulong; index: string; }DataClean在会话结束时请求对给定 Agent 与索引执行deleteByQuery。Manager 侧把index收集到会话上下文的dataCleanIndices集合中(天然去重,兼容重传),并在索引阶段执行清理;缺少index字段会记录错误日志(agentSession.hpp#L429-L495)。
ChecksumModule:完整性校验
table ChecksumModule { session: ulong; index: string; checksum: string; }ChecksumModule用于ModuleCheck模式:Agent 上报自身计算的校验和,Manager 则从已索引文档反算。Manager 的实现是分页拉取该 Agent 在目标索引中的checksum.hash.sha1文档(每批 1000 条,search_after翻页),把所有 SHA1 串接后整体再取 SHA1,与 Agent 值比较;不一致时按固定间隔最多重试 5 次,以容忍近期 delta 尚未完成索引的情况(inventorySyncFacade.hpp#L434-L503、inventorySyncFacade.hpp#L968-L1000)。迟到的ChecksumModule(会话已提交索引后到达)会被直接忽略,避免与索引线程竞争(agentSession.hpp#L314-L323)。
会话收尾与应答:End、EndAck、StartAck
End 与 EndAck
table End { session: ulong; }End关闭会话的上传侧。Manager 只有在收到End且所有预期受序号跟踪的消息都已被GapSet记账后,才完成会话。AgentSessionImpl::handleEnd()的行为(agentSession.hpp#L504-L530):
- 序号齐全:把会话推入索引器队列,立即回
EndAck(Processing); - 序号有缺口:不回 EndAck,而是发出
ReqRet请求缺失序号段重传; - 重复
End:直接回EndAck(Processing),不重复提交。
值得注意的两阶段确认设计:会话在收到End时只承诺“已收齐”,EndAck(Ok)在索引/更新操作真正完成后才由 notify 回调发出(inventorySyncFacade.hpp#L787-L803)。
StartAck
table StartAck { status: Status; session: ulong; }Start 成功时,Manager 在此返回分配的会话 ID。失败时(Agent 被锁、索引器离线、会话数超限、配额不足、size 非法)同样用 StartAck 携带对应Status,且session置为-1,表示没有会话被创建(inventorySyncFacade.hpp#L289-L344)。
table EndAck { status: Status; session: ulong; }EndAck承载会话的最终结果:Ok(完成)、Processing(已受理、异步处理中)、Error(处理失败)、ChecksumMismatch(校验不一致)。
重传机制:Pair 与 ReqRet + GapSet
table Pair { begin: ulong; end: ulong; } table ReqRet { seq: [Pair]; session: ulong; }ReqRet用于请求 Manager 检测到的缺失序号区间。发送路径在 responseDispatcher.hpp#L152-L175:把缺口区间列表逐个CreatePair后组装进ReqRet并经Message封装发出。
缺口计算由 gapSet.hpp 中的GapSet完成,实现上有几个值得注意的点:
- 用有序区间集合表示“已收到”的闭区间
[start, end],observe(seq)为 O(log n),并把相邻/重叠区间自动合并; seq >= size(Start 声明的总量)直接抛std::out_of_range,与AgentSessionImpl的越界拒绝逻辑配合,保证孤儿数据不落盘;- 区间合并覆盖
[0, size-1]时置m_allObserved,empty()检查变为 O(1); ranges()反向导出未覆盖的缺口区间列表(O(k)),正是ReqRet中Pair序列的数据来源(gapSet.hpp#L206-L240);- 另维护
lastUpdate()时间戳,Facade 用它判断会话是否超时(AgentSession::isAlive,agentSession.hpp#L537-L540)。
完整的重传往返可以在 QA 集成测试中验证,例如 reqret_end_flow 与 simple_reqret_test,预期结果在 expected_data/reqret_end_flow.json。
协议帧的发送格式与准入控制(源码补充)
Manager → Agent 的应答不是裸 FlatBuffer:responseDispatcher.hpp 会在 FlatBuffer 字节前拼上文本帧头(msg_to_agent) [] N!s <agentId> <payloadLen> <moduleName>_sync,再经queue/sockets/ar队列(ARQUEUE)发回。理解这一点对抓取/模拟 Manager 响应流量的工具很重要。
Manager 在受理 Start 前还有多层准入控制,都会以StartAck状态体现:
- Agent 锁:同一 Agent 的 metadata/group 更新进行中新会话被拒绝(
Error); - 索引器可用性:不可用时返回
Offline; - 会话上限:
maxSessions达上限返回Offline(默认 1000,可配置,inventorySyncFacade.hpp#L634-L639); - DataValue 配额:以
Start.size预留全局配额(CAS 循环扣减),不足则Offline,会话结束时归还(inventorySyncFacade.hpp#L322-L344)。
协议工具链:如何生成与验证流量
围绕这份 schema,仓库提供了完整的生成与验证链路:
- C/C++ 生成:构建系统调用
flatc -c生成 inventorySync_generated.h; - Python 侧(QA 集成测试):generate_flatbuffers.py 用
flatc生成 Python 类,然后用真实 FlatBuffer 协议模拟 Agent-Manager 交互。运行方式(需要可访问的 Manager 与flatc):
# 在 qa 目录下安装依赖并生成协议类 pip install -r requirements.txt python3 generate_flatbuffers.py # 对本地 Manager 跑全部协议集成测试 python run_tests.py --manager 127.0.0.1 # 只跑指定用例 python run_tests.py --manager 127.0.0.1 --test basic_flow用例覆盖 start/data/end 基础流、无数据会话、ReqRet 重传、DataClean、仅 DataContext 会话、checksum 匹配/不匹配、metadata delta、group delta 等(QA 目录);
- 端到端测试工具:
inventory_sync_testtool位于 testtool/,用单个 JSON 文件描述 Start、DataValue、DataContext 及VDFirst/VDSync选项,模拟完整会话以验证 Inventory Sync、Indexer 与 Vulnerability Scanner 三者的集成。
按 test-tools 文档 的建议:修改会话逻辑/队列/解析时用单元测试(tests/unit/);修改协议语义、ACK、重传、metadata/group 对账或 checksum 逻辑时用 QA 集成套件;验证真实索引或漏洞扫描触发的会话时用 testtool。协议变更时,qa/与testtool/的 schema 消费方需要同步更新。
实用要点
Start.size对MetadataDelta、MetadataCheck、GroupDelta、GroupCheck、ModuleCheck会话可以为 0;其余模式下 size 为 0 会被StartAck(Error)拒绝;DataContext属于 live 协议的一部分,即使它不会被回放进索引器——任何生成/校验流量的工具都必须按GapSet的口径把它计入序号跟踪;DataBatch同样是 live 协议,工具应当支持;Manager 侧会把它展开为逐条DataValue处理;- 越界
seq(≥ size)的DataValue/DataContext/DataClean都会被拒绝,不会产生孤儿存储; EndAck(Processing)不等于最终成功,最终Ok由索引完成回调补发;- 所有
index字段都会经过wazuh-states-*家族白名单过滤,Agent 侧不应假设任意索引名会被接受。
参考路径
| 主题 | 路径 |
|---|---|
| 协议文档(本文主体) | docs/ref/modules/inventory-sync/flatbuffers.md |
| schema 源文件 | src/shared_modules/utils/flatbuffers/schemas/inventorySync.fbs |
| 模块总览 | docs/ref/modules/inventory-sync/README.md |
| Facade / 消息分发 | src/wazuh_modules/inventory_sync/src/inventorySyncFacade.hpp |
| 会话状态机 | src/wazuh_modules/inventory_sync/src/agentSession.hpp |
| 缺口跟踪 | src/wazuh_modules/inventory_sync/src/gapSet.hpp |
| 应答发送 | src/wazuh_modules/inventory_sync/src/responseDispatcher.hpp |
| QA 集成测试 | src/wazuh_modules/inventory_sync/qa/README.md |
| 端到端测试工具 | src/wazuh_modules/inventory_sync/testtool/README.md |
【免费下载链接】wazuhWazuh - The Open Source Security Platform. Unified XDR and SIEM protection for endpoints and cloud workloads.项目地址: https://gitcode.com/GitHub_Trending/wa/wazuh
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考