summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-03-31 19:03:55 -0400
committerAndrew Opalach <andrew@akon.city> 2025-03-31 19:03:55 -0400
commitde6f6090f94b14ed98dac71956ff1bede48f1d8e (patch)
treeeac747f2c7620a8a360f88f83306856527ae2988 /src
parentb4e45dd89e6dd53d8cb940e5609b812e19ade913 (diff)
downloadlibnaunet-de6f6090f94b14ed98dac71956ff1bede48f1d8e.tar.gz
libnaunet-de6f6090f94b14ed98dac71956ff1bede48f1d8e.tar.bz2
libnaunet-de6f6090f94b14ed98dac71956ff1bede48f1d8e.zip
Apply style changes, abide by strict aliasing
- Remove internal sending step in packet pool. Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src')
-rw-r--r--src/curl/curl.c30
-rw-r--r--src/curl/http.c66
-rw-r--r--src/curl/websocket.c15
-rw-r--r--src/ev_embed_compat.c38
-rw-r--r--src/evwrap.h40
-rw-r--r--src/fs_event/fs_event_inotify.c6
-rw-r--r--src/line_processor.c8
-rw-r--r--src/multiplex.c16
-rw-r--r--src/packet_cache.c11
-rw-r--r--src/packet_pool.c58
-rw-r--r--src/packet_pool.h2
-rw-r--r--src/packet_stream.c44
-rw-r--r--src/poll.c2
-rw-r--r--src/rpc.c11
-rw-r--r--src/signal.c2
-rw-r--r--src/socket/socket.h3
-rw-r--r--src/socket/socket_internal.h11
-rw-r--r--src/socket/socket_linux.c41
-rw-r--r--src/socket/socket_windows.c16
-rw-r--r--src/timer.c2
-rw-r--r--src/util/buffer.c4
-rw-r--r--src/util/buffer.h4
-rw-r--r--src/util/file/file.h10
-rw-r--r--src/util/file/file_linux.c16
-rw-r--r--src/util/file/file_stdio.c18
-rw-r--r--src/util/packet.c22
-rw-r--r--src/util/packet.h17
-rw-r--r--src/util/thread/thread_linux.c19
-rw-r--r--src/util/thread/thread_windows.c2
-rw-r--r--src/winwrap.c3
30 files changed, 220 insertions, 317 deletions
diff --git a/src/curl/curl.c b/src/curl/curl.c
index 568f82e..5a5ac01 100644
--- a/src/curl/curl.c
+++ b/src/curl/curl.c
@@ -1,3 +1,4 @@
+#define AL_LOG_SECTION "curl"
#include <al/log.h>
#include "curl.h"
@@ -10,11 +11,11 @@ static void curl_socket_action_callback(struct ev_loop *loop, ev_io *w, s32 reve
{
(void)loop;
struct nn_curl *curl = (struct nn_curl *)w->data;
- s32 action = (revents & EV_READ ? CURL_CSELECT_IN : 0) | (revents & EV_WRITE ? CURL_CSELECT_OUT : 0);
+ s32 action = ((revents & EV_READ) ? CURL_CSELECT_IN : 0) | ((revents & EV_WRITE) ? CURL_CSELECT_OUT : 0);
s32 running;
CURLMcode mc = curl_multi_socket_action(curl->multi_handle, curl->sock, action, &running);
if (mc != CURLM_OK) {
- al_log_error("curl", "curl_multi_socket_action() failed (%s).", curl_multi_strerror(mc));
+ error("curl_multi_socket_action() failed (%s).", curl_multi_strerror(mc));
return;
}
curl->handle_events(curl->userdata, curl);
@@ -26,8 +27,8 @@ static void set_sock(struct nn_curl *curl, s32 what)
if (curl->started) {
ev_io_stop(curl->loop->ev, &curl->event);
}
- s32 action = (what & CURL_POLL_IN ? EV_READ : 0) | (what & CURL_POLL_OUT ? EV_WRITE : 0);
- ev_io_init(&curl->event, curl_socket_action_callback, curl->sock, action);
+ s32 action = ((what & CURL_POLL_IN) ? EV_READ : 0) | ((what & CURL_POLL_OUT) ? EV_WRITE : 0);
+ ev_io_init_n(&curl->event, curl_socket_action_callback, curl->sock, action);
ev_io_start(curl->loop->ev, &curl->event);
curl->started = true;
}
@@ -55,12 +56,10 @@ static void timeout_callback(struct ev_loop *loop, ev_timer *w, s32 revents)
(void)loop;
struct nn_curl *curl = (struct nn_curl *)w->data;
(void)revents;
-
if (curl->timer_started) {
ev_timer_stop(curl->loop->ev, w);
curl->timer_started = false;
}
-
s32 running;
curl_multi_socket_action(curl->multi_handle, CURL_SOCKET_TIMEOUT, 0, &running);
curl->handle_events(curl->userdata, curl);
@@ -70,22 +69,15 @@ static s32 timer_callback(CURLM *multi, s64 timeout_ms, void *userp)
{
(void)multi;
struct nn_curl *curl = (struct nn_curl *)userp;
-
if (curl->timer_started) {
ev_timer_stop(curl->loop->ev, &curl->timer);
curl->timer_started = false;
}
-
- if (timeout_ms == -1) {
- return 0;
- }
-
+ if (timeout_ms == -1) return 0;
curl->timer.data = curl;
- ev_timer_init(&curl->timer, timeout_callback, timeout_ms / 1000.0, 0.0);
-
+ ev_timer_init_n(&curl->timer, timeout_callback, timeout_ms / 1000.0, 0.0);
curl->timer_started = true;
ev_timer_start(curl->loop->ev, &curl->timer);
-
return 0;
}
@@ -95,24 +87,20 @@ bool nn_curl_init(struct nn_curl *curl)
if (!curl->handle) goto err;
//curl_easy_setopt(curl->handle, CURLOPT_VERBOSE, 1L);
curl_easy_setopt(curl->handle, CURLOPT_NOPROGRESS, 1L);
-
curl->multi_handle = curl_multi_init();
if (!curl->multi_handle) goto err;
curl_multi_setopt(curl->multi_handle, CURLMOPT_SOCKETFUNCTION, sock_callback);
curl_multi_setopt(curl->multi_handle, CURLMOPT_SOCKETDATA, curl);
curl_multi_setopt(curl->multi_handle, CURLMOPT_TIMERFUNCTION, timer_callback);
curl_multi_setopt(curl->multi_handle, CURLMOPT_TIMERDATA, curl);
-
curl->loop = NULL;
curl->event.data = curl;
curl->started = false;
curl->timer_started = false;
curl->added = false;
-
return true;
err:
nn_curl_close(curl);
-
return false;
}
@@ -130,7 +118,7 @@ bool nn_curl_add_handle(struct nn_curl *curl)
// a callback is running on another thread.
CURLMcode mc = curl_multi_add_handle(curl->multi_handle, curl->handle);
if (mc != CURLM_OK) {
- al_log_error("curl", "curl_multi_add_handle() failed (%s).", curl_multi_strerror(mc));
+ error("curl_multi_add_handle() failed (%s).", curl_multi_strerror(mc));
return false;
}
curl->added = true;
@@ -142,7 +130,7 @@ bool nn_curl_remove_handle(struct nn_curl *curl)
al_assert(curl->added);
CURLMcode mc = curl_multi_remove_handle(curl->multi_handle, curl->handle);
if (mc != CURLM_OK) {
- al_log_error("curl", "curl_multi_remove_handle() failed (%s).", curl_multi_strerror(mc));
+ error("curl_multi_remove_handle() failed (%s).", curl_multi_strerror(mc));
return false;
}
curl->added = false;
diff --git a/src/curl/http.c b/src/curl/http.c
index ad2a156..993aef7 100644
--- a/src/curl/http.c
+++ b/src/curl/http.c
@@ -1,3 +1,4 @@
+#define AL_LOG_SECTION "curl_http"
#include <al/log.h>
#include "http.h"
@@ -5,17 +6,15 @@
static void check_status_codes(struct nn_curl *curl)
{
struct nn_http *http = (struct nn_http *)curl;
-
- if (http->status_code <= 0L) {
+ if (http->status_code <= 0) {
curl_easy_getinfo(curl->handle, CURLINFO_RESPONSE_CODE, &http->status_code);
- if (http->status_code > 0L) {
+ if (http->status_code > 0) {
http->callback(http->userdata, NNWT_HTTP_RESPONSE_CODE, NULL, &http->status_code);
}
}
-
- if (http->content_length < 0L && http->status_code == 200L) {
+ if (http->content_length < 0) {
curl_easy_getinfo(curl->handle, CURLINFO_CONTENT_LENGTH_DOWNLOAD_T, &http->content_length);
- if (http->content_length >= 0L) {
+ if (http->content_length >= 0) {
http->callback(http->userdata, NNWT_HTTP_CONTENT_LENGTH, NULL, &http->content_length);
}
}
@@ -31,26 +30,19 @@ static void http_handle_events(void *userdata, struct nn_curl *curl)
switch (msg->msg) {
case CURLMSG_DONE: {
al_assert(msg->easy_handle == curl->handle);
-
CURLcode code = msg->data.result;
if (code != CURLE_OK) {
- al_log_error("http", "Curl error: %s.", curl_easy_strerror(code));
- http->status_code = -1L;
+ error("Curl error: %s.", curl_easy_strerror(code));
+ http->status_code = -1;
http->callback(http->userdata, NNWT_HTTP_ERROR, NULL, &http->status_code);
return;
}
-
- if (!nn_curl_remove_handle(curl)) {
- return;
- }
-
+ if (!nn_curl_remove_handle(curl)) return;
check_status_codes(curl);
-
- if (http->status_code == 200L) {
+ if (http->status_code == 200) {
http->callback(http->userdata, NNWT_HTTP_FINISHED, NULL, &http->status_code);
- } else if (http->status_code / 100L == 3L) {
+ } else if (http->status_code / 100 == 3) {
http->callback(http->userdata, NNWT_HTTP_REDIRECT, NULL, &http->status_code);
-
char *redirect_url = NULL;
curl_easy_getinfo(curl->handle, CURLINFO_REDIRECT_URL, &redirect_url);
if (!redirect_url) {
@@ -58,23 +50,19 @@ static void http_handle_events(void *userdata, struct nn_curl *curl)
return;
}
nn_curl_set_url(curl, &al_str_cr(redirect_url));
-
- http->status_code = -1L;
- http->content_length = -1L;
- if (!nn_curl_add_handle(curl)) {
- return;
- }
+ http->status_code = -1;
+ http->content_length = -1;
+ if (!nn_curl_add_handle(curl)) return;
} else {
http->callback(http->userdata, NNWT_HTTP_ERROR, NULL, &http->status_code);
}
-
return;
}
case CURLMSG_LAST:
- al_log_debug("http", "Unhandled CURLMSG_LAST.");
+ debug("Unhandled CURLMSG_LAST.");
break;
case CURLMSG_NONE:
- al_log_debug("http", "Unhandled CURLMSG_NONE.");
+ debug("Unhandled CURLMSG_NONE.");
break;
}
}
@@ -82,8 +70,8 @@ static void http_handle_events(void *userdata, struct nn_curl *curl)
bool nn_http_init(struct nn_http *http)
{
- http->status_code = -1L;
- http->content_length = -1L;
+ http->status_code = -1;
+ http->content_length = -1;
http->headers = NULL;
return nn_curl_init(&http->curl);
}
@@ -148,12 +136,9 @@ static void request_callback(void *userdata, u8 op, u8 *buf, void *opaque)
if (request->pointer + n > size) {
n = size - request->pointer;
}
-
nn_buffer_read(&request->payload, buf, request->pointer, n);
request->pointer += n;
-
*(size_t *)opaque = n;
-
break;
}
case NNWT_HTTP_WRITE: {
@@ -182,30 +167,25 @@ static void request_callback(void *userdata, u8 op, u8 *buf, void *opaque)
static bool init_request_internal(struct nn_http *http, u8 method)
{
struct nn_curl *curl = &http->curl;
-
switch (method) {
case NNWT_HTTP_HEAD:
- curl_easy_setopt(curl->handle, CURLOPT_NOBODY, 1L);
+ curl_easy_setopt(curl->handle, CURLOPT_NOBODY, 1);
break;
case NNWT_HTTP_GET:
- curl_easy_setopt(curl->handle, CURLOPT_HTTPGET, 1L);
+ curl_easy_setopt(curl->handle, CURLOPT_HTTPGET, 1);
break;
case NNWT_HTTP_POST:
- curl_easy_setopt(curl->handle, CURLOPT_POST, 1L);
+ curl_easy_setopt(curl->handle, CURLOPT_POST, 1);
break;
}
-
curl_easy_setopt(curl->handle, CURLOPT_WRITEFUNCTION, stream_write_callback);
curl_easy_setopt(curl->handle, CURLOPT_WRITEDATA, http);
curl_easy_setopt(curl->handle, CURLOPT_READFUNCTION, stream_read_callback);
curl_easy_setopt(curl->handle, CURLOPT_READDATA, http);
-
if (http->headers) {
curl_easy_setopt(curl->handle, CURLOPT_HTTPHEADER, http->headers);
}
-
curl->handle_events = http_handle_events;
-
return nn_curl_add_handle(curl);
}
@@ -235,18 +215,14 @@ bool nn_http_request(struct nn_http_request *request, u8 method, struct nn_event
struct nn_curl *curl = &request->http.curl;
al_assert(curl->loop == NULL);
curl->loop = loop;
-
request->http.callback = request_callback;
request->http.userdata = request;
-
long payload = (long)nn_buffer_get_size(&request->payload);
- if (payload > 0L) {
+ if (payload > 0) {
curl_easy_setopt(curl->handle, CURLOPT_POSTFIELDSIZE, payload);
}
-
request->callback = callback;
request->userdata = userdata;
-
return init_request_internal(&request->http, method);
}
diff --git a/src/curl/websocket.c b/src/curl/websocket.c
index 0030d16..8995221 100644
--- a/src/curl/websocket.c
+++ b/src/curl/websocket.c
@@ -1,3 +1,4 @@
+#define AL_LOG_SECTION "curl_ws"
#include <al/log.h>
#include "../socket/net.h"
@@ -30,18 +31,14 @@ static size_t stream_write_callback(char *buffer, size_t size, size_t nmemb, voi
(void)size;
(void)nmemb;
struct nn_websocket *ws = (struct nn_websocket *)userdata;
-
const struct curl_ws_frame *m = curl_ws_meta(ws->curl.handle);
size_t frame_size = m->offset + m->len + m->bytesleft;
- al_log_debug("ws", "frame_size: %zu, bytesleft: %zu.", frame_size, m->bytesleft);
-
+ debug("frame_size: %zu, bytesleft: %zu.", frame_size, m->bytesleft);
nn_buffer_ensure_space(&ws->frame, frame_size);
al_memcpy(nn_buffer_get_ptr(&ws->frame, m->offset), buffer, m->len);
-
if (m->bytesleft == 0) {
ws->callback(ws->userdata, m, nn_buffer_get_ptr(&ws->frame, 0), frame_size);
}
-
return m->len;
}
@@ -62,17 +59,13 @@ bool nn_websocket_connect(struct nn_websocket *ws, struct nn_event_loop *loop,
struct nn_curl *curl = &ws->curl;
al_assert(curl->loop == NULL);
curl->loop = loop;
-
- curl_easy_setopt(curl->handle, CURLOPT_CONNECT_ONLY, 0L);
+ curl_easy_setopt(curl->handle, CURLOPT_CONNECT_ONLY, 0);
ws->callback = callback;
ws->userdata = userdata;
-
curl_easy_setopt(curl->handle, CURLOPT_WRITEFUNCTION, stream_write_callback);
curl_easy_setopt(curl->handle, CURLOPT_WRITEDATA, ws);
-
curl->handle_events = websocket_handle_events;
curl->userdata = ws;
-
return nn_curl_add_handle(&ws->curl);
}
@@ -89,7 +82,7 @@ bool nn_websocket_send_frame(struct nn_websocket *ws, u8 *data, size_t size)
size_t sent;
CURLcode code = curl_ws_send(ws->curl.handle, data, size, &sent, 0, CURLWS_TEXT);
if (code != CURLE_OK) {
- al_log_error("curl", "curl_ws_send() failed (%s).", curl_easy_strerror(code));
+ error("curl_ws_send() failed (%s).", curl_easy_strerror(code));
return false;
}
al_assert(sent == size);
diff --git a/src/ev_embed_compat.c b/src/ev_embed_compat.c
index a907af4..0400aeb 100644
--- a/src/ev_embed_compat.c
+++ b/src/ev_embed_compat.c
@@ -10,27 +10,27 @@
#endif
#include "evwrap.h"
-#if defined(__clang__) || defined(__GNUC__)
-_Pragma("GCC diagnostic push")
-_Pragma("GCC diagnostic ignored \"-Wstrict-aliasing\"")
-_Pragma("GCC diagnostic ignored \"-Wcomment\"")
-_Pragma("GCC diagnostic ignored \"-Wparentheses\"")
-_Pragma("GCC diagnostic ignored \"-Wunused-value\"")
-_Pragma("GCC diagnostic ignored \"-Wunused-function\"")
-_Pragma("GCC diagnostic ignored \"-Wunused-parameter\"")
-_Pragma("GCC diagnostic ignored \"-Wunused-result\"")
-_Pragma("GCC diagnostic ignored \"-Wsign-compare\"")
-_Pragma("GCC diagnostic ignored \"-Wreturn-type\"")
-_Pragma("GCC diagnostic ignored \"-Wdeprecated-declarations\"")
-_Pragma("GCC diagnostic ignored \"-Wdangling-else\"")
-_Pragma("GCC diagnostic ignored \"-Wtype-limits\"")
-_Pragma("GCC diagnostic ignored \"-Wunknown-pragmas\"")
-#elif _MSC_VER
+
+#ifdef _MSC_VER
#pragma warning(push, 0)
+#else
+AL_IGNORE_WARNING("-Wstrict-aliasing")
+AL_IGNORE_WARNING("-Wcomment")
+AL_IGNORE_WARNING("-Wparentheses")
+AL_IGNORE_WARNING("-Wunused-value")
+AL_IGNORE_WARNING("-Wunused-function")
+AL_IGNORE_WARNING("-Wunused-parameter")
+AL_IGNORE_WARNING("-Wunused-result")
+AL_IGNORE_WARNING("-Wsign-compare")
+AL_IGNORE_WARNING("-Wreturn-type")
+AL_IGNORE_WARNING("-Wdeprecated-declarations")
+AL_IGNORE_WARNING("-Wdangling-else")
+AL_IGNORE_WARNING("-Wtype-limits")
+AL_IGNORE_WARNING("-Wunknown-pragmas")
#endif
#include <ev.c>
-#if defined(__clang__) || defined(__GNUC__)
-_Pragma("GCC diagnostic pop")
-#elif _MSC_VER
+#ifdef _MSC_VER
#pragma warning(pop)
+#else
+AL_IGNORE_WARNING_END
#endif
diff --git a/src/evwrap.h b/src/evwrap.h
index 7e2511b..35b72cc 100644
--- a/src/evwrap.h
+++ b/src/evwrap.h
@@ -16,3 +16,43 @@
#define EV_EMBED_ENABLE 0
#include <ev.h>
+
+AL_IGNORE_WARNING("-Wstrict-aliasing")
+
+// Wrap common libev macros so we can ignore aliasing warnings.
+
+#define ev_init_n(type) ev_init_n_##type
+
+static inline void ev_init_n_ev_io(ev_io *ev, void (*callback)(struct ev_loop *, ev_io *, s32))
+{
+ ev_init(ev, callback);
+}
+
+static inline void ev_init_n_ev_timer(ev_timer *ev, void (*callback)(struct ev_loop *, ev_timer *, s32))
+{
+ ev_init(ev, callback);
+}
+
+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);
+}
+
+static inline void ev_timer_init_n(ev_timer *timer, void (*callback)(struct ev_loop *, ev_timer *, s32), s32 fd, s32 events)
+{
+ ev_timer_init(timer, callback, fd, events);
+}
+
+static inline void ev_async_init_n(ev_async *asyn, void (*callback)(struct ev_loop *, ev_async *, s32))
+{
+ 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/fs_event/fs_event_inotify.c b/src/fs_event/fs_event_inotify.c
index c0862c7..a70550a 100644
--- a/src/fs_event/fs_event_inotify.c
+++ b/src/fs_event/fs_event_inotify.c
@@ -1,3 +1,4 @@
+#define AL_LOG_SECTION "fs_event"
#include <al/log.h>
#include "../util/error.h"
@@ -10,7 +11,6 @@ static void event_callback(struct ev_loop *loop, ev_io *w, s32 revents)
(void)loop;
struct nn_fs_event *fs = (struct nn_fs_event *)w->data;
(void)revents;
-
struct inotify_event event = { 0 };
ssize_t ret = nn_read(fs->fd, &event, sizeof(struct inotify_event));
if (ret == sizeof(struct inotify_event)) {
@@ -29,7 +29,7 @@ void nn_fs_event_init(struct nn_fs_event *fs, str *path, u32 mask,
fs->callback = callback;
fs->userdata = userdata;
fs->event.data = fs;
- ev_io_init(&fs->event, event_callback, fs->fd, EV_READ);
+ ev_io_init_n(&fs->event, event_callback, fs->fd, EV_READ);
}
bool nn_fs_event_start(struct nn_fs_event *fs, struct nn_event_loop *loop)
@@ -37,7 +37,7 @@ bool nn_fs_event_start(struct nn_fs_event *fs, struct nn_event_loop *loop)
fs->loop = loop;
fs->wd = inotify_add_watch(fs->fd, fs->path, fs->mask);
if (fs->wd == -1) {
- al_log_error("fs_event", "inotify_add_watch(%s) failed (%s)", fs->path, nn_strerror(errno));
+ error("inotify_add_watch(%s) failed (%s).", fs->path, nn_strerror(errno));
return false;
}
ev_io_start(fs->loop->ev, &fs->event);
diff --git a/src/line_processor.c b/src/line_processor.c
index f4f04c5..1872ce1 100644
--- a/src/line_processor.c
+++ b/src/line_processor.c
@@ -1,5 +1,3 @@
-#include <al/log.h>
-
#include "line_processor.h"
// Read from a socket/fd and callback for each "line" determined by a delimiter string.
@@ -99,7 +97,7 @@ static void socket_connection_callback(struct ev_loop *loop, ev_io *w, s32 reven
cl->event.data = cl;
al_array_push(pro->clients, cl);
- ev_io_init(&cl->event, read_callback, nn_socket_get_fd(&cl->sock), EV_READ);
+ ev_io_init_n(&cl->event, read_callback, nn_socket_get_fd(&cl->sock), EV_READ);
ev_io_start(pro->loop->ev, &cl->event);
}
@@ -128,7 +126,7 @@ void nn_line_processor_run(struct nn_line_processor *pro, struct nn_event_loop *
switch (pro->type) {
case NNWT_LINE_PROCESSOR_SOCKET:
pro->event.data = pro;
- ev_io_init(&pro->event, socket_connection_callback, nn_socket_get_fd(pro->sock), EV_READ);
+ ev_io_init_n(&pro->event, socket_connection_callback, nn_socket_get_fd(pro->sock), EV_READ);
ev_io_start(pro->loop->ev, &pro->event);
break;
case NNWT_LINE_PROCESSOR_FD: { // Not supported on windows.
@@ -139,7 +137,7 @@ void nn_line_processor_run(struct nn_line_processor *pro, struct nn_event_loop *
cl->pro = pro;
cl->event.data = cl;
al_array_push(pro->clients, cl);
- ev_io_init(&cl->event, read_callback, pro->fd, EV_READ);
+ ev_io_init_n(&cl->event, read_callback, pro->fd, EV_READ);
ev_io_start(pro->loop->ev, &cl->event);
break;
}
diff --git a/src/multiplex.c b/src/multiplex.c
index 2eed7df..fa13a08 100644
--- a/src/multiplex.c
+++ b/src/multiplex.c
@@ -1,3 +1,4 @@
+#define AL_LOG_SECTION "multiplex"
#include <al/log.h>
#include "multiplex.h"
@@ -9,10 +10,8 @@ bool nn_multiplex_socket_init(struct nn_multiplex_socket *multi, u8 type,
if (!nn_socket_init(&multi->sock, NNWT_SOCKET_NONBLOCKING)) {
return false;
}
-
multi->connection_callback = connection_callback;
multi->userdata = userdata;
-
return true;
}
@@ -28,9 +27,7 @@ static void socket_read_callback(struct ev_loop *loop, ev_io *w, s32 revents)
ev_io_stop(loop, &conn->event);
nn_socket_close(&conn->sock);
al_free(conn);
- if (ret == 0) {
- al_log_warn("multiplex", "read() returned 0 on POLLIN.");
- }
+ if (ret == 0) warn("read() returned 0 on POLLIN.");
return;
}
@@ -50,16 +47,14 @@ static void socket_connection_callback(struct ev_loop *loop, ev_io *w, s32 reven
{
struct nn_multiplex_socket *multi = (struct nn_multiplex_socket *)w->data;
(void)revents;
-
struct nn_multiplex_connection *conn = al_alloc_object(struct nn_multiplex_connection);
if (!nn_socket_accept(&multi->sock, &conn->sock, NNWT_SOCKET_NONBLOCKING)) {
al_free(conn);
return;
}
-
conn->multi = multi;
conn->event.data = conn;
- ev_io_init(&conn->event, socket_read_callback, nn_socket_get_fd(&conn->sock), EV_READ);
+ ev_io_init_n(&conn->event, socket_read_callback, nn_socket_get_fd(&conn->sock), EV_READ);
ev_io_start(loop, &conn->event);
}
@@ -68,13 +63,10 @@ bool nn_multiplex_socket_listen(struct nn_multiplex_socket *multi, struct nn_eve
if (!nn_socket_bind(&multi->sock, addr, port) || !nn_socket_listen(&multi->sock)) {
return false;
}
-
multi->loop = loop;
-
multi->event.data = multi;
- ev_io_init(&multi->event, socket_connection_callback, nn_socket_get_fd(&multi->sock), EV_READ);
+ ev_io_init_n(&multi->event, socket_connection_callback, nn_socket_get_fd(&multi->sock), EV_READ);
ev_io_start(multi->loop->ev, &multi->event);
-
return true;
}
diff --git a/src/packet_cache.c b/src/packet_cache.c
index d46f2f4..bc230ed 100644
--- a/src/packet_cache.c
+++ b/src/packet_cache.c
@@ -4,7 +4,7 @@ void nn_packet_cache_init(struct nn_packet_cache *cache, u32 size)
{
al_array_init(cache->cache);
al_array_reserve(cache->cache, size);
- cache->flush = size > 0 ? size / 2 : 0;
+ cache->flush = (size > 0) ? size / 2 : 0;
cache->disabled = false;
nn_cond_init(&cache->cond);
nn_mutex_init(&cache->mutex);
@@ -13,21 +13,16 @@ void nn_packet_cache_init(struct nn_packet_cache *cache, u32 size)
bool nn_packet_cache_send_packet(struct nn_packet_cache *cache, struct nn_packet *packet)
{
nn_mutex_lock(&cache->mutex);
-
if (cache->disabled) {
nn_mutex_unlock(&cache->mutex);
return false;
}
-
al_array_push(cache->cache, packet);
-
bool flush = !packet || cache->cache.count >= cache->flush;
if (flush && nn_cond_is_waiting(&cache->cond)) {
nn_cond_signal(&cache->cond);
}
-
nn_mutex_unlock(&cache->mutex);
-
return true;
}
@@ -43,11 +38,9 @@ void nn_packet_cache_flush(struct nn_packet_cache *cache)
bool nn_packet_cache_wait(struct nn_packet_cache *cache, u32 *count)
{
nn_mutex_lock(&cache->mutex);
-
if (!cache->disabled && !cache->cache.count) {
nn_cond_wait(&cache->cond, &cache->mutex);
}
-
// If there was an active wait before calling disable(), expected
// behavior would be that wait() returns false.
if (cache->disabled) {
@@ -55,9 +48,7 @@ bool nn_packet_cache_wait(struct nn_packet_cache *cache, u32 *count)
nn_mutex_unlock(&cache->mutex);
return false;
}
-
*count = cache->cache.count;
-
return true;
}
diff --git a/src/packet_pool.c b/src/packet_pool.c
index 1f8ca33..fb5e715 100644
--- a/src/packet_pool.c
+++ b/src/packet_pool.c
@@ -11,60 +11,43 @@ static void signal_callback(struct ev_loop *loop, ev_async *w, s32 revents)
(void)loop;
struct nn_packet_pool *pool = (struct nn_packet_pool *)w->data;
(void)revents;
-
nn_mutex_lock(&pool->mutex);
-
al_assert(pool->flushing);
pool->flushing = false;
-
if (pool->disabled) {
nn_mutex_unlock(&pool->mutex);
return;
}
-
- // Asserting on this and pool->flushing above is not necessary for
- // correctness but just to assist in catching errors.
+ // Asserting on this and pool->flushing above is not necessary for correctness.
al_assert(pool->ready.count > 0);
-
- al_array_copy(pool->sending, pool->ready);
- pool->ready.count = 0;
-
- nn_mutex_unlock(&pool->mutex);
-
struct nn_packet *packet;
- al_array_foreach(pool->sending, i, packet) {
+ al_array_foreach(pool->ready, i, packet) {
pool->callback(pool->userdata, packet);
}
- pool->sending.count = 0;
+ pool->ready.count = 0;
+ nn_mutex_unlock(&pool->mutex);
}
void nn_packet_pool_init(struct nn_packet_pool *pool, u32 size,
struct nn_event_loop *loop, void (*callback)(void *, struct nn_packet *), void *userdata)
{
pool->loop = loop;
-
pool->size = size;
pool->grow = false;
pool->flushing = false;
pool->disabled = false;
nn_cond_init(&pool->cond);
nn_mutex_init(&pool->mutex);
-
al_array_init(pool->ready);
al_array_reserve(pool->ready, size);
-
al_array_init(pool->empty);
al_array_reserve(pool->empty, size);
for (u32 i = 0; i < size; i++) {
al_array_push(pool->empty, nn_packet_create());
}
-
- al_array_init(pool->sending);
-
pool->signal.data = pool;
- ev_async_init(&pool->signal, signal_callback);
+ ev_async_init_n(&pool->signal, signal_callback);
ev_async_start(pool->loop->ev, &pool->signal);
-
pool->callback = callback;
pool->userdata = userdata;
}
@@ -72,7 +55,6 @@ void nn_packet_pool_init(struct nn_packet_pool *pool, u32 size,
struct nn_packet *nn_packet_pool_get(struct nn_packet_pool *pool)
{
nn_mutex_lock(&pool->mutex);
-
if (!pool->disabled && !pool->empty.count) {
if (pool->grow) {
al_array_push(pool->empty, nn_packet_create());
@@ -80,38 +62,29 @@ struct nn_packet *nn_packet_pool_get(struct nn_packet_pool *pool)
nn_cond_wait(&pool->cond, &pool->mutex);
}
}
-
if (pool->disabled) {
nn_mutex_unlock(&pool->mutex);
return NULL;
}
-
struct nn_packet *packet = al_array_pop(pool->empty);
-
nn_mutex_unlock(&pool->mutex);
-
return packet;
}
void nn_packet_pool_submit(struct nn_packet_pool *pool, struct nn_packet *packet)
{
nn_mutex_lock(&pool->mutex);
-
if (pool->disabled) {
return_internal(pool, packet);
nn_mutex_unlock(&pool->mutex);
return;
}
-
al_array_push(pool->ready, packet);
-
bool flush = pool->ready.count >= pool->size / 2;
-
if (!pool->flushing && flush) {
pool->flushing = true;
ev_async_send(pool->loop->ev, &pool->signal);
}
-
nn_mutex_unlock(&pool->mutex);
}
@@ -146,30 +119,27 @@ void nn_packet_pool_unlock(struct nn_packet_pool *pool)
void nn_packet_pool_disable(struct nn_packet_pool *pool)
{
nn_mutex_lock(&pool->mutex);
-
+ al_assert(!pool->disabled);
pool->disabled = true;
-
- // This makes disable() and enable() not thread-safe.
+ // This makes disable()/enable() required to be called from the loop thread.
ev_async_stop(pool->loop->ev, &pool->signal);
pool->flushing = false;
-
struct nn_packet *packet;
al_array_foreach(pool->ready, i, packet) {
nn_packet_reset(packet);
al_array_push(pool->empty, packet);
}
pool->ready.count = 0;
-
if (nn_cond_is_waiting(&pool->cond)) {
nn_cond_signal(&pool->cond);
}
-
nn_mutex_unlock(&pool->mutex);
}
void nn_packet_pool_enable(struct nn_packet_pool *pool)
{
nn_mutex_lock(&pool->mutex);
+ al_assert(pool->disabled);
pool->disabled = false;
ev_async_start(pool->loop->ev, &pool->signal);
nn_mutex_unlock(&pool->mutex);
@@ -181,13 +151,9 @@ 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(&pool->signal));
+ al_assert(!ev_is_active_n(ev_async)(&pool->signal));
}
nn_mutex_unlock(&pool->mutex);
-
- nn_mutex_destroy(&pool->mutex);
- nn_cond_destroy(&pool->cond);
-
struct nn_packet *packet;
al_array_foreach(pool->empty, i, packet) {
nn_packet_free(packet);
@@ -197,8 +163,6 @@ void nn_packet_pool_free(struct nn_packet_pool *pool)
nn_packet_free(packet);
}
al_array_free(pool->ready);
- al_array_foreach(pool->sending, i, packet) {
- nn_packet_free(packet);
- }
- al_array_free(pool->sending);
+ nn_mutex_destroy(&pool->mutex);
+ nn_cond_destroy(&pool->cond);
}
diff --git a/src/packet_pool.h b/src/packet_pool.h
index 65b53db..a6d700a 100644
--- a/src/packet_pool.h
+++ b/src/packet_pool.h
@@ -12,6 +12,7 @@
// packet_submit (signal async) /
// callback called \ Thread 1
// packet_return (unblock) /
+// -- Close --
// pool_disable()
// join(thread1) && join(thread2)
// pool_free()
@@ -31,7 +32,6 @@ struct nn_packet_pool {
struct nn_mutex mutex;
array(struct nn_packet *) empty;
array(struct nn_packet *) ready;
- array(struct nn_packet *) sending;
ev_async signal;
void (*callback)(void *, struct nn_packet *);
void *userdata;
diff --git a/src/packet_stream.c b/src/packet_stream.c
index 696481c..5eb2695 100644
--- a/src/packet_stream.c
+++ b/src/packet_stream.c
@@ -24,16 +24,13 @@ void nn_packet_stream_init(struct nn_packet_stream *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;
}
@@ -67,7 +64,6 @@ static void discard_out_queue_internal(struct nn_packet_stream *stream)
stream->packet_sent_callback(stream->userdata, stream->out.packet);
stream->out.packet = NULL;
}
-
if (stream->packets_sent_callback) {
stream->packets_sent_callback(stream->userdata,
stream->out.queue.data, stream->out.queue.count);
@@ -77,26 +73,20 @@ static void discard_out_queue_internal(struct nn_packet_stream *stream)
stream->packet_sent_callback(stream->userdata, packet);
}
}
-
stream->out.queue.count = 0;
}
static void stop_internal(struct nn_packet_stream *stream)
{
stream->connect = PACKET_STREAM_DISCONNECTED;
-
nn_socket_close(&stream->sock);
-
// Mark the stream corked so reconnect() can be consistent
// with connect() and from_socket().
stream->corked = true;
-
nn_packet_reset(stream->in.packet);
stream->in.have_header = false;
stream->in.index = 0;
-
discard_out_queue_internal(stream);
-
stream->connection_closed_callback(stream->userdata, stream);
}
@@ -200,19 +190,14 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents)
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);
-
stream->loop = loop;
stream->sock = *sock;
nn_socket_set_blocking(&stream->sock, false);
-
init_io_state(stream);
-
stream->wevent.data = stream;
- ev_io_init(&stream->wevent, stream_write_callback, nn_socket_get_fd(&stream->sock), EV_WRITE);
-
+ ev_io_init_n(&stream->wevent, stream_write_callback, nn_socket_get_fd(&stream->sock), EV_WRITE);
stream->revent.data = stream;
- ev_io_init(&stream->revent, stream_read_callback, nn_socket_get_fd(&stream->sock), EV_READ);
-
+ ev_io_init_n(&stream->revent, stream_read_callback, nn_socket_get_fd(&stream->sock), EV_READ);
if (stream_connected(stream)) {
ev_io_start(stream->loop->ev, &stream->revent);
}
@@ -224,19 +209,15 @@ static void do_connect_internal(struct nn_packet_stream *stream, str *addr, u16
stream->connection_closed_callback(stream->userdata, stream);
return;
}
-
if (!nn_socket_connect(&stream->sock, addr, port)) {
nn_socket_close(&stream->sock);
stream->connection_closed_callback(stream->userdata, stream);
return;
}
-
stream->revent.data = stream;
- ev_io_init(&stream->revent, stream_read_callback, nn_socket_get_fd(&stream->sock), EV_READ);
-
+ ev_io_init_n(&stream->revent, stream_read_callback, nn_socket_get_fd(&stream->sock), EV_READ);
stream->wevent.data = stream;
- ev_io_init(&stream->wevent, stream_write_callback, nn_socket_get_fd(&stream->sock), EV_WRITE);
-
+ ev_io_init_n(&stream->wevent, stream_write_callback, nn_socket_get_fd(&stream->sock), EV_WRITE);
// Start non-blocking connection.
stream->connect = PACKET_STREAM_CONNECTING;
start_write_internal(stream);
@@ -280,23 +261,17 @@ bool nn_packet_stream_send_packet(struct nn_packet_stream *stream, struct nn_pac
if (stream->connect == PACKET_STREAM_DISCONNECTING) {
return false;
}
-
al_assert(stream->connect == PACKET_STREAM_CONNECTED);
-
nn_packet_write_size(packet);
-
if (stream->direct) {
struct nn_packet_stream *direct = stream->direct;
direct->packet_callback(direct->userdata, direct, packet);
return true;
}
-
al_array_push(stream->out.queue, packet);
-
if (!stream->out.running) {
start_write_internal(stream);
}
-
return true;
}
@@ -338,13 +313,9 @@ void nn_packet_stream_return_packets(struct nn_packet_stream *stream, struct nn_
void nn_packet_stream_disconnect(struct nn_packet_stream *stream)
{
al_assert(stream->connect != PACKET_STREAM_DISCONNECTING);
-
- // If reusing this stream, keep in mind it will have been set
- // back to a default state in stop_internal().
- if (stream->connect == PACKET_STREAM_DISCONNECTED) {
- return;
- }
-
+ // If reusing this stream, keep in mind it will have been set back
+ // to a default state in stop_internal().
+ if (stream->connect == PACKET_STREAM_DISCONNECTED) return;
if (stream->direct) {
struct nn_packet_stream *direct = stream->direct;
direct->connect = PACKET_STREAM_DISCONNECTED;
@@ -362,7 +333,6 @@ void nn_packet_stream_disconnect(struct nn_packet_stream *stream)
stream->connection_closed_callback(stream->userdata, stream);
return;
}
-
if (stream->connect == PACKET_STREAM_CONNECTING) {
// wevent is always active when connect = CONNECTING.
stop_write_internal(stream);
diff --git a/src/poll.c b/src/poll.c
index 8ee38e2..caea5d9 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(&poll->event, event_callback);
+ ev_init_n(ev_io)(&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 ae8b689..d57bde5 100644
--- a/src/rpc.c
+++ b/src/rpc.c
@@ -70,20 +70,15 @@ static bool stream_connection_callback(void *userdata, struct nn_packet_stream *
{
struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata;
struct nn_rpc *rpc = conn->rpc;
-
if (stream->sock.type == NNWT_SOCKET_TCP) {
nn_packet_stream_set_nodelay(stream, 1);
}
-
stream->packet_callback = packet_callback;
stream->packet_sent_callback = packet_sent_callback;
stream->userdata = conn;
-
rpc->connection_callback(rpc->userdata, conn);
-
// Only add connections when acting as a server.
if (!rpc->conn) al_array_push(rpc->connections, conn);
-
return true;
}
@@ -106,7 +101,6 @@ void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream)
{
struct nn_rpc_connection *conn = al_alloc_object(struct nn_rpc_connection);
init_rpc_connection(rpc, conn);
-
conn->stream = stream;
conn->stream->connection_callback = stream_connection_callback;
conn->stream->connection_closed_callback = stream_connection_closed_callback;
@@ -116,13 +110,10 @@ void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream)
void nn_rpc_prepare_client(struct nn_rpc *rpc)
{
al_assert(!rpc->conn);
-
struct nn_rpc_connection *conn = al_alloc_object(struct nn_rpc_connection);
init_rpc_connection(rpc, conn);
-
conn->stream = al_alloc_object(struct nn_packet_stream);
nn_packet_stream_init(conn->stream, stream_connection_callback, stream_connection_closed_callback, conn);
-
rpc->conn = conn;
al_array_push(rpc->connections, conn);
}
@@ -166,7 +157,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 = NNWT_PACKET_GET_ID(packet),
+ .id = nn_packet_get_u32(packet, NNWT_PACKET_HEADER_LENGTH + sizeof(s8)),
.callback = callback,
.userdata = userdata
}));
diff --git a/src/signal.c b/src/signal.c
index a5ad885..f223428 100644
--- a/src/signal.c
+++ b/src/signal.c
@@ -15,7 +15,7 @@ void nn_signal_init(struct nn_signal *signal, struct nn_event_loop *loop,
signal->callback = callback;
signal->userdata = userdata;
signal->signal.data = signal;
- ev_async_init(&signal->signal, signal_callback);
+ ev_async_init_n(&signal->signal, signal_callback);
}
void nn_signal_start(struct nn_signal *signal)
diff --git a/src/socket/socket.h b/src/socket/socket.h
index 57e3212..5f16ec6 100644
--- a/src/socket/socket.h
+++ b/src/socket/socket.h
@@ -23,6 +23,9 @@ enum {
NNWT_SOCKET_NODELAY = 1 << 1
};
+// @TODO: Guard against EINTR?
+// https://github.com/kovidgoyal/kitty/blob/master/kitty/safe-wrappers.h
+
struct nn_socket {
u8 type;
#ifdef NAUNET_ON_WINDOWS
diff --git a/src/socket/socket_internal.h b/src/socket/socket_internal.h
index b79e888..41b48ca 100644
--- a/src/socket/socket_internal.h
+++ b/src/socket/socket_internal.h
@@ -10,15 +10,10 @@ static inline str *nn_socket_display_addr(str *addr)
return addr ? addr : &ipany;
}
-#define nn_socket_fmt_addr(addr) al_str_fmt(nn_socket_display_addr(addr))
+#define nn_addr_x(addr) al_str_x(nn_socket_display_addr(addr))
static inline void nn_socket_apply_flags(struct nn_socket *sock, s32 flags)
{
- if (flags & NNWT_SOCKET_NONBLOCKING) {
- nn_socket_set_blocking(sock, false);
- }
-
- if (flags & NNWT_SOCKET_NODELAY) {
- nn_socket_set_nodelay(sock, 1);
- }
+ if (flags & NNWT_SOCKET_NONBLOCKING) nn_socket_set_blocking(sock, false);
+ if (flags & NNWT_SOCKET_NODELAY) nn_socket_set_nodelay(sock, 1);
}
diff --git a/src/socket/socket_linux.c b/src/socket/socket_linux.c
index b8980be..b2d76a8 100644
--- a/src/socket/socket_linux.c
+++ b/src/socket/socket_linux.c
@@ -1,3 +1,4 @@
+#define AL_LOG_SECTION "socket"
#include <al/log.h>
#include <arpa/inet.h>
#include <netdb.h>
@@ -24,7 +25,7 @@ s32 nn_poll_fds(struct nn_pollfd *fds, nn_nfds nfds, s64 timeout_ns)
struct timespec ts;
ts.tv_sec = timeout_ns / NNWT_TIME_S_TO_NS(1);
ts.tv_nsec = timeout_ns % NNWT_TIME_S_TO_NS(1);
- struct timespec *tsp = timeout_ns >= 0 ? &ts : NULL;
+ struct timespec *tsp = (timeout_ns >= 0) ? &ts : NULL;
return ppoll(fds, nfds, tsp, NULL);
}
@@ -45,14 +46,14 @@ bool nn_socket_init(struct nn_socket *sock, s32 flags)
}
if (sock->fd < 0) {
- al_log_error("socket", "socket() failed: %s (%d).", nn_strerror(errno), errno);
+ error("socket() failed: %s (%d).", nn_strerror(errno), errno);
return false;
}
- nn_socket_apply_flags(sock, flags);
-
sock->addrinfo = NULL;
+ nn_socket_apply_flags(sock, flags);
+
return true;
}
@@ -119,7 +120,7 @@ static bool parse_address(struct nn_socket *sock, str *addr, u16 port)
s32 status = getaddrinfo(c_str, port_str, &hints, &sock->addrinfo);
al_free(c_str);
if (status != 0) {
- al_log_error("socket", "Failed to parse address (%s).", gai_strerror(status));
+ error("Failed to parse address (%s).", gai_strerror(status));
return false;
}
@@ -136,7 +137,6 @@ bool nn_socket_bind(struct nn_socket *sock, str *addr, u16 port)
{
struct sockaddr *saddr = NULL;
socklen_t addrlen = 0;
-
switch (sock->type) {
case NNWT_SOCKET_TCP:
case NNWT_SOCKET_UDP:
@@ -146,7 +146,7 @@ bool nn_socket_bind(struct nn_socket *sock, str *addr, u16 port)
saddr = sock->addrinfo->ai_addr;
addrlen = sock->addrinfo->ai_addrlen;
break;
- case NNWT_SOCKET_UNIX: {
+ case NNWT_SOCKET_UNIX:
al_assert(addr);
sock->addr_un.sun_family = AF_UNIX;
char *c_str = al_str_to_c_str(addr);
@@ -156,18 +156,16 @@ bool nn_socket_bind(struct nn_socket *sock, str *addr, u16 port)
saddr = (struct sockaddr *)&sock->addr_un;
addrlen = sizeof(sock->addr_un);
break;
- }
default:
al_assert_and_return(false);
}
if (bind(sock->fd, saddr, addrlen) < 0) {
- al_log_error("socket", "bind(%.*s:%hu) failed: %s (%d).",
- nn_socket_fmt_addr(addr), port, nn_strerror(errno), errno);
+ error("bind(%.*s:%hu) failed: %s (%d).", nn_addr_x(addr), port, nn_strerror(errno), errno);
return false;
}
- al_log_info("socket", "Socket bound to %.*s:%hu.", nn_socket_fmt_addr(addr), port);
+ info("Socket bound to %.*s:%hu.", nn_addr_x(addr), port);
return true;
}
@@ -177,11 +175,11 @@ bool nn_socket_listen(struct nn_socket *sock)
al_assert(sock->type == NNWT_SOCKET_TCP || sock->type == NNWT_SOCKET_UNIX);
if (listen(sock->fd, SOMAXCONN) < 0) {
- al_log_error("socket", "listen() failed: %s (%d).", nn_strerror(errno), errno);
+ error("listen() failed: %s (%d).", nn_strerror(errno), errno);
return false;
}
- al_log_info("socket", "Listening.");
+ info("Listening.");
return true;
}
@@ -189,19 +187,19 @@ bool nn_socket_listen(struct nn_socket *sock)
bool nn_socket_accept(struct nn_socket *sock, struct nn_socket *cl, s32 flags)
{
al_assert(sock->type == NNWT_SOCKET_TCP || sock->type == NNWT_SOCKET_UNIX);
- cl->type = sock->type;
+ cl->type = sock->type;
struct sockaddr_storage addr = { 0 };
socklen_t addrlen = sizeof(addr);
if ((cl->fd = accept(sock->fd, (struct sockaddr *)&addr, &addrlen)) == -1) {
- al_log_error("socket", "accept() failed: %s (%d).", nn_strerror(errno), errno);
+ error("accept() failed: %s (%d).", nn_strerror(errno), errno);
return false;
}
- nn_socket_apply_flags(cl, flags);
-
cl->addrinfo = NULL;
+ nn_socket_apply_flags(cl, flags);
+
switch (cl->type) {
case NNWT_SOCKET_TCP: {
u16 port = 0;
@@ -215,11 +213,11 @@ bool nn_socket_accept(struct nn_socket *sock, struct nn_socket *cl, s32 flags)
port = ntohs(saddr->sin6_port);
inet_ntop(AF_INET6, &saddr->sin6_addr, addr_str, sizeof(addr_str));
}
- al_log_info("socket", "Connection from %s:%hu.", addr_str, port);
+ info("Connection from %s:%hu.", addr_str, port);
break;
}
case NNWT_SOCKET_UNIX:
- al_log_info("socket", "Connection on UNIX socket.");
+ info("Connection on UNIX socket.");
break;
}
@@ -246,13 +244,10 @@ bool nn_socket_connect(struct nn_socket *sock, str *addr, u16 port)
default:
al_assert_and_return(false);
}
-
if (ret != 0 && errno != EINPROGRESS) {
- al_log_error("socket", "connect(%.*s:%hu) failed: %s (%d).",
- nn_socket_fmt_addr(addr), port, nn_strerror(errno), errno);
+ error("connect(%.*s:%hu) failed: %s (%d).", nn_addr_x(addr), port, nn_strerror(errno), errno);
return false;
}
-
return true;
}
diff --git a/src/socket/socket_windows.c b/src/socket/socket_windows.c
index 85c2853..68a7700 100644
--- a/src/socket/socket_windows.c
+++ b/src/socket/socket_windows.c
@@ -1,3 +1,4 @@
+#define AL_LOG_SECTION "socket"
#include <al/log.h>
#include "socket.h"
@@ -15,7 +16,7 @@ bool nn_socket_init(struct nn_socket *sock, s32 flags)
}
if (sock->fd == INVALID_SOCKET) {
- al_log_error("socket", "socket() failed: %d.", WSAGetLastError());
+ error("socket() failed: %d.", WSAGetLastError());
return false;
}
@@ -56,12 +57,11 @@ bool nn_socket_bind(struct nn_socket *sock, str *addr, u16 port)
s32 ret = bind(sock->fd, (SOCKADDR *)&sock->addr_in, sizeof(sock->addr_in));
if (ret == SOCKET_ERROR) {
- al_log_error("socket", "bind(%.*s:%hu) failed: %d.",
- nn_socket_fmt_addr(addr), port, WSAGetLastError());
+ error("bind(%.*s:%hu) failed: %d.", nn_addr_x(addr), port, WSAGetLastError());
return false;
}
- al_log_info("socket", "Socket bound to %.*s:%hu.", nn_socket_fmt_addr(addr), port);
+ info("Socket bound to %.*s:%hu.", nn_addr_x(addr), port);
return true;
}
@@ -69,11 +69,11 @@ bool nn_socket_bind(struct nn_socket *sock, str *addr, u16 port)
bool nn_socket_listen(struct nn_socket *sock)
{
if (listen(sock->fd, SOMAXCONN) == SOCKET_ERROR) {
- al_log_error("socket", "listen() failed: %d.", WSAGetLastError());
+ error("listen() failed: %d.", WSAGetLastError());
return false;
}
- al_log_info("socket", "Listening.");
+ info("Listening.");
return true;
}
@@ -82,7 +82,7 @@ bool nn_socket_accept(struct nn_socket *sock, struct nn_socket *cl, s32 flags)
{
cl->type = sock->type;
if ((cl->fd = accept(sock->fd, NULL, NULL)) == INVALID_SOCKET) {
- al_log_error("socket", "accept() failed: %d.", WSAGetLastError());
+ error("accept() failed: %d.", WSAGetLastError());
return false;
}
@@ -105,7 +105,7 @@ bool nn_socket_connect(struct nn_socket *sock, str *addr, u16 port)
s32 ret = connect(sock->fd, (SOCKADDR *)&sock->addr_in, sizeof(sock->addr_in));
if (ret == SOCKET_ERROR && WSAGetLastError() != WSAEWOULDBLOCK) {
- al_log_error("socket", "connect() failed: %d.", WSAGetLastError());
+ error("connect() failed: %d.", WSAGetLastError());
return false;
}
diff --git a/src/timer.c b/src/timer.c
index 73902ea..2f35354 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(&timer->timer, timer_callback);
+ ev_init_n(ev_timer)(&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 0bd36b4..071d8bf 100644
--- a/src/util/buffer.c
+++ b/src/util/buffer.c
@@ -8,7 +8,7 @@
void nn_buffer_ensure_space(struct nn_buffer *buf, size_t size)
{
if (size > buf->alloc) {
- buf->alloc = size < NNWT_BUFFER_SIZE ? NNWT_BUFFER_SIZE :
+ buf->alloc = (size < NNWT_BUFFER_SIZE) ? NNWT_BUFFER_SIZE :
(size + NNWT_BUFFER_GROW) & ~NNWT_BUFFER_GROW;
buf->data = (u8 *)(buf->data ? al_realloc(buf->data, buf->alloc) : al_malloc(buf->alloc));
}
@@ -48,7 +48,7 @@ void nn_buffer_append(struct nn_buffer *buf, void *data, size_t size)
u8 *nn_buffer_get_ptr(struct nn_buffer *buf, size_t index)
{
- return &buf->data[index];
+ return buf->data + index;
}
void nn_buffer_read(struct nn_buffer *buf, void *ptr, size_t index, size_t size)
diff --git a/src/util/buffer.h b/src/util/buffer.h
index 612aae4..9b254d4 100644
--- a/src/util/buffer.h
+++ b/src/util/buffer.h
@@ -9,6 +9,10 @@ struct nn_buffer {
u8 *data;
};
+// A pun+dereference breaks on 32bit Android but isn't warned about by -Wstrict-aliasing=3.
+#define NNWT_BUFFER_READ_TYPE(buf, index, type, r) \
+ al_memcpy((void *)&(r), (void *)nn_buffer_get_ptr(buf, index), sizeof(type))
+
void nn_buffer_init(struct nn_buffer *buf);
void nn_buffer_ensure_space(struct nn_buffer *buf, size_t size);
size_t nn_buffer_get_size(struct nn_buffer *buf);
diff --git a/src/util/file/file.h b/src/util/file/file.h
index 03b82ba..91aa82c 100644
--- a/src/util/file/file.h
+++ b/src/util/file/file.h
@@ -4,18 +4,22 @@
#include <al/lib.h>
#ifdef NAUNET_NEEDS_STDIO_ASSIST
-#define _FILE_OFFSET_BITS 64
+// It seems meson adds _FILE_OFFSET_BITS=64 on all "linuxlike" compilers even MinGW.
+// https://github.com/mesonbuild/meson/commit/853634a48da025c59eef70161dba0d150833f60d
+// https://github.com/mesonbuild/meson/issues/12931
+// I don't think we necessarily want this on 32-bit Windows.
#include <stdio.h>
#ifdef NAUNET_ON_WINDOWS
+#if _FILE_OFFSET_BITS == 64
#define fseek _fseeki64
#define ftell _ftelli64
+#endif
#else
#define fseek fseeko
#define ftell ftello
#endif
#else
#ifdef NAUNET_ON_WINDOWS
-// @TODO:
#else
#include <dirent.h>
#include <fcntl.h>
@@ -43,7 +47,6 @@ struct nn_file {
FILE *file;
#else
#ifdef NAUNET_ON_WINDOWS
- // @TODO:
#else
s32 fd;
#endif
@@ -56,7 +59,6 @@ struct nn_dir {
#ifdef NAUNET_NEEDS_STDIO_ASSIST
#else
#ifdef NAUNET_ON_WINDOWS
- // @TODO:
#else
DIR *dir;
#endif
diff --git a/src/util/file/file_linux.c b/src/util/file/file_linux.c
index 5a3621b..c87bc9d 100644
--- a/src/util/file/file_linux.c
+++ b/src/util/file/file_linux.c
@@ -38,13 +38,13 @@ static bool open_file_linux(s32 *fd, str *path, s32 flags)
bool nn_file_open(struct nn_file *file, str *path, s32 flags)
{
file->flags = flags;
- s32 oflags = flags & NNWT_FILE_READONLY ? O_RDONLY : O_RDWR;
+ s32 oflags = (flags & NNWT_FILE_READONLY) ? O_RDONLY : O_RDWR;
if (flags & NNWT_FILE_CREATE) {
oflags |= O_CREAT;
}
if (!open_file_linux(&file->fd, path, oflags)) {
- error("open(%.*s) failed (%d: %s).", al_str_fmt(path), errno, nn_strerror(errno));
+ error("open(%.*s) failed (%d: %s).", al_str_x(path), errno, nn_strerror(errno));
return false;
}
@@ -128,10 +128,10 @@ bool nn_file_read(struct nn_file *file, void *buf, size_t size)
void *nn_file_mmap(struct nn_file *file)
{
- s32 flags = PROT_READ | (file->flags & NNWT_FILE_READONLY ? 0 : PROT_WRITE);
+ s32 flags = PROT_READ | ((file->flags & NNWT_FILE_READONLY) ? 0 : PROT_WRITE);
void *map = mmap(NULL, file->size, flags, MAP_SHARED, file->fd, 0L);
if (map == MAP_FAILED) {
- error("mmap(%.*s) failed (%s).", al_str_fmt(&file->path), nn_strerror(errno));
+ error("mmap(%.*s) failed (%s).", al_str_x(&file->path), nn_strerror(errno));
return NULL;
}
return map;
@@ -170,7 +170,7 @@ bool nn_dir_open(struct nn_dir *dir, str *path)
dir->dir = opendir(c_str);
al_free(c_str);
if (!dir->dir) {
- error("opendir(%.*s) failed (%s).", al_str_fmt(path), nn_strerror(errno));
+ error("opendir(%.*s) failed (%s).", al_str_x(path), nn_strerror(errno));
return false;
}
al_str_clone(&dir->path, path);
@@ -191,7 +191,7 @@ bool nn_dir_exists(str *path)
al_free(c_str);
if (!dir) {
if (errno != ENOENT) {
- error("opendir(%.*s) failed (%s).", al_str_fmt(path), nn_strerror(errno));
+ error("opendir(%.*s) failed (%s).", al_str_x(path), nn_strerror(errno));
}
return false;
}
@@ -205,7 +205,7 @@ bool nn_dir_create(str *path)
s32 ret = mkdir(c_str, 0755);
al_free(c_str);
if (ret == -1) {
- error("mkdir(%.*s) failed (%s).", al_str_fmt(path), nn_strerror(errno));
+ error("mkdir(%.*s) failed (%s).", al_str_x(path), nn_strerror(errno));
return false;
}
return true;
@@ -216,7 +216,7 @@ bool nn_dir_read(struct nn_dir *dir, struct nn_dir_entry *entry)
errno = 0;
if (!(entry->entry = readdir(dir->dir))) {
if (errno) {
- error("readdir(%.*s) failed (%s).", al_str_fmt(&dir->path), nn_strerror(errno));
+ error("readdir(%.*s) failed (%s).", al_str_x(&dir->path), nn_strerror(errno));
}
return false;
}
diff --git a/src/util/file/file_stdio.c b/src/util/file/file_stdio.c
index 3f58664..1969477 100644
--- a/src/util/file/file_stdio.c
+++ b/src/util/file/file_stdio.c
@@ -1,13 +1,19 @@
#define AL_LOG_SECTION "file_stdio"
#include <al/log.h>
+#include <al/lib.h>
#include "../error.h"
-#ifndef NAUNET_NEEDS_STDIO_ASSIST
-#define NAUNET_NEEDS_STDIO_ASSIST
-#endif
#include "file.h"
+/*
+#if defined NAUNET_ON_WINDOWS && defined AL_WE_32BIT
+AL_ASSERT_TYPE_SIZE(off_t, 4);
+#else
+AL_ASSERT_TYPE_SIZE(off_t, 8);
+#endif
+*/
+
static inline bool open_stdio_file(FILE **file, str *path, const char *mode)
{
char *c_str = al_str_to_c_str(path);
@@ -15,7 +21,7 @@ static inline bool open_stdio_file(FILE **file, str *path, const char *mode)
// Still sets errno, according to the Microsoft docs.
errno_t ret = fopen_s(file, c_str, mode);
#else
- s32 ret = (*file = fopen(c_str, mode)) != NULL ? 0 : errno;
+ s32 ret = (*file = fopen(c_str, mode)) ? 0 : errno;
#endif
al_free(c_str);
return ret == 0;
@@ -32,9 +38,9 @@ static bool query_filesize_stdio(struct nn_file *file)
bool nn_file_open(struct nn_file *file, str *path, s32 flags)
{
- const char *mode = flags & NNWT_FILE_CREATE ? "wb+" : "rb+";
+ const char *mode = (flags & NNWT_FILE_CREATE) ? "wb+" : "rb+";
if (!open_stdio_file(&file->file, path, mode)) {
- error("fopen(%.*s) failed (%d: %s).", al_str_fmt(path), errno, nn_strerror(errno));
+ error("fopen(%.*s) failed (%d: %s).", al_str_x(path), errno, nn_strerror(errno));
return false;
}
diff --git a/src/util/packet.c b/src/util/packet.c
index ab2fa83..fe20b58 100644
--- a/src/util/packet.c
+++ b/src/util/packet.c
@@ -27,14 +27,21 @@ void nn_packet_reset(struct nn_packet *packet)
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)
{
- *((u32 *)nn_buffer_get_ptr(&packet->buffer, 0)) = packet->windex;
+ nn_buffer_write(&packet->buffer, &packet->windex, 0, sizeof(u32));
}
u32 nn_packet_get_size(struct nn_packet *packet)
{
- return *((u32 *)nn_buffer_get_ptr(&packet->buffer, 0));
+ return nn_packet_get_u32(packet, 0);
}
#define DEFINE_PACKET_WRITE_FUNC(type) \
@@ -81,23 +88,12 @@ void nn_packet_write_buffer(struct nn_packet *packet, struct nn_buffer *buf)
return r; \
}
-#define DEFINE_PACKET_READ_FUNC_EXT(type) \
- type nn_packet_peek_##type(struct nn_packet *packet) \
- { \
- type r; \
- NNWT_PACKET_PEEK_TYPE(packet, type, r); \
- return r; \
- }
-
DEFINE_PACKET_READ_FUNC(u8)
-DEFINE_PACKET_READ_FUNC_EXT(u8)
DEFINE_PACKET_READ_FUNC(s8)
DEFINE_PACKET_READ_FUNC(u16)
DEFINE_PACKET_READ_FUNC(s16)
DEFINE_PACKET_READ_FUNC(u32)
-DEFINE_PACKET_READ_FUNC_EXT(u32)
DEFINE_PACKET_READ_FUNC(s32)
-DEFINE_PACKET_READ_FUNC_EXT(s32)
DEFINE_PACKET_READ_FUNC(u64)
DEFINE_PACKET_READ_FUNC(s64)
DEFINE_PACKET_READ_FUNC(f32)
diff --git a/src/util/packet.h b/src/util/packet.h
index 5659a70..7183837 100644
--- a/src/util/packet.h
+++ b/src/util/packet.h
@@ -7,7 +7,7 @@
#include "../util/buffer.h"
-#define NNWT_PACKET_HEADER_LENGTH (u32)(sizeof(u32))
+#define NNWT_PACKET_HEADER_LENGTH (u32)(sizeof(u32)) // Just size for now.
struct nn_packet {
struct nn_buffer buffer;
@@ -16,24 +16,18 @@ struct nn_packet {
void *opaque;
};
-#define NNWT_PACKET_GET_ID(p) \
- *((u32 *)nn_buffer_get_ptr(&p->buffer, NNWT_PACKET_HEADER_LENGTH + sizeof(s8)))
-
#define NNWT_PACKET_WRITE_TYPE(p, type, v) \
nn_buffer_write(&(p)->buffer, &(v), (p)->windex, sizeof(type)); \
(p)->windex += (u32)sizeof(type)
#define NNWT_PACKET_WRITE_DATA(p, data, length) \
- nn_buffer_write(&(p)->buffer, (data), (p)->windex, (length)); \
+ nn_buffer_write(&(p)->buffer, data, (p)->windex, length); \
(p)->windex += (u32)length
#define NNWT_PACKET_READ_TYPE(p, type, r) \
- r = *((type *)nn_buffer_get_ptr(&(p)->buffer, (p)->rindex)); \
+ NNWT_BUFFER_READ_TYPE(&(p)->buffer, (p)->rindex, type, r); \
(p)->rindex += (u32)sizeof(type)
-#define NNWT_PACKET_PEEK_TYPE(p, type, r) \
- r = *((type *)nn_buffer_get_ptr(&(p)->buffer, (p)->rindex))
-
#define NNWT_PACKET_READ_DATA(p, length, r) \
r = (void *)nn_buffer_get_ptr(&(p)->buffer, (p)->rindex); \
(p)->rindex += (u32)length
@@ -42,6 +36,8 @@ 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);
@@ -61,14 +57,11 @@ void nn_packet_write_buffer(struct nn_packet *packet, struct nn_buffer *buf);
bool nn_packet_read_bool(struct nn_packet *packet);
u8 nn_packet_read_u8(struct nn_packet *packet);
-u8 nn_packet_peek_u8(struct nn_packet *packet);
s8 nn_packet_read_s8(struct nn_packet *packet);
u16 nn_packet_read_u16(struct nn_packet *packet);
s16 nn_packet_read_s16(struct nn_packet *packet);
u32 nn_packet_read_u32(struct nn_packet *packet);
-u32 nn_packet_peek_u32(struct nn_packet *packet);
s32 nn_packet_read_s32(struct nn_packet *packet);
-s32 nn_packet_peek_s32(struct nn_packet *packet);
u64 nn_packet_read_u64(struct nn_packet *packet);
s64 nn_packet_read_s64(struct nn_packet *packet);
f32 nn_packet_read_f32(struct nn_packet *packet);
diff --git a/src/util/thread/thread_linux.c b/src/util/thread/thread_linux.c
index f07fcb4..f5b480e 100644
--- a/src/util/thread/thread_linux.c
+++ b/src/util/thread/thread_linux.c
@@ -14,7 +14,8 @@ void nn_thread_sleep(nn_os_tstamp delay)
void nn_thread_create(struct nn_thread *thread, nn_thread_func func, void *userdata)
{
- pthread_create(&thread->thread, NULL, func, userdata);
+ s32 ret = pthread_create(&thread->thread, NULL, func, userdata);
+ al_assert(ret == 0);
}
void nn_thread_set_priority(s32 policy, s32 priority)
@@ -51,22 +52,27 @@ void nn_thread_detach(struct nn_thread *thread)
void nn_mutex_init(struct nn_mutex *mutex)
{
- pthread_mutex_init(&mutex->mutex, NULL);
+ s32 ret = pthread_mutex_init(&mutex->mutex, NULL);
+ // Linux manpage says this always returns 0.
+ al_assert(ret == 0);
}
void nn_mutex_lock(struct nn_mutex *mutex)
{
- pthread_mutex_lock(&mutex->mutex);
+ s32 ret = pthread_mutex_lock(&mutex->mutex);
+ al_assert(ret != EINVAL);
}
void nn_mutex_unlock(struct nn_mutex *mutex)
{
- pthread_mutex_unlock(&mutex->mutex);
+ s32 ret = pthread_mutex_unlock(&mutex->mutex);
+ al_assert(ret != EINVAL);
}
void nn_mutex_destroy(struct nn_mutex *mutex)
{
- pthread_mutex_destroy(&mutex->mutex);
+ s32 ret = pthread_mutex_destroy(&mutex->mutex);
+ al_assert(ret != EBUSY);
}
void nn_cond_init(struct nn_cond *cond)
@@ -96,5 +102,6 @@ void nn_cond_signal(struct nn_cond *cond)
void nn_cond_destroy(struct nn_cond *cond)
{
- pthread_cond_destroy(&cond->cond);
+ s32 ret = pthread_cond_destroy(&cond->cond);
+ al_assert(ret != EBUSY);
}
diff --git a/src/util/thread/thread_windows.c b/src/util/thread/thread_windows.c
index f1f1dd7..9ce85e2 100644
--- a/src/util/thread/thread_windows.c
+++ b/src/util/thread/thread_windows.c
@@ -4,7 +4,7 @@ NTSTATUS(__stdcall *NtDelayExecution)(BOOL Alertable, PLARGE_INTEGER DelayInterv
void nn_thread_sleep(nn_os_tstamp delay)
{
- NtDelayExecution(false, &((LARGE_INTEGER){ .QuadPart = -(LONGLONG)(delay * 1E7L) }));
+ NtDelayExecution(false, &((LARGE_INTEGER){ .QuadPart = -(LONGLONG)(delay * 1e7L) }));
}
void nn_thread_create(struct nn_thread *thread, nn_thread_func func, void *userdata)
diff --git a/src/winwrap.c b/src/winwrap.c
index c287339..38bd101 100644
--- a/src/winwrap.c
+++ b/src/winwrap.c
@@ -10,10 +10,9 @@ bool nn_windows_init(void)
return false;
}
+ ULONG actual_resolution;
NTSTATUS(__stdcall *ZwSetTimerResolution)(IN ULONG RequestedResolution, IN BOOLEAN Set, OUT PULONG ActualResolution);
ZwSetTimerResolution = (NTSTATUS(__stdcall *)(IN ULONG, IN BOOLEAN, OUT PULONG))(void *)GetProcAddress(GetModuleHandleW(L"ntdll"), "ZwSetTimerResolution");
-
- ULONG actual_resolution;
ZwSetTimerResolution(1, TRUE, &actual_resolution);
NtDelayExecution = (NTSTATUS(__stdcall *)(BOOL, PLARGE_INTEGER))(void *)GetProcAddress(GetModuleHandleW(L"ntdll"), "NtDelayExecution");