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 --- cross/armv7a-linux-android.txt | 2 +- docs/design.txt | 6 +- flake.lock | 6 +- scripts/adb_logcat.sh | 2 + scripts/run_valgrind.sh | 4 +- 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 +- src/codec/ffmpeg/demuxer.c | 2 +- src/liana/server.c | 36 +++++--- src/liana/server.h | 15 ++- src/libsink/sink.c | 11 ++- src/screen/screen.c | 13 +-- src/screen/screen.h | 4 +- src/server/server.c | 2 +- src/sink/input_simulator.c | 4 +- subprojects/packagefiles/ffmpeg/meson.build | 16 +++- 23 files changed, 190 insertions(+), 133 deletions(-) create mode 100755 scripts/adb_logcat.sh diff --git a/cross/armv7a-linux-android.txt b/cross/armv7a-linux-android.txt index 9557dac..6bad7fe 100644 --- a/cross/armv7a-linux-android.txt +++ b/cross/armv7a-linux-android.txt @@ -10,6 +10,6 @@ CMAKE_SYSTEM_VERSION = '23' [host_machine] system = 'android' -cpu_family = 'armv7-a' +cpu_family = 'armv7-a' # Needed for CMake. cpu = 'armv7-a' endian = 'little' diff --git a/docs/design.txt b/docs/design.txt index 9604efb..bd73f6a 100644 --- a/docs/design.txt +++ b/docs/design.txt @@ -14,9 +14,11 @@ User Accounts ^^^^^^^^^^^^^ Users accounts contain 4 main concepts: lists, media lists, history, and settings. -* *Lists*: Transient lists of media connected to a user-selected set of sinks. Resources retrived via browse/search instances get added to a list for playback. A list's active sinks can be changed on the fly. +* *Lists*: Transient lists of media connected to a user-defined set of sinks. Resources retrived via browse/search instances get added to a list for playback. A list's set of active sinks can be changed on the fly. -* *Media Lists*: A more organized way to collect related media. They are stored to the user's account and therefore saved to the disk. +* *Collections*: A more organized way to collect related media. They are stored to the user's account and therefore saved to the disk. + +* *Packs*: [Working Idea] A more defined collection of media. Possibly being less broad, more detailed, for broader consumption, and/or archival. Something that should be easy to share. * *History*: A detailed history that tracks what a user played and when. Also saved to a user's account. diff --git a/flake.lock b/flake.lock index 81ba877..ab2a096 100644 --- a/flake.lock +++ b/flake.lock @@ -91,11 +91,11 @@ }, "nixpkgs_2": { "locked": { - "lastModified": 1737632463, - "narHash": "sha256-38J9QfeGSej341ouwzqf77WIHAScihAKCt8PQJ+NH28=", + "lastModified": 1737746512, + "narHash": "sha256-nU6AezEX4EuahTO1YopzueAXfjFfmCHylYEFCagduHU=", "owner": "NixOS", "repo": "nixpkgs", - "rev": "0aa475546ed21629c4f5bbf90e38c846a99ec9e9", + "rev": "825479c345a7f806485b7f00dbe3abb50641b083", "type": "github" }, "original": { diff --git a/scripts/adb_logcat.sh b/scripts/adb_logcat.sh new file mode 100755 index 0000000..b0da2a0 --- /dev/null +++ b/scripts/adb_logcat.sh @@ -0,0 +1,2 @@ +#! /usr/bin/env sh +adb logcat *:S CTV diff --git a/scripts/run_valgrind.sh b/scripts/run_valgrind.sh index 583327c..bccc34c 100755 --- a/scripts/run_valgrind.sh +++ b/scripts/run_valgrind.sh @@ -1,5 +1,5 @@ #! /usr/bin/env sh source ../scripts/python_env -valgrind --leak-check=full ./src/fruits/cmv/cmv "$@" +#valgrind --leak-check=full ./src/fruits/cmv/cmv "$@" #valgrind --leak-check=full --show-leak-kinds=all ./src/fruits/cmv/cmv "$@" -#valgrind --leak-check=no --show-error-list=yes ./src/fruits/cmv/cmv "$@" +valgrind --leak-check=no --show-error-list=yes ./src/fruits/cmv/cmv "$@" 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; }; diff --git a/src/codec/ffmpeg/demuxer.c b/src/codec/ffmpeg/demuxer.c index 5c37f1a..b3e5873 100644 --- a/src/codec/ffmpeg/demuxer.c +++ b/src/codec/ffmpeg/demuxer.c @@ -139,7 +139,7 @@ static s32 ff_demuxer_get_packet(struct camu_demuxer *demux, struct camu_codec_p } if (ret < 0) { al_log_error("ff_demuxer", "Failed to read frame (%s).", av_err2str(ret)); - continue; + break; } if (!subscribed_to_index(demux, pkt->stream_index)) { av_packet_unref(pkt); diff --git a/src/liana/server.c b/src/liana/server.c index 04b5d96..45200d1 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -14,13 +14,13 @@ bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop) server->loop = loop; server->increment = 0; al_array_init(server->nodes); - al_array_init(server->zombies); + al_array_init(server->dormant_connections); return true; } -static void remove_zombie(struct lia_server *server, struct nn_packet_stream *stream) +static void remove_dormant_connection(struct lia_server *server, struct nn_packet_stream *stream) { - al_array_remove(server->zombies, stream); + al_array_remove(server->dormant_connections, stream); } static void data_packet_sent_callback(void *userdata, struct nn_packet *packet) @@ -54,6 +54,10 @@ static void packet_pool_callback(void *userdata, struct nn_packet *packet) static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + if (conn->seek_pos != LIANA_TIMESTAMP_INVALID) { + conn->handler->seek(conn->handler, conn->seek_pos); + conn->seek_pos = LIANA_TIMESTAMP_INVALID; + } for (;;) { struct nn_packet *packet = nn_packet_pool_get(&conn->pool); if (!packet) break; @@ -168,7 +172,7 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet // Besides being wasteful, seeking to 0 on a new stream can skip data. if (mask != 0 || seek_pos > 0) { - conn->handler->seek(conn->handler, seek_pos); + conn->seek_pos = seek_pos; } if (mask == 0) { @@ -193,8 +197,7 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_server *server = (struct lia_server *)userdata; - // This connection might no longer be a zombie, but that's fine. - remove_zombie(server, stream); + remove_dormant_connection(server, stream); nn_packet_stream_free(stream); al_free(stream); } @@ -211,7 +214,6 @@ static void demote_and_disconnect_stream(struct lia_server *server, struct nn_pa // This should always be the expected behavior but here it's mainly to // not lose packets that belong to the packet pool. nn_packet_stream_discard_queue(stream); - // This connection will now be nothing but a packet stream. stream->userdata = server; stream->connection_closed_callback = connection_closed_callback; stream->packet_callback = discard_packet_callback; @@ -299,8 +301,8 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str { struct lia_server *server = (struct lia_server *)userdata; - // We got a packet, this connection is no longer a zombie. - remove_zombie(server, stream); + // We got a packet, this connection is no longer dormant. + remove_dormant_connection(server, stream); u32 node_id = nn_packet_read_u32(packet); u32 connection_id = nn_packet_read_u32(packet); @@ -320,6 +322,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str conn->errored = false; conn->disconnected = false; conn->ref = false; + conn->seek_pos = LIANA_TIMESTAMP_INVALID; stream->packet_callback = discard_packet_callback; stream->connection_closed_callback = pre_init_connection_closed_callback; nn_thread_create(&conn->thread, init_thread, conn); @@ -329,6 +332,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str // Cleanup the existing connection's handler and demote it's stream. // The stream was likely already disconnected client-side but it's still safe // to disconnect it here to be sure. + // An alternative to this could be to create a new connection if conn->ref. disable_connection_and_wait(conn); demote_and_disconnect_stream(server, conn->stream); conn->ref = false; @@ -352,7 +356,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) struct lia_server *server = (struct lia_server *)userdata; stream->packet_callback = packet_callback; stream->packet_sent_callback = packet_sent_callback; - al_array_push(server->zombies, stream); + al_array_push(server->dormant_connections, stream); return true; } @@ -400,6 +404,8 @@ static void duration_signal_callback(void *userdata) void lia_node_get_duration(struct lia_node *node) { + // It's vital we don't block during handler->init(), this is very wasteful though. + // We should be able to reuse this initialized handler with a sort of "connection pool". nn_signal_init(&node->signal, node->server->loop, duration_signal_callback, node); nn_signal_start(&node->signal); cch_entry_get_handle(node->entry, &node->handle); @@ -427,9 +433,9 @@ void lia_node_close(struct lia_node *node) void lia_server_close(struct lia_server *server) { - struct nn_packet_stream *zombie; - al_array_foreach_rev(server->zombies, i, zombie) { - nn_packet_stream_disconnect(zombie); + struct nn_packet_stream *conn; + al_array_foreach_rev(server->dormant_connections, i, conn) { + nn_packet_stream_disconnect(conn); } } @@ -443,6 +449,6 @@ void lia_server_free(struct lia_server *server) } al_array_free(server->nodes); - al_assert(!server->zombies.count); - al_array_free(server->zombies); + al_assert(!server->dormant_connections.count); + al_array_free(server->dormant_connections); } diff --git a/src/liana/server.h b/src/liana/server.h index a9e567a..7951ada 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -3,7 +3,6 @@ #include #include #include -#include #include #include "../cache/entry.h" @@ -17,6 +16,7 @@ struct lia_node_connection { bool ref; struct lia_server_handler *handler; struct cch_handle handle; + u64 seek_pos; struct nn_thread thread; struct nn_signal signal; struct nn_packet_pool pool; @@ -32,24 +32,23 @@ struct lia_node { struct cch_entry *entry; array(struct lia_node_connection *) connections; bool closed; - struct lia_server *server; u64 duration; - // Temporary copy from node_connection. We need to - // figure out a "connection pool" structure. - struct lia_server_handler *handler; + struct lia_server *server; + void (*callback)(void *, u8, u64); + void *userdata; + // Temporary copy-paste from node_connection. bool errored; + struct lia_server_handler *handler; struct cch_handle handle; struct nn_thread thread; struct nn_signal signal; - void (*callback)(void *, u8, u64); - void *userdata; }; struct lia_server { struct nn_event_loop *loop; u32 increment; array(struct lia_node *) nodes; - array(struct nn_packet_stream *) zombies; + array(struct nn_packet_stream *) dormant_connections; }; bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop); diff --git a/src/libsink/sink.c b/src/libsink/sink.c index b5e2a65..54f651b 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -795,11 +795,18 @@ static void audio_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_EOF: { nn_mutex_lock(&sink->mutex); al_log_info("sink", "Audio EOF."); + // @TODO: This is not well synced. EOF can happen at any time + // while other stuff is happening in the sink. For example + // while seeking, if EOF happens right after a REMOVE_BUFFERS, + // we might assert during RECONNECT because the entry is ended. + // + // state can be something other than BUFFER_ADDED here + // because it could have changed while waiting on the lock above. if (entry->audio.state == BUFFER_ADDED) { remove_entry_audio_buffer(sink, entry); } // This assert likely doesn't matter due to the handling of the ENDED state. - al_assert(entry->audio.state == BUFFER_SET_OR_BUFFERED); + //al_assert(entry->audio.state == BUFFER_SET_OR_BUFFERED); entry->audio.state = BUFFER_ENDED; bool end_entry = VIDEO_ENDED_OR_EMPTY(entry); #ifndef CAMU_SINK_NO_VIDEO @@ -847,7 +854,7 @@ static void video_buffer_callback(void *userdata, u8 op) if (entry->video.state == BUFFER_ADDED) { remove_entry_video_buffer(sink, entry); } - al_assert(entry->video.state == BUFFER_SET_OR_BUFFERED); + //al_assert(entry->video.state == BUFFER_SET_OR_BUFFERED); entry->video.state = BUFFER_ENDED; if (AUDIO_ENDED_OR_EMPTY(entry)) { swapped = end_entry_and_advance_queue(sink, entry); diff --git a/src/screen/screen.c b/src/screen/screen.c index 929020e..8e3ecbd 100644 --- a/src/screen/screen.c +++ b/src/screen/screen.c @@ -280,8 +280,6 @@ static void key_immediate_callback(void *userdata, u8 state, u8 button) bool camu_screen_init(struct camu_screen *scr, void *context) { al_atomic_store(s32)(&scr->state, CAMU_SCREEN_PAUSED, AL_ATOMIC_RELAXED); - al_array_init(scr->videos); - scr->renderer = NULL; scr->window = stl_window_create(context); scr->window->pointer_pos_callback = pointer_pos_callback; scr->window->scroll_callback = scroll_callback; @@ -292,15 +290,18 @@ bool camu_screen_init(struct camu_screen *scr, void *context) scr->window->refresh_callback = refresh_callback; scr->window->should_close_callback = should_close_callback; scr->window->userdata = scr; + scr->renderer = NULL; + scr->force_render = false; + scr->flags = CAMU_SCREEN_ZOOM_PAN_SIMPLE; + scr->last_mouse_x = 0.0; + scr->last_mouse_y = 0.0; + al_array_init(scr->videos); #ifdef CAMU_SCREEN_THREADED al_array_init(scr->add_queue); al_array_init(scr->rem_queue); al_atomic_store(u8)(&scr->queued, 0, AL_ATOMIC_RELAXED); nn_mutex_init(&scr->mutex); #endif - scr->last_mouse_x = 0.0; - scr->last_mouse_y = 0.0; - scr->flags = CAMU_SCREEN_ZOOM_PAN_SIMPLE; return true; } @@ -317,7 +318,7 @@ static nn_thread_result NNWT_THREADCALL event_thread(void *userdata) bool camu_screen_create_window(struct camu_screen *scr, const char *name) { - s32 flags = STELA_WINDOW_VSYNC; + u32 flags = STELA_WINDOW_VSYNC; if (!scr->window->create_window(scr->window, CAMU_SCREEN_WIDTH, CAMU_SCREEN_HEIGHT, "", name, flags)) { return false; } diff --git a/src/screen/screen.h b/src/screen/screen.h index c84cfc0..7203018 100644 --- a/src/screen/screen.h +++ b/src/screen/screen.h @@ -20,7 +20,7 @@ enum { CAMU_SCREEN_MODIFIER = 1, CAMU_SCREEN_DRAGGING = 1 << 1, - CAMU_SCREEN_ZOOM_PAN_SIMPLE = 1 << 2, + CAMU_SCREEN_ZOOM_PAN_SIMPLE = 1 << 2 }; enum { @@ -50,13 +50,13 @@ struct camu_screen_video { struct camu_screen { atomic(s32) state; - u32 flags; struct stl_window *window; #ifdef STELA_EVENT_BUFFER struct nn_thread thread; #endif struct camu_renderer *renderer; bool force_render; + u32 flags; s32 width; s32 height; bool last_fullscreen; diff --git a/src/server/server.c b/src/server/server.c index 9a6636d..e687bf8 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -654,7 +654,7 @@ static struct nn_rpc_command commands[] = { static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { - // @TODO: Cleanup zombie connections. + // @TODO: Cleanup alien connections. (void)userdata; (void)conn; } diff --git a/src/sink/input_simulator.c b/src/sink/input_simulator.c index baa7dc8..9cf395f 100644 --- a/src/sink/input_simulator.c +++ b/src/sink/input_simulator.c @@ -39,8 +39,8 @@ static nn_thread_result NNWT_THREADCALL input_simulation_thread(void *userdata) case SEEK: { f64 pos = al_rand() / (f64)AL_RAND_MAX; if (pos < 0.005) pos = 0.0; - if (pos > 0.995) pos = 100.0; - else if (pos > 0.99) pos = 99.9; + if (pos > 0.995) pos = 1.0; + else if (pos > 0.99) pos = .9999; camu_sink_seek(sink, pos); break; } diff --git a/subprojects/packagefiles/ffmpeg/meson.build b/subprojects/packagefiles/ffmpeg/meson.build index 435ee90..33dbd30 100644 --- a/subprojects/packagefiles/ffmpeg/meson.build +++ b/subprojects/packagefiles/ffmpeg/meson.build @@ -9,6 +9,7 @@ is_minsize = get_option('buildtype') == 'minsize' is_windows = host_machine.system() == 'windows' is_msvc = compiler.get_id() == 'msvc' or compiler.get_id() == 'clang-cl' is_mingw = is_windows and not is_msvc +is_android = host_machine.system() == 'android' extra_options = [] @@ -35,6 +36,8 @@ if is_windows extra_options += ['--toolchain=msvc'] endif extra_options += ['--enable-w32threads'] +elif is_android + extra_options += ['--target-os=android'] endif #'--extra-ldflags=-Lbuild/subprojects/libalabaster -l:libalabaster.a', @@ -52,7 +55,6 @@ demuxers = 'flac,mp3,aac,wav,image2,mjpeg,image2pipe,image_jpeg_pipe,gif,matrosk parsers = 'aac,opus,mjpeg,jpeg2000,gif,h264,hevc,av1,vp9,vp8' #parsers += ',mpegaudio,mpegvideo,dvd_nav' -#extra_options += ['--enable-zlib', '--enable-libsoxr'] extra_options += ['--enable-zlib'] decoders += ',png,apng' demuxers += ',image_png_pipe' @@ -61,18 +63,24 @@ parsers += ',png' bsfs = 'extract_extradata,mp3_header_decompress' hwaccels = '' -if is_windows +if is_windows and get_option('sink-use-vulkan') extra_options += ['--enable-vulkan'] hwaccels += 'h264_vulkan,hevc_vulkan,av1_vulkan' #extra_options += ['--enable-d3d11va', '--enable-d3d12va', '--enable-dxva2'] #hwaccels += 'h264_dxva2,' +elif is_android + extra_options += ['--enable-jni', '--enable-mediacodec'] + decoders += ',h264_mediacodec,hevc_mediacodec,av1_mediacodec' else - extra_options += ['--disable-vulkan'] + extra_options += ['--disable-vulkan', '--disable-vdpau', '--disable-vaapi'] endif -extra_options += ['--disable-vdpau', '--disable-vaapi'] protocols = 'file,cache' +if not is_android + extra_options += ['--enable-libsoxr'] +endif + ext_proj = import('unstable-external_project') proj = ext_proj.add_project('configure', -- cgit v1.2.3-101-g0448