Wazuh Agent Sync Protocol 会话生命周期:数据同步的阶段划分、消息类型与状态机全解
2026/9/14 3:43:56 网站建设 项目流程

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 之间数据的一致性。一次完整的同步会话经历以下阶段:

  1. Idle 阶段:无活跃同步,协议实例处于空闲状态;
  2. Session Establishment(会话建立):Agent 发送Start消息,Manager 回复StartAck并分配会话 ID;
  3. Data Transfer(数据传输):逐条发送差异数据(differences),Manager 可请求重传丢失的序号区间;
  4. 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)、待发送差异的总数量
  • 状态转移IdleWaitingStartAck

文档给出的基础 Schema:

table Start { mode: Mode; size: uint64; }

实际仓库中的 Schema(inventorySync.fbs)字段远比基础版丰富,除modesize外还携带模块名、同步选项、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
  • 状态转移WaitingStartAckDataTransfer
table StartAck { status: Status; session: uint64; }

session字段是后续所有数据消息和结束消息关联会话的关键。仓库中的Status枚举(见 inventorySync.fbs)实际包含 5 个取值:OkErrorOfflineChecksumMismatchProcessing,其中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枚举与文档一致,仅有UpsertDelete两种。

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
  • 状态转移DataTransferWaitingEndAck
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不删除数据,保留供下次重试
ProcessingManager 已收到End但会话仍在处理中(例如已排队等待索引入库)。Agent 必须继续等待,不得重发End,也不消耗重试次数

状态转移

  • 收到OkErrorWaitingEndAckIdle
  • 收到ProcessingWaitingEndAckWaitingEndAck(重置等待计时,不消耗重试)

Processing状态是异步索引场景下的关键设计:Manager 收到End后并不立刻完成入库,而是先把会话入队再异步索引。若无此状态,Agent 只能靠超时间来兜底,既浪费重试配额又拉长恢复时间。

8. ChecksumModule(校验和消息)

  • 方向:Agent → Manager
  • 内容:会话 ID、索引标识、校验和值
  • 使用时机:完整性校验模式(integrity check mode)下,在 StartAck 与 End 之间发送
table ChecksumModule { session: ulong; index: string; checksum: string; }

消息类型速查表

消息方向关键内容状态影响
StartAgent → Manager模式、差异总数(实际 Schema 还含模块、Agent 元数据、集群信息)IdleWaitingStartAck
StartAckManager → Agent状态、会话 IDWaitingStartAckDataTransfer
DataValueAgent → Managerseq、session、操作类型、id、index、version、data保持DataTransfer
DataCleanAgent → Managerseq、session、index保持DataTransfer
EndAgent → ManagersessionDataTransferWaitingEndAck
ReqRetManager → Agent重传区间列表、session保持传输期(Agent 重发 DataValue)
EndAckManager → Agent状态(Ok/Error/Processing)、sessionOk/ErrorIdleProcessing→ 继续等待
ChecksumModuleAgent → Managersession、index、checksum校验模式专用

此外,Schema 中还定义了DataContextDataBatch两种消息以及统一的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) |

流程说明:

  1. Agent 以 CHECK 模式发送Start
  2. Agent 发送校验和消息(ChecksumModule);
  3. Agent 发送End
  4. Manager 比对后以EndAck回应——校验和一致返回Ok,不一致返回Error,Agent 据此触发全量同步。

仓库中Mode枚举与该模式的对应关系可直接确认:ModuleCheckMetadataCheckGroupCheck三种校验模式(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_versioncluster_namecluster_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) |

每个待清理索引对应一条带递增seqDataClean消息;收到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 ----------------- |

流程

  1. 同步开始前调用clearInMemoryData(),确保内存状态干净;
  2. 通过persistDifferenceInMemory()将恢复数据逐条存入内存;
  3. 触发全量同步;
  4. 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时置位对应标志(startAckReceivedendAckReceivedreqRetReceived)并唤醒等待线程。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),从而决定日志级别。

七、错误处理

协议错误分类

  1. 无效会话 ID:Manager 发来的消息携带了错误的 session ID。Agent 记录错误日志并继续等待,不影响当前同步的进行;
  2. 意外消息类型:收到乱序/不属于当前阶段的消息时,记录为 warning,保持当前阶段不变
  3. 畸形消息: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.hppSyncPhaseSyncState、公开 API 与线程协调机制
结果类型定义agent_sync_protocol_types.hppSyncResultSyncModuleResult
C 接口agent_sync_protocol_c_interface.hC 语言模块调用入口
模块 READMEREADME.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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询