ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Wazuh Agent Sync Protocol 会话生命周期:数据同步的阶段划分、消息类型与状态机全解

Wazuh Agent Sync Protocol 会话生命周期:数据同步的阶段划分、消息类型与状态机全解 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/wazuhWazuh 的 Agent Sync Protocolsrc/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并分配会话 IDData Transfer数据传输逐条发送差异数据differencesManager 可请求重传丢失的序号区间Session Completion会话完成Agent 发送EndManager 以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 managers StartAck. WaitingStartAck, /// brief An end message has been sent, waiting for the managers EndAck. WaitingEndAck };值得注意的是文档中DataTransfer这一“阶段”在SyncPhase中并不单独存在——数据传输期间协议仍处于等待响应ack的状态机上下文中靠m_syncState.phase配合条件变量推进而不是一个独立的等待态。二、消息类型与 FlatBuffer 定义所有消息都序列化为 FlatBuffer并通过 MQueue 消息队列传输。协议定义了 8 种消息类型逐一说明如下。1. Start开始消息方向Agent → Manager内容同步模式Full/Delta、待发送差异的总数量状态转移Idle→WaitingStartAck文档给出的基础 Schematable Start { mode: Mode; size: uint64; }实际仓库中的 SchemainventorySync.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→DataTransfertable 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、数据标识、目标索引、版本号、数据负载状态保持DataTransfertable 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、需要清理的索引名状态保持DataTransfertable DataClean { seq: ulong; session: ulong; index: string; }用途在数据清理同步模式Data Clean Mode下当某模块被禁用例如关闭了 FIM时Agent 用该消息通知 Manager 清除对应索引下的历史数据避免 Manager 侧残留无效数据。5. End结束消息方向Agent → Manager内容会话 ID状态转移DataTransfer→WaitingEndAcktable End { session: ulong; }6. ReqRet请求重传消息方向Manager → Agent内容需要重传的序列号区间列表、会话 IDtable 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::vectorPersistedData filterDataByRanges( const std::vectorPersistedData sourceData, const std::vectorstd::pairuint64_t, uint64_t ranges);即 Agent 保留本次会话的完整待发数据副本收到ReqRet后按闭区间[begin, end]过滤并重发而不是从持久队列重新读取。7. EndAck结束应答方向Manager → Agent内容状态、会话 ID同一会话可能收到多次状态不同的EndAcktable EndAck { status: Status; session: ulong; }状态取值及 Agent 侧语义文档定义与 Schema 中Status枚举吻合状态语义Ok会话成功完成Agent 从持久队列删除已同步的数据Error会话失败如校验和不匹配Agent不删除数据保留供下次重试ProcessingManager 已收到End但会话仍在处理中例如已排队等待索引入库。Agent 必须继续等待不得重发End也不消耗重试次数状态转移收到Ok或ErrorWaitingEndAck→Idle收到ProcessingWaitingEndAck→WaitingEndAck重置等待计时不消耗重试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 元数据、集群信息Idle→WaitingStartAckStartAckManager → Agent状态、会话 IDWaitingStartAck→DataTransferDataValueAgent → Managerseq、session、操作类型、id、index、version、data保持DataTransferDataCleanAgent → Managerseq、session、index保持DataTransferEndAgent → ManagersessionDataTransfer→WaitingEndAckReqRetManager → Agent重传区间列表、session保持传输期Agent 重发 DataValueEndAckManager → Agent状态Ok/Error/Processing、sessionOk/Error→IdleProcessing→ 继续等待ChecksumModuleAgent → Managersession、index、checksum校验模式专用此外Schema 中还定义了DataContext与DataBatch两种消息以及统一的MessageTypeunion 包装inventorySync.fbs分别用于同步上下文的辅助数据与批量值传输属于数据通道的扩展能力。三、特殊同步模式除标准数据传输外协议实现了四种特殊流程分别对应 agent_sync_protocol.hpp 中的公开方法requiresFullSync()、synchronizeMetadataOrGroups()、notifyDataClean()和persistDifferenceInMemory()。1. 完整性校验模式Integrity CheckrequiresFullSync用于在正式同步前比对 Agent 与 Manager 的数据校验和决定是否需要全量同步Agent Manager | | |-------------- Start ---------------- | | (modeCHECK, checksum) | | | |------------ StartAck ---------------- | | (session_id) | | | |---------- ChecksumModule ----------- | | (index, checksum) | | | |--------------- End ------------------ | | (session_id) | | | |------------- EndAck ----------------- | | (status: match/mismatch) |流程说明Agent 以 CHECK 模式发送StartAgent 发送校验和消息ChecksumModuleAgent 发送EndManager 比对后以EndAck回应——校验和一致返回Ok不一致返回ErrorAgent 据此触发全量同步。仓库中Mode枚举与该模式的对应关系可直接确认ModuleCheck、MetadataCheck、GroupCheck三种校验模式inventorySync.fbs。2. 元数据/分组同步模式Metadata/GroupssynchronizeMetadataOrGroups这是无数据传输的简化流程Agent 仅通过握手通知 Manager 处理元数据或分组变更Agent Manager | | |-------------- Start ---------------- | | (modeMETADATA_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 CleannotifyDataClean当模块被禁用时Agent 通知 Manager 清理相应索引并同步清理本地数据库条目Agent Manager | | |-------------- Start ---------------- | | (modeDELTA, sizeN, indices[...]) | | | |------------ StartAck ---------------- | | (session_id) | | | |---------- DataClean[0] -------------- | | (seq0, session, indexfim_files) | | | |---------- DataClean[1] -------------- | | (seq1, session, indexfim_registry)| | | |---------- DataClean[N-1] ------------ | | (seqN-1, session, index...) | | | |--------------- End ------------------ | | | |------------- EndAck ----------------- | | (status: Ok) | | | | clearItemsByIndex() for each index | | (cleanup local database entries) |每个待清理索引对应一条带递增seq的DataClean消息收到EndAck(Ok)后Agent 对每个索引执行clearItemsByIndex()清理本地持久队列条目实现两端数据的同步收敛。4. 内存恢复模式In-Memory RecoverypersistDifferenceInMemory面向恢复场景recovery scenarios差异数据暂存内存而非直接落库Agent (Recovery) Manager | | | clearInMemoryData() | | (cleanup before sync) | | | | persistDifferenceInMemory() × N | | (storing recovery data in memory) | | | |-------------- Start ---------------- | | (modeFULL) | | | |------------ 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数据发送阶段不设超时通过EPSevents per second限流做流量控制防止压垮 Manager3. WaitingEndAck默认超时30 秒收到EndAck(Ok/Error)会话立即完成或失败收到EndAck(Processing)重置等待计时继续等待不重发End、不消耗重试收到ReqRetAgent 重传缺失序号重传行为不消耗重试次数发送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从而决定日志级别。七、错误处理协议错误分类无效会话 IDManager 发来的消息携带了错误的 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.hppSyncPhase、SyncState、公开 API 与线程协调机制结果类型定义agent_sync_protocol_types.hppSyncResult、SyncModuleResultC 接口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),仅供参考
返回列表