summaryrefslogtreecommitdiff
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
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>
-rw-r--r--cross/armv7a-linux-android.txt2
-rw-r--r--docs/design.txt6
-rw-r--r--flake.lock6
-rwxr-xr-xscripts/adb_logcat.sh2
-rwxr-xr-xscripts/run_valgrind.sh4
-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
-rw-r--r--src/codec/ffmpeg/demuxer.c2
-rw-r--r--src/liana/server.c36
-rw-r--r--src/liana/server.h15
-rw-r--r--src/libsink/sink.c11
-rw-r--r--src/screen/screen.c13
-rw-r--r--src/screen/screen.h4
-rw-r--r--src/server/server.c2
-rw-r--r--src/sink/input_simulator.c4
-rw-r--r--subprojects/packagefiles/ffmpeg/meson.build16
23 files changed, 190 insertions, 133 deletions
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 <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;
};
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 <nnwt/event_loop.h>
#include <nnwt/packet_stream.h>
#include <nnwt/packet_pool.h>
-#include <nnwt/socket.h>
#include <nnwt/signal.h>
#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',