diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/curl/curl.c | 30 | ||||
| -rw-r--r-- | src/curl/http.c | 66 | ||||
| -rw-r--r-- | src/curl/websocket.c | 15 | ||||
| -rw-r--r-- | src/ev_embed_compat.c | 38 | ||||
| -rw-r--r-- | src/evwrap.h | 40 | ||||
| -rw-r--r-- | src/fs_event/fs_event_inotify.c | 6 | ||||
| -rw-r--r-- | src/line_processor.c | 8 | ||||
| -rw-r--r-- | src/multiplex.c | 16 | ||||
| -rw-r--r-- | src/packet_cache.c | 11 | ||||
| -rw-r--r-- | src/packet_pool.c | 58 | ||||
| -rw-r--r-- | src/packet_pool.h | 2 | ||||
| -rw-r--r-- | src/packet_stream.c | 44 | ||||
| -rw-r--r-- | src/poll.c | 2 | ||||
| -rw-r--r-- | src/rpc.c | 11 | ||||
| -rw-r--r-- | src/signal.c | 2 | ||||
| -rw-r--r-- | src/socket/socket.h | 3 | ||||
| -rw-r--r-- | src/socket/socket_internal.h | 11 | ||||
| -rw-r--r-- | src/socket/socket_linux.c | 41 | ||||
| -rw-r--r-- | src/socket/socket_windows.c | 16 | ||||
| -rw-r--r-- | src/timer.c | 2 | ||||
| -rw-r--r-- | src/util/buffer.c | 4 | ||||
| -rw-r--r-- | src/util/buffer.h | 4 | ||||
| -rw-r--r-- | src/util/file/file.h | 10 | ||||
| -rw-r--r-- | src/util/file/file_linux.c | 16 | ||||
| -rw-r--r-- | src/util/file/file_stdio.c | 18 | ||||
| -rw-r--r-- | src/util/packet.c | 22 | ||||
| -rw-r--r-- | src/util/packet.h | 17 | ||||
| -rw-r--r-- | src/util/thread/thread_linux.c | 19 | ||||
| -rw-r--r-- | src/util/thread/thread_windows.c | 2 | ||||
| -rw-r--r-- | src/winwrap.c | 3 |
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); @@ -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) @@ -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"); |