summaryrefslogtreecommitdiff
path: root/src/cache
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-01-27 21:56:39 -0500
committerAndrew Opalach <andrew@akon.city> 2025-01-27 21:56:39 -0500
commit457a3cc1a04e45e31370d9083186436b0d12ab1d (patch)
tree7ffb9d676edd4a1b88e8f1020dd905efcfe6b83a /src/cache
parentf760ecedb619a55ec8ee989639ac385f27e82d98 (diff)
downloadcamu-457a3cc1a04e45e31370d9083186436b0d12ab1d.tar.gz
camu-457a3cc1a04e45e31370d9083186436b0d12ab1d.tar.bz2
camu-457a3cc1a04e45e31370d9083186436b0d12ab1d.zip
Refactor threaded waits
- Move seek to handler thread. - Cleanup and comment some stuff. Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/cache')
-rw-r--r--src/cache/entry.c2
-rw-r--r--src/cache/handler.h5
-rw-r--r--src/cache/handlers/cdio.c12
-rw-r--r--src/cache/handlers/cdio.h2
-rw-r--r--src/cache/handlers/http.c14
-rw-r--r--src/cache/handlers/http.h4
-rw-r--r--src/cache/threaded_waits.c137
-rw-r--r--src/cache/threaded_waits.h22
-rw-r--r--src/cache/wait.h2
9 files changed, 116 insertions, 84 deletions
diff --git a/src/cache/entry.c b/src/cache/entry.c
index 8930ea0..3cae147 100644
--- a/src/cache/entry.c
+++ b/src/cache/entry.c
@@ -6,9 +6,9 @@ bool cch_entry_get_handle(struct cch_entry *entry, struct cch_handle *handle)
{
handle->entry = entry;
handle->pointer = 0;
+ handle->wait.disabled = false;
nn_cond_init(&handle->wait.cond);
nn_mutex_init(&handle->wait.mutex);
- handle->wait.disabled = false;
nn_mutex_lock(&entry->mutex);
entry->ref_count++;
nn_mutex_unlock(&entry->mutex);
diff --git a/src/cache/handler.h b/src/cache/handler.h
index b4c6975..0a64573 100644
--- a/src/cache/handler.h
+++ b/src/cache/handler.h
@@ -1,9 +1,7 @@
#pragma once
#include <al/types.h>
-#include <al/array.h>
#include <al/str.h>
-#include <nnwt/thread.h>
#include "handle.h"
@@ -14,7 +12,4 @@ struct cch_handler {
void (*free)(struct cch_handler **);
struct cch_entry *entry;
str liana;
- bool disabled;
- struct nn_mutex mutex;
- array(struct cch_handler_wait *) waits;
};
diff --git a/src/cache/handlers/cdio.c b/src/cache/handlers/cdio.c
index ac90c03..1b4c44b 100644
--- a/src/cache/handlers/cdio.c
+++ b/src/cache/handlers/cdio.c
@@ -68,7 +68,7 @@ static nn_thread_result NNWT_THREADCALL cd_read_thread(void *userdata)
cch_backing_fill_range(cdio->backing, pointer, BYTES_PER_STEP);
cdio->sector += ret;
cdio->backing->unlock(cdio->backing);
- cch_threaded_waits_evaluate(&cdio->handler, cdio->backing);
+ cch_threaded_waits_evaluate(&cdio->waits, cdio->backing);
if (ret < SECTORS_PER_STEP || cdio->sector > cdio->end) {
al_log_debug("cdio", "Normal CD EOF.");
break;
@@ -78,7 +78,7 @@ static nn_thread_result NNWT_THREADCALL cd_read_thread(void *userdata)
if (cdio->sector <= cdio->end) {
al_log_warn("cdio", "Disc read cut short (sectors: %zd, end: %zd).", cdio->sector, cdio->end);
// We didn't read the full disc, disable waiters.
- cch_threaded_waits_disable_all(&cdio->handler);
+ cch_threaded_waits_disable_all(&cdio->waits);
}
return 0;
@@ -101,13 +101,13 @@ static void handler_cdio_maybe_spawn_worker(struct cch_handler *handler, size_t
static bool handler_cdio_wait_for_range(struct cch_handler *handler, struct cch_handler_wait *wait)
{
struct cch_handler_cdio *cdio = (struct cch_handler_cdio *)handler;
- return cch_threaded_wait_for_range(&cdio->handler, cdio->backing, wait);
+ return cch_threaded_wait_for_range(&cdio->waits, cdio->handler.entry, cdio->backing, wait);
}
static void handler_cdio_free(struct cch_handler **handler)
{
struct cch_handler_cdio *cdio = (struct cch_handler_cdio *)*handler;
- cch_threaded_waits_disable_all(&cdio->handler);
+ cch_threaded_waits_disable_all(&cdio->waits);
if (al_atomic_load(s32)(&cdio->running, AL_ATOMIC_RELAXED)) {
al_atomic_store(s32)(&cdio->running, 0, AL_ATOMIC_RELAXED);
nn_thread_join(&cdio->thread);
@@ -116,7 +116,7 @@ static void handler_cdio_free(struct cch_handler **handler)
cdio_log_messages(cdio);
cdio_cddap_close(cdio->drive);
}
- cch_threaded_waits_close(&cdio->handler);
+ cch_threaded_waits_close(&cdio->waits);
al_str_free(&cdio->handler.liana);
al_free(cdio);
*handler = NULL;
@@ -208,6 +208,6 @@ struct cch_entry *cch_handler_cdio_create(void)
nn_buffer_init(&cdio->buffer);
nn_buffer_ensure_space(&cdio->buffer, BYTES_PER_STEP);
nn_buffer_set_size(&cdio->buffer, BYTES_PER_STEP);
- cch_threaded_waits_init(&cdio->handler);
+ cch_threaded_waits_init(&cdio->waits);
return entry;
}
diff --git a/src/cache/handlers/cdio.h b/src/cache/handlers/cdio.h
index 3a3966a..650d2e5 100644
--- a/src/cache/handlers/cdio.h
+++ b/src/cache/handlers/cdio.h
@@ -5,6 +5,7 @@
#include "../handler.h"
#include "../entry.h"
+#include "../threaded_waits.h"
struct cch_handler_cdio {
struct cch_handler handler;
@@ -16,6 +17,7 @@ struct cch_handler_cdio {
struct nn_buffer buffer;
atomic(s32) running;
struct nn_thread thread;
+ struct cch_handler_waits waits;
};
struct cch_entry *cch_handler_cdio_create(void);
diff --git a/src/cache/handlers/http.c b/src/cache/handlers/http.c
index b34e4ba..1856014 100644
--- a/src/cache/handlers/http.c
+++ b/src/cache/handlers/http.c
@@ -30,7 +30,7 @@ static void http_callback(void *userdata, u8 op, u8 *buf, void *opaque)
http->backing->write(http->backing, buf, http->pointer, &n);
http->pointer += n;
*(size_t *)opaque = n;
- cch_threaded_waits_evaluate(&http->handler, http->backing);
+ cch_threaded_waits_evaluate(&http->waits, http->backing);
break;
}
case NNWT_HTTP_RESPONSE_CODE: {
@@ -42,7 +42,7 @@ static void http_callback(void *userdata, u8 op, u8 *buf, void *opaque)
curl_off_t length = *(curl_off_t *)opaque;
al_log_debug("cache_handler_http", "Content-Length: %"CURL_FORMAT_CURL_OFF_T".", length);
cch_entry_set_size(http->handler.entry, length);
- cch_threaded_waits_signal_any(&http->handler);
+ cch_threaded_waits_signal_any(&http->waits);
break;
}
case NNWT_HTTP_REDIRECT: {
@@ -54,7 +54,7 @@ static void http_callback(void *userdata, u8 op, u8 *buf, void *opaque)
al_log_debug("cache_handler_http", "Transfer finished.");
break;
case NNWT_HTTP_ERROR:
- cch_threaded_waits_disable_all(&http->handler);
+ cch_threaded_waits_disable_all(&http->waits);
break;
default:
break;
@@ -77,18 +77,18 @@ static void handler_http_maybe_spawn_worker(struct cch_handler *handler, size_t
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;
- return cch_threaded_wait_for_range(&http->handler, http->backing, wait);
+ return cch_threaded_wait_for_range(&http->waits, http->handler.entry, http->backing, wait);
}
static void handler_http_free(struct cch_handler **handler)
{
struct cch_handler_http *http = (struct cch_handler_http *)*handler;
- cch_threaded_waits_disable_all(&http->handler);
+ cch_threaded_waits_disable_all(&http->waits);
struct nn_http *request;
al_array_foreach_ptr(http->requests, i, request) {
nn_http_close(request);
}
- cch_threaded_waits_close(&http->handler);
+ cch_threaded_waits_close(&http->waits);
al_array_free(http->requests);
al_str_free(&http->url);
al_str_free(&http->handler.liana);
@@ -119,6 +119,6 @@ struct cch_entry *cch_handler_http_create(str *url, struct nn_event_loop *loop)
al_str_clone(&http->url, url);
http->pointer = 0;
al_array_init(http->requests);
- cch_threaded_waits_init(&http->handler);
+ cch_threaded_waits_init(&http->waits);
return entry;
}
diff --git a/src/cache/handlers/http.h b/src/cache/handlers/http.h
index ef47173..2b87db9 100644
--- a/src/cache/handlers/http.h
+++ b/src/cache/handlers/http.h
@@ -6,14 +6,16 @@
#include "../handler.h"
#include "../entry.h"
+#include "../threaded_waits.h"
struct cch_handler_http {
struct cch_handler handler;
struct cch_backing *backing;
+ struct nn_event_loop *loop;
str url;
off_t pointer;
array(struct nn_http) requests;
- struct nn_event_loop *loop;
+ struct cch_handler_waits waits;
};
struct cch_entry *cch_handler_http_create(str *url, struct nn_event_loop *loop);
diff --git a/src/cache/threaded_waits.c b/src/cache/threaded_waits.c
index ae43ae9..541bee6 100644
--- a/src/cache/threaded_waits.c
+++ b/src/cache/threaded_waits.c
@@ -1,78 +1,116 @@
#include "threaded_waits.h"
-#include "entry.h"
-void cch_threaded_waits_init(struct cch_handler *handler)
+void cch_threaded_waits_init(struct cch_handler_waits *waits)
{
- handler->disabled = false;
- nn_mutex_init(&handler->mutex);
- al_array_init(handler->waits);
+ waits->disabled = false;
+ nn_mutex_init(&waits->mutex);
+ al_array_init(waits->active);
}
static bool wait_range_satisfied(struct cch_backing *backing, struct cch_handler_wait *wait)
{
- if (wait->end < 0) return true;
+ if (wait->end < 0) {
+ return true;
+ }
+
struct cch_range *range;
al_array_foreach_ptr(backing->available, i, range) {
if (wait->start >= range->start && range->end >= wait->end) {
return true;
}
}
+
return false;
}
-bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, struct cch_handler_wait *wait)
+bool cch_threaded_wait_for_range(struct cch_handler_waits *waits, struct cch_entry *entry,
+ struct cch_backing *backing, struct cch_handler_wait *wait)
{
- nn_mutex_lock(&handler->mutex);
- bool canceled = handler->disabled;
- off_t size = cch_entry_get_size(handler->entry);
+ nn_mutex_lock(&waits->mutex);
+
+ off_t size = cch_entry_get_size(entry);
backing->lock(backing);
bool satisfied = size >= 0 && wait_range_satisfied(backing, wait);
backing->unlock(backing);
+
+ nn_mutex_lock(&wait->mutex);
+
+ bool canceled = waits->disabled || wait->disabled;
+
if (!canceled && !satisfied) {
- nn_mutex_lock(&wait->mutex);
- if (!wait->disabled) {
- al_array_push(handler->waits, wait);
- nn_mutex_unlock(&handler->mutex);
- nn_cond_wait(&wait->cond, &wait->mutex);
- } else {
- nn_mutex_unlock(&wait->mutex);
- nn_mutex_unlock(&handler->mutex);
- return false;
- }
- nn_mutex_lock(&handler->mutex);
+ al_array_push(waits->active, wait);
+
+ // Unlock handler to wait.
+ nn_mutex_unlock(&waits->mutex);
+ nn_cond_wait(&wait->cond, &wait->mutex);
+ // A successful wait will have been removed before signaling.
+
if (wait->disabled) {
- al_array_remove(handler->waits, wait);
canceled = true;
+ // Re-lock to remove.
+ nn_mutex_lock(&waits->mutex);
+ al_array_remove(waits->active, wait);
}
- nn_mutex_unlock(&wait->mutex);
}
- nn_mutex_unlock(&handler->mutex);
+
+ nn_mutex_unlock(&wait->mutex);
+ if (canceled || satisfied) {
+ // Else, we already unlocked to wait.
+ nn_mutex_unlock(&waits->mutex);
+ }
+
return !canceled;
}
-void cch_threaded_waits_signal_any(struct cch_handler *handler)
+void cch_threaded_waits_signal_any(struct cch_handler_waits *waits)
{
- nn_mutex_lock(&handler->mutex);
+ nn_mutex_lock(&waits->mutex);
+
struct cch_handler_wait *wait;
- al_array_foreach_rev(handler->waits, i, wait) {
+ al_array_foreach_rev(waits->active, i, wait) {
if (wait->end < 0) {
- al_array_remove_at(handler->waits, i);
+ // This will be considered a successful wait.
+ al_array_remove_at(waits->active, i);
+ nn_mutex_lock(&wait->mutex);
+ nn_cond_signal(&wait->cond);
+ nn_mutex_unlock(&wait->mutex);
+ }
+ }
+
+ nn_mutex_unlock(&waits->mutex);
+}
+
+void cch_threaded_waits_evaluate(struct cch_handler_waits *waits, struct cch_backing *backing)
+{
+ nn_mutex_lock(&waits->mutex);
+ backing->lock(backing);
+
+ struct cch_handler_wait *wait;
+ al_array_foreach_rev(waits->active, i, wait) {
+ if (wait_range_satisfied(backing, wait)) {
+ al_array_remove_at(waits->active, i);
nn_mutex_lock(&wait->mutex);
nn_cond_signal(&wait->cond);
nn_mutex_unlock(&wait->mutex);
}
}
- nn_mutex_unlock(&handler->mutex);
+
+ backing->unlock(backing);
+ nn_mutex_unlock(&waits->mutex);
}
void cch_threaded_wait_disable(struct cch_handler_wait *wait)
{
nn_mutex_lock(&wait->mutex);
+
+ // If this wait is in waits->active, cond_is_waiting() will be true.
+ // This should mean we can be assured that successive calls to
+ // disable() -> enable() won't leave this wait in waits->active.
wait->disabled = true;
- // cond_is_waiting() could be false if the handler is disabled.
if (nn_cond_is_waiting(&wait->cond)) {
nn_cond_signal(&wait->cond);
}
+
nn_mutex_unlock(&wait->mutex);
}
@@ -83,44 +121,31 @@ void cch_threaded_wait_enable(struct cch_handler_wait *wait)
nn_mutex_unlock(&wait->mutex);
}
-void cch_threaded_waits_disable_all(struct cch_handler *handler)
+void cch_threaded_waits_disable_all(struct cch_handler_waits *waits)
{
- nn_mutex_lock(&handler->mutex);
- if (handler->disabled) {
- nn_mutex_unlock(&handler->mutex);
+ nn_mutex_lock(&waits->mutex);
+
+ if (waits->disabled) {
+ nn_mutex_unlock(&waits->mutex);
return;
}
+
// Disallow any further waits.
- handler->disabled = true;
+ waits->disabled = true;
+
struct cch_handler_wait *wait;
- al_array_foreach(handler->waits, i, wait) {
+ al_array_foreach_rev(waits->active, i, wait) {
nn_mutex_lock(&wait->mutex);
wait->disabled = true;
nn_cond_signal(&wait->cond);
nn_mutex_unlock(&wait->mutex);
}
- nn_mutex_unlock(&handler->mutex);
-}
-void cch_threaded_waits_evaluate(struct cch_handler *handler, struct cch_backing *backing)
-{
- nn_mutex_lock(&handler->mutex);
- backing->lock(backing);
- struct cch_handler_wait *wait;
- al_array_foreach_rev(handler->waits, i, wait) {
- nn_mutex_lock(&wait->mutex);
- if (wait_range_satisfied(backing, wait)) {
- nn_cond_signal(&wait->cond);
- al_array_remove_at(handler->waits, i);
- }
- nn_mutex_unlock(&wait->mutex);
- }
- backing->unlock(backing);
- nn_mutex_unlock(&handler->mutex);
+ nn_mutex_unlock(&waits->mutex);
}
-void cch_threaded_waits_close(struct cch_handler *handler)
+void cch_threaded_waits_close(struct cch_handler_waits *waits)
{
- nn_mutex_destroy(&handler->mutex);
- al_array_free(handler->waits);
+ nn_mutex_destroy(&waits->mutex);
+ al_array_free(waits->active);
}
diff --git a/src/cache/threaded_waits.h b/src/cache/threaded_waits.h
index f8c05a2..1893395 100644
--- a/src/cache/threaded_waits.h
+++ b/src/cache/threaded_waits.h
@@ -1,15 +1,23 @@
#pragma once
#include <al/types.h>
+#include <nnwt/thread.h>
-#include "handler.h"
#include "backing.h"
+#include "entry.h"
-void cch_threaded_waits_init(struct cch_handler *handler);
-bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, struct cch_handler_wait *wait);
-void cch_threaded_waits_signal_any(struct cch_handler *handler);
+struct cch_handler_waits {
+ bool disabled;
+ struct nn_mutex mutex;
+ array(struct cch_handler_wait *) active;
+};
+
+void cch_threaded_waits_init(struct cch_handler_waits *waits);
+bool cch_threaded_wait_for_range(struct cch_handler_waits *waits, struct cch_entry *entry,
+ struct cch_backing *backing, struct cch_handler_wait *wait);
+void cch_threaded_waits_signal_any(struct cch_handler_waits *waits);
+void cch_threaded_waits_evaluate(struct cch_handler_waits *waits, struct cch_backing *backing);
void cch_threaded_wait_disable(struct cch_handler_wait *wait);
void cch_threaded_wait_enable(struct cch_handler_wait *wait);
-void cch_threaded_waits_disable_all(struct cch_handler *handler);
-void cch_threaded_waits_evaluate(struct cch_handler *handler, struct cch_backing *backing);
-void cch_threaded_waits_close(struct cch_handler *handler);
+void cch_threaded_waits_disable_all(struct cch_handler_waits *waits);
+void cch_threaded_waits_close(struct cch_handler_waits *waits);
diff --git a/src/cache/wait.h b/src/cache/wait.h
index 4c9b228..00c5749 100644
--- a/src/cache/wait.h
+++ b/src/cache/wait.h
@@ -4,7 +4,7 @@
struct cch_handler_wait {
off_t start, end;
+ bool disabled;
struct nn_cond cond;
struct nn_mutex mutex;
- bool disabled;
};