From 5af96b0507128e47c8045588c8a42351303f3e6c Mon Sep 17 00:00:00 2001 From: fengzero Date: Tue, 19 Aug 2025 03:16:17 +0000 Subject: [PATCH] event name --- include/neuron/event/event.h | 4 +- plugins/modbus/modbus_rtu.c | 8 +-- plugins/modbus/modbus_tcp.c | 16 ++--- simulator/modbus_simulator.c | 18 +++--- src/adapter/adapter.c | 24 +++---- src/adapter/driver/driver.c | 77 ++++++++++++----------- src/connection/connection_eth.c | 8 +-- src/connection/mqtt_client.c | 107 ++++++++++++++++---------------- src/core/manager.c | 40 ++++++------ src/event/event_linux.c | 50 ++++++++------- 10 files changed, 182 insertions(+), 170 deletions(-) diff --git a/include/neuron/event/event.h b/include/neuron/event/event.h index 9f173136c..631f0d712 100644 --- a/include/neuron/event/event.h +++ b/include/neuron/event/event.h @@ -35,7 +35,7 @@ typedef struct neu_events neu_events_t; * thread. * @return the newly created event. */ -neu_events_t *neu_event_new(void); +neu_events_t *neu_event_new(const char *name); /** * @brief Close a event. @@ -71,7 +71,7 @@ typedef struct neu_event_timer_param { * @param[in] timer Parameters when creating timer. * @return The added event timer. */ -neu_event_timer_t *neu_event_add_timer(neu_events_t * events, +neu_event_timer_t *neu_event_add_timer(neu_events_t *events, neu_event_timer_param_t timer); /** diff --git a/plugins/modbus/modbus_rtu.c b/plugins/modbus/modbus_rtu.c index 184335db8..083cc82d3 100644 --- a/plugins/modbus/modbus_rtu.c +++ b/plugins/modbus/modbus_rtu.c @@ -102,10 +102,10 @@ static int driver_init(neu_plugin_t *plugin, bool load) { (void) load; plugin->protocol = MODBUS_PROTOCOL_RTU; - plugin->events = neu_event_new(); + plugin->events = neu_event_new(plugin->common.name); plugin->stack = modbus_stack_create((void *) plugin, MODBUS_PROTOCOL_RTU, - modbus_send_msg, modbus_value_handle, - modbus_write_resp); + modbus_send_msg, modbus_value_handle, + modbus_write_resp); plog_notice(plugin, "%s init success", plugin->common.name); return 0; @@ -146,7 +146,7 @@ static int driver_stop(neu_plugin_t *plugin) static int driver_config(neu_plugin_t *plugin, const char *config) { int ret = 0; - char * err_param = NULL; + char *err_param = NULL; neu_conn_param_t param = { 0 }; neu_json_elem_t link = { .name = "link", .t = NEU_JSON_INT }; diff --git a/plugins/modbus/modbus_tcp.c b/plugins/modbus/modbus_tcp.c index dfb6d3c99..5c4bf32d8 100644 --- a/plugins/modbus/modbus_tcp.c +++ b/plugins/modbus/modbus_tcp.c @@ -102,10 +102,10 @@ static int driver_init(neu_plugin_t *plugin, bool load) { (void) load; plugin->protocol = MODBUS_PROTOCOL_TCP; - plugin->events = neu_event_new(); + plugin->events = neu_event_new(plugin->common.name); plugin->stack = modbus_stack_create((void *) plugin, MODBUS_PROTOCOL_TCP, - modbus_send_msg, modbus_value_handle, - modbus_write_resp); + modbus_send_msg, modbus_value_handle, + modbus_write_resp); plog_notice(plugin, "%s init success", plugin->common.name); return 0; @@ -146,20 +146,20 @@ static int driver_stop(neu_plugin_t *plugin) static int driver_config(neu_plugin_t *plugin, const char *config) { int ret = 0; - char * err_param = NULL; + char *err_param = NULL; neu_json_elem_t port = { .name = "port", .t = NEU_JSON_INT }; neu_json_elem_t timeout = { .name = "timeout", .t = NEU_JSON_INT }; neu_json_elem_t host = { .name = "host", - .t = NEU_JSON_STR, - .v.val_str = NULL }; + .t = NEU_JSON_STR, + .v.val_str = NULL }; neu_json_elem_t interval = { .name = "interval", .t = NEU_JSON_INT }; neu_json_elem_t mode = { .name = "connection_mode", .t = NEU_JSON_INT }; neu_conn_param_t param = { 0 }; neu_json_elem_t max_retries = { .name = "max_retries", .t = NEU_JSON_INT }; neu_json_elem_t retry_interval = { .name = "retry_interval", - .t = NEU_JSON_INT }; + .t = NEU_JSON_INT }; neu_json_elem_t endianess_64 = { .name = "endianess_64", - .t = NEU_JSON_INT }; + .t = NEU_JSON_INT }; ret = neu_parse_param((char *) config, &err_param, 5, &port, &host, &mode, &timeout, &interval); diff --git a/simulator/modbus_simulator.c b/simulator/modbus_simulator.c index 7df19c9e6..30352797f 100644 --- a/simulator/modbus_simulator.c +++ b/simulator/modbus_simulator.c @@ -28,9 +28,9 @@ #include "modbus_s.h" zlog_category_t *neuron = NULL; -neu_events_t * events = NULL; -neu_event_io_t * tcp_server_event = NULL; -neu_conn_t * conn = NULL; +neu_events_t *events = NULL; +neu_event_io_t *tcp_server_event = NULL; +neu_conn_t *conn = NULL; int64_t global_timestamp = 0; bool exiting = false; @@ -45,7 +45,7 @@ static bool mode_tcp = true; struct client_event { neu_event_io_t *client; int fd; - void * user_data; + void *user_data; }; struct client_event c_events[10] = { 0 }; @@ -146,7 +146,7 @@ int main(int argc, char *argv[]) .params.tcp_server.stop_listen = stop_listen, }; - events = neu_event_new(); + events = neu_event_new("modbus_simulator"); conn = neu_conn_new(¶m, NULL, connected, disconnected); signal(SIGINT, sig_handler); @@ -202,11 +202,11 @@ static int new_client(enum neu_event_io_type type, int fd, void *usr_data) case NEU_EVENT_IO_READ: { int client_fd = neu_conn_tcp_server_accept(conn); if (client_fd > 0) { - struct cycle_buf * buf = calloc(1, sizeof(struct cycle_buf)); + struct cycle_buf *buf = calloc(1, sizeof(struct cycle_buf)); neu_event_io_param_t io = { - .fd = client_fd, - .usr_data = (void *) buf, - .cb = recv_msg, + .fd = client_fd, + .usr_data = (void *) buf, + .cb = recv_msg, }; neu_event_io_t *client = neu_event_add_io(events, io); diff --git a/src/adapter/adapter.c b/src/adapter/adapter.c index 5e6d4d218..fd7110d96 100644 --- a/src/adapter/adapter.c +++ b/src/adapter/adapter.c @@ -47,7 +47,7 @@ static int adapter_command(neu_adapter_t *adapter, neu_reqresp_head_t header, void *data); static int adapter_response(neu_adapter_t *adapter, neu_reqresp_head_t *header, void *data); -static int adapter_responseto(neu_adapter_t * adapter, +static int adapter_responseto(neu_adapter_t *adapter, neu_reqresp_head_t *header, void *data, struct sockaddr_in dst); static int adapter_register_metric(neu_adapter_t *adapter, const char *name, @@ -115,7 +115,7 @@ static void *adapter_consumer(void *arg) { neu_adapter_t *adapter = (neu_adapter_t *) arg; while (1) { - neu_msg_t * msg = NULL; + neu_msg_t *msg = NULL; uint32_t n = adapter_msg_q_pop(adapter->msg_q, &msg); neu_reqresp_head_t *header = neu_msg_get_header(msg); @@ -145,7 +145,7 @@ neu_adapter_t *neu_adapter_create(neu_adapter_info_t *info, bool load) { int rv = 0; int init_rv = 0; - neu_adapter_t * adapter = NULL; + neu_adapter_t *adapter = NULL; neu_event_io_param_t param = { 0 }; switch (info->module->type) { @@ -173,7 +173,7 @@ neu_adapter_t *neu_adapter_create(neu_adapter_info_t *info, bool load) } adapter->name = strdup(info->name); - adapter->events = neu_event_new(); + adapter->events = neu_event_new("adapter"); adapter->state = NEU_NODE_RUNNING_STATE_INIT; adapter->handle = info->handle; adapter->cb_funs.command = callback_funs.command; @@ -563,7 +563,7 @@ static int adapter_response(neu_adapter_t *adapter, neu_reqresp_head_t *header, return ret; } -static int adapter_responseto(neu_adapter_t * adapter, +static int adapter_responseto(neu_adapter_t *adapter, neu_reqresp_head_t *header, void *data, struct sockaddr_in dst) { @@ -878,7 +878,7 @@ static int adapter_loop(enum neu_event_io_type type, int fd, void *usr_data) case NEU_REQ_GET_TAG: { neu_req_get_tag_t *cmd = (neu_req_get_tag_t *) &header[1]; neu_resp_error_t error = { .error = 0 }; - UT_array * tags = NULL; + UT_array *tags = NULL; if (adapter->module->type == NEU_NA_TYPE_DRIVER) { error.error = neu_adapter_driver_query_tag( @@ -1229,7 +1229,7 @@ static int adapter_loop(enum neu_event_io_type type, int fd, void *usr_data) neu_req_get_ndriver_tags_t *cmd = (neu_req_get_ndriver_tags_t *) &header[1]; neu_resp_error_t error = { .error = 0 }; - UT_array * tags = NULL; + UT_array *tags = NULL; // TODO (void) cmd; @@ -1585,7 +1585,7 @@ int neu_adapter_validate_tag(neu_adapter_t *adapter, neu_datatag_t *tag) return error; } -neu_event_timer_t *neu_adapter_add_timer(neu_adapter_t * adapter, +neu_event_timer_t *neu_adapter_add_timer(neu_adapter_t *adapter, neu_event_timer_param_t param) { return neu_event_add_timer(adapter->events, param); @@ -1610,7 +1610,7 @@ int neu_adapter_register_group_metric(neu_adapter_t *adapter, } int neu_adapter_update_group_metric(neu_adapter_t *adapter, - const char * group_name, + const char *group_name, const char *metric_name, uint64_t n) { if (NULL == adapter->metrics) { @@ -1622,8 +1622,8 @@ int neu_adapter_update_group_metric(neu_adapter_t *adapter, } int neu_adapter_metric_update_group_name(neu_adapter_t *adapter, - const char * group_name, - const char * new_group_name) + const char *group_name, + const char *new_group_name) { if (NULL == adapter->metrics) { return -1; @@ -1634,7 +1634,7 @@ int neu_adapter_metric_update_group_name(neu_adapter_t *adapter, } void neu_adapter_del_group_metrics(neu_adapter_t *adapter, - const char * group_name) + const char *group_name) { if (NULL != adapter->metrics) { neu_node_metrics_del_group(adapter->metrics, group_name); diff --git a/src/adapter/driver/driver.c b/src/adapter/driver/driver.c index d7a5fda8f..631e44d05 100644 --- a/src/adapter/driver/driver.c +++ b/src/adapter/driver/driver.c @@ -39,8 +39,8 @@ typedef struct to_be_write_tag { bool single; neu_datatag_t *tag; neu_value_u value; - UT_array * tvs; - void * req; + UT_array *tvs; + void *req; } to_be_write_tag_t; typedef struct { @@ -52,16 +52,16 @@ typedef struct group { char *name; int64_t timestamp; - neu_group_t * group; - UT_array * static_tags; - UT_array * wt_tags; + neu_group_t *group; + UT_array *static_tags; + UT_array *wt_tags; pthread_mutex_t wt_mtx; neu_event_timer_t *report; neu_event_timer_t *read; neu_event_timer_t *write; - UT_array * apps; // sub_app_t array + UT_array *apps; // sub_app_t array pthread_mutex_t apps_mtx; neu_plugin_group_t grp; @@ -74,7 +74,7 @@ struct neu_adapter_driver { neu_adapter_t adapter; neu_driver_cache_t *cache; - neu_events_t * driver_events; + neu_events_t *driver_events; size_t tag_cnt; struct group *groups; @@ -106,10 +106,10 @@ static void update_with_meta(neu_adapter_t *adapter, const char *group, const char *tag, neu_dvalue_t value, neu_tag_meta_t *metas, int n_meta); static void write_response(neu_adapter_t *adapter, void *r, neu_error error); -static group_t * find_group(neu_adapter_driver_t *driver, const char *name); +static group_t *find_group(neu_adapter_driver_t *driver, const char *name); static void store_write_tag(group_t *group, to_be_write_tag_t *tag); static inline void start_group_timer(neu_adapter_driver_t *driver, - group_t * grp); + group_t *grp); static inline void stop_group_timer(neu_adapter_driver_t *driver, group_t *grp); static void format_tag_value(neu_dvalue_t *value) @@ -173,7 +173,7 @@ static void update_with_meta(neu_adapter_t *adapter, const char *group, const char *tag, neu_dvalue_t value, neu_tag_meta_t *metas, int n_meta) { - neu_adapter_driver_t * driver = (neu_adapter_driver_t *) adapter; + neu_adapter_driver_t *driver = (neu_adapter_driver_t *) adapter; neu_adapter_update_metric_cb_t update_metric = driver->adapter.cb_funs.update_metric; @@ -337,9 +337,9 @@ neu_adapter_driver_t *neu_adapter_driver_create() { neu_adapter_driver_t *driver = calloc(1, sizeof(neu_adapter_driver_t)); - driver->cache = neu_driver_cache_new(); - driver->driver_events = neu_event_new(); - driver->adapter.cb_funs.driver.update = update; + driver->cache = neu_driver_cache_new(); + driver->driver_events = neu_event_new("adapter_driver"); + driver->adapter.cb_funs.driver.update = update; driver->adapter.cb_funs.driver.write_response = write_response; driver->adapter.cb_funs.driver.update_im = update_im; driver->adapter.cb_funs.driver.update_with_meta = update_with_meta; @@ -472,14 +472,17 @@ void neu_adapter_driver_stop_group_timer(neu_adapter_driver_t *driver) { group_t *el = NULL, *tmp = NULL; - HASH_ITER(hh, driver->groups, el, tmp) { stop_group_timer(driver, el); } + HASH_ITER(hh, driver->groups, el, tmp) + { + stop_group_timer(driver, el); + } } void neu_adapter_driver_read_group(neu_adapter_driver_t *driver, - neu_reqresp_head_t * req) + neu_reqresp_head_t *req) { neu_req_read_group_t *cmd = (neu_req_read_group_t *) &req[1]; - group_t * g = find_group(driver, cmd->group); + group_t *g = find_group(driver, cmd->group); if (g == NULL) { neu_resp_error_t error = { .error = NEU_ERR_GROUP_NOT_EXIST }; req->type = NEU_RESP_ERROR; @@ -489,7 +492,7 @@ void neu_adapter_driver_read_group(neu_adapter_driver_t *driver, } neu_resp_read_group_t resp = { 0 }; - neu_group_t * group = g->group; + neu_group_t *group = g->group; UT_array *tags = neu_group_query_read_tag(group, cmd->name, cmd->desc); utarray_new(resp.tags, neu_resp_tag_value_meta_icd()); @@ -549,7 +552,7 @@ void neu_adapter_driver_read_group(neu_adapter_driver_t *driver, } void neu_adapter_driver_read_group_paginate(neu_adapter_driver_t *driver, - neu_reqresp_head_t * req) + neu_reqresp_head_t *req) { neu_req_read_group_paginate_t *cmd = (neu_req_read_group_paginate_t *) &req[1]; @@ -563,8 +566,8 @@ void neu_adapter_driver_read_group_paginate(neu_adapter_driver_t *driver, } neu_resp_read_group_paginate_t resp = { 0 }; - neu_group_t * group = g->group; - UT_array * tags; + neu_group_t *group = g->group; + UT_array *tags; if (cmd->is_error != true && cmd->current_page > 0 && cmd->page_size > 0) { tags = neu_group_query_read_tag_paginate( @@ -837,7 +840,7 @@ static void cal_decimal(neu_type_e tag_type, neu_type_e value_type, } void neu_adapter_driver_write_tags(neu_adapter_driver_t *driver, - neu_reqresp_head_t * req) + neu_reqresp_head_t *req) { neu_req_write_tags_t *cmd = (neu_req_write_tags_t *) &req[1]; @@ -911,10 +914,10 @@ void neu_adapter_driver_write_tags(neu_adapter_driver_t *driver, } void neu_adapter_driver_write_gtags(neu_adapter_driver_t *driver, - neu_reqresp_head_t * req) + neu_reqresp_head_t *req) { neu_req_write_gtags_t *cmd = (neu_req_write_gtags_t *) &req[1]; - group_t * first_g = NULL; + group_t *first_g = NULL; if (driver->adapter.state != NEU_NODE_RUNNING_STATE_RUNNING) { driver->adapter.cb_funs.driver.write_response( @@ -1008,7 +1011,7 @@ void neu_adapter_driver_write_gtags(neu_adapter_driver_t *driver, } void neu_adapter_driver_write_tag(neu_adapter_driver_t *driver, - neu_reqresp_head_t * req) + neu_reqresp_head_t *req) { if (driver->adapter.state != NEU_NODE_RUNNING_STATE_RUNNING) { driver->adapter.cb_funs.driver.write_response( @@ -1017,7 +1020,7 @@ void neu_adapter_driver_write_tag(neu_adapter_driver_t *driver, } neu_req_write_tag_t *cmd = (neu_req_write_tag_t *) &req[1]; - group_t * g = find_group(driver, cmd->group); + group_t *g = find_group(driver, cmd->group); if (g == NULL) { driver->adapter.cb_funs.driver.write_response(&driver->adapter, req, @@ -1268,7 +1271,7 @@ 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) + neu_req_add_gtag_t *cmd) { uint16_t new_groups_count = 0; for (int i = 0; i < cmd->n_group; i++) { @@ -1290,7 +1293,7 @@ static group_t *find_group(neu_adapter_driver_t *driver, const char *name) return find; } int neu_adapter_driver_group_exist(neu_adapter_driver_t *driver, - const char * name) + const char *name) { group_t *find = NULL; int ret = NEU_ERR_GROUP_NOT_EXIST; @@ -1305,7 +1308,7 @@ int neu_adapter_driver_group_exist(neu_adapter_driver_t *driver, UT_array *neu_adapter_driver_get_group(neu_adapter_driver_t *driver) { - group_t * el = NULL, *tmp = NULL; + group_t *el = NULL, *tmp = NULL; UT_array *groups = NULL; UT_icd icd = { sizeof(neu_resp_group_info_t), NULL, NULL, NULL }; @@ -1558,9 +1561,9 @@ void neu_adapter_driver_get_value_tag(neu_adapter_driver_t *driver, } UT_array *neu_adapter_driver_get_read_tag(neu_adapter_driver_t *driver, - const char * group) + const char *group) { - group_t * find = NULL; + group_t *find = NULL; UT_array *tags = NULL; HASH_FIND_STR(driver->groups, group, find); @@ -1574,7 +1577,7 @@ UT_array *neu_adapter_driver_get_read_tag(neu_adapter_driver_t *driver, UT_array *neu_adapter_driver_get_ptag(neu_adapter_driver_t *driver, const char *group, const char *tag) { - group_t * find = NULL; + group_t *find = NULL; UT_array *tags = NULL; HASH_FIND_STR(driver->groups, group, find); @@ -1640,7 +1643,7 @@ static void report_to_app(neu_adapter_driver_t *driver, group_t *group, static int report_callback(void *usr_data) { - group_t * group = (group_t *) usr_data; + group_t *group = (group_t *) usr_data; neu_node_running_state_e state = group->driver->adapter.state; if (state != NEU_NODE_RUNNING_STATE_RUNNING) { return 0; @@ -1777,7 +1780,7 @@ static void group_change(void *arg, int64_t timestamp, UT_array *static_tags, static int write_callback(void *usr_data) { - group_t * group = (group_t *) usr_data; + group_t *group = (group_t *) usr_data; neu_node_running_state_e state = group->driver->adapter.state; if (state != NEU_NODE_RUNNING_STATE_RUNNING) { return 0; @@ -1809,7 +1812,7 @@ static int write_callback(void *usr_data) static int read_callback(void *usr_data) { - group_t * group = (group_t *) usr_data; + group_t *group = (group_t *) usr_data; neu_node_running_state_e state = group->driver->adapter.state; if (state != NEU_NODE_RUNNING_STATE_RUNNING) { return 0; @@ -2377,10 +2380,10 @@ static void store_write_tag(group_t *group, to_be_write_tag_t *tag) } void neu_adapter_driver_subscribe(neu_adapter_driver_t *driver, - neu_req_subscribe_t * req) + neu_req_subscribe_t *req) { sub_app_t sub_app = { 0 }; - group_t * find = NULL; + group_t *find = NULL; HASH_FIND_STR(driver->groups, req->group, find); if (find == NULL) { @@ -2411,7 +2414,7 @@ void neu_adapter_driver_subscribe(neu_adapter_driver_t *driver, report_to_app(driver, find, sub_app.addr); } -void neu_adapter_driver_unsubscribe(neu_adapter_driver_t * driver, +void neu_adapter_driver_unsubscribe(neu_adapter_driver_t *driver, neu_req_unsubscribe_t *req) { group_t *find = NULL; diff --git a/src/connection/connection_eth.c b/src/connection/connection_eth.c index 1e8ec6203..9cf1b4524 100644 --- a/src/connection/connection_eth.c +++ b/src/connection/connection_eth.c @@ -60,7 +60,7 @@ typedef struct { } lldp_callback_elem_t; typedef struct interface_conn { - char * interface; + char *interface; neu_events_t *events; uint8_t count; @@ -82,7 +82,7 @@ static interface_conn_t in_conns[interface_max]; struct neu_conn_eth { interface_conn_t *ic; - void * ctx; + void *ctx; }; static int get_mac(neu_conn_eth_t *conn, const char *interface); @@ -137,7 +137,7 @@ neu_conn_eth_t *neu_conn_eth_init(const char *interface, void *ctx) conn_eth->ic = &in_conns[i]; in_conns[i].interface = strdup(interface); - in_conns[i].events = neu_event_new(); + in_conns[i].events = neu_event_new("eth_conn"); in_conns[i].count += 1; get_mac(conn_eth, interface); @@ -225,7 +225,7 @@ neu_conn_eth_sub_t *neu_conn_eth_register(neu_conn_eth_t *conn, uint8_t xmac[6], { char tmp_mac[20] = { 0 }; neu_conn_eth_sub_t *sub = NULL; - callback_elem_t * elem = NULL; + callback_elem_t *elem = NULL; snprintf(tmp_mac, sizeof(tmp_mac), "%X:%X:%X-%X:%X:%X", xmac[0], xmac[1], xmac[2], xmac[3], xmac[4], xmac[5]); diff --git a/src/connection/mqtt_client.c b/src/connection/mqtt_client.c index 5fb2ec7b1..e8f1fdc13 100644 --- a/src/connection/mqtt_client.c +++ b/src/connection/mqtt_client.c @@ -63,9 +63,9 @@ typedef struct { size_t ref; bool ack; neu_mqtt_qos_e qos; - char * topic; + char *topic; neu_mqtt_client_subscribe_cb_t cb; - void * data; + void *data; UT_hash_handle hh; } subscription_t; @@ -80,10 +80,10 @@ typedef enum { struct { \ neu_mqtt_client_publish_cb_t cb; \ neu_mqtt_qos_e qos; \ - char * topic; \ - uint8_t * payload; \ + char *topic; \ + uint8_t *payload; \ uint32_t len; \ - void * data; \ + void *data; \ } pub; \ subscription_t *sub; \ struct { \ @@ -96,7 +96,7 @@ typedef union { typedef struct task_s { task_kind_e kind; - nng_aio * aio; + nng_aio *aio; union { TASK_UNION_FIELDS; }; @@ -106,31 +106,31 @@ typedef struct task_s { struct neu_mqtt_client_s { nng_socket sock; - nng_mtx * mtx; - neu_events_t * events; - neu_event_timer_t * timer; + nng_mtx *mtx; + neu_events_t *events; + neu_event_timer_t *timer; neu_mqtt_version_e version; - char * host; + char *host; uint16_t port; - char * url; - nng_tls_config * tls_cfg; - nng_msg * conn_msg; + char *url; + nng_tls_config *tls_cfg; + nng_msg *conn_msg; nng_duration retry; bool open; bool connected; neu_mqtt_client_connection_cb_t connect_cb; - void * connect_cb_data; + void *connect_cb_data; neu_mqtt_client_connection_cb_t disconnect_cb; - void * disconnect_cb_data; - nng_mqtt_sqlite_option * sqlite_cfg; + void *disconnect_cb_data; + nng_mqtt_sqlite_option *sqlite_cfg; bool receiving; - nng_aio * recv_aio; - subscription_t * subscriptions; + nng_aio *recv_aio; + subscription_t *subscriptions; size_t suback_count; size_t task_count; size_t task_limit; - task_t * task_free_list; - zlog_category_t * log; + task_t *task_free_list; + zlog_category_t *log; }; static inline task_t *task_new(neu_mqtt_client_t *client); @@ -142,10 +142,10 @@ static void task_handle_sub(task_t *task, neu_mqtt_client_t *client); static void task_handle_unsub(task_t *task, neu_mqtt_client_t *client); static void task_handle_recv(task_t *task, neu_mqtt_client_t *client); -static subscription_t * subscription_new(neu_mqtt_client_t *client, +static subscription_t *subscription_new(neu_mqtt_client_t *client, neu_mqtt_qos_e qos, const char *topic, neu_mqtt_client_subscribe_cb_t cb, - void * data); + void *data); static inline void subscription_free(subscription_t *subscription); static inline subscription_t *subscription_ref(subscription_t *subscription); static inline void subscriptions_free(subscription_t *subscriptions); @@ -159,9 +159,9 @@ static inline task_t *client_alloc_task(neu_mqtt_client_t *client); static inline void client_free_task(neu_mqtt_client_t *client, task_t *task); static inline size_t client_task_free_list_len(neu_mqtt_client_t *client); static inline void client_add_subscription(neu_mqtt_client_t *client, - subscription_t * sub); + subscription_t *sub); static int client_send_sub_msg(neu_mqtt_client_t *client, - subscription_t * subscription); + subscription_t *subscription); static inline void client_start_recv(neu_mqtt_client_t *client); static inline int client_start_timer(neu_mqtt_client_t *client); static inline int client_make_url(neu_mqtt_client_t *client); @@ -307,8 +307,8 @@ static inline void tasks_free(task_t *tasks) static void task_cb(void *arg) { - task_t * task = arg; - nng_aio * aio = task->aio; + task_t *task = arg; + nng_aio *aio = task->aio; neu_mqtt_client_t *client = nng_aio_get_input(aio, 0); if (TASK_PUB == task->kind) { @@ -443,7 +443,7 @@ static void task_handle_recv(task_t *task, neu_mqtt_client_t *client) static subscription_t *subscription_new(neu_mqtt_client_t *client, neu_mqtt_qos_e qos, const char *topic, neu_mqtt_client_subscribe_cb_t cb, - void * data) + void *data) { (void) client; subscription_t *subscription = NULL; @@ -516,7 +516,7 @@ subscription_find_match(subscription_t *subscriptions, const char *topic_name, static int resub_cb(void *data) { neu_mqtt_client_t *client = data; - subscription_t * sub = NULL; + subscription_t *sub = NULL; nng_mtx_lock(client->mtx); if (client->connected && @@ -538,9 +538,9 @@ static void connect_cb(nng_pipe p, nng_pipe_ev ev, void *arg) { (void) p; (void) ev; - neu_mqtt_client_t * client = arg; + neu_mqtt_client_t *client = arg; neu_mqtt_client_connection_cb_t cb = NULL; - void * data = NULL; + void *data = NULL; int reason = 0; nng_pipe_get_int(p, NNG_OPT_MQTT_CONNECT_REASON, &reason); @@ -565,10 +565,10 @@ static void disconnect_cb(nng_pipe p, nng_pipe_ev ev, void *arg) { (void) p; (void) ev; - neu_mqtt_client_t * client = arg; - subscription_t * sub = NULL; + neu_mqtt_client_t *client = arg; + subscription_t *sub = NULL; neu_mqtt_client_connection_cb_t cb = NULL; - void * data = NULL; + void *data = NULL; int reason = 0; nng_pipe_get_int(p, NNG_OPT_MQTT_DISCONNECT_REASON, &reason); @@ -578,7 +578,10 @@ static void disconnect_cb(nng_pipe p, nng_pipe_ev ev, void *arg) client->connected = false; cb = client->disconnect_cb; data = client->disconnect_cb_data; - HASH_LOOP(hh, client->subscriptions, sub) { sub->ack = false; } + HASH_LOOP(hh, client->subscriptions, sub) + { + sub->ack = false; + } client->suback_count = 0; nng_mtx_unlock(client->mtx); @@ -590,9 +593,9 @@ static void disconnect_cb(nng_pipe p, nng_pipe_ev ev, void *arg) static void recv_cb(void *arg) { int rv = 0; - subscription_t * subscription = arg; + subscription_t *subscription = arg; neu_mqtt_client_t *client = arg; - nng_aio * aio = client->recv_aio; + nng_aio *aio = client->recv_aio; if (0 != (rv = nng_aio_result(aio))) { log(error, "mqtt client recv error: %s", nng_strerror(rv)); @@ -621,7 +624,7 @@ static void recv_cb(void *arg) } uint32_t payload_len; - uint8_t * payload = nng_mqtt_msg_get_publish_payload(msg, &payload_len); + uint8_t *payload = nng_mqtt_msg_get_publish_payload(msg, &payload_len); uint32_t topic_len; const char *topic = nng_mqtt_msg_get_publish_topic(msg, &topic_len); uint8_t qos = nng_mqtt_msg_get_publish_qos(msg); @@ -712,7 +715,7 @@ static inline size_t client_task_free_list_len(neu_mqtt_client_t *client) } static inline void client_add_subscription(neu_mqtt_client_t *client, - subscription_t * sub) + subscription_t *sub) { subscription_t *old = NULL; @@ -725,7 +728,7 @@ static inline void client_add_subscription(neu_mqtt_client_t *client, } static inline void client_del_subscription(neu_mqtt_client_t *client, - subscription_t * sub) + subscription_t *sub) { HASH_DEL(client->subscriptions, sub); if (sub->ack) { @@ -735,7 +738,7 @@ static inline void client_del_subscription(neu_mqtt_client_t *client, } static int client_send_sub_msg(neu_mqtt_client_t *client, - subscription_t * subscription) + subscription_t *subscription) { int rv = 0; nng_msg *sub_msg = NULL; @@ -772,7 +775,7 @@ static int client_send_sub_msg(neu_mqtt_client_t *client, } static int client_send_unsub_msg(neu_mqtt_client_t *client, - subscription_t * subscription) + subscription_t *subscription) { int rv = 0; nng_msg *sub_msg = NULL; @@ -813,7 +816,7 @@ static inline void client_start_recv(neu_mqtt_client_t *client) static inline int client_start_timer(neu_mqtt_client_t *client) { - neu_events_t * events = NULL; + neu_events_t *events = NULL; neu_event_timer_t *timer = NULL; if (client->events) { @@ -821,7 +824,7 @@ static inline int client_start_timer(neu_mqtt_client_t *client) return 0; } - events = neu_event_new(); + events = neu_event_new("mqtt_client"); if (NULL == events) { return -1; } @@ -846,7 +849,7 @@ static inline int client_start_timer(neu_mqtt_client_t *client) static inline int client_make_url(neu_mqtt_client_t *client) { - char * url = NULL; + char *url = NULL; const char *fmt = NULL; if (client->tls_cfg) { @@ -906,7 +909,7 @@ static inline nng_mqtt_sqlite_option * alloc_sqlite_config(neu_mqtt_client_t *client) { int rv; - char * db = NULL; + char *db = NULL; nng_mqtt_sqlite_option *cfg = NULL; const mqtt_buf client_id = nng_mqtt_msg_get_connect_client_id(client->conn_msg); @@ -1113,9 +1116,9 @@ int neu_mqtt_client_set_user(neu_mqtt_client_t *client, const char *username, return 0; } -int neu_mqtt_client_set_connect_cb(neu_mqtt_client_t * client, +int neu_mqtt_client_set_connect_cb(neu_mqtt_client_t *client, neu_mqtt_client_connection_cb_t cb, - void * data) + void *data) { nng_mtx_lock(client->mtx); return_failure_if_open(); @@ -1127,9 +1130,9 @@ int neu_mqtt_client_set_connect_cb(neu_mqtt_client_t * client, return 0; } -int neu_mqtt_client_set_disconnect_cb(neu_mqtt_client_t * client, +int neu_mqtt_client_set_disconnect_cb(neu_mqtt_client_t *client, neu_mqtt_client_connection_cb_t cb, - void * data) + void *data) { nng_mtx_lock(client->mtx); return_failure_if_open(); @@ -1287,7 +1290,7 @@ int neu_mqtt_client_set_cache_sync_interval(neu_mqtt_client_t *client, } int neu_mqtt_client_set_zlog_category(neu_mqtt_client_t *client, - zlog_category_t * cat) + zlog_category_t *cat) { nng_mtx_lock(client->mtx); return_failure_if_open(); @@ -1393,7 +1396,7 @@ int neu_mqtt_client_open(neu_mqtt_client_t *client) int neu_mqtt_client_close(neu_mqtt_client_t *client) { int rv = 0; - neu_events_t * events = NULL; + neu_events_t *events = NULL; neu_event_timer_t *timer = NULL; nng_mtx_lock(client->mtx); @@ -1440,7 +1443,7 @@ int neu_mqtt_client_publish(neu_mqtt_client_t *client, neu_mqtt_qos_e qos, { int rv = 0; nng_msg *pub_msg = NULL; - task_t * task = NULL; + task_t *task = NULL; if (0 != (rv = nng_mqtt_msg_alloc(&pub_msg, 0))) { log(error, "nng_mqtt_msg_alloc fail: %s", nng_strerror(rv)); diff --git a/src/core/manager.c b/src/core/manager.c index 58f578881..8369495e1 100644 --- a/src/core/manager.c +++ b/src/core/manager.c @@ -54,11 +54,11 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data); inline static void reply(neu_manager_t *manager, neu_reqresp_head_t *header, void *data); -inline static void forward_msg(neu_manager_t * manager, +inline static void forward_msg(neu_manager_t *manager, neu_reqresp_head_t *header, const char *node); -inline static void forward_msg_copy(neu_manager_t * manager, +inline static void forward_msg_copy(neu_manager_t *manager, neu_reqresp_head_t *header, - const char * node); + const char *node); static void start_static_adapter(neu_manager_t *manager, const char *name); static int update_timestamp(void *usr_data); @@ -85,10 +85,10 @@ uint16_t neu_manager_get_port() neu_manager_t *neu_manager_create() { int rv = 0; - neu_manager_t * manager = calloc(1, sizeof(neu_manager_t)); + neu_manager_t *manager = calloc(1, sizeof(neu_manager_t)); neu_event_io_param_t param = { - .usr_data = (void *) manager, - .cb = manager_loop, + .usr_data = (void *) manager, + .cb = manager_loop, }; neu_event_timer_param_t timestamp_timer_param = { @@ -98,7 +98,7 @@ neu_manager_t *neu_manager_create() .type = NEU_EVENT_TIMER_NOBLOCK, }; - manager->events = neu_event_new(); + manager->events = neu_event_new("manager"); manager->plugin_manager = neu_plugin_manager_create(); manager->node_manager = neu_node_manager_create(); manager->subscribe_manager = neu_subscribe_manager_create(); @@ -206,9 +206,9 @@ void neu_manager_destroy(neu_manager_t *manager) static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data) { int rv = 0; - neu_manager_t * manager = (neu_manager_t *) usr_data; + neu_manager_t *manager = (neu_manager_t *) usr_data; struct sockaddr_in src_addr = { 0 }; - neu_msg_t * msg = NULL; + neu_msg_t *msg = NULL; neu_reqresp_head_t *header = NULL; if (type == NEU_EVENT_IO_CLOSED || type == NEU_EVENT_IO_HUP) { @@ -706,7 +706,7 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data) break; } case NEU_REQ_GET_PLUGIN: { - UT_array * plugins = neu_manager_get_plugins(manager); + UT_array *plugins = neu_manager_get_plugins(manager); neu_resp_get_plugin_t resp = { .plugins = plugins }; header->type = NEU_RESP_GET_PLUGIN; @@ -1095,7 +1095,7 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data) } case NEU_REQ_GET_NODE: { neu_req_get_node_t *cmd = (neu_req_get_node_t *) &header[1]; - UT_array * nodes = + UT_array *nodes = neu_manager_get_nodes(manager, cmd->type, cmd->plugin, cmd->node); neu_resp_get_node_t resp = { .nodes = nodes }; @@ -1229,7 +1229,7 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data) utarray_foreach(groups, neu_resp_subscribe_info_t *, info) { neu_resp_get_sub_driver_tags_info_t in = { 0 }; - neu_adapter_t * driver = + neu_adapter_t *driver = neu_node_manager_find(manager->node_manager, info->driver); assert(driver != NULL); @@ -1669,7 +1669,7 @@ static int manager_loop(enum neu_event_io_type type, int fd, void *usr_data) return 0; } -inline static void forward_msg(neu_manager_t * manager, +inline static void forward_msg(neu_manager_t *manager, neu_reqresp_head_t *header, const char *node) { struct sockaddr_in addr = @@ -1690,9 +1690,9 @@ inline static void forward_msg(neu_manager_t * manager, } } -inline static void forward_msg_copy(neu_manager_t * manager, +inline static void forward_msg_copy(neu_manager_t *manager, neu_reqresp_head_t *header, - const char * node) + const char *node) { neu_msg_t *msg = neu_msg_copy((neu_msg_t *) header); forward_msg(manager, neu_msg_get_header(msg), node); @@ -1700,10 +1700,10 @@ inline static void forward_msg_copy(neu_manager_t * manager, static void start_static_adapter(neu_manager_t *manager, const char *name) { - neu_adapter_t * adapter = NULL; + neu_adapter_t *adapter = NULL; neu_plugin_instance_t instance = { 0 }; neu_adapter_info_t adapter_info = { - .name = name, + .name = name, }; neu_plugin_manager_load_static(manager->plugin_manager, name, &instance); @@ -1719,10 +1719,10 @@ 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) { - neu_adapter_t * adapter = NULL; + neu_adapter_t *adapter = NULL; neu_plugin_instance_t instance = { 0 }; neu_adapter_info_t adapter_info = { - .name = name, + .name = name, }; if (0 != @@ -1752,7 +1752,7 @@ inline static void reply(neu_manager_t *manager, neu_reqresp_head_t *header, neu_node_manager_get_addr(manager->node_manager, header->receiver); neu_reqresp_type_e t = header->type; - void * ctx = header->ctx; + void *ctx = header->ctx; char receiver[NEU_NODE_NAME_LEN] = { 0 }; strncpy(receiver, header->receiver, sizeof(receiver)); diff --git a/src/event/event_linux.c b/src/event/event_linux.c index a70059e1c..bd365c6c2 100644 --- a/src/event/event_linux.c +++ b/src/event/event_linux.c @@ -35,7 +35,7 @@ struct neu_event_timer { int fd; - struct event_data * event_data; + struct event_data *event_data; struct itimerspec value; neu_event_timer_type_e type; pthread_mutex_t mtx; @@ -73,6 +73,8 @@ struct neu_events { pthread_t thread; bool stop; + char *name; + pthread_mutex_t mtx; int n_event; struct event_data event_datas[EVENT_SIZE]; @@ -122,8 +124,8 @@ static void *event_loop(void *arg) } if (ret == -1 || events->stop) { - zlog_warn(neuron, "event loop exit, errno: %s(%d), stop: %d", - strerror(errno), errno, events->stop); + zlog_warn(neuron, "event(%s) loop exit, errno: %s(%d), stop: %d", + events->name, strerror(errno), errno, events->stop); break; } @@ -177,17 +179,18 @@ static void *event_loop(void *arg) return NULL; }; -neu_events_t *neu_event_new(void) +neu_events_t *neu_event_new(const char *name) { neu_events_t *events = calloc(1, sizeof(struct neu_events)); events->epoll_fd = epoll_create(1); - nlog_notice("create epoll: %d(%d)", events->epoll_fd, errno); + nlog_notice("create epoll: %d(%d) %s", events->epoll_fd, errno, name); assert(events->epoll_fd > 0); events->stop = false; events->n_event = 0; + events->name = strdup(name); pthread_mutex_init(&events->mtx, NULL); pthread_create(&events->thread, NULL, event_loop, events); @@ -199,6 +202,7 @@ int neu_event_close(neu_events_t *events) { events->stop = true; close(events->epoll_fd); + free(events->name); pthread_join(events->thread, NULL); pthread_mutex_destroy(&events->mtx); @@ -207,20 +211,21 @@ int neu_event_close(neu_events_t *events) return 0; } -neu_event_timer_t *neu_event_add_timer(neu_events_t * events, +neu_event_timer_t *neu_event_add_timer(neu_events_t *events, neu_event_timer_param_t timer) { int ret = 0; int timer_fd = timerfd_create(CLOCK_MONOTONIC, 0); struct itimerspec value = { - .it_value.tv_sec = timer.second, - .it_value.tv_nsec = timer.millisecond * 1000 * 1000, - .it_interval.tv_sec = timer.second, - .it_interval.tv_nsec = timer.millisecond * 1000 * 1000, + .it_value.tv_sec = timer.second, + .it_value.tv_nsec = timer.millisecond * 1000 * 1000, + .it_interval.tv_sec = timer.second, + .it_interval.tv_nsec = timer.millisecond * 1000 * 1000, }; int index = get_free_event(events); if (index < 0) { - zlog_fatal(neuron, "no free event: %d", events->epoll_fd); + zlog_fatal(neuron, "no free event(%s): %d", events->name, + events->epoll_fd); } assert(index >= 0); @@ -251,18 +256,19 @@ neu_event_timer_t *neu_event_add_timer(neu_events_t * events, zlog_notice(neuron, "add timer, second: %" PRId64 ", millisecond: %" PRId64 - ", timer: %d in epoll %d, " + ", timer: %d in epoll(%s) %d, " "ret: %d, index: %d", - timer.second, timer.millisecond, timer_fd, events->epoll_fd, - ret, index); + timer.second, timer.millisecond, timer_fd, events->name, + events->epoll_fd, ret, index); return timer_ctx; } int neu_event_del_timer(neu_events_t *events, neu_event_timer_t *timer) { - zlog_notice(neuron, "del timer: %d from epoll: %d, index: %d", timer->fd, - events->epoll_fd, timer->event_data->index); + zlog_notice(neuron, "del timer: %d from epoll(%s): %d, index: %d", + timer->fd, events->name, events->epoll_fd, + timer->event_data->index); timer->stop = true; epoll_ctl(events->epoll_fd, EPOLL_CTL_DEL, timer->fd, NULL); @@ -281,8 +287,8 @@ neu_event_io_t *neu_event_add_io(neu_events_t *events, neu_event_io_param_t io) int ret = 0; int index = get_free_event(events); - nlog_notice("add io, fd: %d, epoll: %d, index: %d", io.fd, events->epoll_fd, - index); + nlog_notice("add io, fd: %d, epoll(%s): %d, index: %d", io.fd, events->name, + events->epoll_fd, index); assert(index >= 0); neu_event_io_t *io_ctx = &events->event_datas[index].ctx.io; @@ -303,8 +309,8 @@ neu_event_io_t *neu_event_add_io(neu_events_t *events, neu_event_io_param_t io) ret = epoll_ctl(events->epoll_fd, EPOLL_CTL_ADD, io.fd, &event); - nlog_notice("add io, fd: %d, epoll: %d, ret: %d(%d), index: %d", io.fd, - events->epoll_fd, ret, errno, index); + nlog_notice("add io, fd: %d, epoll(%s): %d, ret: %d(%d), index: %d", io.fd, + events->name, events->epoll_fd, ret, errno, index); assert(ret == 0); return io_ctx; @@ -316,8 +322,8 @@ int neu_event_del_io(neu_events_t *events, neu_event_io_t *io) return 0; } - zlog_notice(neuron, "del io: %d from epoll: %d, index: %d", io->fd, - events->epoll_fd, io->event_data->index); + zlog_notice(neuron, "del io: %d from epoll(%s): %d, index: %d", io->fd, + events->name, events->epoll_fd, io->event_data->index); epoll_ctl(events->epoll_fd, EPOLL_CTL_DEL, io->fd, NULL); free_event(events, io->event_data->index);