summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/curl/http.c10
-rw-r--r--src/curl/http.h2
-rw-r--r--src/curl/websocket.c3
-rw-r--r--src/ev_embed_compat.c1
-rw-r--r--src/packet_cache.c2
-rw-r--r--src/packet_pool.c71
-rw-r--r--src/packet_pool.h27
-rw-r--r--src/packet_stream.c18
-rw-r--r--src/packet_stream.h2
-rw-r--r--src/socket/socket.h4
-rw-r--r--src/socket/socket_linux.c2
-rw-r--r--src/util/thread/thread_windows.c4
-rw-r--r--src/util/timer/timer_windows.c29
13 files changed, 56 insertions, 119 deletions
diff --git a/src/curl/http.c b/src/curl/http.c
index cdd740b..b3c13a2 100644
--- a/src/curl/http.c
+++ b/src/curl/http.c
@@ -92,12 +92,12 @@ void aki_http_add_header(struct aki_http *http, str *header)
al_free(c_str);
}
-void aki_http_set_range(struct aki_http *http, ssize_t start, ssize_t end)
+void aki_http_set_range(struct aki_http *http, size_t start, size_t end)
{
- char range[64];
- if (start < 0) al_sprintf(range, "-%zd", end);
- else if (end < 0) al_sprintf(range, "%zd-", start);
- else al_sprintf(range, "%zd-%zd", start, end);
+ char range[44];
+ if (start == 0u) al_snprintf(range, sizeof(range), "-%zu", end);
+ else if (end == 0u) al_snprintf(range, sizeof(range), "%zu-", start);
+ else al_snprintf(range, sizeof(range), "%zu-%zu", start, end);
curl_easy_setopt(http->curl.handle, CURLOPT_RANGE, range);
}
diff --git a/src/curl/http.h b/src/curl/http.h
index 28942b9..3d0b796 100644
--- a/src/curl/http.h
+++ b/src/curl/http.h
@@ -46,7 +46,7 @@ bool aki_http_init(struct aki_http *http);
void aki_http_set_url(struct aki_http *http, str *url);
void aki_http_set_user_agent(struct aki_http *http, str *user_agent);
void aki_http_add_header(struct aki_http *http, str *header);
-void aki_http_set_range(struct aki_http *http, ssize_t start, ssize_t end);
+void aki_http_set_range(struct aki_http *http, size_t start, size_t end);
void aki_http_close(struct aki_http *http);
bool aki_http_request_stream(struct aki_http *http, u8 method, struct aki_event_loop *loop,
size_t (*callback)(void *, u8, u8 *, s64), void *userdata);
diff --git a/src/curl/websocket.c b/src/curl/websocket.c
index aede3de..8016981 100644
--- a/src/curl/websocket.c
+++ b/src/curl/websocket.c
@@ -54,9 +54,6 @@ void aki_websocket_set_url(struct aki_websocket *ws, str *url)
bool aki_websocket_connect(struct aki_websocket *ws, struct aki_event_loop *loop,
void (*callback)(void *, const struct curl_ws_frame *, u8 *, size_t), void *userdata)
{
- // Connect blocking.
- //curl_easy_setopt(ws->curl.handle, CURLOPT_CONNECT_ONLY, 2L);
- //curl_easy_perform(ws->curl.handle);
curl_easy_setopt(ws->curl.handle, CURLOPT_CONNECT_ONLY, 0L);
ws->curl.loop = loop;
ws->callback = callback;
diff --git a/src/ev_embed_compat.c b/src/ev_embed_compat.c
index 4c7a1f2..552e25b 100644
--- a/src/ev_embed_compat.c
+++ b/src/ev_embed_compat.c
@@ -18,6 +18,7 @@ _Pragma("GCC diagnostic ignored \"-Wparentheses\"")
_Pragma("GCC diagnostic ignored \"-Wunused-value\"")
_Pragma("GCC diagnostic ignored \"-Wunused-function\"")
_Pragma("GCC diagnostic ignored \"-Wunused-parameter\"")
+_Pragma("GCC diagnostic ignored \"-Wunused-result\"")
_Pragma("GCC diagnostic ignored \"-Wsign-compare\"")
_Pragma("GCC diagnostic ignored \"-Wreturn-type\"")
_Pragma("GCC diagnostic ignored \"-Wdeprecated-declarations\"")
diff --git a/src/packet_cache.c b/src/packet_cache.c
index 209b3a4..c279982 100644
--- a/src/packet_cache.c
+++ b/src/packet_cache.c
@@ -4,7 +4,7 @@ void aki_packet_cache_init(struct aki_packet_cache *cache, u32 size)
{
al_array_init(cache->cache);
al_array_reserve(cache->cache, size);
- cache->flush = size / 2;
+ cache->flush = size > 0 ? size / 2 : 0;
aki_cond_init(&cache->cond);
aki_mutex_init(&cache->mutex);
cache->disabled = false;
diff --git a/src/packet_pool.c b/src/packet_pool.c
index a3ca57b..1e3117f 100644
--- a/src/packet_pool.c
+++ b/src/packet_pool.c
@@ -22,44 +22,30 @@ static void signal_callback(struct ev_loop *loop, ev_async *w, s32 revents)
while (pool->ready.size > 0) {
struct aki_packet *packet;
al_array_pop_at(pool->ready, 0, packet);
- aki_mutex_unlock(&pool->mutex);
u8 ret = pool->callback(pool->userdata, packet);
- al_assert(ret != AKI_PACKET_POOL_CLOSED);
- aki_mutex_lock(&pool->mutex);
switch (ret) {
- case AKI_PACKET_POOL_KEEP: // AKI_PACKET_POOL_NOP
+ case AKI_PACKET_POOL_KEEP:
break;
case AKI_PACKET_POOL_RETURN:
- al_assert(pool->mode == AKI_PACKET_POOL_MODE_POOL);
return_internal(pool, packet);
break;
- case AKI_PACKET_POOL_DISABLE:
- al_assert(pool->mode == AKI_PACKET_POOL_MODE_POOL);
- pool->disabled = true;
- return_internal(pool, packet);
- aki_mutex_unlock(&pool->mutex);
- ret = pool->callback(pool->userdata, NULL);
- al_assert(ret == AKI_PACKET_POOL_CLOSED);
- // signal_callback must never be called again.
- return;
}
}
aki_mutex_unlock(&pool->mutex);
}
-static void init_internal(struct aki_packet_pool *pool, u32 size, struct aki_event_loop *loop,
- u8 mode, u8 (*callback)(void *, struct aki_packet *), void *userdata)
+void aki_packet_pool_init(struct aki_packet_pool *pool, u32 size, struct aki_event_loop *loop,
+ u8 (*callback)(void *, struct aki_packet *), void *userdata)
{
- pool->mode = mode;
pool->loop = loop;
+ al_assert(size > 0);
+ pool->size = size;
al_array_init(pool->ready);
al_array_reserve(pool->ready, size);
- if (pool->mode == AKI_PACKET_POOL_MODE_POOL) {
- al_array_init(pool->empty);
- al_array_reserve(pool->empty, size);
- for (u32 i = 0; i < size; i++) {
- al_array_push(pool->empty, aki_packet_create());
- }
+ al_array_init(pool->empty);
+ al_array_reserve(pool->empty, size);
+ for (u32 i = 0; i < size; i++) {
+ al_array_push(pool->empty, aki_packet_create());
}
pool->disabled = false;
aki_cond_init(&pool->cond);
@@ -71,23 +57,10 @@ static void init_internal(struct aki_packet_pool *pool, u32 size, struct aki_eve
pool->userdata = userdata;
}
-void aki_packet_pool_init(struct aki_packet_pool *pool, u32 size, struct aki_event_loop *loop,
- u8 (*callback)(void *, struct aki_packet *), void *userdata)
-{
- init_internal(pool, size, loop, AKI_PACKET_POOL_MODE_POOL, callback, userdata);
-}
-
-void aki_packet_pool_init_ex(struct aki_packet_pool *pool, u32 size, struct aki_event_loop *loop,
- u8 mode, u8 (*callback)(void *, struct aki_packet *), void *userdata)
-{
- init_internal(pool, size, loop, mode, callback, userdata);
-}
-
static const bool grow = false;
struct aki_packet *aki_packet_pool_get(struct aki_packet_pool *pool)
{
- al_assert(pool->mode == AKI_PACKET_POOL_MODE_POOL);
aki_mutex_lock(&pool->mutex);
if (!pool->disabled && !pool->empty.size) {
if (grow) {
@@ -113,13 +86,20 @@ void aki_packet_pool_submit(struct aki_packet_pool *pool, struct aki_packet *pac
return;
}
al_array_push(pool->ready, packet);
+ bool flush = pool->ready.size >= pool->size / 2;
aki_mutex_unlock(&pool->mutex);
+ if (flush) {
+ ev_async_send(pool->loop->ev, &pool->signal);
+ }
+}
+
+void aki_packet_pool_flush(struct aki_packet_pool *pool)
+{
ev_async_send(pool->loop->ev, &pool->signal);
}
void aki_packet_pool_return(struct aki_packet_pool *pool, struct aki_packet *packet)
{
- al_assert(pool->mode == AKI_PACKET_POOL_MODE_POOL);
aki_mutex_lock(&pool->mutex);
return_internal(pool, packet);
aki_mutex_unlock(&pool->mutex);
@@ -127,7 +107,6 @@ void aki_packet_pool_return(struct aki_packet_pool *pool, struct aki_packet *pac
void aki_packet_pool_disable(struct aki_packet_pool *pool)
{
- al_assert(pool->mode == AKI_PACKET_POOL_MODE_POOL);
aki_mutex_lock(&pool->mutex);
pool->disabled = true;
struct aki_packet *packet;
@@ -144,7 +123,6 @@ void aki_packet_pool_disable(struct aki_packet_pool *pool)
void aki_packet_pool_enable(struct aki_packet_pool *pool)
{
- al_assert(pool->mode == AKI_PACKET_POOL_MODE_POOL);
aki_mutex_lock(&pool->mutex);
pool->disabled = false;
aki_mutex_unlock(&pool->mutex);
@@ -152,7 +130,6 @@ void aki_packet_pool_enable(struct aki_packet_pool *pool)
struct aki_packet *aki_packet_pool_pop(struct aki_packet_pool *pool)
{
- al_assert(pool->mode == AKI_PACKET_POOL_MODE_PASSTHROUGH);
struct aki_packet *packet = NULL;
if (pool->ready.size > 0) {
al_array_pop_at(pool->ready, 0, packet);
@@ -166,14 +143,12 @@ void aki_packet_pool_free(struct aki_packet_pool *pool)
aki_mutex_destroy(&pool->mutex);
aki_cond_destroy(&pool->cond);
struct aki_packet *packet;
- if (pool->mode == AKI_PACKET_POOL_MODE_POOL) {
- al_array_foreach(pool->empty, i, packet) {
- aki_packet_free(packet);
- }
- al_array_free(pool->empty);
- al_array_foreach(pool->ready, i, packet) {
- aki_packet_free(packet);
- }
+ al_array_foreach(pool->empty, i, packet) {
+ aki_packet_free(packet);
+ }
+ al_array_free(pool->empty);
+ al_array_foreach(pool->ready, i, packet) {
+ aki_packet_free(packet);
}
al_array_free(pool->ready);
}
diff --git a/src/packet_pool.h b/src/packet_pool.h
index 2a2e608..c31ac9e 100644
--- a/src/packet_pool.h
+++ b/src/packet_pool.h
@@ -7,7 +7,7 @@
#include "loop.h"
-// Steps (MODE_POOL):
+// Steps:
// packet_get (block) \ Thread 2
// packet_submit (signal async) /
// callback called \ Thread 1
@@ -16,32 +16,16 @@
// join(thread1) && join(thread2)
// pool_free()
-// Steps (MODE_PASSTHROUGH):
-// packet_submit (signal async) -> Thread 2
-// callback called -> Thread 1
-// join(thread1) && join(thread2)
-// while (pop()) {} (empty pool)
-// pool_free()
-
-enum {
- AKI_PACKET_POOL_MODE_POOL = 0,
- AKI_PACKET_POOL_MODE_PASSTHROUGH
-};
-
enum {
AKI_PACKET_POOL_KEEP = 0,
- AKI_PACKET_POOL_RETURN,
- AKI_PACKET_POOL_DISABLE,
- AKI_PACKET_POOL_CLOSED
+ AKI_PACKET_POOL_RETURN
};
-#define AKI_PACKET_POOL_NOP AKI_PACKET_POOL_KEEP
-
struct aki_packet_pool {
- u8 mode;
+ struct aki_event_loop *loop;
+ u32 size;
array(struct aki_packet *) empty;
array(struct aki_packet *) ready;
- struct aki_event_loop *loop;
struct aki_cond cond;
struct aki_mutex mutex;
ev_async signal;
@@ -52,10 +36,9 @@ struct aki_packet_pool {
void aki_packet_pool_init(struct aki_packet_pool *pool, u32 size, struct aki_event_loop *loop,
u8 (*callback)(void *, struct aki_packet *), void *userdata);
-void aki_packet_pool_init_ex(struct aki_packet_pool *pool, u32 size, struct aki_event_loop *loop,
- u8 mode, u8 (*callback)(void *, struct aki_packet *), void *userdata);
struct aki_packet *aki_packet_pool_get(struct aki_packet_pool *pool);
void aki_packet_pool_submit(struct aki_packet_pool *pool, struct aki_packet *packet);
+void aki_packet_pool_flush(struct aki_packet_pool *pool);
void aki_packet_pool_return(struct aki_packet_pool *pool, struct aki_packet *packet);
void aki_packet_pool_disable(struct aki_packet_pool *pool);
void aki_packet_pool_enable(struct aki_packet_pool *pool);
diff --git a/src/packet_stream.c b/src/packet_stream.c
index 5ae91f1..8b5c611 100644
--- a/src/packet_stream.c
+++ b/src/packet_stream.c
@@ -112,17 +112,16 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents)
al_assert(revents & EV_WRITE); // Assert on EV_ERROR.
// EV_WRITE can mean any of POLLOUT, POLLERR, or POLLHUP.
if (UNLIKELY(stream->connect == PACKET_STREAM_CONNECTING)) {
- // Shove in the multiplex ID byte here assuming POLLOUT means
- // we can write a minimum of 1 byte.
ssize_t ret = aki_socket_write(&stream->sock, &stream->id, sizeof(u8));
if (ret < 0) {
// It seems POLLOUT can be signaled even if any write() will error.
stop_internal(stream);
return;
} else if (ret == 0) {
- // Try again on next POLLOUT. Not sure if this can happen.
+ // Try again on next POLLOUT, not sure if this can actually happen.
return;
}
+ al_assert(ret == sizeof(u8));
stream->connect = PACKET_STREAM_CONNECTED;
// Callbacks must be set in connection_callback.
stream->connection_callback(stream->userdata, stream);
@@ -215,9 +214,9 @@ void aki_packet_stream_reconnect(struct aki_packet_stream *stream, str *addr, s3
void aki_packet_stream_cork(struct aki_packet_stream *stream, bool cork)
{
if (stream->connect == PACKET_STREAM_CONNECTED) {
- if (cork && !stream->corked) {
+ if (cork && stream->revent.active) {
ev_io_stop(stream->loop->ev, &stream->revent);
- } else if (!cork && stream->corked) {
+ } else if (!cork && !stream->revent.active) {
ev_io_start(stream->loop->ev, &stream->revent);
}
}
@@ -225,11 +224,10 @@ void aki_packet_stream_cork(struct aki_packet_stream *stream, bool cork)
stream->corked = cork;
}
-void aki_packet_stream_send_packet(struct aki_packet_stream *stream, struct aki_packet *packet)
+bool aki_packet_stream_send_packet(struct aki_packet_stream *stream, struct aki_packet *packet)
{
if (stream->connect == PACKET_STREAM_DISCONNECTING) {
- stream->packet_sent_callback(stream->userdata, packet);
- return;
+ return false;
}
al_assert(stream->connect == PACKET_STREAM_CONNECTED);
aki_packet_write_size(packet);
@@ -237,6 +235,7 @@ void aki_packet_stream_send_packet(struct aki_packet_stream *stream, struct aki_
if (!stream->out.have_data) {
start_write_internal(stream);
}
+ return true;
}
void aki_packet_stream_disconnect(struct aki_packet_stream *stream)
@@ -249,7 +248,8 @@ void aki_packet_stream_disconnect(struct aki_packet_stream *stream)
}
stream->connect = PACKET_STREAM_DISCONNECTING;
// If we are not corked, this relies on stream_read_callback() to resolve and call stop_internal().
- // NOTE: read() can still return data after shutdown().
+ // NOTE: read() can still return data after shutdown(). This is why we have to check if
+ // connect = DISCONNECTING in send_packet().
aki_socket_shutdown(&stream->sock);
if (stream->corked || stream->connect == PACKET_STREAM_CONNECTING) {
stop_internal(stream);
diff --git a/src/packet_stream.h b/src/packet_stream.h
index f0f9cc2..80d8a2f 100644
--- a/src/packet_stream.h
+++ b/src/packet_stream.h
@@ -42,6 +42,6 @@ void aki_packet_stream_from_socket(struct aki_packet_stream *stream, struct aki_
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);
+bool aki_packet_stream_send_packet(struct aki_packet_stream *stream, struct aki_packet *packet);
void aki_packet_stream_disconnect(struct aki_packet_stream *stream);
void aki_packet_stream_free(struct aki_packet_stream *stream);
diff --git a/src/socket/socket.h b/src/socket/socket.h
index 18e3040..36b101a 100644
--- a/src/socket/socket.h
+++ b/src/socket/socket.h
@@ -59,8 +59,10 @@ bool aki_socket_init(struct aki_socket *sock);
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 *sock, u32 sndbuf);
u32 aki_socket_get_send_buf(struct aki_socket *sock);
+void aki_socket_set_send_buf(struct aki_socket *sock, u32 sndbuf);
+u32 aki_socket_get_recv_buf(struct aki_socket *sock);
+void aki_socket_set_recv_buf(struct aki_socket *sock, u32 rcvbuf);
bool aki_socket_set(struct aki_socket *sock, str *addr, u16 port);
bool aki_socket_bind(struct aki_socket *sock, str *addr, u16 port);
diff --git a/src/socket/socket_linux.c b/src/socket/socket_linux.c
index 753f6b9..b89aac6 100644
--- a/src/socket/socket_linux.c
+++ b/src/socket/socket_linux.c
@@ -112,7 +112,7 @@ static bool parse_address(struct aki_socket *sock, str *addr, u16 port)
char *c_str = addr ? al_str_to_c_str(addr) : NULL;
if (!c_str) hints.ai_flags = AI_PASSIVE;
char port_str[6];
- al_sprintf(port_str, "%hu", port);
+ al_snprintf(port_str, sizeof(port_str), "%hu", port);
s32 status = getaddrinfo(c_str, port_str, &hints, &sock->addrinfo);
al_free(c_str);
if (status != 0) {
diff --git a/src/util/thread/thread_windows.c b/src/util/thread/thread_windows.c
index a58e1b1..b948d3e 100644
--- a/src/util/thread/thread_windows.c
+++ b/src/util/thread/thread_windows.c
@@ -4,9 +4,7 @@ NTSTATUS(__stdcall *NtDelayExecution)(BOOL Alertable, PLARGE_INTEGER DelayInterv
void aki_thread_sleep(aki_os_tstamp delay)
{
- LARGE_INTEGER li;
- li.QuadPart = -(s64)(delay * 1e7L);
- NtDelayExecution(false, &li);
+ NtDelayExecution(false, &((LARGE_INTEGER){ .QuadPart = -(s64)(delay * 1e7L) }));
}
void aki_thread_create(struct aki_thread *thread, aki_thread_func func, void *userdata)
diff --git a/src/util/timer/timer_windows.c b/src/util/timer/timer_windows.c
index 007efed..edd1297 100644
--- a/src/util/timer/timer_windows.c
+++ b/src/util/timer/timer_windows.c
@@ -1,5 +1,3 @@
-#include <al/lib.h>
-
#include "../../winwrap.h"
#include "timer.h"
@@ -11,29 +9,12 @@ f64 aki_get_tick(void)
return result.QuadPart / (f64)performance_frequency.QuadPart;
}
-static SYSTEMTIME unix_epoch = {
- .wYear = 1970,
- .wMonth = 1,
- .wDayOfWeek = 4,
- .wDay = 1,
- .wHour = 0,
- .wMinute = 0,
- .wSecond = 0,
- .wMilliseconds = 0
-};
+// https://stackoverflow.com/a/46024468
+static const s64 UNIX_TIME_START = 0x019DB1DED53E8000; // January 1st, 1970 in "ticks"
u64 aki_get_timestamp(void)
{
- FILETIME result0, result1;
- SystemTimeToFileTime(&unix_epoch, &result0);
- ULARGE_INTEGER li0 = {
- .LowPart = result0.dwLowDateTime,
- .HighPart = result0.dwHighDateTime
- };
- GetSystemTimePreciseAsFileTime(&result1);
- ULARGE_INTEGER li1 = {
- .LowPart = result1.dwLowDateTime,
- .HighPart = result1.dwHighDateTime
- };
- return (li1.QuadPart - li0.QuadPart) / 10L;
+ FILETIME ft;
+ GetSystemTimePreciseAsFileTime(&ft);
+ return (((LARGE_INTEGER){ .LowPart = ft.dwLowDateTime, .HighPart = ft.dwHighDateTime }).QuadPart - UNIX_TIME_START) / 10Lu;
}