
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.txtManager 侧解析入口在 inventorySyncFacade.hppQA 集成测试与inventory_sync_testtool也复用同一份 live schema见 test-tools 文档。从源码结构看该模块对协议的消费分为两层InventorySyncFacadeImpl::run()按MessageType分发inventorySyncFacade.hpp#L125-L431每个会话的状态机由AgentSessionImpl承担agentSession.hpp。此外Facade 在工作线程中会先用flatbuffers::VerifierVerifyMessageBuffer做缓冲区完整性校验校验失败的消息直接丢弃并记录日志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 → ManagerStart、DataValue、DataBatch、DataContext、DataClean、ChecksumModule、EndManager → AgentStartAck、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、OptionMode会话工作模式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-L796ErrorAgent 被锁、或Start.size非法等参数级错误Offline索引器不可用、会话数达到上限、DataValue 配额耗尽时Manager 用StartAck(Offline, session-1)拒绝新会话Agent 稍后重试inventorySyncFacade.hpp#L296-L344Processing收到End且序号齐全时Manager 先回EndAck(Processing)表示“会话已提交处理”索引完成后另发最终EndAck(Ok)agentSession.hpp#L504-L524ChecksumMismatchModuleCheck模式校验不一致时返回。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_vdmodefull、delta、integrity-check、metadata 或 group 模式size本次会话预期收到的、受序号跟踪的消息数。Manager 用它初始化GapSet并作为 DataValue 配额的预留量inventorySyncFacade.hpp#L322-L344index本次会话的目标索引列表。Manager 会过滤只保留属于 Agent 作用域的wazuh-states-*家族索引不匹配的索引会被告警并忽略白名单函数isAgentScopedStateIndex接受wazuh-states-inventory-*、wazuh-states-vulnerabilities、wazuh-states-fim-*、wazuh-states-sca、wazuh-states-sca-*agentSession.hpp#L44-L48option漏洞扫描联动行为global_versionmetadata / group 更新流程使用的版本cluster_name、cluster_nodeManager 侧会话上下文传播的集群元数据agentidManager 侧会自动做零填充——不足 3 位时左侧补 0agentSession.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、ChecksumModuleDataValue主索引载荷table DataValue { seq: ulong; session: ulong; operation: Operation; id: string; index: string; version: ulong; data: [byte]; }这是协议中主要的可索引载荷类型seq序号被GapSet用于缺口跟踪与重传请求sessionStartAck中分配的会话 IDoperationUpsert或Deleteid逻辑文档 ID 片段index目标状态索引version可选的文档版本会传播到索引器dataJSON 载荷字节。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后缀存 RocksDBkey 为{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、StartAckEnd 与 EndAcktable 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。StartAcktable 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 GapSettable 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 sizeStart 声明的总量直接抛std::out_of_range与AgentSessionImpl的越界拒绝逻辑配合保证孤儿数据不落盘区间合并覆盖[0, size-1]时置m_allObservedempty()检查变为 O(1)ranges()反向导出未覆盖的缺口区间列表O(k)正是ReqRet中Pair序列的数据来源gapSet.hpp#L206-L240另维护lastUpdate()时间戳Facade 用它判断会话是否超时AgentSession::isAliveagentSession.hpp#L537-L540。完整的重传往返可以在 QA 集成测试中验证例如 reqret_end_flow 与 simple_reqret_test预期结果在 expected_data/reqret_end_flow.json。协议帧的发送格式与准入控制源码补充Manager → Agent 的应答不是裸 FlatBufferresponseDispatcher.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-L639DataValue 配额以Start.size预留全局配额CAS 循环扣减不足则Offline会话结束时归还inventorySyncFacade.hpp#L322-L344。协议工具链如何生成与验证流量围绕这份 schema仓库提供了完整的生成与验证链路C/C 生成构建系统调用flatc -c生成 inventorySync_generated.hPython 侧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.mdschema 源文件src/shared_modules/utils/flatbuffers/schemas/inventorySync.fbs模块总览docs/ref/modules/inventory-sync/README.mdFacade / 消息分发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.hppQA 集成测试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),仅供参考