diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/ev_embed_compat.c | 1 | ||||
| -rw-r--r-- | src/evwrap.h | 23 | ||||
| -rw-r--r-- | src/loop.c | 2 | ||||
| -rw-r--r-- | src/multiplex.c | 2 | ||||
| -rw-r--r-- | src/packet_pool.c | 2 | ||||
| -rw-r--r-- | src/packet_stream.c | 64 | ||||
| -rw-r--r-- | src/poll.c | 2 | ||||
| -rw-r--r-- | src/rpc.c | 13 | ||||
| -rw-r--r-- | src/rpc2.h | 2 | ||||
| -rw-r--r-- | src/socket/socket.h | 2 | ||||
| -rw-r--r-- | src/socket/socket_linux.c | 10 | ||||
| -rw-r--r-- | src/socket/socket_windows.c | 19 | ||||
| -rw-r--r-- | src/timer.c | 2 | ||||
| -rw-r--r-- | src/util/buffer.c | 7 | ||||
| -rw-r--r-- | src/util/error.c | 17 | ||||
| -rw-r--r-- | src/util/error.h | 7 | ||||
| -rw-r--r-- | src/util/packet.c | 9 | ||||
| -rw-r--r-- | src/util/packet.h | 10 |
18 files changed, 124 insertions, 70 deletions
diff --git a/src/ev_embed_compat.c b/src/ev_embed_compat.c index 0400aeb..9243de2 100644 --- a/src/ev_embed_compat.c +++ b/src/ev_embed_compat.c @@ -17,6 +17,7 @@ AL_IGNORE_WARNING("-Wstrict-aliasing") AL_IGNORE_WARNING("-Wcomment") AL_IGNORE_WARNING("-Wparentheses") +AL_IGNORE_WARNING("-Wunused-variable") AL_IGNORE_WARNING("-Wunused-value") AL_IGNORE_WARNING("-Wunused-function") AL_IGNORE_WARNING("-Wunused-parameter") diff --git a/src/evwrap.h b/src/evwrap.h index 35b72cc..1d3a44b 100644 --- a/src/evwrap.h +++ b/src/evwrap.h @@ -17,22 +17,24 @@ #include <ev.h> -AL_IGNORE_WARNING("-Wstrict-aliasing") - -// Wrap common libev macros so we can ignore aliasing warnings. +// Wrap common libev macros to ignore aliasing warnings. -#define ev_init_n(type) ev_init_n_##type +AL_IGNORE_WARNING("-Wstrict-aliasing") -static inline void ev_init_n_ev_io(ev_io *ev, void (*callback)(struct ev_loop *, ev_io *, s32)) +static inline void ev_init_nn(ev_watcher *ev, void (*callback)(struct ev_loop *, ev_watcher *, s32)) { ev_init(ev, callback); } -static inline void ev_init_n_ev_timer(ev_timer *ev, void (*callback)(struct ev_loop *, ev_timer *, s32)) +#define ev_init_n(w, callback) ev_init_nn((ev_watcher *)w, (void (*)(struct ev_loop *, ev_watcher *, s32))callback) + +static inline bool ev_is_active_nn(ev_watcher *ev) { - ev_init(ev, callback); + return ev_is_active(ev); } +#define ev_is_active_n(w) ev_is_active_nn((ev_watcher *)w) + static inline void ev_io_init_n(ev_io *io, void (*callback)(struct ev_loop *, ev_io *, s32), s32 fd, s32 events) { ev_io_init(io, callback, fd, events); @@ -48,11 +50,4 @@ static inline void ev_async_init_n(ev_async *asyn, void (*callback)(struct ev_lo ev_async_init(asyn, callback); } -#define ev_is_active_n(type) ev_is_active_n_##type - -static inline bool ev_is_active_n_ev_async(ev_async *asyn) -{ - return ev_is_active(asyn); -} - AL_IGNORE_WARNING_END @@ -6,6 +6,8 @@ bool nn_event_loop_init(struct nn_event_loop *loop) { loop->ev = ev_default_loop(EVFLAG_AUTO); //loop->ev = ev_loop_new(EVFLAG_AUTO); + // @TODO: + // https://pod.tst.eu/http://cvs.schmorp.de/libev/ev.pod @ EVBACKEND_SELECT //ev_set_timeout_collect_interval(loop->ev, 0.1); //ev_set_io_collect_interval(loop->ev, 0.05); return true; diff --git a/src/multiplex.c b/src/multiplex.c index 4378884..a9eb5b4 100644 --- a/src/multiplex.c +++ b/src/multiplex.c @@ -24,7 +24,7 @@ static void socket_read_callback(struct ev_loop *loop, ev_io *w, s32 revents) (void)revents; u8 id; - ssize_t ret = nn_socket_read(&conn->sock, &id, sizeof(u8)); + ssize_t ret = nn_socket_read(&conn->sock, &id, 1); if (ret <= 0) { ev_io_stop(loop, &conn->event); nn_socket_close(&conn->sock); diff --git a/src/packet_pool.c b/src/packet_pool.c index 7f52e63..7d9ac4d 100644 --- a/src/packet_pool.c +++ b/src/packet_pool.c @@ -154,7 +154,7 @@ void nn_packet_pool_free(struct nn_packet_pool *pool) if (!pool->disabled) { ev_async_stop(pool->loop->ev, &pool->signal); } else { - al_assert(!ev_is_active_n(ev_async)(&pool->signal)); + al_assert(!ev_is_active_n(&pool->signal)); } nn_mutex_unlock(&pool->mutex); struct nn_packet *packet; diff --git a/src/packet_stream.c b/src/packet_stream.c index cca1b0f..b183df5 100644 --- a/src/packet_stream.c +++ b/src/packet_stream.c @@ -19,26 +19,6 @@ static inline void init_io_state(struct nn_packet_stream *stream) stream->out.running = false; } -void nn_packet_stream_init(struct nn_packet_stream *stream, - bool (*connection_callback)(void *, struct nn_packet_stream *), - void (*connection_closed_callback)(void *, struct nn_packet_stream *), void *userdata) -{ - init_io_state(stream); - stream->direct = NULL; - stream->connection_callback = connection_callback; - stream->connection_closed_callback = connection_closed_callback; - stream->packet_callback = NULL; - stream->packet_sent_callback = NULL; - stream->packets_sent_callback = NULL; - stream->userdata = userdata; - stream->connect = PACKET_STREAM_DISCONNECTED; -} - -void nn_packet_stream_set_nodelay(struct nn_packet_stream *stream, s32 nodelay) -{ - nn_socket_set_nodelay(&stream->sock, nodelay); -} - static void start_write_internal(struct nn_packet_stream *stream) { stream->out.running = true; @@ -60,6 +40,7 @@ static void shutdown_internal(struct nn_packet_stream *stream) static void discard_out_queue_internal(struct nn_packet_stream *stream) { + al_assert(!stream->out.running); if (stream->out.packet) { stream->packet_sent_callback(stream->userdata, stream->out.packet); stream->out.packet = NULL; @@ -148,9 +129,10 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents) al_assert(revents & EV_WRITE); // Assert on EV_ERROR. if (UNLIKELY(stream->connect == PACKET_STREAM_CONNECTING)) { - ssize_t ret = nn_socket_write(&stream->sock, &stream->id, sizeof(u8)); + ssize_t ret = nn_socket_write(&stream->sock, &stream->id, 1); if (ret < 0) { - // It seems POLLOUT can be signaled even if any write() will error. + // EV_WRITE may be signaled even if any subsequent write()/send() will error. + nn_socket_check_error(ret); stop_write_internal(stream); stop_internal(stream); return; @@ -189,6 +171,30 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents) } } +void nn_packet_stream_init(struct nn_packet_stream *stream, + bool (*connection_callback)(void *, struct nn_packet_stream *), + void (*connection_closed_callback)(void *, struct nn_packet_stream *), void *userdata) +{ + init_io_state(stream); + stream->revent.data = stream; + stream->wevent.data = stream; + ev_init_n(&stream->revent, stream_read_callback); + ev_init_n(&stream->wevent, stream_write_callback); + stream->direct = NULL; + stream->connection_callback = connection_callback; + stream->connection_closed_callback = connection_closed_callback; + stream->packet_callback = NULL; + stream->packet_sent_callback = NULL; + stream->packets_sent_callback = NULL; + stream->userdata = userdata; + stream->connect = PACKET_STREAM_DISCONNECTED; +} + +void nn_packet_stream_set_nodelay(struct nn_packet_stream *stream, s32 nodelay) +{ + nn_socket_set_nodelay(&stream->sock, nodelay); +} + void nn_packet_stream_from_socket(struct nn_packet_stream *stream, struct nn_event_loop *loop, struct nn_socket *sock) { al_assert(stream->connection_callback && stream->connection_closed_callback); @@ -196,10 +202,11 @@ void nn_packet_stream_from_socket(struct nn_packet_stream *stream, struct nn_eve stream->sock = *sock; nn_socket_set_blocking(&stream->sock, false); init_io_state(stream); - stream->wevent.data = stream; - ev_io_init_n(&stream->wevent, stream_write_callback, nn_socket_get_fd(&stream->sock), EV_WRITE); stream->revent.data = stream; - ev_io_init_n(&stream->revent, stream_read_callback, nn_socket_get_fd(&stream->sock), EV_READ); + stream->wevent.data = stream; + s32 fd = nn_socket_get_fd(&stream->sock); + ev_io_init_n(&stream->revent, stream_read_callback, fd, EV_READ); + ev_io_init_n(&stream->wevent, stream_write_callback, fd, EV_WRITE); if (stream_connected(stream)) { ev_io_start(stream->loop->ev, &stream->revent); } @@ -216,10 +223,9 @@ static void do_connect_internal(struct nn_packet_stream *stream, str *addr, u16 stream->connection_closed_callback(stream->userdata, stream); return; } - stream->revent.data = stream; - ev_io_init_n(&stream->revent, stream_read_callback, nn_socket_get_fd(&stream->sock), EV_READ); - stream->wevent.data = stream; - ev_io_init_n(&stream->wevent, stream_write_callback, nn_socket_get_fd(&stream->sock), EV_WRITE); + s32 fd = nn_socket_get_fd(&stream->sock); + ev_io_set(&stream->revent, fd, EV_READ); + ev_io_set(&stream->wevent, fd, EV_WRITE); // Start non-blocking connection. stream->connect = PACKET_STREAM_CONNECTING; start_write_internal(stream); @@ -12,7 +12,7 @@ void nn_poll_init(struct nn_poll *poll, void (*callback)(void *, s32 revents), v poll->callback = callback; poll->userdata = userdata; poll->event.data = poll; - ev_init_n(ev_io)(&poll->event, event_callback); + ev_init_n(&poll->event, event_callback); } void nn_poll_set(struct nn_poll *poll, s32 fd, s32 events) @@ -24,8 +24,8 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str { struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata; al_assert(conn->stream == stream); - s8 op = nn_packet_read_s8(packet); u32 id = nn_packet_read_u32(packet); + s8 op = nn_packet_read_s8(packet); if (op == -1) { // Response struct nn_rpc_callback *callback; al_array_foreach_ptr(conn->callbacks, i, callback) { @@ -40,8 +40,8 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str al_array_foreach_ptr(conn->rpc->commands, i, command) { if (command->op == op) { struct nn_packet *rpacket = nn_packet_create(); - nn_packet_write_s8(rpacket, -1); nn_packet_write_u32(rpacket, id); + nn_packet_write_s8(rpacket, -1); if (command->callback(command->userdata, conn, packet, rpacket)) { conn->outgoing++; nn_packet_stream_send_packet(stream, rpacket); @@ -107,7 +107,7 @@ void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream) conn->stream->userdata = conn; } -void nn_rpc_prepare_client(struct nn_rpc *rpc) +struct nn_rpc_connection *nn_rpc_prepare_client(struct nn_rpc *rpc) { al_assert(!rpc->conn); struct nn_rpc_connection *conn = al_alloc_object(struct nn_rpc_connection); @@ -116,6 +116,7 @@ void nn_rpc_prepare_client(struct nn_rpc *rpc) nn_packet_stream_init(conn->stream, stream_connection_callback, stream_connection_closed_callback, conn); rpc->conn = conn; al_array_push(rpc->connections, conn); + return rpc->conn; } void nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port) @@ -130,12 +131,12 @@ void nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port) nn_packet_stream_reconnect(rpc->conn->stream, addr, port); } +// @TODO: Pack opcode into a u32, reduce ID by 8 bits (make struct with bitmask) struct nn_packet *nn_rpc_get_packet(struct nn_rpc *rpc, s8 op) { struct nn_packet *packet = nn_packet_create(); + nn_packet_write_u32(packet, (rpc->increment = al_u32_inc_wrap(rpc->increment))); nn_packet_write_s8(packet, op); - nn_packet_write_u32(packet, rpc->increment); - rpc->increment = al_u32_inc_wrap(rpc->increment); return packet; } @@ -157,7 +158,7 @@ void nn_rpc_connection_command(struct nn_rpc_connection *conn, struct nn_packet { if (callback) { al_array_push(conn->callbacks, ((struct nn_rpc_callback){ - .id = nn_packet_get_u32(packet, NNWT_PACKET_HEADER_LENGTH + sizeof(s8)), + .id = nn_packet_get_u32(packet, NNWT_PACKET_HEADER_LENGTH), .callback = callback, .userdata = userdata })); @@ -43,7 +43,7 @@ bool nn_rpc_init(struct nn_rpc *rpc, struct nn_event_loop *loop, void (*connection_closed_callback)(void *, struct nn_rpc_connection *), void *userdata); void nn_rpc_add_command(struct nn_rpc *rpc, struct nn_rpc_command *command); void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream); -void nn_rpc_prepare_client(struct nn_rpc *rpc); +struct nn_rpc_connection *nn_rpc_prepare_client(struct nn_rpc *rpc); void nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port); void nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port); struct nn_packet *nn_rpc_get_packet(struct nn_rpc *rpc, s8 op); diff --git a/src/socket/socket.h b/src/socket/socket.h index e1b203a..3872bd4 100644 --- a/src/socket/socket.h +++ b/src/socket/socket.h @@ -80,5 +80,7 @@ ssize_t nn_socket_write(struct nn_socket *sock, void *buf, size_t size); ssize_t nn_socket_sendto(struct nn_socket *sock, void *buf, size_t size); ssize_t nn_socket_recvfrom(struct nn_socket *sock, void *buf, size_t size); +bool nn_socket_check_error(ssize_t ret); + void nn_socket_shutdown(struct nn_socket *sock); void nn_socket_close(struct nn_socket *sock); diff --git a/src/socket/socket_linux.c b/src/socket/socket_linux.c index d854644..70d9184 100644 --- a/src/socket/socket_linux.c +++ b/src/socket/socket_linux.c @@ -277,6 +277,16 @@ ssize_t nn_socket_recvfrom(struct nn_socket *sock, void *buf, size_t size) return recvfrom(sock->fd, buf, size, 0, sock->addrinfo->ai_addr, &sock->addrinfo->ai_addrlen); } +// This would normally be a log_warn() but we rely on expected errors for control flow. +bool nn_socket_check_error(ssize_t ret) +{ + if (ret < 0) { + log_debug("Socket error: %s (%d).", nn_strerror(errno), errno); + return true; + } + return false; +} + s32 nn_socket_get_fd(struct nn_socket *sock) { return sock->fd; diff --git a/src/socket/socket_windows.c b/src/socket/socket_windows.c index 123640e..69fa55b 100644 --- a/src/socket/socket_windows.c +++ b/src/socket/socket_windows.c @@ -1,6 +1,8 @@ #define AL_LOG_SECTION "socket" #include <al/log.h> +#include "../util/error.h" + #include "socket.h" #include "socket_internal.h" #include "net.h" @@ -125,12 +127,12 @@ s32 nn_socket_get_fd(struct nn_socket *sock) ssize_t nn_socket_read(struct nn_socket *sock, void *buf, size_t size) { - return recv(sock->fd, buf, (s32)size, 0); + return (ssize_t)recv(sock->fd, buf, (s32)size, 0); } ssize_t nn_socket_write(struct nn_socket *sock, void *buf, size_t size) { - return send(sock->fd, buf, (s32)size, 0); + return (ssize_t)send(sock->fd, buf, (s32)size, 0); } ssize_t nn_socket_sendto(struct nn_socket *sock, void *buf, size_t size) @@ -149,6 +151,17 @@ ssize_t nn_socket_recvfrom(struct nn_socket *sock, void *buf, size_t size) return 0; } +bool nn_socket_check_error(ssize_t ret) +{ + if (ret == SOCKET_ERROR) { + s32 err = WSAGetLastError(); + char *strerror = nn_win32_error_message(err); + log_debug("Socket error: %s (%d).", strerror ? strerror : "(None)", err); + return true; + } + return false; +} + void nn_socket_shutdown(struct nn_socket *sock) { if (shutdown(sock->fd, SD_BOTH) == SOCKET_ERROR) {} @@ -157,5 +170,7 @@ void nn_socket_shutdown(struct nn_socket *sock) void nn_socket_close(struct nn_socket *sock) { _close(sock->internal_fd); + // https://learn.microsoft.com/en-us/cpp/c-runtime-library/reference/close?view=msvc-170 + // Based on these docs, I assume calling closesocket() is not necessary. //closesocket(sock->fd); } diff --git a/src/timer.c b/src/timer.c index 2f35354..c0759cc 100644 --- a/src/timer.c +++ b/src/timer.c @@ -16,7 +16,7 @@ void nn_timer_init(struct nn_timer *timer, struct nn_event_loop *loop, timer->userdata = userdata; timer->timer.data = timer; timer->disabled = false; - ev_init_n(ev_timer)(&timer->timer, timer_callback); + ev_init_n(&timer->timer, timer_callback); } void nn_timer_set_repeat(struct nn_timer *timer, nn_os_tstamp repeat) diff --git a/src/util/buffer.c b/src/util/buffer.c index 4c8b3c2..ab7adac 100644 --- a/src/util/buffer.c +++ b/src/util/buffer.c @@ -30,9 +30,12 @@ void nn_buffer_shrink_to_size(struct nn_buffer *buffer) void nn_buffer_write(struct nn_buffer *buffer, void *data, size_t index, size_t size) { + al_assert(buffer->size <= buffer->alloc); size_t reach = index + size; - nn_buffer_ensure_space(buffer, reach); - if (reach > buffer->size) buffer->size = reach; + if (reach > buffer->size) { + nn_buffer_ensure_space(buffer, reach); + buffer->size = reach; + } al_memcpy(&buffer->data[index], data, size); } diff --git a/src/util/error.c b/src/util/error.c index db2d001..61c6136 100644 --- a/src/util/error.c +++ b/src/util/error.c @@ -15,3 +15,20 @@ char *nn_strerror(s32 errnum) #endif return errorbuf; } + +#ifdef NAUNET_ON_WINDOWS +#include "../winwrap.h" + +// https://stackoverflow.com/a/46104456 +char *nn_win32_error_message(s32 err) +{ + DWORD ret = FormatMessage(FORMAT_MESSAGE_FROM_SYSTEM | FORMAT_MESSAGE_IGNORE_INSERTS, + NULL, + err, + MAKELANGID(LANG_NEUTRAL, SUBLANG_DEFAULT), + errorbuf, + NNWT_ERROR_STR_MAXLEN, + NULL); + return !ret ? NULL : errorbuf; +} +#endif diff --git a/src/util/error.h b/src/util/error.h index 3731a01..ef6fd46 100644 --- a/src/util/error.h +++ b/src/util/error.h @@ -3,6 +3,13 @@ #include <al/lib.h> #include <errno.h> +#ifdef NAUNET_ON_WINDOWS +#define NNWT_ERROR_STR_MAXLEN 256 +#else #define NNWT_ERROR_STR_MAXLEN 94 +#endif char *nn_strerror(s32 errnum); +#ifdef NAUNET_ON_WINDOWS +char *nn_win32_error_message(s32 err); +#endif diff --git a/src/util/packet.c b/src/util/packet.c index d7df1da..c991a7f 100644 --- a/src/util/packet.c +++ b/src/util/packet.c @@ -17,7 +17,6 @@ struct nn_packet *nn_packet_clone(struct nn_packet *packet) al_memcpy(nn_buffer_get_ptr(&c->buffer, 0), nn_buffer_get_ptr(&packet->buffer, 0), size); c->windex = packet->windex; c->rindex = packet->rindex; - c->opaque = packet->opaque; return c; } @@ -26,14 +25,6 @@ void nn_packet_reset(struct nn_packet *packet) nn_buffer_ensure_space(&packet->buffer, NNWT_PACKET_HEADER_LENGTH); packet->rindex = NNWT_PACKET_HEADER_LENGTH; packet->windex = NNWT_PACKET_HEADER_LENGTH; - packet->opaque = NULL; -} - -u32 nn_packet_get_u32(struct nn_packet *packet, u32 index) -{ - u32 value; - NNWT_BUFFER_READ_TYPE(&packet->buffer, index, u32, value); - return value; } void nn_packet_write_size(struct nn_packet *packet) diff --git a/src/util/packet.h b/src/util/packet.h index a2837f1..1f76559 100644 --- a/src/util/packet.h +++ b/src/util/packet.h @@ -13,7 +13,6 @@ struct nn_packet { struct nn_buffer buffer; u32 rindex; u32 windex; - void *opaque; }; #define NNWT_PACKET_WRITE_TYPE(p, type, v) \ @@ -32,12 +31,17 @@ struct nn_packet { r = (__typeof__(r))nn_buffer_get_ptr(&(p)->buffer, (p)->rindex); \ (p)->rindex += (u32)length +static inline u32 nn_packet_get_u32(struct nn_packet *packet, u32 index) +{ + u32 value; + NNWT_BUFFER_READ_TYPE(&packet->buffer, index, u32, value); + return value; +} + struct nn_packet *nn_packet_create(void); struct nn_packet *nn_packet_clone(struct nn_packet *packet); void nn_packet_reset(struct nn_packet *packet); -u32 nn_packet_get_u32(struct nn_packet *packet, u32 index); - void nn_packet_write_size(struct nn_packet *packet); u32 nn_packet_get_size(struct nn_packet *packet); |