summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--src/ev_embed_compat.c1
-rw-r--r--src/evwrap.h23
-rw-r--r--src/loop.c2
-rw-r--r--src/multiplex.c2
-rw-r--r--src/packet_pool.c2
-rw-r--r--src/packet_stream.c64
-rw-r--r--src/poll.c2
-rw-r--r--src/rpc.c13
-rw-r--r--src/rpc2.h2
-rw-r--r--src/socket/socket.h2
-rw-r--r--src/socket/socket_linux.c10
-rw-r--r--src/socket/socket_windows.c19
-rw-r--r--src/timer.c2
-rw-r--r--src/util/buffer.c7
-rw-r--r--src/util/error.c17
-rw-r--r--src/util/error.h7
-rw-r--r--src/util/packet.c9
-rw-r--r--src/util/packet.h10
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
diff --git a/src/loop.c b/src/loop.c
index 07d55a4..774feaf 100644
--- a/src/loop.c
+++ b/src/loop.c
@@ -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);
diff --git a/src/poll.c b/src/poll.c
index caea5d9..dd5c9e2 100644
--- a/src/poll.c
+++ b/src/poll.c
@@ -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)
diff --git a/src/rpc.c b/src/rpc.c
index d57bde5..1b00e63 100644
--- a/src/rpc.c
+++ b/src/rpc.c
@@ -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
}));
diff --git a/src/rpc2.h b/src/rpc2.h
index ecfb91a..e6c2098 100644
--- a/src/rpc2.h
+++ b/src/rpc2.h
@@ -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);