Wazuh Agent Sync Protocol 会话生命周期:数据同步的阶段划分、消息类型与状态机全解
【免费下载链接】wazuhWazuh - The Open Source Security Platform. Unified XDR and SIEM protection for endpoints and cloud workloads.项目地址: https://gitcode.com/GitHub_Trending/wa/wazuh
Wazuh 的 Agent Sync Protocol(src/shared_modules/sync_protocol)是 Agent 内部模块(FIM、SCA、Inventory/Syscollector 等)向 Manager 可靠同步数据的共享组件。本文基于 Protocol Lifecycle 文档 完整解析一次同步会话的四个阶段、全部 8 种 FlatBuffer 消息类型、四种特殊同步模式,以及状态机、超时重试与错误处理机制,并结合仓库中的实际 FlatBuffer 模式文件与协议实现头文件,说明每个设计决策背后的源码依据。
一、协议总览与同步阶段划分
Agent Sync Protocol 采用基于会话(session-based)的同步机制,通过明确定义的消息类型和状态转移保证 Agent 与 Manager 之间数据的一致性。一次完整的同步会话经历以下阶段:
- Idle 阶段:无活跃同步,协议实例处于空闲状态;
- Session Establishment(会话建立):Agent 发送
Start消息,Manager 回复StartAck并分配会话 ID; - Data Transfer(数据传输):逐条发送差异数据(differences),Manager 可请求重传丢失的序号区间;
- Session Completion(会话完成):Agent 发送
End,Manager 以EndAck确认会话成功或失败。
在源码层面,阶段由 agent_sync_protocol.hpp 中的SyncPhase枚举表达,与文档的阶段划分一一对应:
/// @brief Defines the possible phases of a synchronization process. enum class SyncPhase { /// @brief The protocol is not in an active synchronization process. Idle, /// @brief A start message has been sent, waiting for the manager's StartAck. WaitingStartAck, /// @brief An end message has been sent, waiting for the manager's EndAck. WaitingEndAck };值得注意的是,文档中DataTransfer这一“阶段”在SyncPhase中并不单独存在——数据传输期间协议仍处于等待响应(ack)的状态机上下文中,靠m_syncState.phase配合条件变量推进,而不是一个独立的等待态。
二、消息类型与 FlatBuffer 定义
所有消息都序列化为 FlatBuffer,并通过 MQueue 消息队列传输。协议定义了 8 种消息类型,逐一说明如下。
1. Start(开始消息)
- 方向:Agent → Manager
- 内容:同步模式(Full/Delta)、待发送差异的总数量
- 状态转移:
Idle→WaitingStartAck
文档给出的基础 Schema:
table Start { mode: Mode; size: uint64; }实际仓库中的 Schema(inventorySync.fbs)字段远比基础版丰富,除mode和size外还携带模块名、同步选项、Agent 元数据与集群信息:
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; }其中index: [string]与global_version正是后文“元数据/分组同步”和“数据清理”模式所依赖的字段——从Start消息的设计可以推断,Manager 端在一次握手中即可获知本次会话涉及哪些索引、以及全局版本号,用于判定是否需要全量同步。
2. StartAck(开始应答)
- 方向:Manager → Agent
- 内容:应答状态、本次同步会话的唯一会话 ID
- 状态转移:
WaitingStartAck→DataTransfer
table StartAck { status: Status; session: uint64; }session字段是后续所有数据消息和结束消息关联会话的关键。仓库中的Status枚举(见 inventorySync.fbs)实际包含 5 个取值:Ok、Error、Offline、ChecksumMismatch、Processing,其中Offline表示 Manager 当前无法服务该 Agent(如 Manager 重启窗口期、无可用 Indexer),这也是 Agent 端 SyncModuleResult 中managerNotReady标志的触发条件之一。
3. DataValue(数据消息)
- 方向:Agent → Manager
- 内容:序列号、会话 ID、操作类型(Upsert/Delete)、数据标识、目标索引、版本号、数据负载
- 状态:保持
DataTransfer
table DataValue { seq: ulong; session: ulong; operation: Operation; id: string; index: string; version: ulong; data: [byte]; }每条差异数据都分配递增的seq,这是后续 Manager 通过ReqRet按序号区间重传的前提。仓库中Operation枚举与文档一致,仅有Upsert和Delete两种。
4. DataClean(清理通知消息)
- 方向:Agent → Manager
- 内容:序列号、会话 ID、需要清理的索引名
- 状态:保持
DataTransfer
table DataClean { seq: ulong; session: ulong; index: string; }用途:在数据清理同步模式(Data Clean Mode)下,当某模块被禁用(例如关闭了 FIM)时,Agent 用该消息通知 Manager 清除对应索引下的历史数据,避免 Manager 侧残留无效数据。
5. End(结束消息)
- 方向:Agent → Manager
- 内容:会话 ID
- 状态转移:
DataTransfer→WaitingEndAck
table End { session: ulong; }6. ReqRet(请求重传消息)
- 方向:Manager → Agent
- 内容:需要重传的序列号区间列表、会话 ID
table ReqRet { seq: [Pair]; session: ulong; } table Pair { begin: ulong; end: ulong; }Agent 行为:重传指定区间内的DataValue消息。在源码中,这一过程对应 agent_sync_protocol.hpp 中的filterDataByRanges()方法:
/// @brief Filters a vector of persisted data based on a list of sequence number ranges. /// @param sourceData The complete vector of `PersistedData` items. /// @param ranges A vector of pairs [begin, end] inclusive range of sequence numbers. /// @return A new vector containing only the `PersistedData` items that match the requested ranges. std::vector<PersistedData> filterDataByRanges( const std::vector<PersistedData>& sourceData, const std::vector<std::pair<uint64_t, uint64_t>>& ranges);即 Agent 保留本次会话的完整待发数据副本,收到ReqRet后按闭区间[begin, end]过滤并重发,而不是从持久队列重新读取。
7. EndAck(结束应答)
- 方向:Manager → Agent
- 内容:状态、会话 ID
- 同一会话可能收到多次状态不同的
EndAck
table EndAck { status: Status; session: ulong; }状态取值及 Agent 侧语义(文档定义,与 Schema 中Status枚举吻合):
| 状态 | 语义 |
|---|---|
Ok | 会话成功完成,Agent 从持久队列删除已同步的数据 |
Error | 会话失败(如校验和不匹配),Agent不删除数据,保留供下次重试 |
Processing | Manager 已收到End但会话仍在处理中(例如已排队等待索引入库)。Agent 必须继续等待,不得重发End,也不消耗重试次数 |
状态转移:
- 收到
Ok或Error:WaitingEndAck→Idle - 收到
Processing:WaitingEndAck→WaitingEndAck(重置等待计时,不消耗重试)
Processing状态是异步索引场景下的关键设计:Manager 收到End后并不立刻完成入库,而是先把会话入队再异步索引。若无此状态,Agent 只能靠超时间来兜底,既浪费重试配额又拉长恢复时间。
8. ChecksumModule(校验和消息)
- 方向:Agent → Manager
- 内容:会话 ID、索引标识、校验和值
- 使用时机:完整性校验模式(integrity check mode)下,在 StartAck 与 End 之间发送
table ChecksumModule { session: ulong; index: string; checksum: string; }消息类型速查表
| 消息 | 方向 | 关键内容 | 状态影响 |
|---|---|---|---|
| Start | Agent → Manager | 模式、差异总数(实际 Schema 还含模块、Agent 元数据、集群信息) | Idle→WaitingStartAck |
| StartAck | Manager → Agent | 状态、会话 ID | WaitingStartAck→DataTransfer |
| DataValue | Agent → Manager | seq、session、操作类型、id、index、version、data | 保持DataTransfer |
| DataClean | Agent → Manager | seq、session、index | 保持DataTransfer |
| End | Agent → Manager | session | DataTransfer→WaitingEndAck |
| ReqRet | Manager → Agent | 重传区间列表、session | 保持传输期(Agent 重发 DataValue) |
| EndAck | Manager → Agent | 状态(Ok/Error/Processing)、session | Ok/Error→Idle;Processing→ 继续等待 |
| ChecksumModule | Agent → Manager | session、index、checksum | 校验模式专用 |
此外,Schema 中还定义了DataContext与DataBatch两种消息以及统一的MessageTypeunion 包装(inventorySync.fbs),分别用于同步上下文的辅助数据与批量值传输,属于数据通道的扩展能力。
三、特殊同步模式
除标准数据传输外,协议实现了四种特殊流程,分别对应 agent_sync_protocol.hpp 中的公开方法requiresFullSync()、synchronizeMetadataOrGroups()、notifyDataClean()和persistDifferenceInMemory()。
1. 完整性校验模式(Integrity Check,requiresFullSync)
用于在正式同步前比对 Agent 与 Manager 的数据校验和,决定是否需要全量同步:
Agent Manager | | |-------------- Start ----------------> | | (mode=CHECK, checksum) | | | |<------------ StartAck ---------------- | | (session_id) | | | |---------- ChecksumModule -----------> | | (index, checksum) | | | |--------------- End ------------------> | | (session_id) | | | |<------------- EndAck ----------------- | | (status: match/mismatch) |流程说明:
- Agent 以 CHECK 模式发送
Start; - Agent 发送校验和消息(
ChecksumModule); - Agent 发送
End; - Manager 比对后以
EndAck回应——校验和一致返回Ok,不一致返回Error,Agent 据此触发全量同步。
仓库中Mode枚举与该模式的对应关系可直接确认:ModuleCheck、MetadataCheck、GroupCheck三种校验模式(inventorySync.fbs)。
2. 元数据/分组同步模式(Metadata/Groups,synchronizeMetadataOrGroups)
这是无数据传输的简化流程,Agent 仅通过握手通知 Manager 处理元数据或分组变更:
Agent Manager | | |-------------- Start ----------------> | | (mode=METADATA_DELTA/GROUP_DELTA) | | | |<------------ StartAck ---------------- | | (session_id) | | | |--------------- End ------------------> | | (session_id) | | | |<------------- EndAck ----------------- | | (success) |支持的模式:METADATA_DELTA(元数据增量)、METADATA_CHECK(元数据校验)、GROUP_DELTA(分组增量)、GROUP_CHECK(分组校验)。
流程:Agent 以对应模式发送Start→ 不发送任何DataValue→ 立即发送End→ Manager 处理元数据/分组信息后回复EndAck。实际 Schema 中Start消息携带的groups: [string]、global_version、cluster_name、cluster_node等字段正是这类元数据同步的载体。
3. 数据清理模式(Data Clean,notifyDataClean)
当模块被禁用时,Agent 通知 Manager 清理相应索引,并同步清理本地数据库条目:
Agent Manager | | |-------------- Start ----------------> | | (mode=DELTA, size=N, indices=[...]) | | | |<------------ StartAck ---------------- | | (session_id) | | | |---------- DataClean[0] --------------> | | (seq=0, session, index="fim_files") | | | |---------- DataClean[1] --------------> | | (seq=1, session, index="fim_registry")| | | |---------- DataClean[N-1] ------------> | | (seq=N-1, session, index=...) | | | |--------------- End ------------------> | | | |<------------- EndAck ----------------- | | (status: Ok) | | | | clearItemsByIndex() for each index | | (cleanup local database entries) |每个待清理索引对应一条带递增seq的DataClean消息;收到EndAck(Ok)后,Agent 对每个索引执行clearItemsByIndex()清理本地持久队列条目,实现两端数据的同步收敛。
4. 内存恢复模式(In-Memory Recovery,persistDifferenceInMemory)
面向恢复场景(recovery scenarios),差异数据暂存内存而非直接落库:
Agent (Recovery) Manager | | | clearInMemoryData() | | (cleanup before sync) | | | | persistDifferenceInMemory() × N | | (storing recovery data in memory) | | | |-------------- Start ----------------> | | (mode=FULL) | | | |<------------ StartAck ---------------- | | | |------- DataValue (from memory) ------> | | (from memory) ... | | | |--------------- End ------------------> | | | |<------------- EndAck ----------------- |流程:
- 同步开始前调用
clearInMemoryData(),确保内存状态干净; - 通过
persistDifferenceInMemory()将恢复数据逐条存入内存; - 触发全量同步;
DataValue消息直接从内存向量(m_inMemoryData,见 agent_sync_protocol.hpp)发出,而非从 SQLite 持久队列读取。
源码中该成员的定义印证了这一设计:“In-memory vector to store PersistedData for recovery scenarios”。
四、完整同步流程
成功同步
Agent Manager | | |-------------- Start ----------------> | | (mode, count) | | | |<------------ StartAck ---------------- | | (session_id) | | | |----------- DataValue[0] -------------> | |----------- DataValue[1] -------------> | | ... | |----------- DataValue[N] -------------> | | | |--------------- End ------------------> | | (session_id) | | | |<------------- EndAck ----------------- | | (success) |EndAck(Ok)之后,Agent 删除持久队列中本次会话已同步的数据;EndAck(Error)则保留数据等待下一轮重试。
含 EndAck(Processing) 的同步
当 Manager 收到End后把会话入队、尚未完成索引时,会先发EndAck(Processing)让 Agent 继续等待:
Agent Manager | | |-------------- Start ----------------> | | | |<------------ StartAck ---------------- | | | |----------- DataValue[0..N] ----------> | | | |--------------- End ------------------> | | | |<------- EndAck(Processing) ----------- | (session queued, not yet indexed) | (wait again, no retry consumed, | | End NOT resent) | | | |<------------- EndAck(Ok) ------------- | (indexing complete)规则要点:不重发End、不消耗重试次数、重置等待计时。实现上,SyncState 结构体中设有专门的processingAckReceived标志位来区分这一状态,避免与普通等待混同。
含重传的同步
当某条DataValue在网络中丢失时,Manager 用ReqRet指明缺失区间,Agent 精准补发:
Agent Manager | | |-------------- Start ----------------> | | | |<------------ StartAck ---------------- | | | |----------- DataValue[0] -------------> | |----------- DataValue[1] -------------> | |----------- DataValue[2] -----X (lost) | |----------- DataValue[3] -------------> | |----------- DataValue[4] -------------> | | | |--------------- End ------------------> | | | |<------------- ReqRet ----------------- | | (ranges: [[2,2]]) | | | |----------- DataValue[2] -------------> | (retransmission) | | |<------------- EndAck ----------------- |重传依赖 Agent 端保存的完整会话数据副本与filterDataByRanges()的闭区间过滤,因此重传是精确补洞而非整批重发,带宽开销与丢失数据量成正比。
五、状态机
文档给出的完整状态机如下:
实现上,状态转移由SyncState结构体内的互斥锁 + 条件变量(std::mutex mtx; std::condition_variable cv;)协调主同步线程与响应处理线程:主线程发完消息后在cv上带超时等待,parseResponseBuffer()解析到StartAck/EndAck/ReqRet时置位对应标志(startAckReceived、endAckReceived、reqRetReceived)并唤醒等待线程。SyncState析构函数还会主动notify_all(),防止条件变量销毁时线程仍在等待造成死锁(见 agent_sync_protocol.hpp 的注释)。
另外,m_syncInProgress原子标志防止并发调用synchronizeModule():当后台刷盘线程与模块定时器线程同时触发同步时,第二个调用方直接跳过本轮——正在进行的同步会排空共享队列,并发调用是冗余的。
六、超时与重试机制
各阶段的超时行为如下(默认值来自 lifecycle.md 文档):
1. WaitingStartAck
- 默认超时:30 秒
- 超时动作:重发
Start消息 - 最大重试:可配置,默认 3 次
- 超过最大重试:中止本次同步
2. DataTransfer
- 数据发送阶段不设超时
- 通过EPS(events per second)限流做流量控制,防止压垮 Manager
3. WaitingEndAck
- 默认超时:30 秒
- 收到
EndAck(Ok/Error):会话立即完成或失败 - 收到
EndAck(Processing):重置等待计时,继续等待,不重发End、不消耗重试 - 收到
ReqRet:Agent 重传缺失序号;重传行为不消耗重试次数 - 发送
End后超时:重发End,消耗 1 次重试 - 最大重试:可配置,默认 3 次;耗尽后中止同步
这些可配置项在源码中有着落:AgentSyncProtocol构造函数显式接收timeout(默认超时)、retries(默认重试次数)、maxEps(默认 EPS 上限)与syncEndDelay(结束消息延迟)四个参数(agent_sync_protocol.hpp),即文档所述“30 秒/3 次”是各模块构造协议实例时注入的默认值,而非硬编码常量。
失败结果由 SyncResult 枚举分类上报,与超时策略直接对应:
enum class SyncResult { SUCCESS, ///< Operation completed successfully COMMUNICATION_ERROR, ///< Failed to communicate with the manager CHECKSUM_ERROR, ///< Checksum validation failed START_TIMEOUT_ERROR, ///< Exceeded maximum retries waiting for Start END_TIMEOUT_ERROR, ///< Exceeded maximum retries waiting for End PROTOCOL_ERROR, ///< Manager sent an unexpected or invalid response NO_GROUPS_ERROR, ///< No groups available in metadata. };SyncModuleResult还额外提供consecutiveFailures(连续失败计数,首次成功后归零)与managerNotReady标志,帮助调用模块区分“重启后的一次性抖动”与“持续故障”(如 Manager 无可用 Indexer),从而决定日志级别。
七、错误处理
协议错误分类
- 无效会话 ID:Manager 发来的消息携带了错误的 session ID。Agent 记录错误日志并继续等待,不影响当前同步的进行;
- 意外消息类型:收到乱序/不属于当前阶段的消息时,记录为 warning,保持当前阶段不变;
- 畸形消息:FlatBuffer 解析失败时记录 error 日志,丢弃该消息。
这套“宽松处理 + 不破坏状态”的策略与源码中validatePhaseAndSession()的校验职责一致:先校验消息所属阶段与会话 ID 是否匹配,不匹配的消息被忽略而非让状态机回退。
设计权衡
- 校验类错误(invalid session / unexpected type)只打日志不回退状态,是因为会话 ID 本身已能隔离并发或残留消息;
- 真正会导致会话失败的只有超时耗尽、
EndAck(Error)(如 checksum mismatch,对应Status::ChecksumMismatch)和通信中断; Error状态不删除持久队列数据,保证下一轮同步可以整批重发,这是“宁可多传、不可丢数据”的可靠性取向。
八、延伸阅读:相关源码与文档
| 资源 | 路径 | 说明 |
|---|---|---|
| 协议生命周期文档(本文主体) | lifecycle.md | 阶段、消息、状态机、超时重试的权威描述 |
| FlatBuffer 模式定义 | inventorySync.fbs | 全部消息表、Mode/Operation/Status/Option枚举与MessageTypeunion |
| 协议实现头文件 | agent_sync_protocol.hpp | SyncPhase、SyncState、公开 API 与线程协调机制 |
| 结果类型定义 | agent_sync_protocol_types.hpp | SyncResult、SyncModuleResult |
| C 接口 | agent_sync_protocol_c_interface.h | C 语言模块调用入口 |
| 模块 README | README.md | 架构总览、FIM/SCA/Inventory 各自独立 SQLite 库(如fim_sync.db)的设计 |
| API 参考 | api-reference.md | 完整函数签名 |
| 集成指南 | integration-guide.md | 模块接入步骤 |
| 序列图 | sequence-diagrams.md | 协议交互的可视化表示 |
| 单元测试 | test_agent_sync_protocol.cpp | 会话流程的测试用例 |
| 集成测试 | test_sync_protocol_integration.cpp | 端到端集成验证 |
从架构角度看,README 文档 说明每个内部模块(FIM、SCA、Inventory)持有独立的协议实例与独立的 SQLite 持久库,共享同一条 MQueue 消息队列。本文所述的生命周期、状态机与重试机制,正是运行在这些相互隔离的实例之上,使各模块的同步会话互不干扰,同时通过统一的 EPS 限流保证 Manager 侧的整体负载可控。
【免费下载链接】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),仅供参考