Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 28 additions & 3 deletions plugins/mqtt/schema.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
3 changes: 3 additions & 0 deletions plugins/mqtt/schema.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
12 changes: 8 additions & 4 deletions src/core/manager.c
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand All @@ -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);
}
Expand Down Expand Up @@ -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 = {
Expand All @@ -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);
Expand Down
12 changes: 12 additions & 0 deletions src/core/node_manager.c
Original file line number Diff line number Diff line change
Expand Up @@ -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 };
Expand Down
1 change: 1 addition & 0 deletions src/core/node_manager.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
42 changes: 42 additions & 0 deletions tests/ut/mqtt_schema_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down
Loading