diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/curl/http.c | 10 | ||||
| -rw-r--r-- | src/curl/http.h | 2 | ||||
| -rw-r--r-- | src/curl/websocket.c | 3 | ||||
| -rw-r--r-- | src/ev_embed_compat.c | 1 | ||||
| -rw-r--r-- | src/packet_cache.c | 2 | ||||
| -rw-r--r-- | src/packet_pool.c | 71 | ||||
| -rw-r--r-- | src/packet_pool.h | 27 | ||||
| -rw-r--r-- | src/packet_stream.c | 18 | ||||
| -rw-r--r-- | src/packet_stream.h | 2 | ||||
| -rw-r--r-- | src/socket/socket.h | 4 | ||||
| -rw-r--r-- | src/socket/socket_linux.c | 2 | ||||
| -rw-r--r-- | src/util/thread/thread_windows.c | 4 | ||||
| -rw-r--r-- | src/util/timer/timer_windows.c | 29 |
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; } |