From 457a3cc1a04e45e31370d9083186436b0d12ab1d Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 27 Jan 2025 21:56:39 -0500 Subject: Refactor threaded waits - Move seek to handler thread. - Cleanup and comment some stuff. Signed-off-by: Andrew Opalach --- src/cache/entry.c | 2 +- src/cache/handler.h | 5 -- src/cache/handlers/cdio.c | 12 ++-- src/cache/handlers/cdio.h | 2 + src/cache/handlers/http.c | 14 ++--- src/cache/handlers/http.h | 4 +- src/cache/threaded_waits.c | 137 +++++++++++++++++++++++++++------------------ src/cache/threaded_waits.h | 22 +++++--- src/cache/wait.h | 2 +- 9 files changed, 116 insertions(+), 84 deletions(-) (limited to 'src/cache') 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 -#include #include -#include #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 +#include -#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; }; -- cgit v1.2.3-101-g0448