summaryrefslogtreecommitdiff
path: root/src/cache/handlers
diff options
context:
space:
mode:
Diffstat (limited to 'src/cache/handlers')
-rw-r--r--src/cache/handlers/file.c49
-rw-r--r--src/cache/handlers/file.h12
-rw-r--r--src/cache/handlers/http.c207
-rw-r--r--src/cache/handlers/http.h24
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);