summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-06-24 14:34:21 -0400
committerAndrew Opalach <andrew@akon.city> 2024-06-24 14:34:21 -0400
commita320ee672ccf52297dc509b25ed263f165709b31 (patch)
treeef78116173aa19c28ee4a46f861861f700528c50 /src
parent3d1ab87859a291cf8965b56703931193d5e6ae25 (diff)
downloadlibnaunet-a320ee672ccf52297dc509b25ed263f165709b31.tar.gz
libnaunet-a320ee672ccf52297dc509b25ed263f165709b31.tar.bz2
libnaunet-a320ee672ccf52297dc509b25ed263f165709b31.zip
Add multiplex, a bunch of cleanup
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src')
-rw-r--r--src/common.c3
-rw-r--r--src/curl/http.c1
-rw-r--r--src/curl/websocket.c2
-rw-r--r--src/ev_embed_compat.c3
-rw-r--r--src/line_processor.c17
-rw-r--r--src/line_processor.h6
-rw-r--r--src/multiplex.c47
-rw-r--r--src/multiplex.h20
-rw-r--r--src/packet_stream.c184
-rw-r--r--src/packet_stream.h24
-rw-r--r--src/rpc.c118
-rw-r--r--src/rpc2.h19
-rw-r--r--src/socket/socket.h35
-rw-r--r--src/socket/socket_linux.c169
-rw-r--r--src/socket/socket_windows.c4
-rw-r--r--src/util/file/file_linux.c13
-rw-r--r--src/util/packet.c28
-rw-r--r--src/util/packet.h7
18 files changed, 385 insertions, 315 deletions
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 <al/lib.h>
#include <al/random.h>
#include <locale.h>
+#include <signal.h>
#ifndef _WIN32
#include <unistd.h>
@@ -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 <al/array.h>
+
+#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 <al/log.h>
#include <unistd.h>
#include <sys/fcntl.h>
#include <netinet/tcp.h>
#include <netdb.h>
-#include <signal.h>
-#include <al/log.h>
#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)
+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 = 0
- };
+ 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);