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
2 changes: 1 addition & 1 deletion include/neuron/event/event.h
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion plugins/modbus/modbus_rtu.c
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,7 @@ 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);
Expand Down
2 changes: 1 addition & 1 deletion plugins/modbus/modbus_tcp.c
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ 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);
Expand Down
2 changes: 1 addition & 1 deletion plugins/mqtt/mqtt_plugin_intf.c
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ static int start_hearbeat_timer(neu_plugin_t *plugin, uint64_t interval)
}

if (NULL == plugin->events) {
plugin->events = neu_event_new();
plugin->events = neu_event_new(plugin->common.name);
if (NULL == plugin->events) {
plog_error(plugin, "neu_event_new fail");
return NEU_ERR_EINTERNAL;
Expand Down
2 changes: 1 addition & 1 deletion simulator/modbus_simulator.c
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,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(&param, NULL, connected, disconnected);

signal(SIGINT, sig_handler);
Expand Down
2 changes: 1 addition & 1 deletion simulator/modbus_tty_simulator.c
Original file line number Diff line number Diff line change
Expand Up @@ -267,7 +267,7 @@ int main(int argc, char *argv[])
.params.tty_client.timeout = 3000,
};

events = neu_event_new();
events = neu_event_new("modbus_tty_simulator");
conn = neu_conn_new(&param, NULL, connected, disconnected);

neu_conn_start(conn);
Expand Down
2 changes: 1 addition & 1 deletion src/adapter/adapter.c
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,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;
Expand Down
2 changes: 1 addition & 1 deletion src/adapter/driver/driver.c
Original file line number Diff line number Diff line change
Expand Up @@ -434,7 +434,7 @@ 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->driver_events = neu_event_new("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;
Expand Down
2 changes: 1 addition & 1 deletion src/connection/connection_eth.c
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
2 changes: 1 addition & 1 deletion src/connection/mqtt_client.c
Original file line number Diff line number Diff line change
Expand Up @@ -907,7 +907,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;
}
Expand Down
2 changes: 1 addition & 1 deletion src/core/manager.c
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,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();
Expand Down
37 changes: 21 additions & 16 deletions src/event/event_linux.c
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ struct neu_events {
int epoll_fd;
pthread_t thread;
bool stop;
char * name;

pthread_mutex_t mtx;
int n_event;
Expand Down Expand Up @@ -122,8 +123,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 loop(%s) exit, errno: %s(%d), stop: %d",
events->name, strerror(errno), errno, events->stop);
break;
}

Expand Down Expand Up @@ -177,17 +178,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(%s): %d(%d)", name, events->epoll_fd, errno);
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);
Expand All @@ -199,6 +201,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);
Expand All @@ -220,7 +223,8 @@ neu_event_timer_t *neu_event_add_timer(neu_events_t * events,
};
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);

Expand Down Expand Up @@ -251,18 +255,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);
Expand All @@ -281,8 +286,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;
Expand All @@ -303,8 +308,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;
Expand All @@ -316,8 +321,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);
Expand Down
2 changes: 1 addition & 1 deletion src/otel/otel_manager.c
Original file line number Diff line number Diff line change
Expand Up @@ -1113,7 +1113,7 @@ static int otel_timer_cb(void *data)
void neu_otel_start()
{
if (otel_event == NULL) {
otel_event = neu_event_new();
otel_event = neu_event_new("otel_event");
}

if (otel_timer == NULL) {
Expand Down
Loading