From a320ee672ccf52297dc509b25ed263f165709b31 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 24 Jun 2024 14:34:21 -0400 Subject: Add multiplex, a bunch of cleanup Signed-off-by: Andrew Opalach --- include/aki/multiplex.h | 3 + meson.build | 23 +++--- meson_options.txt | 1 + src/common.c | 3 + src/curl/http.c | 1 + src/curl/websocket.c | 2 +- src/ev_embed_compat.c | 3 +- src/line_processor.c | 17 ++-- src/line_processor.h | 6 +- src/multiplex.c | 47 +++++++++++ src/multiplex.h | 20 +++++ src/packet_stream.c | 184 ++++++++++++++++++++++---------------------- src/packet_stream.h | 24 +++--- src/rpc.c | 118 +++++++++++++++------------- src/rpc2.h | 19 ++--- src/socket/socket.h | 35 +++++---- src/socket/socket_linux.c | 169 +++++++++++++++++++--------------------- src/socket/socket_windows.c | 4 +- src/util/file/file_linux.c | 15 ++-- src/util/packet.c | 28 +++---- src/util/packet.h | 7 +- 21 files changed, 403 insertions(+), 326 deletions(-) create mode 100644 include/aki/multiplex.h create mode 100644 src/multiplex.c create mode 100644 src/multiplex.h diff --git a/include/aki/multiplex.h b/include/aki/multiplex.h new file mode 100644 index 0000000..1cf1dbe --- /dev/null +++ b/include/aki/multiplex.h @@ -0,0 +1,3 @@ +#pragma once + +#include "../../src/multiplex.h" diff --git a/meson.build b/meson.build index 2f2e51b..f6636fa 100644 --- a/meson.build +++ b/meson.build @@ -35,6 +35,7 @@ if get_option('event-loop').enabled() 'src/packet_pool.c', 'src/packet_cache.c', 'src/packet_stream.c', + 'src/multiplex.c', 'src/rpc.c', 'src/line_processor.c' ] @@ -75,17 +76,19 @@ if get_option('event-loop').enabled() endif endif -jansson = dependency('jansson', required: false) -if not jansson.found() - jansson_opts = cmake.subproject_options() - jansson_opts.add_cmake_defines({ 'CMAKE_BUILD_TYPE': is_debug ? 'Debug' : 'Release' }) - jansson_opts.add_cmake_defines({ 'JANSSON_EXAMPLES': false }) - jansson_opts.add_cmake_defines({ 'JANSSON_WITHOUT_TESTS': true }) - jansson_opts.add_cmake_defines({ 'JANSSON_BUILD_DOCS': false }) - jansson_opts.add_cmake_defines({ 'JANSSON_INSTALL': false }) - jansson = cmake.subproject('jansson', options: jansson_opts).dependency('jansson') +if get_option('json').enabled() + jansson = dependency('jansson', required: false) + if not jansson.found() + jansson_opts = cmake.subproject_options() + jansson_opts.add_cmake_defines({ 'CMAKE_BUILD_TYPE': is_debug ? 'Debug' : 'Release' }) + jansson_opts.add_cmake_defines({ 'JANSSON_EXAMPLES': false }) + jansson_opts.add_cmake_defines({ 'JANSSON_WITHOUT_TESTS': true }) + jansson_opts.add_cmake_defines({ 'JANSSON_BUILD_DOCS': false }) + jansson_opts.add_cmake_defines({ 'JANSSON_INSTALL': false }) + jansson = cmake.subproject('jansson', options: jansson_opts).dependency('jansson') + endif + akiyo_deps += [jansson] endif -akiyo_deps += [jansson] if not is_windows akiyo_src += [ diff --git a/meson_options.txt b/meson_options.txt index e7aef06..1156b3d 100644 --- a/meson_options.txt +++ b/meson_options.txt @@ -1,3 +1,4 @@ option('event-loop', type: 'feature', value: 'enabled') option('curl', type: 'feature', value: 'enabled') +option('json', type: 'feature', value: 'enabled') option('tests', type: 'boolean', value: true) diff --git a/src/common.c b/src/common.c index 803ce97..ae9edf7 100644 --- a/src/common.c +++ b/src/common.c @@ -1,6 +1,7 @@ #include #include #include +#include #ifndef _WIN32 #include @@ -28,6 +29,8 @@ bool aki_common_init(void) al_set_page_size(info.dwPageSize); #endif + signal(SIGPIPE, SIG_IGN); + return true; } diff --git a/src/curl/http.c b/src/curl/http.c index 5674a18..c58d5a6 100644 --- a/src/curl/http.c +++ b/src/curl/http.c @@ -31,6 +31,7 @@ static void http_handle_events(struct aki_curl *curl) CURLcode code = msg->data.result; if (code != CURLE_OK) { al_log_error("http", "Curl error: %s.", curl_easy_strerror(code)); + http->callback(http->userdata, AKI_HTTP_ERROR, NULL, -1); return; } if (!aki_curl_remove_handle(curl)) return; diff --git a/src/curl/websocket.c b/src/curl/websocket.c index 7e1d9b9..4d480c3 100644 --- a/src/curl/websocket.c +++ b/src/curl/websocket.c @@ -28,7 +28,7 @@ static size_t stream_write_callback(char *buffer, size_t size, size_t nmemb, voi struct aki_websocket *ws = (struct aki_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); + al_log_debug("ws", "frame_size: %zu, bytesleft: %zu.", frame_size, m->bytesleft); aki_buffer_ensure_space(&ws->frame, frame_size); al_memcpy(aki_buffer_get_ptr(&ws->frame, m->offset), buffer, m->len); if (m->bytesleft == 0) { diff --git a/src/ev_embed_compat.c b/src/ev_embed_compat.c index 8e36572..ab98423 100644 --- a/src/ev_embed_compat.c +++ b/src/ev_embed_compat.c @@ -14,8 +14,9 @@ _Pragma("GCC diagnostic push") \ _Pragma("GCC diagnostic ignored \"-Wstrict-aliasing\"") _Pragma("GCC diagnostic ignored \"-Wcomment\"") -_Pragma("GCC diagnostic ignored \"-Wunused-value\"") _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 \"-Wsign-compare\"") _Pragma("GCC diagnostic ignored \"-Wreturn-type\"") diff --git a/src/line_processor.c b/src/line_processor.c index 158c5cb..b7d4ef0 100644 --- a/src/line_processor.c +++ b/src/line_processor.c @@ -12,7 +12,7 @@ static void remove_client(struct aki_line_processor *pro, struct aki_line_proces { ev_io_stop(pro->loop->ev, &client->event); if (client->pro->type == AKI_LINE_PROCESSOR_SOCKET) { - aki_socket_close(&client->s); + aki_socket_close(&client->sock); } aki_buffer_free(&client->buf); struct aki_line_processor_client *rclient; @@ -34,7 +34,7 @@ static void read_callback(struct ev_loop *loop, ev_io *w, s32 revents) u8 *ptr = aki_buffer_get_ptr(&client->buf, client->index); ssize_t ret; if (client->pro->type == AKI_LINE_PROCESSOR_SOCKET) { - ret = aki_socket_read(&client->s, ptr, CHUNK_SIZE); + ret = aki_socket_read(&client->sock, ptr, CHUNK_SIZE); } else { ret = aki_read(client->pro->fd, ptr, CHUNK_SIZE); } @@ -75,19 +75,18 @@ static void socket_connection_callback(struct ev_loop *loop, ev_io *w, s32 reven struct aki_line_processor *pro = (struct aki_line_processor *)w->data; (void)revents; struct aki_line_processor_client *client = al_alloc_object(struct aki_line_processor_client); - if (!aki_socket_accept(pro->s, &client->s)) { + if (!aki_socket_accept(pro->sock, &client->sock)) { al_free(client); return; } - aki_socket_set_blocking(&client->s, false); - aki_socket_set_no_delay(&client->s, 1); + aki_socket_set_blocking(&client->sock, false); client->index = 0; aki_buffer_init(&client->buf); client->mark = 0; client->pro = pro; client->event.data = client; al_array_push(pro->clients, client); - ev_io_init(&client->event, read_callback, aki_socket_get_fd(&client->s), EV_READ); + ev_io_init(&client->event, read_callback, aki_socket_get_fd(&client->sock), EV_READ); ev_io_start(pro->loop->ev, &client->event); } @@ -98,10 +97,10 @@ void aki_line_processor_init(struct aki_line_processor *pro, str *delim) al_array_init(pro->clients); } -void aki_line_processor_open_socket(struct aki_line_processor *pro, struct aki_socket *socket) +void aki_line_processor_open_socket(struct aki_line_processor *pro, struct aki_socket *sock) { pro->type = AKI_LINE_PROCESSOR_SOCKET; - pro->s = socket; + pro->sock = sock; } void aki_line_processor_open_fd(struct aki_line_processor *pro, s32 fd) @@ -116,7 +115,7 @@ void aki_line_processor_run(struct aki_line_processor *pro, struct aki_event_loo switch (pro->type) { case AKI_LINE_PROCESSOR_SOCKET: pro->event.data = pro; - ev_io_init(&pro->event, socket_connection_callback, aki_socket_get_fd(pro->s), EV_READ); + ev_io_init(&pro->event, socket_connection_callback, aki_socket_get_fd(pro->sock), EV_READ); ev_io_start(pro->loop->ev, &pro->event); break; case AKI_LINE_PROCESSOR_FD: { // Not supported on windows. diff --git a/src/line_processor.h b/src/line_processor.h index cc391ca..20f1101 100644 --- a/src/line_processor.h +++ b/src/line_processor.h @@ -24,7 +24,7 @@ enum { struct aki_line_processor_client { ev_io event; - struct aki_socket s; + struct aki_socket sock; size_t index; struct aki_buffer buf; size_t mark; @@ -35,7 +35,7 @@ struct aki_line_processor { struct aki_event_loop *loop; u8 type; ev_io event; - struct aki_socket *s; + struct aki_socket *sock; s32 fd; str delim; array(struct aki_line_processor_client *) clients; @@ -44,7 +44,7 @@ struct aki_line_processor { }; void aki_line_processor_init(struct aki_line_processor *pro, str *delim); -void aki_line_processor_open_socket(struct aki_line_processor *pro, struct aki_socket *socket); +void aki_line_processor_open_socket(struct aki_line_processor *pro, struct aki_socket *sock); void aki_line_processor_open_fd(struct aki_line_processor *pro, s32 fd); void aki_line_processor_run(struct aki_line_processor *pro, struct aki_event_loop *loop); void aki_line_processor_stop(struct aki_line_processor *pro); diff --git a/src/multiplex.c b/src/multiplex.c new file mode 100644 index 0000000..d3bb14e --- /dev/null +++ b/src/multiplex.c @@ -0,0 +1,47 @@ +#include "multiplex.h" + +bool aki_multiplex_socket_init(struct aki_multiplex_socket *multi, u8 type, + bool (*connection_callback)(void *, u8, struct aki_socket *), void *userdata) +{ + multi->sock.type = type; + if (!aki_socket_init(&multi->sock)) return false; + aki_socket_set_blocking(&multi->sock, false); + multi->connection_callback = connection_callback; + multi->userdata = userdata; + return true; +} + +static void socket_connection_callback(struct ev_loop *loop, ev_io *w, s32 revents) +{ + (void)loop; + struct aki_multiplex_socket *multi = (struct aki_multiplex_socket *)w->data; + (void)revents; + struct aki_socket sock; + if (!aki_socket_accept(&multi->sock, &sock)) return; + u8 id; + ssize_t ret = aki_socket_read(&sock, &id, sizeof(u8)); + if (ret > 0) { + al_assert(ret == sizeof(u8)); + if (multi->connection_callback(multi->userdata, id, &sock)) { + return; + } + } + aki_socket_close(&sock); +} + +bool aki_multiplex_socket_listen(struct aki_multiplex_socket *multi, struct aki_event_loop *loop, str *addr, s32 port) +{ + if (!aki_socket_listen(&multi->sock, addr, port)) return false; + multi->loop = loop; + multi->event.data = multi; + ev_io_init(&multi->event, socket_connection_callback, aki_socket_get_fd(&multi->sock), EV_READ); + ev_io_start(multi->loop->ev, &multi->event); + return true; +} + +void aki_multiplex_socket_close(struct aki_multiplex_socket *multi) +{ + aki_socket_shutdown(&multi->sock); + aki_socket_close(&multi->sock); + ev_io_stop(multi->loop->ev, &multi->event); +} diff --git a/src/multiplex.h b/src/multiplex.h new file mode 100644 index 0000000..cd20c29 --- /dev/null +++ b/src/multiplex.h @@ -0,0 +1,20 @@ +#pragma once + +#include + +#include "socket/socket.h" + +#include "loop.h" + +struct aki_multiplex_socket { + struct aki_socket sock; + struct aki_event_loop *loop; + ev_io event; + bool (*connection_callback)(void *, u8 id, struct aki_socket *); + void *userdata; +}; + +bool aki_multiplex_socket_init(struct aki_multiplex_socket *multi, u8 type, + bool (*connection_callback)(void *, u8, struct aki_socket *), void *userdata); +bool aki_multiplex_socket_listen(struct aki_multiplex_socket *multi, struct aki_event_loop *loop, str *addr, s32 port); +void aki_multiplex_socket_close(struct aki_multiplex_socket *multi); diff --git a/src/packet_stream.c b/src/packet_stream.c index 9064a1f..f653c54 100644 --- a/src/packet_stream.c +++ b/src/packet_stream.c @@ -6,33 +6,38 @@ static void init_internal(struct aki_packet_stream *stream) stream->out.packet = NULL; stream->out.have_data = false; stream->in.packet = aki_packet_create(); - stream->in.index = 0; stream->in.have_size = false; + stream->in.index = 0; } bool aki_packet_stream_init(struct aki_packet_stream *stream, u8 type, void (*connection_callback)(void *, struct aki_packet_stream *), - void (*connection_closed_callback)(void *, struct aki_packet_stream *), - void (*packet_callback)(void *, struct aki_packet_stream *, struct aki_packet *), - void (*packet_sent_callback)(void *, struct aki_packet *), void *userdata) + void (*connection_closed_callback)(void *, struct aki_packet_stream *), void *userdata) { + stream->sock.type = type; + if (!aki_socket_init(&stream->sock)) { + return false; + } + aki_socket_set_blocking(&stream->sock, false); + stream->corked = false; + stream->connected = false; stream->connection_callback = connection_callback; stream->connection_closed_callback = connection_closed_callback; - stream->packet_callback = packet_callback; - stream->packet_sent_callback = packet_sent_callback; - stream->flushed_callback = NULL; + stream->packet_callback = NULL; + stream->packet_sent_callback = NULL; stream->userdata = userdata; - stream->s.type = type; - if (!aki_socket_init(&stream->s)) return false; - aki_socket_set_blocking(&stream->s, false); init_internal(stream); - stream->connected = false; return true; } -void aki_packet_stream_set_no_delay(struct aki_packet_stream *stream, s32 no_delay) +void aki_packet_stream_set_nodelay(struct aki_packet_stream *stream, s32 nodelay) +{ + aki_socket_set_nodelay(&stream->sock, nodelay); +} + +void aki_packet_stream_set_multiplex(struct aki_packet_stream *stream, u8 id) { - aki_socket_set_no_delay(&stream->s, no_delay); + stream->id = id; } static void stop_write_internal(struct aki_packet_stream *stream) @@ -48,13 +53,13 @@ static void stop_internal(struct aki_packet_stream *stream) // ev_timer_stop(stream->loop->ev, &stream->timer); // stream->timer_started = false; //} - if (stream->revent.active) { + if (!stream->corked) { ev_io_stop(stream->loop->ev, &stream->revent); } if (stream->out.have_data) { stop_write_internal(stream); } - aki_socket_close(&stream->s); + aki_socket_close(&stream->sock); stream->connection_closed_callback(stream->userdata, stream); } @@ -63,9 +68,10 @@ static void stream_read_callback(struct ev_loop *loop, ev_io *w, s32 revents) (void)loop; struct aki_packet_stream *stream = (struct aki_packet_stream *)w->data; (void)revents; + // EV_READ can mean any of POLLIN, POLLERR, or POLLHUP. u8 *ptr = aki_buffer_get_ptr(&stream->in.packet->buffer, stream->in.index); u32 size = ((stream->in.have_size) ? aki_packet_size(stream->in.packet) : AKI_PACKET_HEADER_LENGTH); - ssize_t ret = aki_socket_read(&stream->s, ptr, size - stream->in.index); + ssize_t ret = aki_socket_read(&stream->sock, ptr, size - stream->in.index); if (ret <= 0) { stop_internal(stream); return; @@ -86,24 +92,29 @@ static void stream_read_callback(struct ev_loop *loop, ev_io *w, s32 revents) } } -static void start_read_internal(struct aki_packet_stream *stream) -{ - ev_io_start(stream->loop->ev, &stream->revent); -} - static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents) { (void)loop; struct aki_packet_stream *stream = (struct aki_packet_stream *)w->data; (void)revents; + // EV_WRITE can mean any of POLLOUT, POLLERR, or POLLHUP. if (UNLIKELY(!stream->connected)) { stream->connected = true; - stop_write_internal(stream); - // This needs to be called before starting to read because setting - // connection_closed_callback and/or packet_callback here should be valid. + // Shove in the multiplex ID byte here assuming POLLOUT means + // we can write a minimum of 1 byte. + // If there is a case where POLLOUT can end up writing 0 bytes, we + // should wait to set stream->connected until write() returns >0 and + // return before connection_callback(). + ssize_t ret = aki_socket_write(&stream->sock, &stream->id, sizeof(u8)); + // Assume any error (ret < 0) will be resolved in the read callback. + al_assert(ret <= 0 || ret == sizeof(u8)); + // Packet callbacks are set in connection_callback. stream->connection_callback(stream->userdata, stream); - start_read_internal(stream); - return; + al_assert(stream->packet_callback); + al_assert(stream->packet_sent_callback); + if (!stream->corked) { + ev_io_start(stream->loop->ev, &stream->revent); + } } if (!stream->out.packet) { if (stream->out.queue.size > 0) { @@ -116,7 +127,7 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents) } u8 *ptr = aki_buffer_get_ptr(&stream->out.packet->buffer, stream->out.index); u32 size = aki_packet_size(stream->out.packet); - ssize_t ret = aki_socket_write(&stream->s, ptr, size - stream->out.index); + ssize_t ret = aki_socket_write(&stream->sock, ptr, size - stream->out.index); if (ret <= 0) { stop_write_internal(stream); return; @@ -125,46 +136,25 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents) if (stream->out.index >= size) { stream->packet_sent_callback(stream->userdata, stream->out.packet); stream->out.packet = NULL; - if (stream->out.queue.size == 0 && stream->flushed_callback) { - stream->flushed_callback(stream->userdata); - } } } -static void socket_connection_callback(struct ev_loop *loop, ev_io *w, s32 revents) -{ - (void)loop; - struct aki_packet_stream *server = (struct aki_packet_stream *)w->data; - (void)revents; - struct aki_packet_stream *client = al_alloc_object(struct aki_packet_stream); - if (!aki_socket_accept(&server->s, &client->s)) { - al_free(client); - return; - } - aki_socket_set_blocking(&client->s, false); - client->connected = true; - client->loop = server->loop; - init_internal(client); - server->connection_callback(server->userdata, client); - client->wevent.data = client; - ev_io_init(&client->wevent, stream_write_callback, - aki_socket_get_fd(&client->s), EV_WRITE); - client->revent.data = client; - ev_io_init(&client->revent, stream_read_callback, - aki_socket_get_fd(&client->s), EV_READ); - ev_io_start(client->loop->ev, &client->revent); -} - -void aki_packet_stream_listen(struct aki_packet_stream *stream, - struct aki_event_loop *loop, str *addr, s32 port) +void aki_packet_stream_from_socket(struct aki_packet_stream *stream, struct aki_event_loop *loop, struct aki_socket *sock) { stream->loop = loop; - // wevent not used for listening. + stream->corked = false; + stream->connected = true; + stream->sock = *sock; + aki_socket_set_blocking(&stream->sock, false); + init_internal(stream); + stream->connection_callback(stream->userdata, stream); + stream->wevent.data = stream; + ev_io_init(&stream->wevent, stream_write_callback, + aki_socket_get_fd(&stream->sock), EV_WRITE); stream->revent.data = stream; - ev_io_init(&stream->revent, socket_connection_callback, - aki_socket_get_fd(&stream->s), EV_READ); - ev_io_start(loop->ev, &stream->revent); - aki_socket_listen(&stream->s, addr, port); + ev_io_init(&stream->revent, stream_read_callback, + aki_socket_get_fd(&stream->sock), EV_READ); + ev_io_start(stream->loop->ev, &stream->revent); } /* @@ -183,35 +173,54 @@ static void start_write_internal(struct aki_packet_stream *stream) ev_io_start(stream->loop->ev, &stream->wevent); } -void aki_packet_stream_connect(struct aki_packet_stream *stream, - struct aki_event_loop *loop, str *addr, s32 port) +static void connect_internal(struct aki_packet_stream *stream, str *addr, s32 port) { - stream->loop = loop; + if (!aki_socket_connect(&stream->sock, addr, port)) { + stream->connection_closed_callback(stream->userdata, stream); + return; + } stream->revent.data = stream; ev_io_init(&stream->revent, stream_read_callback, - aki_socket_get_fd(&stream->s), EV_READ); + aki_socket_get_fd(&stream->sock), EV_READ); stream->wevent.data = stream; ev_io_init(&stream->wevent, stream_write_callback, - aki_socket_get_fd(&stream->s), EV_WRITE); + aki_socket_get_fd(&stream->sock), EV_WRITE); start_write_internal(stream); //stream->timer.data = stream; //ev_timer_init(&stream->timer, timeout_callback, 5.0, 0.0); - if (!aki_socket_connect(&stream->s, addr, port)) { - stop_internal(stream); - return; - } //stream->timer_started = true; //ev_timer_start(loop->ev, &stream->timer); } +void aki_packet_stream_connect(struct aki_packet_stream *stream, struct aki_event_loop *loop, str *addr, s32 port) +{ + stream->loop = loop; + connect_internal(stream, addr, port); +} + +void aki_packet_stream_reconnect(struct aki_packet_stream *stream, str *addr, s32 port) +{ + if (!aki_socket_init(&stream->sock)) { + stream->connection_closed_callback(stream->userdata, stream); + return; + } + aki_packet_reset(stream->in.packet); + stream->in.have_size = false; + stream->in.index = 0; + aki_socket_set_blocking(&stream->sock, false); + connect_internal(stream, addr, port); +} + void aki_packet_stream_cork(struct aki_packet_stream *stream, bool cork) { - if (!stream->connected) return; - if (cork && stream->revent.active) { - ev_io_stop(stream->loop->ev, &stream->revent); - } else if (!cork && !stream->revent.active) { - ev_io_start(stream->loop->ev, &stream->revent); + if (stream->connected) { + if (cork && !stream->corked) { + ev_io_stop(stream->loop->ev, &stream->revent); + } else if (!cork && stream->corked) { + ev_io_start(stream->loop->ev, &stream->revent); + } } + stream->corked = cork; } void aki_packet_stream_send_packet(struct aki_packet_stream *stream, struct aki_packet *packet) @@ -224,30 +233,19 @@ void aki_packet_stream_send_packet(struct aki_packet_stream *stream, struct aki_ } } -void aki_packet_stream_flush(struct aki_packet_stream *stream, void (*flushed_callback)(void *)) -{ - stream->flushed_callback = flushed_callback; - if (!stream->out.have_data) { - stream->flushed_callback(stream->userdata); - } -} - void aki_packet_stream_disconnect(struct aki_packet_stream *stream) { - if (stream->connected) { - aki_socket_shutdown(&stream->s); - if (stream->revent.active) { - ev_io_stop(stream->loop->ev, &stream->revent); - } - stop_internal(stream); - } + if (!stream->connected) return; + aki_socket_shutdown(&stream->sock); + stop_internal(stream); } void aki_packet_stream_free(struct aki_packet_stream *stream) { - if (stream->in.packet) { - aki_packet_free(stream->in.packet); - } + // Loosely assert that stop_internal has ran. + al_assert(!stream->out.have_data); + al_assert(stream->in.packet); + aki_packet_free(stream->in.packet); if (stream->out.packet) { stream->packet_sent_callback(stream->userdata, stream->out.packet); } diff --git a/src/packet_stream.h b/src/packet_stream.h index 964bb9b..b678ba7 100644 --- a/src/packet_stream.h +++ b/src/packet_stream.h @@ -8,18 +8,19 @@ #include "loop.h" struct aki_packet_stream { - struct aki_socket s; + struct aki_socket sock; struct aki_event_loop *loop; ev_io revent; ev_io wevent; - ev_timer timer; - bool timer_started; + bool corked; + u8 id; // multiplex + //ev_timer timer; + //bool timer_started; bool connected; void (*connection_callback)(void *, struct aki_packet_stream *); void (*connection_closed_callback)(void *, struct aki_packet_stream *); void (*packet_callback)(void *, struct aki_packet_stream *, struct aki_packet *); void (*packet_sent_callback)(void *, struct aki_packet *); - void (*flushed_callback)(void *); void *userdata; struct { array(struct aki_packet *) queue; @@ -36,16 +37,13 @@ struct aki_packet_stream { bool aki_packet_stream_init(struct aki_packet_stream *stream, u8 type, void (*connection_callback)(void *, struct aki_packet_stream *), - void (*connection_closed_callback)(void *, struct aki_packet_stream *), - void (*packet_callback)(void *, struct aki_packet_stream *, struct aki_packet *), - void (*packet_sent_callback)(void *, struct aki_packet *), void *userdata); -void aki_packet_stream_set_no_delay(struct aki_packet_stream *stream, s32 no_delay); -void aki_packet_stream_listen(struct aki_packet_stream *stream, - struct aki_event_loop *loop, str *addr, s32 port); -void aki_packet_stream_connect(struct aki_packet_stream *stream, - struct aki_event_loop *loop, str *addr, s32 port); + void (*connection_closed_callback)(void *, struct aki_packet_stream *), void *userdata); +void aki_packet_stream_set_nodelay(struct aki_packet_stream *stream, s32 no_delay); +void aki_packet_stream_set_multiplex(struct aki_packet_stream *stream, u8 id); +void aki_packet_stream_from_socket(struct aki_packet_stream *stream, struct aki_event_loop *loop, struct aki_socket *sock); +void aki_packet_stream_connect(struct aki_packet_stream *stream, struct aki_event_loop *loop, str *addr, s32 port); +void aki_packet_stream_reconnect(struct aki_packet_stream *stream, str *addr, s32 port); void aki_packet_stream_cork(struct aki_packet_stream *stream, bool cork); void aki_packet_stream_send_packet(struct aki_packet_stream *stream, struct aki_packet *packet); -void aki_packet_stream_flush(struct aki_packet_stream *stream, void (*flushed_callback)(void *)); void aki_packet_stream_disconnect(struct aki_packet_stream *stream); void aki_packet_stream_free(struct aki_packet_stream *stream); diff --git a/src/rpc.c b/src/rpc.c index 6891ea3..0c0046a 100644 --- a/src/rpc.c +++ b/src/rpc.c @@ -1,13 +1,26 @@ #include "rpc2.h" -static void connection_closed_callback(void *userdata, struct aki_packet_stream *s) +bool aki_rpc_init(struct aki_rpc *rpc, struct aki_event_loop *loop, + void (*connection_callback)(void *, struct aki_rpc_connection *), + void (*connection_closed_callback)(void *, struct aki_rpc_connection *), void *userdata) { - struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata; - al_assert(s == conn->stream); - conn->rpc->connection_closed_callback(conn->rpc->userdata, conn); + rpc->loop = loop; + rpc->increment = 0; + al_array_init(rpc->commands); + rpc->conn = NULL; + al_array_init(rpc->connections); + rpc->connection_callback = connection_callback; + rpc->connection_closed_callback = connection_closed_callback; + rpc->userdata = userdata; + return true; +} + +void aki_rpc_add_command(struct aki_rpc *rpc, struct aki_rpc_command *command) +{ + al_array_push(rpc->commands, *command); } -static void packet_callback(void *userdata, struct aki_packet_stream *s, struct aki_packet *packet) +static void packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) { struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata; s8 op = aki_packet_read_s8(packet); @@ -29,7 +42,7 @@ static void packet_callback(void *userdata, struct aki_packet_stream *s, struct aki_packet_write_s8(rpacket, -1); aki_packet_write_u32(rpacket, id); if (command->callback(command->userdata, conn, packet, rpacket)) { - aki_packet_stream_send_packet(s, rpacket); + aki_packet_stream_send_packet(stream, rpacket); } else { aki_packet_free(rpacket); } @@ -46,84 +59,74 @@ static void packet_sent_callback(void *userdata, struct aki_packet *packet) aki_packet_free(packet); } -static void connection_callback(void *userdata, struct aki_packet_stream *s) +static void stream_connection_callback(void *userdata, struct aki_packet_stream *stream) { - struct aki_rpc *rpc = (struct aki_rpc *)userdata; - aki_packet_stream_set_no_delay(s, true); - struct aki_rpc_connection *conn = al_alloc_object(struct aki_rpc_connection); - al_array_push(rpc->connections, conn); - conn->rpc = rpc; - conn->stream = s; - // connection_callback will never be called twice for the same stream. - s->connection_closed_callback = connection_closed_callback; - s->packet_callback = packet_callback; - s->packet_sent_callback = packet_sent_callback; - s->flushed_callback = NULL; - // For aki_rpc_connect this overwrites rpc as userdata, whereas for aki_rpc_listen - // this is a new aki_packet_stream object. - s->userdata = conn; - al_array_init(conn->callbacks); + struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata; + struct aki_rpc *rpc = conn->rpc; + aki_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 in server mode. + if (!rpc->conn) al_array_push(rpc->connections, conn); } -bool aki_rpc_init(struct aki_rpc *rpc, u8 type, - void (*connection_callback)(void *, struct aki_rpc_connection *), - void (*connection_closed_callback)(void *, struct aki_rpc_connection *), void *userdata) +static void stream_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) { - rpc->type = type; - rpc->id = 0; - rpc->connection_callback = connection_callback; - rpc->connection_closed_callback = connection_closed_callback; - rpc->userdata = userdata; - al_array_init(rpc->connections); - al_array_init(rpc->commands); - return true; + struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata; + al_assert(stream == &conn->stream); + conn->rpc->connection_closed_callback(conn->rpc->userdata, conn); } -void aki_rpc_add_command(struct aki_rpc *rpc, struct aki_rpc_command *command) +void aki_rpc_add_socket(struct aki_rpc *rpc, struct aki_socket *sock) { - al_array_push(rpc->commands, *command); + struct aki_rpc_connection *conn = al_alloc_object(struct aki_rpc_connection); + conn->rpc = rpc; + al_array_init(conn->callbacks); + conn->stream.connection_callback = stream_connection_callback; + conn->stream.connection_closed_callback = stream_connection_closed_callback; + conn->stream.userdata = conn; + aki_packet_stream_from_socket(&conn->stream, rpc->loop, sock); } -bool aki_rpc_listen(struct aki_rpc *rpc, struct aki_event_loop *loop, str *addr, s32 port) +bool aki_rpc_prepare_client(struct aki_rpc *rpc, u8 type, u8 id) { - if (!aki_packet_stream_init(&rpc->stream, rpc->type, connection_callback, NULL, NULL, NULL, rpc)) { + al_assert(!rpc->conn); + struct aki_rpc_connection *conn = al_alloc_object(struct aki_rpc_connection); + rpc->conn = conn; + al_array_push(rpc->connections, conn); + if (!aki_packet_stream_init(&conn->stream, type, stream_connection_callback, + stream_connection_closed_callback, conn)) { return false; } - aki_packet_stream_set_no_delay(&rpc->stream, true); - aki_packet_stream_listen(&rpc->stream, loop, addr, port); + conn->rpc = rpc; + al_array_init(conn->callbacks); + aki_packet_stream_set_nodelay(&conn->stream, 1); + aki_packet_stream_set_multiplex(&conn->stream, id); return true; } -bool aki_rpc_connect(struct aki_rpc *rpc, struct aki_event_loop *loop, str *addr, s32 port) +void aki_rpc_connect(struct aki_rpc *rpc, str *addr, s32 port) { - if (!aki_packet_stream_init(&rpc->stream, rpc->type, connection_callback, NULL, NULL, NULL, rpc)) { - return false; - } - aki_packet_stream_set_no_delay(&rpc->stream, true); - aki_packet_stream_connect(&rpc->stream, loop, addr, port); - return true; + al_assert(rpc->conn); + aki_packet_stream_connect(&rpc->conn->stream, rpc->loop, addr, port); } struct aki_packet *aki_rpc_get_packet(struct aki_rpc *rpc, s8 op) { struct aki_packet *packet = aki_packet_create(); aki_packet_write_s8(packet, op); - aki_packet_write_u32(packet, rpc->id); - rpc->id = al_u32_inc_wrap(rpc->id); + aki_packet_write_u32(packet, rpc->increment); + rpc->increment = al_u32_inc_wrap(rpc->increment); return packet; } -void aki_rpc_disconnect(struct aki_rpc *rpc) -{ - aki_packet_stream_disconnect(&rpc->stream); -} - void aki_rpc_free(struct aki_rpc *rpc) { - aki_packet_stream_free(&rpc->stream); struct aki_rpc_connection *conn; al_array_foreach(rpc->connections, i, conn) { + aki_packet_stream_free(&conn->stream); al_array_free(conn->callbacks); al_free(conn); } @@ -141,5 +144,10 @@ void aki_rpc_connection_command(struct aki_rpc_connection *conn, struct aki_pack .userdata = userdata })); } - aki_packet_stream_send_packet(conn->stream, packet); + aki_packet_stream_send_packet(&conn->stream, packet); +} + +void aki_rpc_conn_disconnect(struct aki_rpc_connection *conn) +{ + aki_packet_stream_disconnect(&conn->stream); } diff --git a/src/rpc2.h b/src/rpc2.h index 482ad56..fd98c2d 100644 --- a/src/rpc2.h +++ b/src/rpc2.h @@ -13,9 +13,9 @@ struct aki_rpc_callback { }; struct aki_rpc_connection { - struct aki_rpc *rpc; - struct aki_packet_stream *stream; + struct aki_packet_stream stream; array(struct aki_rpc_callback) callbacks; + struct aki_rpc *rpc; }; struct aki_rpc_command { @@ -25,25 +25,26 @@ struct aki_rpc_command { }; struct aki_rpc { - u8 type; - struct aki_packet_stream stream; + struct aki_event_loop *loop; + u32 increment; array(struct aki_rpc_command) commands; + struct aki_rpc_connection *conn; // client array(struct aki_rpc_connection *) connections; - u32 id; void (*connection_callback)(void *, struct aki_rpc_connection *); void (*connection_closed_callback)(void *, struct aki_rpc_connection *); void *userdata; }; -bool aki_rpc_init(struct aki_rpc *rpc, u8 type, +bool aki_rpc_init(struct aki_rpc *rpc, struct aki_event_loop *loop, void (*connection_callback)(void *, struct aki_rpc_connection *), void (*connection_closed_callback)(void *, struct aki_rpc_connection *), void *userdata); void aki_rpc_add_command(struct aki_rpc *rpc, struct aki_rpc_command *command); -bool aki_rpc_listen(struct aki_rpc *rpc, struct aki_event_loop *loop, str *addr, s32 port); -bool aki_rpc_connect(struct aki_rpc *rpc, struct aki_event_loop *loop, str *addr, s32 port); +void aki_rpc_add_socket(struct aki_rpc *rpc, struct aki_socket *sock); +bool aki_rpc_prepare_client(struct aki_rpc *rpc, u8 type, u8 id); +void aki_rpc_connect(struct aki_rpc *rpc, str *addr, s32 port); struct aki_packet *aki_rpc_get_packet(struct aki_rpc *rpc, s8 op); -void aki_rpc_disconnect(struct aki_rpc *rpc); void aki_rpc_free(struct aki_rpc *rpc); void aki_rpc_connection_command(struct aki_rpc_connection *conn, struct aki_packet *packet, void (*callback)(void *, struct aki_packet *), void *userdata); +void aki_rpc_conn_disconnect(struct aki_rpc_connection *conn); diff --git a/src/socket/socket.h b/src/socket/socket.h index c0cf796..bb42a3d 100644 --- a/src/socket/socket.h +++ b/src/socket/socket.h @@ -25,11 +25,11 @@ enum { struct aki_socket { u8 type; #ifndef _WIN32 - s32 sock; + s32 fd; struct sockaddr_in addr_in; struct sockaddr_un addr_un; #else - SOCKET sock; + SOCKET fd; SOCKADDR_IN addr_in; #endif }; @@ -45,6 +45,7 @@ struct aki_socket { #define aki_htons htons #define aki_htonl htonl #define aki_ntohs ntohs +#define aki_ntohl ntohl #ifndef _WIN32 // Posix-only helpers. #define aki_pollfd pollfd @@ -53,24 +54,24 @@ void aki_fd_set_blocking(s32 fd, bool blocking); s32 aki_poll_fds(struct aki_pollfd *fds, aki_nfds nfds, s64 timeout_ns); #endif -bool aki_socket_init(struct aki_socket *s); +bool aki_socket_init(struct aki_socket *sock); -void aki_socket_set_blocking(struct aki_socket *s, bool blocking); -void aki_socket_set_no_delay(struct aki_socket *s, s32 no_delay); +void aki_socket_set_blocking(struct aki_socket *sock, bool blocking); +void aki_socket_set_nodelay(struct aki_socket *sock, s32 nodelay); -void aki_socket_set_send_buf(struct aki_socket *s, u32 sndbuf); -u32 aki_socket_get_send_buf(struct aki_socket *s); +void aki_socket_set_send_buf(struct aki_socket *sock, u32 sndbuf); +u32 aki_socket_get_send_buf(struct aki_socket *sock); -bool aki_socket_listen(struct aki_socket *s, str *addr, s32 port); -bool aki_socket_accept(struct aki_socket *s, struct aki_socket *c); -bool aki_socket_connect(struct aki_socket *s, str *addr, s32 port); +bool aki_socket_listen(struct aki_socket *sock, str *addr, s32 port); +bool aki_socket_accept(struct aki_socket *sock, struct aki_socket *c); +bool aki_socket_connect(struct aki_socket *sock, str *addr, s32 port); -s32 aki_socket_get_fd(struct aki_socket *s); +s32 aki_socket_get_fd(struct aki_socket *sock); -ssize_t aki_socket_read(struct aki_socket *s, void *buf, size_t size); -ssize_t aki_socket_write(struct aki_socket *s, void *buf, size_t size); -ssize_t aki_socket_sendto(struct aki_socket *s, void *buf, size_t size); -ssize_t aki_socket_recvfrom(struct aki_socket *s, void *buf, size_t size); +ssize_t aki_socket_read(struct aki_socket *sock, void *buf, size_t size); +ssize_t aki_socket_write(struct aki_socket *sock, void *buf, size_t size); +ssize_t aki_socket_sendto(struct aki_socket *sock, void *buf, size_t size); +ssize_t aki_socket_recvfrom(struct aki_socket *sock, void *buf, size_t size); -void aki_socket_shutdown(struct aki_socket *s); -void aki_socket_close(struct aki_socket *s); +void aki_socket_shutdown(struct aki_socket *sock); +void aki_socket_close(struct aki_socket *sock); diff --git a/src/socket/socket_linux.c b/src/socket/socket_linux.c index b442806..adab91b 100644 --- a/src/socket/socket_linux.c +++ b/src/socket/socket_linux.c @@ -1,9 +1,8 @@ +#include #include #include #include #include -#include -#include #include "../util/error.h" @@ -11,40 +10,38 @@ void aki_fd_set_blocking(s32 fd, bool blocking) { - s32 flags = fcntl(fd, F_GETFL, 0); - al_assert(flags != -1); - flags = blocking ? (flags & ~O_NONBLOCK) : (flags | O_NONBLOCK); - s32 ret = fcntl(fd, F_SETFL, flags); + s32 ret = fcntl(fd, F_GETFL, 0); // ret = flags + al_assert(ret != -1); + ret = fcntl(fd, F_SETFL, blocking ? (ret & ~O_NONBLOCK) : (ret | O_NONBLOCK)); // ret = success/error al_assert(ret == 0); } -#define AKI_TIME_S_TO_NS(s) ((s) * INT64_C(1000000000)) - // https://github.com/mpv-player/mpv/blob/e575ec4fc3654387c7358bd3640877ef32628d2c/osdep/poll_wrapper.c#L29 +#define AKI_TIME_S_TO_NS(s) ((s) * INT64_C(1000000000)) s32 aki_poll_fds(struct aki_pollfd *fds, aki_nfds nfds, s64 timeout_ns) { struct timespec ts; - ts.tv_sec = timeout_ns / AKI_TIME_S_TO_NS(1); + ts.tv_sec = timeout_ns / AKI_TIME_S_TO_NS(1); ts.tv_nsec = timeout_ns % AKI_TIME_S_TO_NS(1); struct timespec *tsp = timeout_ns >= 0 ? &ts : NULL; return ppoll(fds, nfds, tsp, NULL); } -bool aki_socket_init(struct aki_socket *s) +bool aki_socket_init(struct aki_socket *sock) { - switch (s->type) { + switch (sock->type) { case AKI_SOCKET_TCP: - s->sock = socket(AF_INET, SOCK_STREAM, 0); + sock->fd = socket(AF_INET, SOCK_STREAM, 0); break; case AKI_SOCKET_UDP: - s->sock = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); + sock->fd = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); break; case AKI_SOCKET_UNIX: - s->sock = socket(AF_UNIX, SOCK_STREAM, 0); + sock->fd = socket(AF_UNIX, SOCK_STREAM, 0); break; } - if (s->sock < 0) { + if (sock->fd < 0) { al_log_error("socket", "socket() failed: %s (%d).", aki_strerror(errno), errno); return false; } @@ -52,125 +49,123 @@ bool aki_socket_init(struct aki_socket *s) return true; } -void aki_socket_set_blocking(struct aki_socket *s, bool blocking) +void aki_socket_set_blocking(struct aki_socket *sock, bool blocking) { - aki_fd_set_blocking(s->sock, blocking); + aki_fd_set_blocking(sock->fd, blocking); } -void aki_socket_set_no_delay(struct aki_socket *s, s32 no_delay) +void aki_socket_set_nodelay(struct aki_socket *sock, s32 nodelay) { - if (s->type != AKI_SOCKET_TCP) return; - setsockopt(s->sock, IPPROTO_TCP, TCP_NODELAY, &no_delay, sizeof(no_delay)); + //al_assert(sock->type == AKI_SOCKET_TCP); + if (sock->type != AKI_SOCKET_TCP) return; // TODO: + setsockopt(sock->fd, IPPROTO_TCP, TCP_NODELAY, &nodelay, sizeof(nodelay)); } -u32 aki_socket_get_send_buf(struct aki_socket *s) +u32 aki_socket_get_send_buf(struct aki_socket *sock) { u32 send_queue_size; socklen_t optlen = sizeof(send_queue_size); - getsockopt(s->sock, SOL_SOCKET, SO_SNDBUF, &send_queue_size, &optlen); + getsockopt(sock->fd, SOL_SOCKET, SO_SNDBUF, &send_queue_size, &optlen); return send_queue_size / 2; } -void aki_socket_set_send_buf(struct aki_socket *s, u32 sndbuf) +void aki_socket_set_send_buf(struct aki_socket *sock, u32 sndbuf) { u32 send_queue_size = sndbuf; - setsockopt(s->sock, SOL_SOCKET, SO_SNDBUF, &send_queue_size, sizeof(send_queue_size)); - al_assert(aki_socket_get_send_buf(s) == send_queue_size); + setsockopt(sock->fd, SOL_SOCKET, SO_SNDBUF, &send_queue_size, sizeof(send_queue_size)); + al_assert(aki_socket_get_send_buf(sock) == send_queue_size); } -u32 aki_socket_get_recv_buf(struct aki_socket *s) +u32 aki_socket_get_recv_buf(struct aki_socket *sock) { u32 receive_queue_size; socklen_t optlen = sizeof(receive_queue_size); - getsockopt(s->sock, SOL_SOCKET, SO_RCVBUF, &receive_queue_size, &optlen); + getsockopt(sock->fd, SOL_SOCKET, SO_RCVBUF, &receive_queue_size, &optlen); return receive_queue_size / 2; } -void aki_socket_set_recv_buf(struct aki_socket *s, u32 rcvbuf) +void aki_socket_set_recv_buf(struct aki_socket *sock, u32 rcvbuf) { u32 receive_queue_size = rcvbuf; - setsockopt(s->sock, SOL_SOCKET, SO_RCVBUF, &receive_queue_size, sizeof(receive_queue_size)); - al_assert(aki_socket_get_recv_buf(s) == receive_queue_size); + setsockopt(sock->fd, SOL_SOCKET, SO_RCVBUF, &receive_queue_size, sizeof(receive_queue_size)); + al_assert(aki_socket_get_recv_buf(sock) == receive_queue_size); } #define al_log_socket_err(func) \ al_log_error("socket", #func"(%.*s:%d) failed: %s (%d).", \ AL_STR_PRINTF(addr), port, aki_strerror(errno), errno); -bool aki_socket_listen(struct aki_socket *s, str *addr, s32 port) +bool aki_socket_listen(struct aki_socket *sock, str *addr, s32 port) { struct sockaddr *saddr = NULL; socklen_t addrlen = 0; - switch (s->type) { + switch (sock->type) { case AKI_SOCKET_TCP: { - s->addr_in.sin_family = AF_INET; - s->addr_in.sin_addr.s_addr = htonl(INADDR_ANY); - s->addr_in.sin_port = htons(port); - saddr = (struct sockaddr *)&s->addr_in; - addrlen = sizeof(s->addr_in); + sock->addr_in.sin_family = AF_INET; + sock->addr_in.sin_addr.s_addr = htonl(INADDR_ANY); + sock->addr_in.sin_port = htons(port); + saddr = (struct sockaddr *)&sock->addr_in; + addrlen = sizeof(sock->addr_in); break; } case AKI_SOCKET_UNIX: { + sock->addr_un.sun_family = AF_UNIX; char *c_str = al_str_to_c_str(addr); - s->addr_un.sun_family = AF_UNIX; - size_t len = al_strlen(c_str); - al_memcpy(s->addr_un.sun_path, c_str, len); - s->addr_un.sun_path[len] = '\0'; + al_memcpy(sock->addr_un.sun_path, c_str, al_strlen(c_str) + 1); al_free(c_str); - saddr = (struct sockaddr *)&s->addr_un; - addrlen = sizeof(s->addr_un); + saddr = (struct sockaddr *)&sock->addr_un; + addrlen = sizeof(sock->addr_un); break; } case AKI_SOCKET_UDP: al_assert(false); } - if (bind(s->sock, saddr, addrlen) < 0) { + if (bind(sock->fd, saddr, addrlen) < 0) { al_log_socket_err(bind) return false; } - if (listen(s->sock, SOMAXCONN) < 0) { + if (listen(sock->fd, SOMAXCONN) < 0) { al_log_socket_err(listen) return false; } al_log_info("socket", "Listening on %.*s:%d.", AL_STR_PRINTF(addr), port); - signal(SIGPIPE, SIG_IGN); - return true; } -bool aki_socket_accept(struct aki_socket *s, struct aki_socket *c) +bool aki_socket_accept(struct aki_socket *sock, struct aki_socket *cl) { - c->type = s->type; + cl->type = sock->type; struct sockaddr_storage addr; socklen_t len = sizeof(addr); - if ((c->sock = accept(s->sock, (struct sockaddr *)&addr, &len)) == -1) { + if ((cl->fd = accept(sock->fd, (struct sockaddr *)&addr, &len)) == -1) { + al_log_error("socket", "accept() failed: %s (%d).", aki_strerror(errno), errno); return false; } - switch (s->type) { + switch (cl->type) { case AKI_SOCKET_TCP: { s32 port = 0; - char ipstr[INET6_ADDRSTRLEN] = {0}; + char ipstr[INET6_ADDRSTRLEN] = { 0 }; if (addr.ss_family == AF_INET) { - struct sockaddr_in *s = (struct sockaddr_in *)&addr; - port = ntohs(s->sin_port); - inet_ntop(AF_INET, &s->sin_addr, ipstr, sizeof(ipstr)); + struct sockaddr_in *saddr = (struct sockaddr_in *)&addr; + port = ntohs(saddr->sin_port); + inet_ntop(AF_INET, &saddr->sin_addr, ipstr, sizeof(ipstr)); } else if (addr.ss_family == AF_INET6) { - struct sockaddr_in6 *s = (struct sockaddr_in6 *)&addr; - port = ntohs(s->sin6_port); - inet_ntop(AF_INET6, &s->sin6_addr, ipstr, sizeof(ipstr)); + struct sockaddr_in6 *saddr = (struct sockaddr_in6 *)&addr; + port = ntohs(saddr->sin6_port); + inet_ntop(AF_INET6, &saddr->sin6_addr, ipstr, sizeof(ipstr)); } al_log_info("socket", "Connection from %s:%d.", ipstr, port); break; } case AKI_SOCKET_UNIX: - al_log_info("socket", "Connection from a UNIX socket."); + al_log_info("socket", "Connection on UNIX socket."); break; case AKI_SOCKET_UDP: al_assert(false); @@ -179,9 +174,9 @@ bool aki_socket_accept(struct aki_socket *s, struct aki_socket *c) return true; } -bool aki_socket_connect(struct aki_socket *s, str *addr, s32 port) +bool aki_socket_connect(struct aki_socket *sock, str *addr, s32 port) { - switch (s->type) { + switch (sock->type) { case AKI_SOCKET_TCP: case AKI_SOCKET_UDP: { char *c_str = al_str_to_c_str(addr); @@ -198,9 +193,9 @@ bool aki_socket_connect(struct aki_socket *s, str *addr, s32 port) return false; } - s->addr_in.sin_family = AF_INET; - s->addr_in.sin_addr.s_addr = ((struct in_addr *)hptr->h_addr_list[0])->s_addr; - s->addr_in.sin_port = htons(port); + sock->addr_in.sin_family = AF_INET; + sock->addr_in.sin_addr.s_addr = ((struct in_addr *)hptr->h_addr_list[0])->s_addr; + sock->addr_in.sin_port = htons(port); break; } @@ -208,21 +203,21 @@ bool aki_socket_connect(struct aki_socket *s, str *addr, s32 port) break; } - switch (s->type) { + switch (sock->type) { case AKI_SOCKET_TCP: case AKI_SOCKET_UDP: { - s32 ret = connect(s->sock, (struct sockaddr *)&s->addr_in, sizeof(s->addr_in)); - if (ret != 0 && errno != 115) { + s32 ret = connect(sock->fd, (struct sockaddr *)&sock->addr_in, sizeof(sock->addr_in)); + if (ret != 0 && errno != EINPROGRESS) { al_log_socket_err(connect) return false; } break; } case AKI_SOCKET_UNIX: { - s->addr_un.sun_family = AF_UNIX; - al_memcpy(s->addr_un.sun_path, addr->data, addr->len); - s->addr_un.sun_path[addr->len] = '\0'; - if (connect(s->sock, (struct sockaddr *)&s->addr_un, sizeof(s->addr_un)) < 0) { + sock->addr_un.sun_family = AF_UNIX; + al_memcpy(sock->addr_un.sun_path, addr->data, addr->len); + sock->addr_un.sun_path[addr->len] = '\0'; + if (connect(sock->fd, (struct sockaddr *)&sock->addr_un, sizeof(sock->addr_un)) < 0) { al_log_socket_err(connect) return false; } @@ -233,43 +228,41 @@ bool aki_socket_connect(struct aki_socket *s, str *addr, s32 port) break; } - signal(SIGPIPE, SIG_IGN); - return true; } -ssize_t aki_socket_read(struct aki_socket *s, void *buf, size_t size) +ssize_t aki_socket_read(struct aki_socket *sock, void *buf, size_t size) { - return recv(s->sock, buf, size, 0); + return recv(sock->fd, buf, size, 0); } -ssize_t aki_socket_write(struct aki_socket *s, void *buf, size_t size) +ssize_t aki_socket_write(struct aki_socket *sock, void *buf, size_t size) { - return send(s->sock, buf, size, MSG_NOSIGNAL); + return send(sock->fd, buf, size, MSG_NOSIGNAL); } -ssize_t aki_socket_sendto(struct aki_socket *s, void *buf, size_t size) +ssize_t aki_socket_sendto(struct aki_socket *sock, void *buf, size_t size) { - return sendto(s->sock, buf, size, 0, (const struct sockaddr *)&s->addr_in, sizeof(s->addr_in)); + return sendto(sock->fd, buf, size, 0, (const struct sockaddr *)&sock->addr_in, sizeof(sock->addr_in)); } -ssize_t aki_socket_recvfrom(struct aki_socket *s, void *buf, size_t size) +ssize_t aki_socket_recvfrom(struct aki_socket *sock, void *buf, size_t size) { - socklen_t addrlen = sizeof(s->addr_in); - return recvfrom(s->sock, buf, size, 0, (struct sockaddr *)&s->addr_in, &addrlen); + socklen_t addrlen = sizeof(sock->addr_in); + return recvfrom(sock->fd, buf, size, 0, (struct sockaddr *)&sock->addr_in, &addrlen); } -s32 aki_socket_get_fd(struct aki_socket *s) +s32 aki_socket_get_fd(struct aki_socket *sock) { - return s->sock; + return sock->fd; } -void aki_socket_shutdown(struct aki_socket *s) +void aki_socket_shutdown(struct aki_socket *sock) { - shutdown(s->sock, SHUT_RDWR); + shutdown(sock->fd, SHUT_RDWR); } -void aki_socket_close(struct aki_socket *s) +void aki_socket_close(struct aki_socket *sock) { - close(s->sock); + close(sock->fd); } diff --git a/src/socket/socket_windows.c b/src/socket/socket_windows.c index dcb7e08..a7192d8 100644 --- a/src/socket/socket_windows.c +++ b/src/socket/socket_windows.c @@ -30,10 +30,10 @@ void aki_socket_set_blocking(struct aki_socket *s, bool blocking) ioctlsocket(s->sock, FIONBIO, &mode); } -void aki_socket_set_no_delay(struct aki_socket *s, s32 no_delay) +void aki_socket_set_nodelay(struct aki_socket *s, s32 nodelay) { (void)s; - (void)no_delay; + (void)nodelay; } #define al_log_socket_err(func) \ diff --git a/src/util/file/file_linux.c b/src/util/file/file_linux.c index 403d685..ac78bd8 100644 --- a/src/util/file/file_linux.c +++ b/src/util/file/file_linux.c @@ -5,14 +5,9 @@ #include "file.h" -static bool lock_file_internal(s32 fd) -{ - struct flock l = { - .l_type = F_WRLCK, - .l_whence = SEEK_SET, - .l_start = 0, - .l_len = 0 - }; +static bool lock_file_internal(s32 fd, size_t size) +{ + struct flock l = { .l_type = F_WRLCK, .l_whence = SEEK_SET, .l_start = 0, .l_len = size }; if (fcntl(fd, F_SETLKW, &l) == -1) { al_log_error("file", "fcntl(F_SETLKW) failed (%s).", aki_strerror(errno)); close(fd); @@ -27,7 +22,7 @@ bool aki_file_open(struct aki_file *file, str *path, s32 flags) char *c_str = al_str_to_c_str(path); s32 oflags = (flags & AKI_FILE_READONLY) ? O_RDONLY : O_RDWR; if (flags & AKI_FILE_CREATE) oflags |= O_CREAT; - file->fd = open(c_str, oflags, 0666); + file->fd = open(c_str, oflags, 0644); al_free(c_str); if (file->fd == -1) { al_log_error("file", "open(%.*s) failed (%s).", AL_STR_PRINTF(path), aki_strerror(errno)); @@ -41,7 +36,7 @@ bool aki_file_open(struct aki_file *file, str *path, s32 flags) return false; } file->filesize = sb.st_size; - if (flags & AKI_FILE_LOCK) lock_file_internal(file->fd); + if (flags & AKI_FILE_LOCK) lock_file_internal(file->fd, file->filesize); return true; } diff --git a/src/util/packet.c b/src/util/packet.c index 2d8b440..373a3a3 100644 --- a/src/util/packet.c +++ b/src/util/packet.c @@ -58,12 +58,6 @@ void aki_packet_write_str(struct aki_packet *packet, str *s) AKI_PACKET_WRITE_DATA(packet, s->data, s->len); } -void aki_packet_write_wstr(struct aki_packet *packet, wstr *w) -{ - AKI_PACKET_WRITE_TYPE(packet, u32, w->len); - AKI_PACKET_WRITE_DATA(packet, w->data, w->len * sizeof(wchar_t)); -} - void aki_packet_write_buffer(struct aki_packet *packet, struct aki_buffer *buf) { AKI_PACKET_WRITE_TYPE(packet, size_t, buf->size); @@ -96,13 +90,6 @@ void aki_packet_read_str(struct aki_packet *packet, str *s) s->alloc = 0; } -void aki_packet_read_wstr(struct aki_packet *packet, wstr *w) -{ - AKI_PACKET_READ_TYPE(packet, u32, w->len); - AKI_PACKET_READ_DATA(packet, w->len * sizeof(wchar_t), w->data); - w->alloc = 0; -} - void aki_packet_read_buffer(struct aki_packet *packet, struct aki_buffer *buf) { AKI_PACKET_READ_TYPE(packet, size_t, buf->size); @@ -110,6 +97,21 @@ void aki_packet_read_buffer(struct aki_packet *packet, struct aki_buffer *buf) buf->alloc = 0; } +#ifdef AL_HAVE_WIDE_STRING +void aki_packet_write_wstr(struct aki_packet *packet, wstr *w) +{ + AKI_PACKET_WRITE_TYPE(packet, u32, w->len); + AKI_PACKET_WRITE_DATA(packet, w->data, w->len * sizeof(wchar_t)); +} + +void aki_packet_read_wstr(struct aki_packet *packet, wstr *w) +{ + AKI_PACKET_READ_TYPE(packet, u32, w->len); + AKI_PACKET_READ_DATA(packet, w->len * sizeof(wchar_t), w->data); + w->alloc = 0; +} +#endif + void aki_packet_free(struct aki_packet *packet) { aki_buffer_free(&packet->buffer); diff --git a/src/util/packet.h b/src/util/packet.h index f877c72..cacd82c 100644 --- a/src/util/packet.h +++ b/src/util/packet.h @@ -53,7 +53,6 @@ void aki_packet_write_u64(struct aki_packet *packet, u64 v); void aki_packet_write_f32(struct aki_packet *packet, f32 v); void aki_packet_write_f64(struct aki_packet *packet, f64 v); void aki_packet_write_str(struct aki_packet *packet, str *s); -void aki_packet_write_wstr(struct aki_packet *packet, wstr *w); void aki_packet_write_buffer(struct aki_packet *packet, struct aki_buffer *buf); s8 aki_packet_read_s8(struct aki_packet *packet); @@ -67,7 +66,11 @@ u64 aki_packet_read_u64(struct aki_packet *packet); f32 aki_packet_read_f32(struct aki_packet *packet); f64 aki_packet_read_f64(struct aki_packet *packet); void aki_packet_read_str(struct aki_packet *packet, str *s); -void aki_packet_read_wstr(struct aki_packet *packet, wstr *w); void aki_packet_read_buffer(struct aki_packet *packet, struct aki_buffer *buf); +#ifdef AL_HAVE_WIDE_STRING +void aki_packet_write_wstr(struct aki_packet *packet, wstr *w); +void aki_packet_read_wstr(struct aki_packet *packet, wstr *w); +#endif + void aki_packet_free(struct aki_packet *packet); -- cgit v1.2.3-101-g0448