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
15 changes: 13 additions & 2 deletions include/neuron/msg.h
Original file line number Diff line number Diff line change
Expand Up @@ -326,11 +326,21 @@ typedef struct neu_req_get_node {
neu_node_type_e type;
char plugin[NEU_PLUGIN_NAME_LEN];
char node[NEU_NODE_NAME_LEN];

struct {
bool q_state;
int state;
bool q_link;
int link;
bool s_delay;
char q_group_name[NEU_GROUP_NAME_LEN];
} query;
} neu_req_get_node_t;

typedef struct neu_resp_node_info {
char node[NEU_NODE_NAME_LEN];
char plugin[NEU_PLUGIN_NAME_LEN];
int64_t delay;
char node[NEU_NODE_NAME_LEN];
char plugin[NEU_PLUGIN_NAME_LEN];
} neu_resp_node_info_t;

typedef struct neu_resp_get_node {
Expand Down Expand Up @@ -368,6 +378,7 @@ typedef struct neu_req_del_group {

typedef struct neu_req_get_group {
char driver[NEU_NODE_NAME_LEN];
char q_group[NEU_GROUP_NAME_LEN];
} neu_req_get_group_t;

typedef struct neu_resp_group_info {
Expand Down
28 changes: 28 additions & 0 deletions plugins/restful/adapter_handle.c
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,7 @@ void handle_get_adapter(nng_aio *aio)
int ret = 0;
char plugin_name[NEU_PLUGIN_NAME_LEN] = { 0 };
char node_name[NEU_NODE_NAME_LEN] = { 0 };
char group_name[NEU_GROUP_NAME_LEN] = { 0 };

neu_plugin_t * plugin = neu_rest_get_plugin();
neu_node_type_e node_type = { 0 };
Expand Down Expand Up @@ -173,6 +174,33 @@ void handle_get_adapter(nng_aio *aio)
strcpy(cmd.node, node_name);
}

if (neu_http_get_param_str(aio, "group", group_name, sizeof(group_name)) >
0) {
strcpy(cmd.query.q_group_name, group_name);
}

uintmax_t query_state = 0;
if (neu_http_get_param_uintmax(aio, "state", &query_state) == 0) {
if (query_state == 1 || query_state == 2 || query_state == 3 ||
query_state == 4) {
cmd.query.q_state = true;
cmd.query.state = (int) query_state;
}
}

uintmax_t query_link = 0;
if (neu_http_get_param_uintmax(aio, "link", &query_link) == 0) {
if (query_link == 0 || query_link == 1) {
cmd.query.q_link = true;
cmd.query.link = (int) query_link;
}
}

uintmax_t sort_delay = 0;
if (neu_http_get_param_uintmax(aio, "delay", &sort_delay) == 0) {
cmd.query.s_delay = true;
}

cmd.type = node_type;
ret = neu_plugin_op(plugin, header, &cmd);
if (ret != 0) {
Expand Down
16 changes: 11 additions & 5 deletions plugins/restful/group_config_handle.c
Original file line number Diff line number Diff line change
Expand Up @@ -160,11 +160,12 @@ void handle_del_group_config(nng_aio *aio)

void handle_get_group_config(nng_aio *aio)
{
neu_plugin_t * plugin = neu_rest_get_plugin();
char node_name[NEU_NODE_NAME_LEN] = { 0 };
int ret = 0;
neu_req_get_group_t cmd = { 0 };
neu_reqresp_head_t header = {
neu_plugin_t * plugin = neu_rest_get_plugin();
char node_name[NEU_NODE_NAME_LEN] = { 0 };
char group_name[NEU_GROUP_NAME_LEN] = { 0 };
int ret = 0;
neu_req_get_group_t cmd = { 0 };
neu_reqresp_head_t header = {
.ctx = aio,
.type = NEU_REQ_GET_GROUP,
.otel_trace_type = NEU_OTEL_TRACE_TYPE_REST_COMM,
Expand All @@ -178,6 +179,11 @@ void handle_get_group_config(nng_aio *aio)
header.type = NEU_REQ_GET_DRIVER_GROUP;
}

if (neu_http_get_param_str(aio, "group", group_name, sizeof(group_name)) >
0) {
strcpy(cmd.q_group, group_name);
}

ret = neu_plugin_op(plugin, header, &cmd);
if (ret != 0) {
NEU_JSON_RESPONSE_ERROR(NEU_ERR_IS_BUSY, {
Expand Down
15 changes: 14 additions & 1 deletion src/adapter/adapter.c
Original file line number Diff line number Diff line change
Expand Up @@ -1033,12 +1033,15 @@ static int adapter_loop(enum neu_event_io_type type, int fd, void *usr_data)
break;
}
case NEU_REQ_GET_GROUP: {

neu_req_get_group_t *cmd = (neu_req_get_group_t *) &header[1];

neu_msg_exchange(header);

if (adapter->module->type == NEU_NA_TYPE_DRIVER) {
neu_resp_get_group_t resp = {
.groups = neu_adapter_driver_get_group(
(neu_adapter_driver_t *) adapter)
(neu_adapter_driver_t *) adapter, cmd->q_group)
};
header->type = NEU_RESP_GET_GROUP;
reply(adapter, header, &resp);
Expand Down Expand Up @@ -1736,6 +1739,16 @@ neu_node_state_t neu_adapter_get_state(neu_adapter_t *adapter)
return state;
}

UT_array *neu_adapter_get_groups(neu_adapter_t *adapter, const char *filter)
{
if (adapter->module->type != NEU_NA_TYPE_DRIVER) {
return NULL;
} else {
return neu_adapter_driver_get_group((neu_adapter_driver_t *) adapter,
filter);
}
}

neu_event_timer_t *neu_adapter_add_timer(neu_adapter_t * adapter,
neu_event_timer_param_t param)
{
Expand Down
1 change: 1 addition & 0 deletions src/adapter/adapter_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ void neu_adapter_del_timer(neu_adapter_t *adapter, neu_event_timer_t *timer);
int neu_adapter_set_setting(neu_adapter_t *adapter, const char *config);
int neu_adapter_get_setting(neu_adapter_t *adapter, char **config);
neu_node_state_t neu_adapter_get_state(neu_adapter_t *adapter);
UT_array *neu_adapter_get_groups(neu_adapter_t *adapter, const char *filter);

static inline void neu_adapter_reset_metrics(neu_adapter_t *adapter)
{
Expand Down
16 changes: 10 additions & 6 deletions src/adapter/driver/driver.c
Original file line number Diff line number Diff line change
Expand Up @@ -1757,7 +1757,8 @@ int neu_adapter_driver_group_exist(neu_adapter_driver_t *driver,
return ret;
}

UT_array *neu_adapter_driver_get_group(neu_adapter_driver_t *driver)
UT_array *neu_adapter_driver_get_group(neu_adapter_driver_t *driver,
const char * filter)
{
group_t * el = NULL, *tmp = NULL;
UT_array *groups = NULL;
Expand All @@ -1767,13 +1768,16 @@ UT_array *neu_adapter_driver_get_group(neu_adapter_driver_t *driver)

HASH_ITER(hh, driver->groups, el, tmp)
{
neu_resp_group_info_t info = { 0 };

info.interval = neu_group_get_interval(el->group);
info.tag_count = neu_group_tag_size(el->group);
strncpy(info.name, el->name, sizeof(info.name));
if (strlen(filter) == 0 || strstr(el->name, filter) != NULL) {
neu_resp_group_info_t info = { 0 };

utarray_push_back(groups, &info);
info.interval = neu_group_get_interval(el->group);
info.tag_count = neu_group_tag_size(el->group);
strncpy(info.name, el->name, sizeof(info.name));

utarray_push_back(groups, &info);
}
}

return groups;
Expand Down
3 changes: 2 additions & 1 deletion src/adapter/driver/driver_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,8 @@ int neu_adapter_driver_del_group(neu_adapter_driver_t *driver,
const char * name);
int neu_adapter_driver_group_exist(neu_adapter_driver_t *driver,
const char * name);
UT_array *neu_adapter_driver_get_group(neu_adapter_driver_t *driver);
UT_array *neu_adapter_driver_get_group(neu_adapter_driver_t *driver,
const char * filter);
uint16_t neu_adapter_driver_group_count(neu_adapter_driver_t *driver);
uint16_t neu_adapter_driver_new_group_count(neu_adapter_driver_t *driver,
neu_req_add_gtag_t * cmd);
Expand Down
14 changes: 9 additions & 5 deletions src/core/manager.c
Original file line number Diff line number Diff line change
Expand Up @@ -435,7 +435,8 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data)
neu_resp_error_t e = { 0 };

UT_array *nodes = neu_manager_get_nodes(
manager, NEU_NA_TYPE_DRIVER | NEU_NA_TYPE_APP, cmd->plugin, "");
manager, NEU_NA_TYPE_DRIVER | NEU_NA_TYPE_APP, cmd->plugin, "",
false, false, 0, false, 0, "");

if (nodes != NULL) {
if (utarray_len(nodes) > 0) {
Expand Down Expand Up @@ -616,7 +617,8 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data)
}

UT_array *nodes = neu_manager_get_nodes(
manager, NEU_NA_TYPE_DRIVER | NEU_NA_TYPE_APP, module_name, "");
manager, NEU_NA_TYPE_DRIVER | NEU_NA_TYPE_APP, module_name, "",
false, false, 0, false, 0, "");

if (nodes != NULL) {
if (utarray_len(nodes) > 0) {
Expand Down Expand Up @@ -866,9 +868,11 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data)
break;
}
case NEU_REQ_GET_NODE: {
neu_req_get_node_t *cmd = (neu_req_get_node_t *) &header[1];
UT_array * nodes =
neu_manager_get_nodes(manager, cmd->type, cmd->plugin, cmd->node);
neu_req_get_node_t *cmd = (neu_req_get_node_t *) &header[1];
UT_array * nodes = neu_manager_get_nodes(
manager, cmd->type, cmd->plugin, cmd->node, cmd->query.s_delay,
cmd->query.q_state, cmd->query.state, cmd->query.q_link,
cmd->query.link, cmd->query.q_group_name);
neu_resp_get_node_t resp = { .nodes = nodes };

header->type = NEU_RESP_GET_NODE;
Expand Down
10 changes: 7 additions & 3 deletions src/core/manager_internal.c
Original file line number Diff line number Diff line change
Expand Up @@ -114,9 +114,13 @@ int neu_manager_del_node(neu_manager_t *manager, const char *node_name)
}

UT_array *neu_manager_get_nodes(neu_manager_t *manager, int type,
const char *plugin, const char *node)
const char *plugin, const char *node,
bool sort_delay, bool q_state, int state,
bool q_link, int link, const char *q_group_name)
{
return neu_node_manager_filter(manager->node_manager, type, plugin, node);
return neu_node_manager_filter(manager->node_manager, type, plugin, node,
sort_delay, q_state, state, q_link, link,
q_group_name);
}

int neu_manager_update_node_name(neu_manager_t *manager, const char *node,
Expand Down Expand Up @@ -182,7 +186,7 @@ UT_array *neu_manager_get_driver_group(neu_manager_t *manager)
neu_adapter_t *adapter =
neu_node_manager_find(manager->node_manager, driver->node);
UT_array *groups =
neu_adapter_driver_get_group((neu_adapter_driver_t *) adapter);
neu_adapter_driver_get_group((neu_adapter_driver_t *) adapter, "");

utarray_foreach(groups, neu_resp_group_info_t *, g)
{
Expand Down
5 changes: 4 additions & 1 deletion src/core/manager_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,10 @@ int neu_manager_add_node(neu_manager_t *manager, const char *node_name,
neu_node_running_state_e state, bool load);
int neu_manager_del_node(neu_manager_t *manager, const char *node_name);
UT_array *neu_manager_get_nodes(neu_manager_t *manager, int type,
const char *plugin, const char *node);
const char *plugin, const char *node,
bool sort_delay, bool q_state, int state,
bool q_link, int link,
const char *q_group_name);
int neu_manager_update_node_name(neu_manager_t *manager, const char *node,
const char *new_name);
int neu_manager_update_group_name(neu_manager_t *manager, const char *driver,
Expand Down
61 changes: 60 additions & 1 deletion src/core/node_manager.c
Original file line number Diff line number Diff line change
Expand Up @@ -194,8 +194,24 @@ UT_array *neu_node_manager_get(neu_node_manager_t *mgr, int type)
return array;
}

static int neu_resp_node_info_sort(const void *a, const void *b)
{
const neu_resp_node_info_t *info_a = (const neu_resp_node_info_t *) a;
const neu_resp_node_info_t *info_b = (const neu_resp_node_info_t *) b;

if (info_a->delay < info_b->delay) {
return -1;
} else if (info_a->delay > info_b->delay) {
return 1;
}
return 0;
}

UT_array *neu_node_manager_filter(neu_node_manager_t *mgr, int type,
const char *plugin, const char *node)
const char *plugin, const char *node,
bool sort_delay, bool q_state, int state,
bool q_link, int link,
const char *q_group_name)
{
UT_array * array = NULL;
UT_icd icd = { sizeof(neu_resp_node_info_t), NULL, NULL, NULL };
Expand All @@ -215,14 +231,57 @@ UT_array *neu_node_manager_filter(neu_node_manager_t *mgr, int type,
strstr(el->adapter->name, node) == NULL) {
continue;
}
if (q_state) {
if (el->adapter->state !=
(neu_node_running_state_e) state) {
continue;
}
}
if (q_link) {
neu_node_state_t node_state =
neu_adapter_get_state(el->adapter);
if (node_state.link != (neu_node_link_state_e) link) {
continue;
}
}

if (strlen(q_group_name) > 0) {
if (strstr(el->adapter->name, q_group_name) == NULL) {
UT_array *groups =
neu_adapter_get_groups(el->adapter, q_group_name);
if (NULL == groups || utarray_len(groups) == 0) {
if (groups != NULL) {
utarray_free(groups);
}
continue;
} else {
utarray_free(groups);
}
}
}

neu_resp_node_info_t info = { 0 };
strcpy(info.node, el->adapter->name);
strcpy(info.plugin, el->adapter->module->module_name);
if (sort_delay) {
neu_metric_entry_t *e = NULL;
if (NULL != el->adapter->metrics) {
HASH_FIND_STR(el->adapter->metrics->entries,
NEU_METRIC_LAST_RTT_MS, e);
}
info.delay = NULL != e ? e->value : 0;
} else {
info.delay = 0;
}
utarray_push_back(array, &info);
}
}
}

if (sort_delay) {
utarray_sort(array, neu_resp_node_info_sort);
}

return array;
}

Expand Down
5 changes: 4 additions & 1 deletion src/core/node_manager.h
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,10 @@ uint16_t neu_node_manager_size(neu_node_manager_t *mgr);
// neu_resp_node_info array
UT_array *neu_node_manager_get(neu_node_manager_t *mgr, int type);
UT_array *neu_node_manager_filter(neu_node_manager_t *mgr, int type,
const char *plugin, const char *node);
const char *plugin, const char *node,
bool sort_delay, bool q_state, int state,
bool q_link, int link,
const char *q_group_name);
UT_array *neu_node_manager_get_all(neu_node_manager_t *mgr);

// neu_adapter_t array
Expand Down
Loading