Wazuh Schema Validator 集成指南:Syscollector、SCA 与 FIM 的同步数据预校验实践
Wazuh Schema Validator 集成指南Syscollector、SCA 与 FIM 的同步数据预校验实践【免费下载链接】wazuhWazuh - The Open Source Security Platform. Unified XDR and SIEM protection for endpoints and cloud workloads.项目地址: https://gitcode.com/GitHub_Trending/wa/wazuh本文围绕 Wazuh 仓库中的 Schema Validator 集成指南 展开讲解如何将schema_validator库接入 C 模块Syscollector、SCA与 C 模块FIM确保在数据进入同步协议sync protocol并写入索引器之前完成 JSON Schema 预校验。读完本文你将掌握模块启动时的校验器工厂初始化流程、验证辅助函数的标准写法、失败数据的“延迟批量删除”deferred deletion模式、错误处理与优雅降级策略、单元测试与 Mock 注入方法以及 CMake 层面的集成细节。Wazuh 的 Syscollector、SCA、FIM 等模块都会向索引器发送状态数据如wazuh-states-inventory-packages、wazuh-states-sca、wazuh-states-fim-file等索引。如果消息结构与索引模板映射index template mapping不一致索引端会拒绝写入甚至可能触发一致性同步的反复重试。Schema Validator 的作用就是在发送前用嵌入的 JSON Schema 做本地校验把不合规数据尽早丢弃并清理本地数据库。一、Schema Validator 库的结构总览在进入集成步骤之前先从源码确认这个库对外暴露了什么这决定了后续每一步集成的调用方式。核心接口位于 schemaValidator.hpp定义了三层结构ValidationResultL27-L36校验结果包含isValid布尔标志和errors错误信息向量。默认构造时isValid为true只有发现错误才会被置为false。ISchemaValidatorEngineL43-L73抽象校验引擎接口提供validate(const std::string)、validate(const nlohmann::json)与getSchemaName()三个虚方法。该抽象接口的意义在于支持依赖注入——测试时可以用 GMock 实现一个 Mock 校验器替换真实引擎。SchemaValidatorFactoryL137-L192单例工厂提供getInstance()、initialize(customValidators)、getValidator(indexPattern)、isInitialized()、reset()五个方法。initialize()支持传入自定义校验器映射用于测试注入为空时默认从编译期嵌入的资源加载 Schemareset()则专门用于测试中清空单例状态。C 语言模块则通过 schemaValidator_c.h 中的三个导出函数使用同一个工厂schema_validator_initialize()从嵌入资源初始化工厂schema_validator_is_initialized()检查工厂是否已初始化schema_validator_validate(indexPattern, message, errorMessage)按索引模式校验消息失败时输出错误信息调用方需free。校验规则的实际实现从 schemaValidator.cpp 的实现看引擎解析的是Elasticsearch/OpenSearch 索引模板格式的 JSON它读取template.mappings.properties作为字段定义读取index_patterns[0]去掉*通配符后缀作为 Schema 名称也就是getValidator(index)所使用的索引 keyL139-L178。字段级校验规则如下对应 L266-L370 的validateFieldSchema 类型合法取值keyword/text/match_only_text字符串integer/long/short/unsigned_long整数scaled_float数值dateepoch 数字或符合 ISO8601 的字符串正则校验ip经inet_pton验证的 IPv4/IPv6 字符串兼容 OpenSearch 行为%zone后缀会被截断后再解析如fe80::1%eth0boolean布尔值含properties的对象递归校验子字段几条值得注意的宽松规则与 OpenSearch 的实际索引行为对齐null值在任何字段上都被放行包括嵌套对象字段L243-L247数组按元素类型逐个校验路径形如field[0]Strict 模式当模板的dynamic字段为strict时未定义的额外字段会被拒绝——但值为null的额外字段例外放行因为 null 不参与索引映射L206-L223字段缺省是允许的字段默认可选只有出现在消息中的字段才会被校验类型。嵌入资源与线程安全工厂的initialize()实现L457-L499有两个对集成方很重要的特性幂等性工厂是进程级单例Syscollector 与 SCA 可能并发启动并各自调用initialize()。一旦已初始化后续调用直接返回成功不会重复加载。这也是集成指南中“初始化前先检查isInitialized()”这一步的底层原因。内部使用std::shared_mutex对m_validators做读写锁保护写独占、读共享。Schema 在编译期嵌入schema_validator/CMakeLists.txt 会 Globsrc/external/indexer-plugins/目录下的全部 JSON 模板文件排除仅供 Python 框架运行时使用的metrics-*.json生成schemaResources.cpp把每个模板以 raw string literal 形式嵌入 C 源码并由 schemaResources.hpp 声明的Resources::getEmbeddedSchemas()以“文件名 → JSON 内容”的 map 返回。运行时不需要磁盘上的 Schema 文件这也是日志中 “initialized successfully from embedded resources” 说法的来源。二、C 模块集成Syscollector、SCA以下 5 个步骤完整覆盖一个 C 模块接入校验器的过程与仓库中 syscollectorImp.cpp、sca_event_handler.cpp 的实际用法一致。Step 1包含头文件#include schemaValidator.hppStep 2模块启动时初始化一次void YourModule::initialize() { // Initialize schema validator from embedded resources auto validatorFactory SchemaValidator::SchemaValidatorFactory::getInstance(); if (!validatorFactory.isInitialized()) { if (validatorFactory.initialize()) { m_logFunction(LOG_INFO, Schema validator initialized successfully from embedded resources); } else { m_logFunction(LOG_WARNING, Failed to initialize schema validator. Schema validation will be disabled.); } } }注意initialize()本身是幂等且线程安全的见上文这里的isInitialized()检查主要是为了避免无谓的调用并保证只在启动期打印一次成功/失败日志。Step 3创建校验辅助函数将校验逻辑封装成一个可复用的辅助函数统一处理“工厂未初始化 → 跳过校验”“无对应索引的校验器 → 跳过校验”两种降级路径bool YourModule::validateSchemaAndLog(const std::string data, const std::string index, const std::string context) const { auto validatorFactory SchemaValidator::SchemaValidatorFactory::getInstance(); if (!validatorFactory.isInitialized()) { return true; // Validation disabled } auto validator validatorFactory.getValidator(index); if (!validator) { return true; // No validator for this index } auto validationResult validator-validate(data); if (validationResult.isValid) { return true; } // Validation failed - log errors std::string errorMsg Schema validation failed for message ( context , index: index ). Errors: ; for (const auto error : validationResult.errors) { errorMsg - error; } if (m_logFunction) { m_logFunction(LOG_ERROR, errorMsg); m_logFunction(LOG_ERROR, Raw event that failed validation: data); } return false; }辅助函数的设计要点返回true的两种降级情形工厂未初始化、该索引没有校验器都表示“放行”保证校验器故障不会阻断正常的同步流程失败时输出两条LOG_ERROR一条是结构化错误列表一条是原始事件 JSON。原始事件日志是排障的关键——排查“校验总是失败”时需要拿它和对应索引的 Schema 模板逐字段比对context参数用于标注调用点例如table: tableName使多条错误日志可定位到具体处理路径。Step 4发送数据前调用校验void YourModule::processEvent(const std::string data, const std::string index) { // Validate data std::string context event processing; bool validationPassed validateSchemaAndLog(data, index, context); if (!validationPassed) { // Discard invalid data if (m_logFunction) { m_logFunction(LOG_ERROR, Discarding invalid message); } // Mark for deletion from database markForDeletion(data); return; } // Send valid data to sync protocol m_spSyncProtocol-persistDifference(id, operation, index, data, version); }校验失败的数据会被丢弃并标记为待删除而非立即删除删除动作延迟到扫描周期结束原因见后文“延迟删除模式”一节。Step 5实现批量删除void YourModule::deleteFailedItemsFromDB( const std::vectorstd::pairstd::string, nlohmann::json failedItems) const { if (failedItems.empty() || !m_spDBSync) { return; } try { // Create a transaction DBSyncTxn deleteTxn(m_spDBSync-handle(), nlohmann::json::array(), 0, 1, [](ReturnTypeCallback, const nlohmann::json) {}); // Delete all failed items for (const auto [tableName, data] : failedItems) { if (m_logFunction) { m_logFunction(LOG_DEBUG, Deleting entry from table tableName due to validation failure); } try { auto deleteQuery DeleteQuery::builder() .table(tableName) .data(data) .rowFilter() .build(); m_spDBSync-deleteRows(deleteQuery.query()); } catch (const std::exception e) { if (m_logFunction) { m_logFunction(LOG_ERROR, Failed to delete from DBSync: std::string(e.what())); } } } // Finalize transaction deleteTxn.getDeletedRows([](ReturnTypeCallback, const nlohmann::json) {}); if (m_logFunction) { m_logFunction(LOG_DEBUG, Deleted std::to_string(failedItems.size()) item(s) from DBSync due to validation failure); } } catch (const std::exception e) { if (m_logFunction) { m_logFunction(LOG_ERROR, Failed to create DBSync transaction for deletion: std::string(e.what())); } } }这个函数把所有失败项放进一个 DBSync 事务中批量删除DBSyncTxn单条删除抛异常不会中断整体流程只是记录LOG_ERROR继续处理剩余条目。三、C 模块集成FIMFIMsyscheckd是 C 语言实现通过 schemaValidator_c.h 的 C 包装函数接入同一个工厂。仓库中 syscheck.c、run_check.c 等文件可见实际调用点。Step 1包含头文件#include schemaValidator_c.hStep 2模块启动时初始化void fim_initialize(void) { // Initialize schema validator from embedded resources if (!schema_validator_is_initialized()) { if (schema_validator_initialize()) { minfo(Schema validator initialized successfully from embedded resources); } else { mwarn(Failed to initialize schema validator. Schema validation will be disabled.); } } }Step 3发送数据前校验FIM 的校验函数多了一个与 C 模块不同的条件只有同步功能启用时才校验并且失败项通过OSList累积用于延迟删除bool fim_validate_and_queue(const char* index, const char* data, void* item_data, OSList* failed_list) { bool validation_passed true; // Only validate if synchronization is enabled and schema validator is initialized if (syscheck.enable_synchronization schema_validator_is_initialized()) { char* errorMessage NULL; if (!schema_validator_validate(index, data, errorMessage)) { // Validation failed - log errors if (errorMessage) { mdebug2(Schema validation failed for FIM message (index: %s). Error: %s, index, errorMessage); mdebug2(Raw event that failed validation: %s, data); free(errorMessage); } // Mark for deferred deletion from database if (failed_list item_data) { mdebug1(Marking FIM entry for deferred deletion due to validation failure); OSList_AddData(failed_list, item_data); } validation_passed false; } } return validation_passed; }C 接口与 C 接口在语义上完全对应schema_validator_validate内部就是取出对应索引的ISchemaValidatorEngine调用validate并将ValidationResult::errors拼接为单条错误字符串输出。注意 C 侧由调用方负责free(errorMessage)避免内存泄漏。Step 4批量删除失败项void fim_delete_failed_items(OSList* failed_list) { if (!failed_list || OSList_GetSize(failed_list) 0) { return; } mdebug1(Deleting %d FIM item(s) from database due to validation failure, OSList_GetSize(failed_list)); OSListNode* node; OSList_foreach(node, failed_list) { void* item_data node-data; // Delete item from database fim_db_remove_path(syscheck.database, item_data); } }FIM 的本地存储是 SQLite 路径数据库因此直接按条目调用fim_db_remove_path逐条删除与 C 模块的 DBSync 事务批量删除在目标上一致把无法通过索引 Schema 的数据从本地状态库中清掉防止下一轮同步再次发送同样的坏数据形成一致性死循环。四、两个模块的具体辅助函数模式集成指南给出了两种已经落地的模式分别对应 Syscollector 和 SCA 的实现差异。模式 1带上下文日志的校验Syscollectorbool Syscollector::validateSchemaAndLog(const std::string data, const std::string index, const std::string context) const { auto validatorFactory SchemaValidator::SchemaValidatorFactory::getInstance(); if (!validatorFactory.isInitialized()) { return true; } auto validator validatorFactory.getValidator(index); if (!validator) { return true; } auto validationResult validator-validate(data); if (validationResult.isValid) { return true; } // Validation failed - log errors std::string errorMsg Schema validation failed for Syscollector message ( context , index: index ). Errors: ; for (const auto error : validationResult.errors) { errorMsg - error; } if (m_logFunction) { m_logFunction(LOG_ERROR, errorMsg); m_logFunction(LOG_ERROR, Raw event that failed validation: data); } return false; }调用方式对应 syscollectorImp.cpp 中的用法bool validationPassed validateSchemaAndLog(statefulToSend, index, table: tableName); if (!validationPassed) { // Discard and mark for deletion }Syscollector 每次同步的数据携带具体表名hardware、os、packages 等对应wazuh-states-inventory-*系列索引因此context传入table: tableName让错误日志直接定位到出问题的表和索引。模式 2校验加延迟删除SCAbool SCAEventHandler::ValidateAndHandleStatefulMessage( const nlohmann::json statefulEvent, const std::string context, const nlohmann::json checkData, std::vectornlohmann::json* failedChecks) const { if (statefulEvent.empty()) { return true; } auto validatorFactory SchemaValidator::SchemaValidatorFactory::getInstance(); if (!validatorFactory.isInitialized()) { return true; } auto validator validatorFactory.getValidator(SCA_SYNC_INDEX); if (!validator) { return true; } std::string statefulData statefulEvent.dump(); auto validationResult validator-validate(statefulData); if (validationResult.isValid) { return true; } // Validation failed - log errors std::string errorMsg Schema validation failed for SCA message ( context , index: std::string(SCA_SYNC_INDEX) ). Errors: ; for (const auto error : validationResult.errors) { errorMsg - error; } LoggingHelper::getInstance().log(LOG_ERROR, errorMsg); LoggingHelper::getInstance().log(LOG_ERROR, Raw event that failed validation: statefulData); // Handle deletion from DBSync to prevent integrity sync loops if (!checkData.empty() failedChecks) { // Deferred deletion: accumulate for batch deletion with transaction LoggingHelper::getInstance().log(LOG_DEBUG, Marking SCA check for deferred deletion due to validation failure); failedChecks-push_back(checkData); } return false; }调用方式对应 sca_event_handler.cpp 中的事件循环std::vectornlohmann::json failedChecks; // Process events for (const auto event : events) { bool validationPassed ValidateAndHandleStatefulMessage( event, context, checkData, failedChecks); if (validationPassed) { PushStateful(event, operation, version); } } // Batch delete DeleteFailedChecksFromDB(failedChecks);与模式 1 的关键差异在于SCA 的校验函数同时承担“累积失败项”的职责——它接收一个failedChecks输出参数把失败 check 的数据库记录checkData压入其中供循环结束后统一批量删除。代码注释点明了动机“prevent integrity sync loops”——如果只丢弃消息而不删除本地 DBSync 中的记录下一轮同步会重新发送同一条非法数据造成无休止的重发循环。五、延迟删除模式Deferred Deletion这是集成指南中最核心的工程模式值得单独理解。为什么需要延迟删除问题校验回调发生在 DBSync 的处理流程内部此时如果立即对数据库执行删除会形成嵌套事务nested transactions破坏 DBSync 的事务模型。解决方案在处理过程中只“累积”失败项处理全部结束后再用单个批量事务一次性删除。实现四步Step 1创建累积容器std::vectorstd::pairstd::string, nlohmann::json failedItems; m_failedItems failedItems; // Make accessible to callbacks容器元素是(表名, 数据)对通过成员指针m_failedItems让回调链路中的深层函数也能访问到当前周期累积的列表。Step 2处理过程中累积if (!validationPassed) { if (m_failedItems) { m_failedItems-push_back({tableName, data}); } }Step 3处理完成后清理指针m_failedItems nullptr;置空指针防止后续误写——这是一个生命周期明确的裸指针模式nullptr即“当前不在扫描周期内”。Step 4批量删除deleteFailedItemsFromDB(failedItems);完整示例Syscollectorvoid Syscollector::scan() { // Vector to accumulate items that fail validation std::vectorstd::pairstd::string, nlohmann::json failedItems; m_failedItems failedItems; // Run scans scanHardware(); scanOs(); scanPackages(); // ... etc // Clean up after all scans m_failedItems nullptr; // Delete all items that failed schema validation deleteFailedItemsFromDB(failedItems); }整个扫描周期内各子扫描hardware、os、packages 等独立触发回调并累积失败项scan()退出回调上下文后才执行删除从而彻底避开嵌套事务。六、错误处理策略优雅降级Graceful Degradation校验器不可用时永远放行绝不阻断主流程auto factory SchemaValidator::SchemaValidatorFactory::getInstance(); // Check if initialized if (!factory.isInitialized()) { // Log warning once during startup m_logFunction(LOG_WARNING, Schema validator not initialized. Validation disabled.); return true; // Continue without validation } // Check if validator exists for index auto validator factory.getValidator(index); if (!validator) { // No validator for this index - continue without validation return true; } // Proceed with validation auto result validator-validate(data);这一策略的合理性在于校验器只是防御层索引端OpenSearch最终仍会按映射做类型检查。校验器初始化失败时降级为“无校验直发”比让模块停止同步更合理。日志分级策略初始化阶段// During startup LOG_INFO: Schema validator initialized successfully LOG_WARNING: Schema validator not initialized. Validation disabled.校验失败时// When validation fails LOG_ERROR: Schema validation failed for module message (context, index: index). Errors: details LOG_ERROR: Raw event that failed validation: json LOG_DEBUG: Marking entry for deferred deletion due to validation failure删除阶段// After batch deletion LOG_DEBUG: Deleted N item(s) from database due to validation failure LOG_ERROR: Failed to delete from database: error // If deletion fails规律是错误详情与原始事件用LOG_ERROR必须留痕是排障核心输入累积标记与批量删除成功用LOG_DEBUG高频且可接受删除失败才升级为LOG_ERROR。七、测试集成单元测试结构仓库自带完整测试schemaValidator_test.cpp 与 schemaValidatorConcurrency_test.cpp后者专门验证工厂在多线程initialize()竞争下的正确性。集成指南给出的标准测试骨架如下——注意SetUp/TearDown中必须调用factory.reset()否则单例状态会跨用例污染#include gtest/gtest.h #include schemaValidator.hpp class SchemaValidatorTest : public ::testing::Test { protected: void SetUp() override { // Reset factory for clean state auto factory SchemaValidator::SchemaValidatorFactory::getInstance(); factory.reset(); } void TearDown() override { // Clean up auto factory SchemaValidator::SchemaValidatorFactory::getInstance(); factory.reset(); } }; TEST_F(SchemaValidatorTest, ValidMessage) { auto factory SchemaValidator::SchemaValidatorFactory::getInstance(); ASSERT_TRUE(factory.initialize()); auto validator factory.getValidator(wazuh-states-inventory-packages); ASSERT_NE(validator, nullptr); std::string validJson R({ agent: {id: 001}, package: {name: nginx, version: 1.18.0} }); auto result validator-validate(validJson); EXPECT_TRUE(result.isValid); EXPECT_TRUE(result.errors.empty()); } TEST_F(SchemaValidatorTest, InvalidMessage) { auto factory SchemaValidator::SchemaValidatorFactory::getInstance(); ASSERT_TRUE(factory.initialize()); auto validator factory.getValidator(wazuh-states-inventory-packages); ASSERT_NE(validator, nullptr); std::string invalidJson R({ agent: {id: 001}, package: {name: 123} }); auto result validator-validate(invalidJson); EXPECT_FALSE(result.isValid); EXPECT_FALSE(result.errors.empty()); }这两个用例正好覆盖上文提到的类型规则package.name在 Schema 中是字符串类型传入数字123会触发Expected string, got number错误并附带路径信息。用 Mock 校验器测试模块的失败处理ISchemaValidatorEngine抽象接口 initialize(customValidators)注入参数见 schemaValidator.hpp 的注释“This allows dependency injection for testing”就是为这类测试设计的class MockSchemaValidator : public SchemaValidator::ISchemaValidatorEngine { public: MOCK_METHOD(ValidationResult, validate, (const std::string), (override)); MOCK_METHOD(ValidationResult, validate, (const nlohmann::json), (override)); MOCK_METHOD(std::string, getSchemaName, (), (const, override)); }; TEST_F(YourModuleTest, ValidationFailureHandling) { // Create mock validator that always fails auto mockValidator std::make_sharedMockSchemaValidator(); ON_CALL(*mockValidator, validate(testing::_)) .WillByDefault(testing::Return(ValidationResult{false, {Test error}})); // Inject mock std::mapstd::string, std::shared_ptrSchemaValidator::ISchemaValidatorEngine mocks; mocks[test-index] mockValidator; auto factory SchemaValidator::SchemaValidatorFactory::getInstance(); factory.reset(); factory.initialize(mocks); // Test your modules handling of validation failure bool result yourModule-processData(testData, test-index); EXPECT_FALSE(result); // Should handle validation failure correctly }这个模式允许在不依赖真实 Schema 文件的情况下精确模拟“校验总是失败”的场景验证模块是否正确丢弃数据、累积删除项、输出正确的日志——这正是集成指南“Testing Integration”一节要覆盖的行为。八、CMake 集成模块的CMakeLists.txt只需链接schema_validator库target_link_libraries(your_module PRIVATE schema_validator )从 schema_validator/CMakeLists.txt 看schema_validator是一个 C17 共享库由schemaValidator.cpp、schemaValidator_c.cpp与构建时生成的schemaResources.cpp三部分组成它把include/目录作为PUBLIC包含路径导出因此链接方直接#include schemaValidator.hpp/#include schemaValidator_c.h即可。库还处理了平台差异Windows 下以WIN_EXPORT宏导出符号并生成 DLL链接时附加-Wl,--add-stdcall-aliasLinux 下设置-rpath$ORIGIN使运行时能从模块旁找到.so。如果构建时src/external/indexer-plugins/下没有 Schema 文件CMake 只会告警“Schemas will not be embedded”模块侧则表现为initialize()失败、按“优雅降级”策略禁用校验——这也解释了排障表中“工厂返回 nullptr”的第一种成因。九、常见故障排查问题 1getValidator返回 nullptr原因工厂未初始化或该索引模式没有对应的嵌入 Schema。解决if (!factory.isInitialized()) { factory.initialize(); } auto validator factory.getValidator(index); if (!validator) { m_logFunction(LOG_WARNING, No validator found for index: index); // Continue without validation }排查时先确认模块发送的索引名如wazuh-states-sca、wazuh-states-inventory-hardware与src/external/indexer-plugins/下模板文件的index_patterns是否一致——工厂的 map key 取自模板的第一个 pattern 去掉*后缀两者必须完全匹配。问题 2校验总是失败原因数据与 Schema 结构不匹配。解决查看错误日志中输出的原始事件 JSON“Raw event that failed validation”对照该索引对应的 Schema 模板文件逐字段比对确认字段名与类型严格一致注意 strict 模式下额外非 null 字段会被拒绝检查是否有缺失的必填结构以及date/ip字段的格式是否满足 ISO8601 /inet_pton的要求。问题 3性能退化原因在循环内反复调用getValidator或者本可复用的校验器被重复获取。解决缓存校验器实例循环内复用// Cache validator auto validator factory.getValidator(index); // Reuse in loop for (const auto item : items) { validator-validate(item); // Fast }getValidator内部走的是std::shared_mutex共享锁的 map 查找L501-L513单次开销不大但批量事件处理时缓存仍然能省去重复加锁与查找且缓存的shared_ptr保证校验器实例在整个处理周期内稳定有效。十、相关文档Schema Validator API Reference完整 API 文档Syscollector ArchitectureSyscollector 中的 Schema 校验集成SCA ArchitectureSCA 中的 Schema 校验集成FIM ArchitectureFIM 中的 Schema 校验集成READMESchema Validator 模块总览。【免费下载链接】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),仅供参考