From 0e4bf245e1d82a85f5ba88c6a9e6c4c5691a5c98 Mon Sep 17 00:00:00 2001 From: fengzero Date: Tue, 11 Aug 2026 06:46:06 +0000 Subject: [PATCH 1/2] schema support int/bool/float --- plugins/mqtt/schema.c | 31 +++++++++++++++++++++++--- plugins/mqtt/schema.h | 3 +++ tests/ut/mqtt_schema_test.cc | 42 ++++++++++++++++++++++++++++++++++++ 3 files changed, 73 insertions(+), 3 deletions(-) diff --git a/plugins/mqtt/schema.c b/plugins/mqtt/schema.c index 69458560c..997293db9 100644 --- a/plugins/mqtt/schema.c +++ b/plugins/mqtt/schema.c @@ -49,10 +49,31 @@ static int mqtt_schema_parse(json_t *root, mqtt_schema_vt_t **vts, str_val[strlen(str_val) - 1] == '}') { return -1; } else { - vt->vt = MQTT_SCHEMA_UD; + vt->vt = MQTT_SCHEMA_UD; + vt->jtype = NEU_JSON_STR; strncpy(vt->ud, str_val, sizeof(vt->ud) - 1); } } + } else if (json_is_integer(value) || json_is_real(value) || + json_is_boolean(value)) { + *vts_len += 1; + *vts = realloc(*vts, *vts_len * sizeof(mqtt_schema_vt_t)); + mqtt_schema_vt_t *vt = &(*vts)[*vts_len - 1]; + + memset(vt, 0, sizeof(mqtt_schema_vt_t)); + strncpy(vt->name, key, sizeof(vt->name) - 1); + vt->vt = MQTT_SCHEMA_UD; + + if (json_is_integer(value)) { + vt->jtype = NEU_JSON_INT; + vt->jvalue.val_int = json_integer_value(value); + } else if (json_is_real(value)) { + vt->jtype = NEU_JSON_DOUBLE; + vt->jvalue.val_double = json_real_value(value); + } else { + vt->jtype = NEU_JSON_BOOL; + vt->jvalue.val_bool = json_boolean_value(value); + } } else if (json_is_object(value)) { if (deep >= 3) { return -1; @@ -280,8 +301,12 @@ static void *schema_encode(char *driver, char *group, break; } case MQTT_SCHEMA_UD: - elem.t = NEU_JSON_STR; - elem.v.val_str = vts[i].ud; + elem.t = vts[i].jtype; + if (vts[i].jtype == NEU_JSON_STR) { + elem.v.val_str = vts[i].ud; + } else { + elem.v = vts[i].jvalue; + } break; case MQTT_SCHEMA_OBJECT: { void *sub_root = schema_encode(driver, group, tags, vts[i].sub_vts, diff --git a/plugins/mqtt/schema.h b/plugins/mqtt/schema.h index ebeece9ab..b51a50825 100644 --- a/plugins/mqtt/schema.h +++ b/plugins/mqtt/schema.h @@ -53,6 +53,9 @@ typedef struct mqtt_schema_vt { char ud[128]; + neu_json_type_e jtype; + neu_json_value_u jvalue; + mqtt_schema_vt_t *sub_vts; size_t n_sub_vts; } mqtt_schema_vt_t; diff --git a/tests/ut/mqtt_schema_test.cc b/tests/ut/mqtt_schema_test.cc index 5d94dda0e..754a84cde 100644 --- a/tests/ut/mqtt_schema_test.cc +++ b/tests/ut/mqtt_schema_test.cc @@ -89,6 +89,48 @@ TEST(validate_deep_success, schema) free(vts); } +TEST(validate_constants, schema) +{ + const char *schema = + "{\"string\":\"value\",\"integer\":1234,\"float\":11.44," + "\"boolean\":true,\"object\":{\"integer\":-42,\"float\":0.5," + "\"boolean\":false}}"; + mqtt_schema_vt_t *vts = NULL; + size_t n_vts; + int ret = mqtt_schema_validate(schema, &vts, &n_vts); + + ASSERT_EQ(ret, 0); + ASSERT_EQ(n_vts, 5); + EXPECT_EQ(vts[0].jtype, NEU_JSON_STR); + EXPECT_STREQ(vts[0].ud, "value"); + EXPECT_EQ(vts[1].jtype, NEU_JSON_INT); + EXPECT_EQ(vts[1].jvalue.val_int, 1234); + EXPECT_EQ(vts[2].jtype, NEU_JSON_DOUBLE); + EXPECT_EQ(vts[2].jvalue.val_double, 11.44); + EXPECT_EQ(vts[3].jtype, NEU_JSON_BOOL); + EXPECT_TRUE(vts[3].jvalue.val_bool); + + ASSERT_EQ(vts[4].vt, MQTT_SCHEMA_OBJECT); + ASSERT_EQ(vts[4].n_sub_vts, 3); + EXPECT_EQ(vts[4].sub_vts[0].jvalue.val_int, -42); + EXPECT_EQ(vts[4].sub_vts[1].jvalue.val_double, 0.5); + EXPECT_FALSE(vts[4].sub_vts[2].jvalue.val_bool); + + char * result = NULL; + neu_json_read_resp_t tags = { 0 }; + ret = mqtt_schema_encode((char *) "driver", (char *) "group", &tags, vts, + n_vts, NULL, 0, &result); + ASSERT_EQ(ret, 0); + EXPECT_STREQ(result, + "{\"string\": \"value\", \"integer\": 1234, \"float\": " + "11.44, \"boolean\": true, \"object\": {\"integer\": -42, " + "\"float\": 0.5, \"boolean\": false}}"); + + free(result); + free(vts[4].sub_vts); + free(vts); +} + TEST(validate_failed, schema) { const char *schema = From 4e9df9cc9fd94d02a85661c4d28be3f045375ece Mon Sep 17 00:00:00 2001 From: fengzero Date: Wed, 12 Aug 2026 02:35:53 +0000 Subject: [PATCH 2/2] store single node --- src/core/manager.c | 12 ++++++++---- src/core/node_manager.c | 12 ++++++++++++ src/core/node_manager.h | 1 + 3 files changed, 21 insertions(+), 4 deletions(-) diff --git a/src/core/manager.c b/src/core/manager.c index 60d4ef33c..20c7cde76 100644 --- a/src/core/manager.c +++ b/src/core/manager.c @@ -152,6 +152,8 @@ neu_manager_t *neu_manager_create() nlog_warn("load plugin error"); } + manager_load_node(manager); + UT_array *single_plugins = neu_plugin_manager_get_single(manager->plugin_manager); @@ -162,7 +164,6 @@ neu_manager_t *neu_manager_create() } utarray_free(single_plugins); - manager_load_node(manager); while (neu_node_manager_exist_uninit(manager->node_manager)) { usleep(1000 * 100); } @@ -1866,6 +1867,11 @@ static void start_static_adapter(neu_manager_t *manager, const char *name) static void start_single_adapter(neu_manager_t *manager, const char *name, const char *plugin_name, bool display) { + if (neu_node_manager_is_exist(manager->node_manager, name)) { + nlog_warn("adapter %s already exist", name); + return; + } + neu_adapter_t * adapter = NULL; neu_plugin_instance_t instance = { 0 }; neu_adapter_info_t adapter_info = { @@ -1883,9 +1889,7 @@ static void start_single_adapter(neu_manager_t *manager, const char *name, adapter = neu_adapter_create(&adapter_info, true); neu_node_manager_add_single(manager->node_manager, adapter, display); - if (display) { - manager_storage_add_node(manager, name, ""); - } + manager_storage_add_node(manager, name, ""); neu_adapter_init(adapter, false); neu_adapter_start_single(adapter); diff --git a/src/core/node_manager.c b/src/core/node_manager.c index f3fd20747..4b6f63147 100644 --- a/src/core/node_manager.c +++ b/src/core/node_manager.c @@ -459,6 +459,18 @@ bool neu_node_manager_is_driver(neu_node_manager_t *mgr, const char *name) return false; } +bool neu_node_manager_is_exist(neu_node_manager_t *mgr, const char *name) +{ + node_entity_t *node = NULL; + + HASH_FIND_STR(mgr->nodes, name, node); + if (node != NULL) { + return true; + } + + return false; +} + UT_array *neu_node_manager_get_addrs(neu_node_manager_t *mgr, int type) { UT_icd icd = { sizeof(struct sockaddr_un), NULL, NULL, NULL }; diff --git a/src/core/node_manager.h b/src/core/node_manager.h index 7ac2a2854..087420d64 100644 --- a/src/core/node_manager.h +++ b/src/core/node_manager.h @@ -63,6 +63,7 @@ UT_array * neu_node_manager_get_adapter(neu_node_manager_t *mgr, int type); neu_adapter_t *neu_node_manager_find(neu_node_manager_t *mgr, const char *name); bool neu_node_manager_is_single(neu_node_manager_t *mgr, const char *name); bool neu_node_manager_is_driver(neu_node_manager_t *mgr, const char *name); +bool neu_node_manager_is_exist(neu_node_manager_t *mgr, const char *name); // addr array UT_array *neu_node_manager_get_addrs(neu_node_manager_t *mgr, int type);