From a107c0aebb21536863edc719948203c359c4ffc6 Mon Sep 17 00:00:00 2001 From: fengzero Date: Wed, 20 Aug 2025 06:11:30 +0000 Subject: [PATCH] app: get all groups --- include/neuron/msg.h | 36 +++++-- plugins/restful/group_config_handle.c | 107 ++++++++++++++++++ plugins/restful/group_config_handle.h | 4 + plugins/restful/handle.c | 6 ++ plugins/restful/rest.c | 7 ++ src/adapter/adapter.c | 1 + src/base/msg_internal.h | 149 +++++++++++++------------- src/core/manager.c | 12 +++ src/core/manager_internal.c | 49 +++++++++ src/core/manager_internal.h | 2 + src/parser/neu_json_group_config.c | 44 ++++++++ src/parser/neu_json_group_config.h | 18 ++++ 12 files changed, 353 insertions(+), 82 deletions(-) diff --git a/include/neuron/msg.h b/include/neuron/msg.h index 5c588c6e9..7f6652ec7 100644 --- a/include/neuron/msg.h +++ b/include/neuron/msg.h @@ -60,6 +60,8 @@ typedef enum neu_reqresp_type { NEU_REQ_SUBSCRIBE_GROUPS, NEU_REQ_GET_SUBSCRIBE_GROUP, NEU_RESP_GET_SUBSCRIBE_GROUP, + NEU_REQ_GET_DRIVER_SUBSCRIBE_GROUP, + NEU_RESP_GET_DRIVER_SUBSCRIBE_GROUP, NEU_REQ_GET_SUB_DRIVER_TAGS, NEU_RESP_GET_SUB_DRIVER_TAGS, @@ -143,14 +145,17 @@ static const char *neu_reqresp_type_string_t[] = { [NEU_REQ_WRITE_TAGS] = "NEU_REQ_WRITE_TAGS", [NEU_REQ_WRITE_GTAGS] = "NEU_REQ_WRITE_GTAGS", - [NEU_REQ_SUBSCRIBE_GROUP] = "NEU_REQ_SUBSCRIBE_GROUP", - [NEU_REQ_UNSUBSCRIBE_GROUP] = "NEU_REQ_UNSUBSCRIBE_GROUP", - [NEU_REQ_UPDATE_SUBSCRIBE_GROUP] = "NEU_REQ_UPDATE_SUBSCRIBE_GROUP", - [NEU_REQ_SUBSCRIBE_GROUPS] = "NEU_REQ_SUBSCRIBE_GROUPS", - [NEU_REQ_GET_SUBSCRIBE_GROUP] = "NEU_REQ_GET_SUBSCRIBE_GROUP", - [NEU_RESP_GET_SUBSCRIBE_GROUP] = "NEU_RESP_GET_SUBSCRIBE_GROUP", - [NEU_REQ_GET_SUB_DRIVER_TAGS] = "NEU_REQ_GET_SUB_DRIVER_TAGS", - [NEU_RESP_GET_SUB_DRIVER_TAGS] = "NEU_RESP_GET_SUB_DRIVER_TAGS", + [NEU_REQ_SUBSCRIBE_GROUP] = "NEU_REQ_SUBSCRIBE_GROUP", + [NEU_REQ_UNSUBSCRIBE_GROUP] = "NEU_REQ_UNSUBSCRIBE_GROUP", + [NEU_REQ_UPDATE_SUBSCRIBE_GROUP] = "NEU_REQ_UPDATE_SUBSCRIBE_GROUP", + [NEU_REQ_SUBSCRIBE_GROUPS] = "NEU_REQ_SUBSCRIBE_GROUPS", + [NEU_REQ_GET_SUBSCRIBE_GROUP] = "NEU_REQ_GET_SUBSCRIBE_GROUP", + [NEU_RESP_GET_SUBSCRIBE_GROUP] = "NEU_RESP_GET_SUBSCRIBE_GROUP", + [NEU_REQ_GET_DRIVER_SUBSCRIBE_GROUP] = "NEU_REQ_GET_DRIVER_SUBSCRIBE_GROUP", + [NEU_RESP_GET_DRIVER_SUBSCRIBE_GROUP] = + "NEU_RESP_GET_DRIVER_SUBSCRIBE_GROUP", + [NEU_REQ_GET_SUB_DRIVER_TAGS] = "NEU_REQ_GET_SUB_DRIVER_TAGS", + [NEU_RESP_GET_SUB_DRIVER_TAGS] = "NEU_RESP_GET_SUB_DRIVER_TAGS", [NEU_REQ_NODE_INIT] = "NEU_REQ_NODE_INIT", [NEU_REQ_NODE_UNINIT] = "NEU_REQ_NODE_UNINIT", @@ -547,6 +552,11 @@ typedef struct neu_req_get_subscribe_group { char group[NEU_GROUP_NAME_LEN]; } neu_req_get_subscribe_group_t; +typedef struct neu_req_get_driver_subscribe_group { + char app[NEU_NODE_NAME_LEN]; + char name[NEU_NODE_NAME_LEN]; +} neu_req_get_driver_subscribe_group_t; + typedef struct { char app[NEU_NODE_NAME_LEN]; } neu_req_get_sub_driver_tags_t; @@ -559,6 +569,13 @@ typedef struct neu_resp_subscribe_info { char *static_tags; } neu_resp_subscribe_info_t; +typedef struct neu_resp_driver_subscribe_info { + char driver[NEU_NODE_NAME_LEN]; + char app[NEU_NODE_NAME_LEN]; + char group[NEU_GROUP_NAME_LEN]; + bool subscribed; +} neu_resp_driver_subscribe_info_t; + static inline void neu_resp_subscribe_info_fini(neu_resp_subscribe_info_t *info) { free(info->params); @@ -582,7 +599,8 @@ typedef struct { } neu_resp_get_sub_driver_tags_t; typedef struct neu_resp_get_subscribe_group { - UT_array *groups; // array neu_resp_subscribe_info_t + UT_array *groups; // array neu_resp_subscribe_info_t or + // neu_resp_driver_subscribe_info_t } neu_resp_get_subscribe_group_t; typedef struct neu_req_node_setting { diff --git a/plugins/restful/group_config_handle.c b/plugins/restful/group_config_handle.c index 3fdbd91d2..3f185fa99 100644 --- a/plugins/restful/group_config_handle.c +++ b/plugins/restful/group_config_handle.c @@ -488,4 +488,111 @@ void handle_grp_get_subscribe_resp(nng_aio * aio, free(result); free(sub_grp_configs.groups); utarray_free(groups->groups); +} + +void handle_grp_get_subscribes(nng_aio *aio) +{ + int ret = 0; + neu_plugin_t * plugin = neu_rest_get_plugin(); + neu_req_get_driver_subscribe_group_t cmd = { 0 }; + neu_reqresp_head_t header = { + .ctx = aio, + .type = NEU_REQ_GET_DRIVER_SUBSCRIBE_GROUP, + .otel_trace_type = NEU_OTEL_TRACE_TYPE_REST_COMM, + }; + + NEU_VALIDATE_JWT(aio); + + // required parameter + ret = neu_http_get_param_str(aio, "app", cmd.app, sizeof(cmd.app)); + if (ret <= 0) { + NEU_JSON_RESPONSE_ERROR(NEU_ERR_PARAM_IS_WRONG, { + neu_http_response(aio, error_code.error, result_error); + }) + return; + } + + // optional parameter + ret = neu_http_get_param_str(aio, "name", cmd.name, sizeof(cmd.name)); + if (-1 == ret || (size_t) ret == sizeof(cmd.name)) { + NEU_JSON_RESPONSE_ERROR(NEU_ERR_PARAM_IS_WRONG, { + neu_http_response(aio, error_code.error, result_error); + }) + return; + } + + ret = neu_plugin_op(plugin, header, &cmd); + if (ret != 0) { + NEU_JSON_RESPONSE_ERROR(NEU_ERR_IS_BUSY, { + neu_http_response(aio, NEU_ERR_IS_BUSY, result_error); + }); + } +} + +void handle_grp_get_subscribes_resp(nng_aio * aio, + neu_resp_get_subscribe_group_t *drivers) +{ + char * result = NULL; + neu_json_get_driver_subscribe_resp_t drivers_resp = { 0 }; + + utarray_foreach(drivers->groups, neu_resp_driver_subscribe_info_t *, group) + { + // int index = utarray_eltidx(drivers->groups, group); + bool found = false; + + for (int k = 0; k < drivers_resp.n_driver; k++) { + if (strcmp(drivers_resp.drivers[k].driver, group->driver) == 0) { + found = true; + drivers_resp.drivers[k].n_group++; + drivers_resp.drivers[k].groups = realloc( + drivers_resp.drivers[k].groups, + drivers_resp.drivers[k].n_group * + sizeof(neu_json_get_driver_subscribe_resp_group_t)); + drivers_resp.drivers[k] + .groups[drivers_resp.drivers[k].n_group - 1] + .group = strdup(group->group); + drivers_resp.drivers[k] + .groups[drivers_resp.drivers[k].n_group - 1] + .subscribed = group->subscribed; + break; + } + } + + if (!found) { + drivers_resp.n_driver++; + drivers_resp.drivers = realloc( + drivers_resp.drivers, + drivers_resp.n_driver * + sizeof(neu_json_get_driver_subscribe_resp_driver_t)); + + drivers_resp.drivers[drivers_resp.n_driver - 1].driver = + strdup(group->driver); + drivers_resp.drivers[drivers_resp.n_driver - 1].n_group = 1; + drivers_resp.drivers[drivers_resp.n_driver - 1].groups = + calloc(1, sizeof(neu_json_get_driver_subscribe_resp_group_t)); + + drivers_resp.drivers[drivers_resp.n_driver - 1].groups[0].group = + strdup(group->group); + drivers_resp.drivers[drivers_resp.n_driver - 1] + .groups[0] + .subscribed = group->subscribed; + } + } + + neu_json_encode_by_fn(&drivers_resp, + neu_json_encode_get_driver_subscribe_resp, &result); + + neu_http_ok(aio, result); + free(result); + + for (int i = 0; i < drivers_resp.n_driver; i++) { + free(drivers_resp.drivers[i].driver); + for (int j = 0; j < drivers_resp.drivers[i].n_group; j++) { + free(drivers_resp.drivers[i].groups[j].group); + } + free(drivers_resp.drivers[i].groups); + } + free(drivers_resp.drivers); + + utarray_free(drivers->groups); } \ No newline at end of file diff --git a/plugins/restful/group_config_handle.h b/plugins/restful/group_config_handle.h index 040aba1f6..1ed401857 100644 --- a/plugins/restful/group_config_handle.h +++ b/plugins/restful/group_config_handle.h @@ -39,4 +39,8 @@ void handle_grp_get_subscribe(nng_aio *aio); void handle_grp_get_subscribe_resp(nng_aio * aio, neu_resp_get_subscribe_group_t *groups); +void handle_grp_get_subscribes(nng_aio *aio); +void handle_grp_get_subscribes_resp(nng_aio * aio, + neu_resp_get_subscribe_group_t *groups); + #endif \ No newline at end of file diff --git a/plugins/restful/handle.c b/plugins/restful/handle.c index 7ff39ed99..861286eec 100644 --- a/plugins/restful/handle.c +++ b/plugins/restful/handle.c @@ -350,6 +350,12 @@ static struct neu_http_handler rest_handlers[] = { .url = "/api/v2/subscribe", .value.handler = handle_grp_get_subscribe, }, + { + .method = NEU_HTTP_METHOD_GET, + .type = NEU_HTTP_HANDLER_FUNCTION, + .url = "/api/v2/subscribes", + .value.handler = handle_grp_get_subscribes, + }, { .method = NEU_HTTP_METHOD_POST, .type = NEU_HTTP_HANDLER_FUNCTION, diff --git a/plugins/restful/rest.c b/plugins/restful/rest.c index 41f52cb9a..44121d6a8 100644 --- a/plugins/restful/rest.c +++ b/plugins/restful/rest.c @@ -279,6 +279,13 @@ static int dashb_plugin_request(neu_plugin_t * plugin, neu_otel_scope_set_status_code2(scope, NEU_OTEL_STATUS_OK, 0); } break; + case NEU_RESP_GET_DRIVER_SUBSCRIBE_GROUP: + handle_grp_get_subscribes_resp(header->ctx, + (neu_resp_get_subscribe_group_t *) data); + if (neu_otel_control_is_started() && trace) { + neu_otel_scope_set_status_code2(scope, NEU_OTEL_STATUS_OK, 0); + } + break; case NEU_RESP_GET_NODE_SETTING: handle_get_node_setting_resp(header->ctx, (neu_resp_get_node_setting_t *) data); diff --git a/src/adapter/adapter.c b/src/adapter/adapter.c index aba013b66..8c79d4a79 100644 --- a/src/adapter/adapter.c +++ b/src/adapter/adapter.c @@ -727,6 +727,7 @@ static int adapter_loop(enum neu_event_io_type type, int fd, void *usr_data) case NEU_RESP_GET_NODE_SETTING: case NEU_REQ_UPDATE_GROUP: case NEU_RESP_GET_SUBSCRIBE_GROUP: + case NEU_RESP_GET_DRIVER_SUBSCRIBE_GROUP: case NEU_RESP_ADD_TAG: case NEU_RESP_ADD_GTAG: case NEU_RESP_UPDATE_TAG: diff --git a/src/base/msg_internal.h b/src/base/msg_internal.h index 2d403250b..c8faaba92 100644 --- a/src/base/msg_internal.h +++ b/src/base/msg_internal.h @@ -30,79 +30,82 @@ extern "C" { #include "msg.h" -#define NEU_REQRESP_TYPE_MAP(XX) \ - XX(NEU_RESP_ERROR, neu_resp_error_t) \ - XX(NEU_REQ_READ_GROUP, neu_req_read_group_t) \ - XX(NEU_RESP_READ_GROUP, neu_resp_read_group_t) \ - XX(NEU_REQ_READ_GROUP_PAGINATE, neu_req_read_group_paginate_t) \ - XX(NEU_RESP_READ_GROUP_PAGINATE, neu_resp_read_group_paginate_t) \ - XX(NEU_REQ_TEST_READ_TAG, neu_req_test_read_tag_t) \ - XX(NEU_RESP_TEST_READ_TAG, neu_resp_test_read_tag_t) \ - XX(NEU_REQ_WRITE_TAG, neu_req_write_tag_t) \ - XX(NEU_REQ_WRITE_TAGS, neu_req_write_tags_t) \ - XX(NEU_REQ_WRITE_GTAGS, neu_req_write_gtags_t) \ - XX(NEU_REQ_SUBSCRIBE_GROUP, neu_req_subscribe_t) \ - XX(NEU_REQ_UNSUBSCRIBE_GROUP, neu_req_unsubscribe_t) \ - XX(NEU_REQ_UPDATE_SUBSCRIBE_GROUP, neu_req_subscribe_t) \ - XX(NEU_REQ_SUBSCRIBE_GROUPS, neu_req_subscribe_groups_t) \ - XX(NEU_REQ_GET_SUBSCRIBE_GROUP, neu_req_get_subscribe_group_t) \ - XX(NEU_RESP_GET_SUBSCRIBE_GROUP, neu_resp_get_subscribe_group_t) \ - XX(NEU_REQ_GET_SUB_DRIVER_TAGS, neu_req_get_sub_driver_tags_t) \ - XX(NEU_RESP_GET_SUB_DRIVER_TAGS, neu_resp_get_sub_driver_tags_t) \ - XX(NEU_REQ_NODE_INIT, neu_req_node_init_t) \ - XX(NEU_REQ_NODE_UNINIT, neu_req_node_uninit_t) \ - XX(NEU_RESP_NODE_UNINIT, neu_resp_node_uninit_t) \ - XX(NEU_REQ_ADD_NODE, neu_req_add_node_t) \ - XX(NEU_REQ_UPDATE_NODE, neu_req_update_node_t) \ - XX(NEU_REQ_DEL_NODE, neu_req_del_node_t) \ - XX(NEU_REQ_GET_NODE, neu_req_get_node_t) \ - XX(NEU_RESP_GET_NODE, neu_resp_get_node_t) \ - XX(NEU_REQ_NODE_SETTING, neu_req_node_setting_t) \ - XX(NEU_REQ_GET_NODE_SETTING, neu_req_get_node_setting_t) \ - XX(NEU_RESP_GET_NODE_SETTING, neu_resp_get_node_setting_t) \ - XX(NEU_REQ_GET_NODE_STATE, neu_req_get_node_state_t) \ - XX(NEU_RESP_GET_NODE_STATE, neu_resp_get_node_state_t) \ - XX(NEU_REQ_GET_NODES_STATE, neu_req_get_nodes_state_t) \ - XX(NEU_RESP_GET_NODES_STATE, neu_resp_get_nodes_state_t) \ - XX(NEU_REQ_NODE_CTL, neu_req_node_ctl_t) \ - XX(NEU_REQ_NODE_RENAME, neu_req_node_rename_t) \ - XX(NEU_RESP_NODE_RENAME, neu_resp_node_rename_t) \ - XX(NEU_REQ_ADD_GROUP, neu_req_add_group_t) \ - XX(NEU_REQ_DEL_GROUP, neu_req_del_group_t) \ - XX(NEU_REQ_UPDATE_GROUP, neu_req_update_group_t) \ - XX(NEU_REQ_UPDATE_DRIVER_GROUP, neu_req_update_group_t) \ - XX(NEU_RESP_UPDATE_DRIVER_GROUP, neu_resp_update_group_t) \ - XX(NEU_REQ_GET_GROUP, neu_req_get_group_t) \ - XX(NEU_RESP_GET_GROUP, neu_resp_get_group_t) \ - XX(NEU_REQ_GET_DRIVER_GROUP, neu_req_get_group_t) \ - XX(NEU_RESP_GET_DRIVER_GROUP, neu_resp_get_driver_group_t) \ - XX(NEU_REQ_ADD_TAG, neu_req_add_tag_t) \ - XX(NEU_RESP_ADD_TAG, neu_resp_add_tag_t) \ - XX(NEU_REQ_ADD_GTAG, neu_req_add_gtag_t) \ - XX(NEU_RESP_ADD_GTAG, neu_resp_add_tag_t) \ - XX(NEU_REQ_DEL_TAG, neu_req_del_tag_t) \ - XX(NEU_REQ_UPDATE_TAG, neu_req_update_tag_t) \ - XX(NEU_RESP_UPDATE_TAG, neu_resp_update_tag_t) \ - XX(NEU_REQ_GET_TAG, neu_req_get_tag_t) \ - XX(NEU_RESP_GET_TAG, neu_resp_get_tag_t) \ - XX(NEU_REQ_ADD_PLUGIN, neu_req_add_plugin_t) \ - XX(NEU_REQ_DEL_PLUGIN, neu_req_del_plugin_t) \ - XX(NEU_REQ_UPDATE_PLUGIN, neu_req_update_plugin_t) \ - XX(NEU_REQ_GET_PLUGIN, neu_req_get_plugin_t) \ - XX(NEU_RESP_GET_PLUGIN, neu_resp_get_plugin_t) \ - XX(NEU_REQRESP_TRANS_DATA, neu_reqresp_trans_data_t) \ - XX(NEU_REQRESP_NODES_STATE, neu_reqresp_nodes_state_t) \ - XX(NEU_REQRESP_NODE_DELETED, neu_reqresp_node_deleted_t) \ - XX(NEU_REQ_ADD_DRIVERS, neu_req_driver_array_t) \ - XX(NEU_REQ_UPDATE_LOG_LEVEL, neu_req_update_log_level_t) \ - XX(NEU_REQ_PRGFILE_UPLOAD, neu_req_prgfile_upload_t) \ - XX(NEU_REQ_PRGFILE_PROCESS, neu_req_prgfile_process_t) \ - XX(NEU_RESP_PRGFILE_PROCESS, neu_resp_prgfile_process_t) \ - XX(NEU_REQ_SCAN_TAGS, neu_req_scan_tags_t) \ - XX(NEU_RESP_SCAN_TAGS, neu_resp_scan_tags_t) \ - XX(NEU_REQ_CHECK_SCHEMA, neu_req_check_schema_t) \ - XX(NEU_RESP_CHECK_SCHEMA, neu_resp_check_schema_t) \ - XX(NEU_REQ_DRIVER_ACTION, neu_req_driver_action_t) \ +#define NEU_REQRESP_TYPE_MAP(XX) \ + XX(NEU_RESP_ERROR, neu_resp_error_t) \ + XX(NEU_REQ_READ_GROUP, neu_req_read_group_t) \ + XX(NEU_RESP_READ_GROUP, neu_resp_read_group_t) \ + XX(NEU_REQ_READ_GROUP_PAGINATE, neu_req_read_group_paginate_t) \ + XX(NEU_RESP_READ_GROUP_PAGINATE, neu_resp_read_group_paginate_t) \ + XX(NEU_REQ_TEST_READ_TAG, neu_req_test_read_tag_t) \ + XX(NEU_RESP_TEST_READ_TAG, neu_resp_test_read_tag_t) \ + XX(NEU_REQ_WRITE_TAG, neu_req_write_tag_t) \ + XX(NEU_REQ_WRITE_TAGS, neu_req_write_tags_t) \ + XX(NEU_REQ_WRITE_GTAGS, neu_req_write_gtags_t) \ + XX(NEU_REQ_SUBSCRIBE_GROUP, neu_req_subscribe_t) \ + XX(NEU_REQ_UNSUBSCRIBE_GROUP, neu_req_unsubscribe_t) \ + XX(NEU_REQ_UPDATE_SUBSCRIBE_GROUP, neu_req_subscribe_t) \ + XX(NEU_REQ_SUBSCRIBE_GROUPS, neu_req_subscribe_groups_t) \ + XX(NEU_REQ_GET_SUBSCRIBE_GROUP, neu_req_get_subscribe_group_t) \ + XX(NEU_RESP_GET_SUBSCRIBE_GROUP, neu_resp_get_subscribe_group_t) \ + XX(NEU_REQ_GET_DRIVER_SUBSCRIBE_GROUP, \ + neu_req_get_driver_subscribe_group_t) \ + XX(NEU_RESP_GET_DRIVER_SUBSCRIBE_GROUP, neu_resp_get_subscribe_group_t) \ + XX(NEU_REQ_GET_SUB_DRIVER_TAGS, neu_req_get_sub_driver_tags_t) \ + XX(NEU_RESP_GET_SUB_DRIVER_TAGS, neu_resp_get_sub_driver_tags_t) \ + XX(NEU_REQ_NODE_INIT, neu_req_node_init_t) \ + XX(NEU_REQ_NODE_UNINIT, neu_req_node_uninit_t) \ + XX(NEU_RESP_NODE_UNINIT, neu_resp_node_uninit_t) \ + XX(NEU_REQ_ADD_NODE, neu_req_add_node_t) \ + XX(NEU_REQ_UPDATE_NODE, neu_req_update_node_t) \ + XX(NEU_REQ_DEL_NODE, neu_req_del_node_t) \ + XX(NEU_REQ_GET_NODE, neu_req_get_node_t) \ + XX(NEU_RESP_GET_NODE, neu_resp_get_node_t) \ + XX(NEU_REQ_NODE_SETTING, neu_req_node_setting_t) \ + XX(NEU_REQ_GET_NODE_SETTING, neu_req_get_node_setting_t) \ + XX(NEU_RESP_GET_NODE_SETTING, neu_resp_get_node_setting_t) \ + XX(NEU_REQ_GET_NODE_STATE, neu_req_get_node_state_t) \ + XX(NEU_RESP_GET_NODE_STATE, neu_resp_get_node_state_t) \ + XX(NEU_REQ_GET_NODES_STATE, neu_req_get_nodes_state_t) \ + XX(NEU_RESP_GET_NODES_STATE, neu_resp_get_nodes_state_t) \ + XX(NEU_REQ_NODE_CTL, neu_req_node_ctl_t) \ + XX(NEU_REQ_NODE_RENAME, neu_req_node_rename_t) \ + XX(NEU_RESP_NODE_RENAME, neu_resp_node_rename_t) \ + XX(NEU_REQ_ADD_GROUP, neu_req_add_group_t) \ + XX(NEU_REQ_DEL_GROUP, neu_req_del_group_t) \ + XX(NEU_REQ_UPDATE_GROUP, neu_req_update_group_t) \ + XX(NEU_REQ_UPDATE_DRIVER_GROUP, neu_req_update_group_t) \ + XX(NEU_RESP_UPDATE_DRIVER_GROUP, neu_resp_update_group_t) \ + XX(NEU_REQ_GET_GROUP, neu_req_get_group_t) \ + XX(NEU_RESP_GET_GROUP, neu_resp_get_group_t) \ + XX(NEU_REQ_GET_DRIVER_GROUP, neu_req_get_group_t) \ + XX(NEU_RESP_GET_DRIVER_GROUP, neu_resp_get_driver_group_t) \ + XX(NEU_REQ_ADD_TAG, neu_req_add_tag_t) \ + XX(NEU_RESP_ADD_TAG, neu_resp_add_tag_t) \ + XX(NEU_REQ_ADD_GTAG, neu_req_add_gtag_t) \ + XX(NEU_RESP_ADD_GTAG, neu_resp_add_tag_t) \ + XX(NEU_REQ_DEL_TAG, neu_req_del_tag_t) \ + XX(NEU_REQ_UPDATE_TAG, neu_req_update_tag_t) \ + XX(NEU_RESP_UPDATE_TAG, neu_resp_update_tag_t) \ + XX(NEU_REQ_GET_TAG, neu_req_get_tag_t) \ + XX(NEU_RESP_GET_TAG, neu_resp_get_tag_t) \ + XX(NEU_REQ_ADD_PLUGIN, neu_req_add_plugin_t) \ + XX(NEU_REQ_DEL_PLUGIN, neu_req_del_plugin_t) \ + XX(NEU_REQ_UPDATE_PLUGIN, neu_req_update_plugin_t) \ + XX(NEU_REQ_GET_PLUGIN, neu_req_get_plugin_t) \ + XX(NEU_RESP_GET_PLUGIN, neu_resp_get_plugin_t) \ + XX(NEU_REQRESP_TRANS_DATA, neu_reqresp_trans_data_t) \ + XX(NEU_REQRESP_NODES_STATE, neu_reqresp_nodes_state_t) \ + XX(NEU_REQRESP_NODE_DELETED, neu_reqresp_node_deleted_t) \ + XX(NEU_REQ_ADD_DRIVERS, neu_req_driver_array_t) \ + XX(NEU_REQ_UPDATE_LOG_LEVEL, neu_req_update_log_level_t) \ + XX(NEU_REQ_PRGFILE_UPLOAD, neu_req_prgfile_upload_t) \ + XX(NEU_REQ_PRGFILE_PROCESS, neu_req_prgfile_process_t) \ + XX(NEU_RESP_PRGFILE_PROCESS, neu_resp_prgfile_process_t) \ + XX(NEU_REQ_SCAN_TAGS, neu_req_scan_tags_t) \ + XX(NEU_RESP_SCAN_TAGS, neu_resp_scan_tags_t) \ + XX(NEU_REQ_CHECK_SCHEMA, neu_req_check_schema_t) \ + XX(NEU_RESP_CHECK_SCHEMA, neu_resp_check_schema_t) \ + XX(NEU_REQ_DRIVER_ACTION, neu_req_driver_action_t) \ XX(NEU_RESP_DRIVER_ACTION, neu_resp_driver_action_t) static inline size_t neu_reqresp_size(neu_reqresp_type_e t) diff --git a/src/core/manager.c b/src/core/manager.c index 041f5eb07..4abe87d6c 100644 --- a/src/core/manager.c +++ b/src/core/manager.c @@ -1000,6 +1000,18 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data) reply(manager, header, &resp); break; } + case NEU_REQ_GET_DRIVER_SUBSCRIBE_GROUP: { + neu_req_get_driver_subscribe_group_t *cmd = + (neu_req_get_driver_subscribe_group_t *) &header[1]; + UT_array *groups = + neu_manager_get_driver_groups(manager, cmd->app, cmd->name); + neu_resp_get_subscribe_group_t resp = { .groups = groups }; + + strcpy(header->receiver, header->sender); + header->type = NEU_RESP_GET_DRIVER_SUBSCRIBE_GROUP; + reply(manager, header, &resp); + break; + } case NEU_REQ_GET_SUB_DRIVER_TAGS: { neu_req_get_sub_driver_tags_t *cmd = (neu_req_get_sub_driver_tags_t *) &header[1]; diff --git a/src/core/manager_internal.c b/src/core/manager_internal.c index 9437dc5e5..093d17935 100644 --- a/src/core/manager_internal.c +++ b/src/core/manager_internal.c @@ -376,6 +376,55 @@ UT_array *neu_manager_get_sub_group_deep_copy(neu_manager_t *manager, return subs; } +UT_array *neu_manager_get_driver_groups(neu_manager_t *manager, const char *app, + const char *name) +{ + UT_icd icd = { sizeof(neu_resp_driver_subscribe_info_t), NULL, NULL, NULL }; + UT_array *result = NULL; + + UT_array *sub_groups = + neu_subscribe_manager_get(manager->subscribe_manager, app, NULL, NULL); + UT_array *driver_groups = neu_manager_get_driver_group(manager); + + utarray_new(result, &icd); + utarray_foreach(driver_groups, neu_resp_driver_group_info_t *, grp) + { + neu_resp_driver_subscribe_info_t info = { 0 }; + + if (name != NULL && strlen(name) > 0) { + if (strcmp(grp->driver, name) == 0 || + strcmp(grp->group, name) == 0) { + strcpy(info.app, app); + strcpy(info.driver, grp->driver); + strcpy(info.group, grp->group); + info.subscribed = false; + } + } else { + strcpy(info.app, app); + strcpy(info.driver, grp->driver); + strcpy(info.group, grp->group); + info.subscribed = false; + } + + if (strlen(info.driver) > 0 && strlen(info.group) > 0) { + utarray_foreach(sub_groups, neu_resp_subscribe_info_t *, sub) + { + if (strcmp(sub->driver, info.driver) == 0 && + strcmp(sub->group, info.group) == 0) { + info.subscribed = true; + break; + } + } + + utarray_push_back(result, &info); + } + } + + utarray_free(sub_groups); + utarray_free(driver_groups); + return result; +} + int neu_manager_get_node_info(neu_manager_t *manager, const char *name, neu_persist_node_info_t *info) { diff --git a/src/core/manager_internal.h b/src/core/manager_internal.h index 5ebcaefc8..e637ccc1d 100644 --- a/src/core/manager_internal.h +++ b/src/core/manager_internal.h @@ -85,6 +85,8 @@ UT_array *neu_manager_get_sub_group_deep_copy(neu_manager_t *manager, const char * app, const char * driver, const char * group); +UT_array *neu_manager_get_driver_groups(neu_manager_t *manager, const char *app, + const char *name); int neu_manager_get_node_info(neu_manager_t *manager, const char *name, neu_persist_node_info_t *info); diff --git a/src/parser/neu_json_group_config.c b/src/parser/neu_json_group_config.c index 2aa3a9c88..8c8261701 100644 --- a/src/parser/neu_json_group_config.c +++ b/src/parser/neu_json_group_config.c @@ -362,6 +362,50 @@ int neu_json_encode_get_subscribe_resp(void *object, void *param) return ret; } +int neu_json_encode_get_driver_subscribe_resp(void *object, void *param) +{ + + neu_json_get_driver_subscribe_resp_t *resp = + (neu_json_get_driver_subscribe_resp_t *) param; + + void *driver_array = neu_json_array(); + for (int i = 0; i < resp->n_driver; i++) { + if (resp->drivers[i].n_group == 0) { + continue; + } + + json_t *ob = json_object(); + void * group_array = neu_json_array(); + + json_object_set_new(ob, "name", json_string(resp->drivers[i].driver)); + + for (int j = 0; j < resp->drivers[i].n_group; j++) { + json_t *group_ob = json_object(); + + json_object_set_new(group_ob, "name", + json_string(resp->drivers[i].groups[j].group)); + json_object_set_new( + group_ob, "subscribed", + json_boolean(resp->drivers[i].groups[j].subscribed)); + + json_array_append_new(group_array, group_ob); + } + + json_object_set_new(ob, "groups", group_array); + + json_array_append_new(driver_array, ob); + } + + neu_json_elem_t resp_elems[] = { { + .name = "drivers", + .t = NEU_JSON_OBJECT, + .v.val_object = driver_array, + } }; + + return neu_json_encode_field(object, resp_elems, + NEU_JSON_ELEM_SIZE(resp_elems)); +} + static inline int dump_params(void *root, char **const result) { return neu_json_dump_key(root, "params", result, false); diff --git a/src/parser/neu_json_group_config.h b/src/parser/neu_json_group_config.h index c6f8438ef..1e77517e6 100644 --- a/src/parser/neu_json_group_config.h +++ b/src/parser/neu_json_group_config.h @@ -110,6 +110,24 @@ typedef struct { int neu_json_encode_get_subscribe_resp(void *json_object, void *param); +typedef struct { + char *group; + bool subscribed; +} neu_json_get_driver_subscribe_resp_group_t; + +typedef struct { + char * driver; + int n_group; + neu_json_get_driver_subscribe_resp_group_t *groups; +} neu_json_get_driver_subscribe_resp_driver_t; + +typedef struct { + int n_driver; + neu_json_get_driver_subscribe_resp_driver_t *drivers; +} neu_json_get_driver_subscribe_resp_t; + +int neu_json_encode_get_driver_subscribe_resp(void *json_object, void *param); + typedef struct { char *group; char *app;