diff options
Diffstat (limited to 'src/cache/handlers')
| -rw-r--r-- | src/cache/handlers/file.c | 49 | ||||
| -rw-r--r-- | src/cache/handlers/file.h | 12 | ||||
| -rw-r--r-- | src/cache/handlers/http.c | 207 | ||||
| -rw-r--r-- | src/cache/handlers/http.h | 24 |
4 files changed, 292 insertions, 0 deletions
diff --git a/src/cache/handlers/file.c b/src/cache/handlers/file.c new file mode 100644 index 0000000..ec6a146 --- /dev/null +++ b/src/cache/handlers/file.c @@ -0,0 +1,49 @@ +#include "../backings/file.h" + +#include "file.h" + +static bool handler_file_can_seek(struct cch_handler *handler) +{ + (void)handler; + return true; +} + +static void handler_file_maybe_spawn_worker(struct cch_handler *handler, size_t index) +{ + struct cch_handler_file *file = (struct cch_handler_file *)handler; + (void)file; + (void)index; +} + +static bool handler_file_wait_for_range(struct cch_handler *handler, struct cch_handler_wait *wait) +{ + (void)handler; + (void)wait; + return true; +} + +static void handler_file_free(struct cch_handler **handler) +{ + struct cch_handler_file *file = (struct cch_handler_file *)*handler; + al_free(file); + *handler = NULL; +} + +struct cch_entry *cch_handler_file_create(str *path) +{ + struct cch_backing *backing = cch_backing_file_create(path, 0); + if (!backing) return NULL; + struct cch_entry *entry = al_alloc_object(struct cch_entry); + entry->backing = backing; + entry->ref_count = 0; + entry->hash = 0; + entry->handler = (struct cch_handler *)al_alloc_object(struct cch_handler_file); + aki_mutex_init(&entry->mutex); + entry->handler->can_seek = handler_file_can_seek; + entry->handler->maybe_spawn_worker = handler_file_maybe_spawn_worker; + entry->handler->wait_for_range = handler_file_wait_for_range; + entry->handler->free = handler_file_free; + entry->handler->entry = entry; + entry->handler->backing = entry->backing; + return entry; +} diff --git a/src/cache/handlers/file.h b/src/cache/handlers/file.h new file mode 100644 index 0000000..52e5c94 --- /dev/null +++ b/src/cache/handlers/file.h @@ -0,0 +1,12 @@ +#pragma once + +#include <al/str.h> + +#include "../handler.h" +#include "../entry.h" + +struct cch_handler_file { + struct cch_handler handler; +}; + +struct cch_entry *cch_handler_file_create(str *path); diff --git a/src/cache/handlers/http.c b/src/cache/handlers/http.c new file mode 100644 index 0000000..1daf231 --- /dev/null +++ b/src/cache/handlers/http.c @@ -0,0 +1,207 @@ +#include <al/random.h> +#include <al/log.h> + +#include "../backings/file.h" +#include "../backings/memory.h" + +#include "http.h" + +#define USER_AGENT al_str_c("Mozilla/5.0 (X11; Linux x86_64; rv:96.0) Gecko/20100101 Firefox/96.0") + +static bool handler_http_can_seek(struct cch_handler *handler) +{ + (void)handler; + return false; +} + +static size_t http_callback(void *userdata, u8 op, u8 *buf, s64 int0) +{ + struct cch_handler_http *http = (struct cch_handler_http *)userdata; + size_t ret = (size_t)int0; + switch (op) { + case AKI_HTTP_FINISHED: { + al_log_debug("cache_handler_http", "Transfer finished."); + break; + } + case AKI_HTTP_ERROR: { + al_log_debug("cache_handler_http", "Error."); + break; // Unhandled. + } + case AKI_HTTP_RESPONSE_CODE: { + al_log_debug("cache_handler_http", "HTTP %ld.", int0); + break; + } + case AKI_HTTP_REDIRECT: { + al_log_debug("cache_handler_http", "Redirect."); + break; + } + case AKI_HTTP_CONTENT_LENGTH: { + off_t content_length = (off_t)int0; + al_log_debug("cache_handler_http", "Content-Length: %lld.", content_length); + aki_mutex_lock(&http->mutex); + cch_entry_set_size(http->handler.entry, content_length); + struct cch_handler_wait *wait; + al_array_foreach(http->handler.waits, i, wait) { + if (wait->end < 0) { + aki_mutex_lock(&wait->mutex); + aki_cond_signal(&wait->cond); + al_array_remove_at_iter(http->handler.waits, i); + aki_mutex_unlock(&wait->mutex); + } + } + aki_mutex_unlock(&http->mutex); + break; + } + case AKI_HTTP_WRITE: { + http->handler.backing->write(http->handler.backing, buf, http->pointer, &ret); + aki_mutex_lock(&http->mutex); + http->pointer += ret; + struct cch_handler_wait *wait; + al_array_foreach(http->handler.waits, i, wait) { + aki_mutex_lock(&wait->mutex); + // Assume that if we get data before a length, we won't get a length. + // So, signal the waiter (wait_for_range returns false). + if (wait->end < 0) wait->disabled = true; + if (wait->disabled || http->pointer >= wait->end) { + aki_cond_signal(&wait->cond); + al_array_remove_at_iter(http->handler.waits, i); + } + aki_mutex_unlock(&wait->mutex); + } + aki_mutex_unlock(&http->mutex); + break; + } + default: + break; + } + return ret; +} + +static void signal_callback(void *userdata) +{ + struct cch_handler_http *http = (struct cch_handler_http *)userdata; + do { + u32 size; + camu_queue_size(http->queue, size); + if (size == 0) break; + s32 index; + camu_queue_pop(http->queue, index); + if (index == -1) { + aki_event_loop_break(&http->loop); + break; + } + al_array_push(http->requests, (struct aki_http){0}); + struct aki_http *request = &al_array_last(http->requests); + aki_http_init(request); + aki_http_set_url(request, &http->url); + aki_http_set_user_agent(request, USER_AGENT); + aki_http_request_stream(request, AKI_HTTP_GET, &http->loop, http_callback, http); + } while (1); +} + +static void handler_http_maybe_spawn_worker(struct cch_handler *handler, size_t index) +{ + struct cch_handler_http *http = (struct cch_handler_http *)handler; + camu_queue_push(http->queue, index); + aki_signal_send(&http->signal); +} + +static bool handler_http_wait_for_range(struct cch_handler *handler, struct cch_handler_wait *wait) +{ + struct cch_handler_http *http = (struct cch_handler_http *)handler; + aki_mutex_lock(&http->mutex); + bool canceled = false; + if (cch_entry_get_size(http->handler.entry) < 0 || http->pointer < wait->end) { + aki_mutex_lock(&wait->mutex); + if (!wait->disabled) { + al_array_push(http->handler.waits, wait); + aki_mutex_unlock(&http->mutex); + aki_cond_wait(&wait->cond, &wait->mutex); + } else { + aki_mutex_unlock(&wait->mutex); + aki_mutex_unlock(&http->mutex); + return false; + } + aki_mutex_unlock(&wait->mutex); + aki_mutex_lock(&http->mutex); + if (wait->disabled) { + u16 id = wait->id; + al_array_foreach(http->handler.waits, i, wait) { + if (wait->id == id) { + al_array_remove_at_iter(http->handler.waits, i); + break; + } + } + canceled = true; + } + } + aki_mutex_unlock(&http->mutex); + return !canceled; +} + +static void handler_http_free(struct cch_handler **handler) +{ + struct cch_handler_http *http = (struct cch_handler_http *)*handler; + aki_mutex_lock(&http->mutex); + struct cch_handler_wait *wait; + al_array_foreach(http->handler.waits, i, wait) { + aki_mutex_lock(&wait->mutex); + wait->disabled = true; + aki_cond_signal(&wait->cond); + aki_mutex_unlock(&wait->mutex); + } + aki_mutex_unlock(&http->mutex); + camu_queue_push(http->queue, -1); + aki_signal_send(&http->signal); + aki_thread_join(&http->thread); + aki_mutex_destroy(&http->mutex); + struct aki_http *request; + al_array_foreach_ptr(http->requests, i, request) { + aki_http_close(request); + } + al_array_free(http->requests); + camu_queue_free(http->queue); + al_str_free(&http->url); + al_array_free(http->handler.waits); + al_free(http); + *handler = NULL; +} + +static aki_thread_result AKI_THREADCALL event_loop_thread(void *userdata) +{ + struct cch_handler_http *http = (struct cch_handler_http *)userdata; + aki_event_loop_run(&http->loop); + aki_event_loop_destroy(&http->loop); + return 0; +} + +struct cch_entry *cch_handler_http_create(str *url) +{ + //struct cch_backing *backing = cch_backing_file_create(al_str_c("/tmp/camu_http_data"), 1024 * 128); + struct cch_backing *backing = cch_backing_memory_create(1024 * 128); + if (!backing) return NULL; + struct cch_entry *entry = al_alloc_object(struct cch_entry); + entry->backing = backing; + entry->ref_count = 0; + entry->hash = 0; + entry->handler = (struct cch_handler *)al_alloc_object(struct cch_handler_http); + aki_mutex_init(&entry->mutex); + entry->handler->can_seek = handler_http_can_seek; + entry->handler->maybe_spawn_worker = handler_http_maybe_spawn_worker; + entry->handler->wait_for_range = handler_http_wait_for_range; + entry->handler->free = handler_http_free; + entry->handler->entry = entry; + entry->handler->backing = entry->backing; + al_array_init(entry->handler->waits); + struct cch_handler_http *http = (struct cch_handler_http *)entry->handler; + al_str_clone(&http->url, url); + http->pointer = 0; + al_array_init(http->requests); + camu_queue_init(http->queue); + aki_event_loop_init(&http->loop); + aki_signal_init(&http->signal, signal_callback, http); + aki_signal_start(&http->signal, &http->loop); + aki_mutex_init(&http->mutex); + aki_thread_create(&http->thread, event_loop_thread, http); + return entry; +} diff --git a/src/cache/handlers/http.h b/src/cache/handlers/http.h new file mode 100644 index 0000000..f79a39a --- /dev/null +++ b/src/cache/handlers/http.h @@ -0,0 +1,24 @@ +#pragma once + +#include <aki/http.h> +#include <aki/thread.h> +#include <aki/signal.h> + +#include "../../util/queue.h" + +#include "../handler.h" +#include "../entry.h" + +struct cch_handler_http { + struct cch_handler handler; + off_t pointer; + str url; + array(struct aki_http) requests; + queue(s32) queue; + struct aki_event_loop loop; + struct aki_signal signal; + struct aki_mutex mutex; + struct aki_thread thread; +}; + +struct cch_entry *cch_handler_http_create(str *url); |