Browse Source

uasync: ре-валидация node перед последующими коллбэками (счётчик cb_called) вместо снапшота

v2
evgeny 3 weeks ago
parent
commit
187bcd80b9
  1. 122
      lib/u_async.c

122
lib/u_async.c

@ -1158,22 +1158,11 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
struct socket_node* node = &ua->sockets->sockets[idx];
if (!node->active) continue;
/* Снапшот полей — коллбэки могут реаллоцировать ua->sockets (node станет висячим) */
int local_fd = node->fd;
socket_t local_sock = node->sock;
int local_type = node->type;
void* local_ud = node->user_data;
socket_callback_t local_except = node->except_cbk;
socket_callback_t local_read = node->read_cbk;
socket_callback_t local_write = node->write_cbk;
socket_t_callback_t local_read_sock = node->read_cbk_sock;
socket_t_callback_t local_write_sock = node->write_cbk_sock;
SOCKET s;
if (local_type == SOCKET_NODE_TYPE_SOCK) {
s = local_sock;
if (node->type == SOCKET_NODE_TYPE_SOCK) {
s = node->sock;
} else {
s = (SOCKET)local_fd;
s = (SOCKET)node->fd;
}
int has_read = FD_ISSET(s, &read_fds);
@ -1183,34 +1172,49 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
if (!has_read && !has_write && !has_except) continue;
DEBUG_DEBUG(DEBUG_CATEGORY_SYS, "select→fd=%d r=%d w=%d e=%d", (int)s, has_read, has_write, has_except);
/* Коллбэки могут удалить/добавить сокет (realloc массива) — перед
* каждым последующим коллбэком пере-валидируем node по индексу. */
int cb_called = 0;
if (has_except) {
if (local_except) {
local_except(local_fd, local_ud);
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
cb_called++;
}
if (has_read) {
if (local_type == SOCKET_NODE_TYPE_SOCK) {
if (local_read_sock) {
local_read_sock(local_sock, local_ud);
if (cb_called) {
node = &ua->sockets->sockets[idx];
if (!node->active) continue;
}
if (node->type == SOCKET_NODE_TYPE_SOCK) {
if (node->read_cbk_sock) {
node->read_cbk_sock(node->sock, node->user_data);
}
} else {
if (local_read) {
local_read(local_fd, local_ud);
if (node->read_cbk) {
node->read_cbk(node->fd, node->user_data);
}
}
cb_called++;
}
if (has_write) {
if (local_type == SOCKET_NODE_TYPE_SOCK) {
if (local_write_sock) {
local_write_sock(local_sock, local_ud);
if (cb_called) {
node = &ua->sockets->sockets[idx];
if (!node->active) continue;
}
if (node->type == SOCKET_NODE_TYPE_SOCK) {
if (node->write_cbk_sock) {
node->write_cbk_sock(node->sock, node->user_data);
}
} else {
if (local_write) {
local_write(local_fd, local_ud);
if (node->write_cbk) {
node->write_cbk(node->fd, node->user_data);
}
}
cb_called++;
}
}
}
@ -1254,63 +1258,75 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
}
/* Socket event - lookup by fd */
struct socket_node* node = socket_array_get(ua->sockets, ua->poll_fds[i].fd);
int fd = ua->poll_fds[i].fd;
struct socket_node* node = socket_array_get(ua->sockets, fd);
if (!node) { // Try by socket_t (in case this is a socket)
socket_t lookup_sock = ua->poll_fds[i].fd;
node = socket_array_get_by_sock(ua->sockets, lookup_sock);
node = socket_array_get_by_sock(ua->sockets, fd);
}
if (!node) continue; // Socket may have been removed
/* Снапшот полей — коллбэки могут реаллоцировать ua->sockets (node станет висячим) */
int local_fd = node->fd;
socket_t local_sock = node->sock;
int local_type = node->type;
void* local_ud = node->user_data;
socket_callback_t local_except = node->except_cbk;
socket_callback_t local_read = node->read_cbk;
socket_callback_t local_write = node->write_cbk;
socket_t_callback_t local_read_sock = node->read_cbk_sock;
socket_t_callback_t local_write_sock = node->write_cbk_sock;
/* Коллбэки могут удалить/добавить сокет (realloc массива) — перед
* каждым последующим коллбэком пере-валидируем node по fd. */
int cb_called = 0;
/* Read readiness BEFORE error — avoid losing data on combined IN+ERR/HUP events */
if (ua->poll_fds[i].revents & POLLIN) {
if (local_type == SOCKET_NODE_TYPE_SOCK) {
if (local_read_sock) {
local_read_sock(local_sock, local_ud);
if (node->type == SOCKET_NODE_TYPE_SOCK) {
if (node->read_cbk_sock) {
node->read_cbk_sock(node->sock, node->user_data);
}
} else {
if (local_read) {
local_read(local_fd, local_ud);
if (node->read_cbk) {
node->read_cbk(node->fd, node->user_data);
}
}
cb_called++;
}
/* Write readiness BEFORE error — flush pending writes before handling HUP */
if (ua->poll_fds[i].revents & POLLOUT) {
if (local_type == SOCKET_NODE_TYPE_SOCK) {
if (local_write_sock) {
local_write_sock(local_sock, local_ud);
if (cb_called) {
node = socket_array_get(ua->sockets, fd);
if (!node) node = socket_array_get_by_sock(ua->sockets, fd);
if (!node) continue;
}
if (node->type == SOCKET_NODE_TYPE_SOCK) {
if (node->write_cbk_sock) {
node->write_cbk_sock(node->sock, node->user_data);
}
} else {
if (local_write) {
local_write(local_fd, local_ud);
if (node->write_cbk) {
node->write_cbk(node->fd, node->user_data);
}
}
cb_called++;
}
/* Check for error conditions LAST — I/O handlers drain/process data first */
if (ua->poll_fds[i].revents & (POLLERR | POLLHUP | POLLNVAL)) {
if (cb_called) {
node = socket_array_get(ua->sockets, fd);
if (!node) node = socket_array_get_by_sock(ua->sockets, fd);
if (!node) continue;
}
/* Treat as exceptional condition */
if (local_except) {
local_except(local_fd, local_ud);
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
cb_called++;
}
/* Exceptional data (out-of-band) */
if (ua->poll_fds[i].revents & POLLPRI) {
if (local_except) {
local_except(local_fd, local_ud);
if (cb_called) {
node = socket_array_get(ua->sockets, fd);
if (!node) node = socket_array_get_by_sock(ua->sockets, fd);
if (!node) continue;
}
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
cb_called++;
}
}
}

Loading…
Cancel
Save