From e9475ce94ba69bd437d8cf0cf3f78062be928568 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Sat, 7 Jun 2025 11:54:55 -0400 Subject: Clarify and fix various sink behaviors Signed-off-by: Andrew Opalach --- src/buffer/audio.c | 2 + src/buffer/clock.c | 13 ++++- src/buffer/video.c | 34 ++++++----- src/buffer/video_null.h | 6 +- src/cache/handlers/cdio.c | 1 + src/codec/ffmpeg/decoder.c | 6 +- src/fruits/cmc/cmc.c | 2 +- src/fruits/cmsrv/cmsrv.c | 2 +- src/fruits/cmv/cmv.c | 3 +- src/fruits/ctv/ctv.c | 1 + src/liana/client.c | 2 +- src/liana/handlers/codec_client.c | 2 +- src/liana/list.c | 13 +++-- src/liana/server.c | 17 +++--- src/liana/vcr.c | 7 ++- src/liana/vcr.h | 3 +- src/libsink/input_simulator.c | 3 +- src/libsink/sink.c | 100 ++++++++++++++++++-------------- src/libsink/sink.h | 1 + src/portal/scripts/create_vendor.sh | 2 + src/portal/scripts/fetch_python_deps.sh | 8 +++ src/portal/src/search.c | 1 + src/render/meson.build | 6 +- src/screen/screen.c | 91 +++++++++++++++++++++-------- src/screen/screen.h | 7 ++- src/server/server.c | 2 + src/util/print_time.h | 5 +- 27 files changed, 224 insertions(+), 116 deletions(-) create mode 100755 src/portal/scripts/fetch_python_deps.sh (limited to 'src') diff --git a/src/buffer/audio.c b/src/buffer/audio.c index 37aff4a..fd395e0 100644 --- a/src/buffer/audio.c +++ b/src/buffer/audio.c @@ -250,6 +250,8 @@ void camu_audio_buffer_resync(struct camu_audio_buffer *buf) ptrdiff_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, ptrdiff_t req) { + al_assert(buf->buffered); + struct camu_audio_format *fmt = &buf->fmt.req; ptrdiff_t ret, signal = req; diff --git a/src/buffer/clock.c b/src/buffer/clock.c index 6a30af9..508a899 100644 --- a/src/buffer/clock.c +++ b/src/buffer/clock.c @@ -44,14 +44,23 @@ void camu_clock_seek(struct camu_clock *clock, f64 base, u64 target) al_atomic_store(f64)(&clock->tick, tick, AL_ATOMIC_RELAXED); } } else { - // The value of clock->pause cannot be touched here. A reconnecting entry - // may be relying on a CLOCK_PAUSED callback for sync. + // Don't touch clock->pause here for the sake of sync. + // Theoretically, the sink could rely on a CLOCK_PAUSED event from an entry + // even after it was seeked. That would ultimately maintain sync but the + // clock would report the old position in get_pts() before triggering CLOCK_PAUSED, + // likely confusing the audio/video buffers of an entry. This is mitigated in + // CLIENT_REMOVE_BUFFERS by immediately switching to the target instead of waiting for + // a clock_callback(). clock->paused_at = 0.0; } } void camu_clock_loop(struct camu_clock *clock, f64 last_pts) { + // Seek to 0 but include the time it took to perform the seek. + // Even without skipping a frame this method of looping is likely + // to mess up the frame pacing between the last frame of the previous + // loop and the first frame of this loop. clock->offset += last_pts - clock->base; clock->base = 0.0; } diff --git a/src/buffer/video.c b/src/buffer/video.c index 3eee026..9a33540 100644 --- a/src/buffer/video.c +++ b/src/buffer/video.c @@ -1,15 +1,14 @@ #define AL_LOG_SECTION "video_buffer" #include -#include "video.h" -#include "common.h" -#include "common_internal.h" - #ifdef CAMU_HAVE_FFMPEG #include "../codec/ffmpeg/common.h" +#include "../codec/ffmpeg/scaler.h" #endif -#include "../codec/ffmpeg/scaler.h" +#include "video.h" +#include "common.h" +#include "common_internal.h" #define BUFFER_MARK_LOW ((1.0 / 24.0) * 2) #define BUFFER_MARK_BUFFERED ((1.0 / 24.0) * 4) // Must be >1. @@ -22,14 +21,9 @@ bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *cl buf->latency = 0.0; al_atomic_store(f64)(&buf->pts, -1.0, AL_ATOMIC_RELAXED); buf->reset_pts = -1.0; - // The least confusing behavior for single_frame is that it can't be - // set if the buffer is empty. - buf->single_frame = false; - buf->avg_frame_duration = 0.0; + buf->single_frame = true; buf->queue = NULL; buf->buffered = false; - buf->buffered_with_one_frame = false; - buf->weighted_read = true; al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED); #ifdef CAMU_SCREEN_THREADED al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED); @@ -49,6 +43,7 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_code switch (stream->mode) { case CAMU_NORMAL: { buf->single_frame = true; + buf->avg_frame_duration = 0.0; log_info("Stream: %s (%ux%u) IMAGE.", format_name, fmt->width, fmt->height); break; } @@ -57,15 +52,18 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_code AVRational frame_rate = stream->av.stream->avg_frame_rate; buf->single_frame = stream->duration == 0 || frame_rate.den == 0; if (buf->single_frame) { + buf->avg_frame_duration = 0.0; log_info("Stream: %s (%ux%u) IMAGE.", format_name, fmt->width, fmt->height); } else { - buf->avg_frame_duration = (frame_rate.den > 0) ? av_q2d(av_inv_q(frame_rate)) : 0.0; + al_assert(frame_rate.den > 0); + buf->avg_frame_duration = av_q2d(av_inv_q(frame_rate)); log_info("Stream: %s (%ux%u) VIDEO %.3ffps.", format_name, fmt->width, fmt->height, av_q2d(frame_rate)); } break; } #endif } + buf->weighted_read = !buf->single_frame; struct camu_video_format *in = &buf->fmt.in; struct camu_video_format *req = &buf->fmt.req; camu_video_format_copy(in, fmt); @@ -114,6 +112,9 @@ void camu_video_buffer_set_latency(struct camu_video_buffer *buf, s32 frames) static void after_push_internal(struct camu_video_buffer *buf) { + // Note that queue->reset() must only be called from the read() thread. + // If a reset where to happen from the this thread, the latest read() + // frame could be freed before it was used. s32 count = buf->queue->count(buf->queue); f64 have = count * buf->avg_frame_duration; if (!buf->buffered && (buf->single_frame || have >= BUFFER_MARK_BUFFERED)) { @@ -126,9 +127,6 @@ static void after_push_internal(struct camu_video_buffer *buf) buf->buffered_with_one_frame = count == 1; log_debug("Buffered (mark: %.2fs).", have); buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED); - } else if (have >= BUFFER_MARK_RESET) { - log_warn("Buffer overflow, resetting."); - buf->queue->reset(buf->queue); } else if (have >= BUFFER_MARK_HIGH) { buf->callback(buf->userdata, CAMU_BUFFER_CORK); } @@ -220,6 +218,8 @@ void camu_video_buffer_reset(struct camu_video_buffer *buf, f64 pts) bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weighted) { + al_assert(buf->buffered); + f64 base_pts = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE); if (!buf->single_frame) { f64 pts = camu_clock_get_pts(buf->clock, buf->latency, !buf->weighted_read); @@ -250,6 +250,10 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weig f64 have = buf->queue->count(buf->queue) * buf->avg_frame_duration; if (have <= BUFFER_MARK_LOW) { buf->callback(buf->userdata, CAMU_BUFFER_UNCORK); + } else if (have >= BUFFER_MARK_RESET) { + log_warn("Buffer overflow, resetting."); + ret = CAMU_QUEUE_ERR; // Don't use the frame written to out. + buf->queue->reset(buf->queue); } } diff --git a/src/buffer/video_null.h b/src/buffer/video_null.h index 3de0001..15d2aba 100644 --- a/src/buffer/video_null.h +++ b/src/buffer/video_null.h @@ -26,7 +26,7 @@ AL_IGNORE_WARNING("-Wunused-function") static bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock) { (void)clock; - buf->single_frame = false; + buf->single_frame = true; buf->avg_frame_duration = 0.0; #ifdef CAMU_SCREEN_THREADED al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED); @@ -84,8 +84,8 @@ static bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, boo { (void)buf; (void)out; - *weighted = false; - return false; + (void)weighted; + al_assert_and_return(false); } static void camu_video_buffer_free(struct camu_video_buffer *buf) diff --git a/src/cache/handlers/cdio.c b/src/cache/handlers/cdio.c index ba7898f..82600dc 100644 --- a/src/cache/handlers/cdio.c +++ b/src/cache/handlers/cdio.c @@ -36,6 +36,7 @@ static bool handler_cdio_can_seek(struct cch_handler *handler) static nn_thread_result NNWT_THREADCALL cd_read_thread(void *userdata) { struct cch_handler_cdio *cdio = (struct cch_handler_cdio *)userdata; + nn_thread_set_name("cdio_read"); for (;;) { if (!al_atomic_load(s32)(&cdio->running, AL_ATOMIC_RELAXED)) { diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c index e097e0c..f201ea3 100644 --- a/src/codec/ffmpeg/decoder.c +++ b/src/codec/ffmpeg/decoder.c @@ -12,11 +12,13 @@ #ifdef CAMU_FF_DECODER_HWACCEL #if defined STELA_API_VULKAN static const char *hwdevces[] = { "vulkan" }; -#elif defined NAUNET_ON_WINDOWS +#else +#if defined NAUNET_ON_WINDOWS static const char *hwdevces[] = { "d3d11va" }; -#elif defined STELA_API_OPENGL +#else static const char *hwdevces[] = { "vaapi" }; #endif +#endif static s32 get_buffer2(AVCodecContext *context, AVFrame *pic, s32 flags) { diff --git a/src/fruits/cmc/cmc.c b/src/fruits/cmc/cmc.c index b9810f6..0529f81 100644 --- a/src/fruits/cmc/cmc.c +++ b/src/fruits/cmc/cmc.c @@ -265,7 +265,7 @@ static bool parse_command_line(s32 argc, char *argv[]) s32 main(s32 argc, char *argv[]) { - if (!nn_common_init()) return EXIT_FAILURE; + if (!nn_common_init("cmc_main")) return EXIT_FAILURE; al_array_init(c.cli.args); if (!parse_command_line(argc, argv)) { diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c index ca65beb..431c682 100644 --- a/src/fruits/cmsrv/cmsrv.c +++ b/src/fruits/cmsrv/cmsrv.c @@ -130,7 +130,7 @@ s32 main(s32 argc, char *argv[]) { (void)argc; (void)argv; - if (!nn_common_init()) return EXIT_FAILURE; + if (!nn_common_init("cmsrv_main")) return EXIT_FAILURE; #ifdef CMSRV_USE_UI if (!cmsrv_ui_init(&s.ui, &s.server)) return EXIT_FAILURE; al_set_print(log_callback, &s); diff --git a/src/fruits/cmv/cmv.c b/src/fruits/cmv/cmv.c index 28dcbb5..727261b 100644 --- a/src/fruits/cmv/cmv.c +++ b/src/fruits/cmv/cmv.c @@ -40,6 +40,7 @@ static void exit_callback(void *userdata, struct camu_desktop *desktop) static nn_thread_result NNWT_THREADCALL event_loop_thread(void *userdata) { struct cmv *c = (struct cmv *)userdata; + nn_thread_set_name("event_loop"); nn_event_loop_run(&c->loop); return 0; } @@ -231,7 +232,7 @@ s32 wmain(s32 argc, wchar_t **argv) s32 main(s32 argc, char *argv[]) #endif { - if (!nn_common_init() || !stl_global_init(false)) { + if (!nn_common_init("cmv_main") || !stl_global_init(false)) { return EXIT_FAILURE; } diff --git a/src/fruits/ctv/ctv.c b/src/fruits/ctv/ctv.c index 5becb88..7357177 100644 --- a/src/fruits/ctv/ctv.c +++ b/src/fruits/ctv/ctv.c @@ -18,6 +18,7 @@ struct ctv { static nn_thread_result NNWT_THREADCALL event_loop_thread(void *userdata) { struct ctv *c = (struct ctv *)userdata; + nn_thread_set_name("event_loop"); nn_event_loop_init(&c->loop); diff --git a/src/liana/client.c b/src/liana/client.c index 00eaca6..5222d7a 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -290,7 +290,7 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, client->mask = 0; al_array_init(client->streams); client->reconnect = RECONNECT_NONE; - lia_vcr_init(&client->vcr, client->loop, &client->data); + lia_vcr_init(&client->vcr, client->loop, &client->data, node_id); al_str_clone(&client->addr, addr); client->port = port; client->connection_id = 0; diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index 38b5d3b..b85fc43 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -4,7 +4,6 @@ #include "../../codec/ffmpeg/decoder.h" #include "../../codec/ffmpeg/packet_ext.h" #endif - #include "../../codec/stb_image/decoder.h" #include "../../codec/spng/decoder.h" #include "../../codec/wuffs/decoder.h" @@ -126,6 +125,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc struct camu_codec_stream *stream = codec->handler.stream; AVRational time_base = stream->av.stream->time_base; pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, time_base); + // @TODO: Shift av_packet_free() to vcr. This leaks right now. #else av_packet_free_side_data(pkt); av_packet_free(&pkt); diff --git a/src/liana/list.c b/src/liana/list.c index fd58ea1..2a1458f 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -51,11 +51,17 @@ static void signal_meta(struct lia_list *list, struct lia_list_entry *entry, u8 list->callback(list->userdata, LIANA_LIST_META, entry, META_OPAQUE(meta)); } +static void entry_unload(struct lia_list *list, struct lia_list_entry *entry) +{ + list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL); +} + static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, bool *error) { u8 status; list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &status); if (status == LIANA_ENTRY_ERRORED) { + // If sequence is <0 that must mean entry is not contained in list->entries. if (sequence >= 0) { al_array_remove_at(list->entries, (u32)sequence); if (list->current > sequence) { @@ -74,6 +80,8 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e } *error = true; signal_meta(list, entry, LIANA_META_ENTRY_ERRORED); + al_assert(!al_array_contains(list->entries, entry)); + entry_unload(list, entry); return false; } *error = false; @@ -84,11 +92,6 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e return false; } -static void entry_unload(struct lia_list *list, struct lia_list_entry *entry) -{ - list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL); -} - static void entry_ref(struct lia_list *list, struct lia_list_entry *entry) { list->callback(list->userdata, LIANA_REF_ENTRY, entry, NULL); diff --git a/src/liana/server.c b/src/liana/server.c index cc4ceed..d07cb8d 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -54,6 +54,7 @@ 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; + nn_thread_set_name("liana_handler"); if (conn->seek_pos != LIANA_TIMESTAMP_INVALID) { conn->handler->seek(conn->handler, conn->seek_pos); @@ -107,6 +108,10 @@ static void free_node(struct lia_node *node) { struct lia_server *server = node->server; cch_entry_free(&node->entry); + al_assert(!node->requests.count); + al_array_free(node->requests); + al_assert(!node->connections.count); + al_array_free(node->connections); al_array_remove(server->nodes, node); al_free(node); } @@ -298,6 +303,7 @@ static void signal_callback(void *userdata) static nn_thread_result NNWT_THREADCALL init_thread(void *userdata) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + nn_thread_set_name("liana_node_init"); if (!conn->handler->init(conn->handler, &conn->handle)) { conn->errored = true; } @@ -427,6 +433,7 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en static nn_thread_result NNWT_THREADCALL init_duration_thread(void *userdata) { struct lia_node *node = (struct lia_node *)userdata; + nn_thread_set_name("liana_init_dur"); if (!node->handler->init(node->handler, &node->handle)) { node->errored = true; } else { @@ -496,14 +503,8 @@ void lia_server_close(struct lia_server *server) void lia_server_free(struct lia_server *server) { - struct lia_node *node; - al_array_foreach(server->nodes, i, node) { - al_assert(!node->requests.count); - al_array_free(node->requests); - al_assert(!node->connections.count); - al_array_free(node->connections); - al_free(node); - } + // Assuming we joined on the event loop, server->nodes should be empty. + al_assert(!server->nodes.count); al_array_free(server->nodes); al_assert(!server->dormant_connections.count); al_array_free(server->dormant_connections); diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 270bc38..228d758 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -35,9 +35,10 @@ static void reset_metrics(struct lia_vcr *vcr) vcr->metric.last_report_ts = 0; } -void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data) +void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data, u16 node_id) { vcr->data = data; + vcr->node_id = node_id; al_array_init(vcr->tracks); al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); vcr->mark.buffered = VCR_BUFFER_BUFFERED; @@ -63,9 +64,11 @@ static void return_entire_cache(struct lia_vcr_track *track) static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) { nn_thread_set_priority(NNWT_THREAD_SCHED_FIFO, 32); - struct lia_vcr_track *track = (struct lia_vcr_track *)userdata; struct lia_vcr *vcr = track->vcr; + const char thread_name[16]; // 16 = limit. + al_snprintf((char *)thread_name, sizeof(thread_name), "vcr:%hu", vcr->node_id); + nn_thread_set_name(thread_name); s32 state; bool corked; diff --git a/src/liana/vcr.h b/src/liana/vcr.h index 94d0c47..167ce7f 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -25,6 +25,7 @@ struct lia_vcr_track { struct lia_vcr { struct nn_packet_stream *data; + u16 node_id; array(struct lia_vcr_track *) tracks; atomic(u64) count; struct { u64 buffered, low; } mark; @@ -39,7 +40,7 @@ struct lia_vcr { } metric; }; -void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data); +void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data, u16 node_id); void lia_vcr_start(struct lia_vcr *vcr); void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track); bool lia_vcr_is_empty(struct lia_vcr *vcr); diff --git a/src/libsink/input_simulator.c b/src/libsink/input_simulator.c index 80fde02..871f0ce 100644 --- a/src/libsink/input_simulator.c +++ b/src/libsink/input_simulator.c @@ -21,8 +21,9 @@ enum { static nn_thread_result NNWT_THREADCALL input_simulation_thread(void *userdata) { struct camu_sink *sink = (struct camu_sink *)userdata; + nn_thread_set_name("input_simulator"); while (!quit) { - nn_thread_sleep(NNWT_TS_FROM_USEC(30000)); + nn_thread_sleep(NNWT_TS_FROM_USEC(60000)); switch (al_rand() % MARK) { case SKIP: camu_sink_skip(sink, (al_rand() % 5)); diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 29af2f5..94ee3c6 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -204,7 +204,7 @@ static void remove_entry_video_buffer(struct camu_sink_entry *entry) // It's possible for some of an entry's buffers to be ENDED while others are still ADDED. // This means entry->ended and BUFFER_ENDED have two distinct considerations. -// entry->ended: Completely ignored and needs special consideration in CLIENT_REMOVE_BUFFERS. +// entry->ended: Completely ignored and needs special handling in CLIENT_REMOVE_BUFFERS. // BUFFER_ENDED: No-op'd in remove_entry_buffers() and add_or_queue_entry() but otherwise unchanged. static void remove_entry_buffers(struct camu_sink_entry *entry) { @@ -279,19 +279,15 @@ static void sink_local_pause(struct camu_sink *sink, struct camu_sink_entry *ent } #endif -static inline s32 get_sequence_for_command(struct camu_sink *sink) +static inline s32 get_sequence_for_command(struct camu_sink_entry *entry) { // SEQUENCE_ANY resolves order on the server. s32 sequence = LIANA_SEQUENCE_ANY; // If we're local this could only lead to feeling like your inputs were eaten. #ifndef CAMU_SINK_LOCAL - if (sink->target && sink->target != (struct camu_sink_entry *)0xb00b) { - sequence = sink->target->sequence; - } else if (sink->current) { - sequence = sink->current->sequence; - } + if (entry) sequence = entry->sequence; #else - (void)sink; + (void)entry; #endif return sequence; } @@ -393,7 +389,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); nn_packet_write_str(packet, &sink->default_list); nn_packet_write_u8(packet, CAMU_LIST_SKIP); - nn_packet_write_s32(packet, get_sequence_for_command(sink)); + nn_packet_write_s32(packet, get_sequence_for_command((struct camu_sink_entry *)cmd->opaque)); nn_packet_write_s32(packet, (s32)cmd->value.i); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; @@ -406,7 +402,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); nn_packet_write_str(packet, &sink->default_list); nn_packet_write_u8(packet, CAMU_LIST_TOGGLE_PAUSE); - nn_packet_write_s32(packet, get_sequence_for_command(sink)); + nn_packet_write_s32(packet, get_sequence_for_command((struct camu_sink_entry *)cmd->opaque)); nn_packet_write_f64(packet, cmd->value.f); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); #endif @@ -674,7 +670,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) if (detached) { al_assert(detached == current); sink->detached = NULL; - log_warn("Unset detached as a substitute for remove."); + log_debug("Unset detached as a substitute for remove."); // This should only matter if detached was unconfigured. if (AUDIO_EMPTY(detached)) AUDIO_STATE(detached) = BUFFER_INIT; if (VIDEO_EMPTY(detached)) VIDEO_STATE(detached) = BUFFER_INIT; @@ -690,7 +686,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) remove_previous_if_contains(sink, target); add_or_queue_entry(target); } else { - if (VIDEO_IS_SINGLE_FRAME(target)) { + if (!VIDEO_EMPTY(target) && VIDEO_IS_SINGLE_FRAME(target)) { add_video_if_set_and_buffered(target); } stop_video = true; @@ -698,7 +694,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) } if (ensure_removed) { - if (VIDEO_IS_SINGLE_FRAME(current)) { + if (!VIDEO_EMPTY(current) && VIDEO_IS_SINGLE_FRAME(current)) { remove_entry_video_buffer(current); } al_assert(AUDIO_NOT_ADDED(current)); @@ -849,24 +845,25 @@ static void video_buffer_callback(void *userdata, u8 op) break; case CAMU_BUFFER_EOF: log_debug("Video EOF."); - bool single_frame = VIDEO_IS_SINGLE_FRAME(entry); - bool swapped = false; nn_mutex_lock(&sink->lock); // This buffer's state could be ADDED, SET_OR_BUFFERED, or CONFIGURED. if (!AUDIO_EMPTY(entry)) { camu_audio_buffer_set_no_video(&entry->audio.buf, true); } - if (!single_frame) { - if (VIDEO_STATE(entry) == BUFFER_ADDED) { - remove_entry_video_buffer(entry); - } - VIDEO_STATE(entry) = BUFFER_ENDED; - if (AUDIO_ENDED_OR_EMPTY(entry)) { - swapped = end_entry_and_advance_queue(sink, entry); - } + if (VIDEO_IS_SINGLE_FRAME(entry)) { + nn_mutex_unlock(&sink->lock); + return; + } + if (VIDEO_STATE(entry) == BUFFER_ADDED) { + remove_entry_video_buffer(entry); + } + VIDEO_STATE(entry) = BUFFER_ENDED; + bool swapped = false; + if (AUDIO_ENDED_OR_EMPTY(entry)) { + swapped = end_entry_and_advance_queue(sink, entry); } nn_mutex_unlock(&sink->lock); - if (!swapped && !single_frame) { + if (!swapped) { queue_cmd(sink, (struct camu_sink_cmd){ .op = STOP, .value.i = CAMU_SINK_VIDEO @@ -907,8 +904,8 @@ static void clock_callback(void *userdata, u8 op) static void evaluate_and_set_buffer_params(struct camu_sink *sink, struct camu_sink_entry *entry) { -#ifdef CAMU_SINK_LOCAL bool ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry); +#ifdef CAMU_SINK_LOCAL if (!AUDIO_EMPTY(entry) && !ignore_video) { f64 audio = camu_mixer_get_latency(sink->audio.mixer); s32 frames = audio / entry->video.buf.avg_frame_duration; @@ -923,13 +920,12 @@ static void evaluate_and_set_buffer_params(struct camu_sink *sink, struct camu_s // To sync clients with differing audio latencies our only option is to factor the mixer // latency directly into the audio buffer. f64 audio = camu_mixer_get_latency(sink->audio.mixer); - if (!VIDEO_EMPTY(entry)) { + if (!ignore_video) { s32 frames = audio / entry->video.buf.avg_frame_duration; struct camu_renderer *renderer = sink->video.renderer; if (renderer) frames += renderer->get_latency(renderer); camu_video_buffer_set_latency(&entry->video.buf, frames); } - bool ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry); camu_audio_buffer_set_latency(&entry->audio.buf, audio); camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video); #endif @@ -1051,6 +1047,13 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str nn_mutex_lock(&sink->lock); log_trace("remove_buffers("ENTRY_FMT", %s, %s), entry == current: %s.", ENTRY_ARG(entry), BOOLSTR(rec->reconnect), BOOLSTR(rec->unconfigured), BOOLSTR(entry == sink->current)); + // Immediately switch to a potential target to avoid excessive delay/catchup + // that would be caused by this entry being re-added with an in-between clock state. + if (sink->target) { + switch_to(sink, sink->target); + sink->target = NULL; + } + // This entry might be in previous if it was added to previous then, // 1. it's being cleaned up after ENTRY_MAX_AGE - 1 entries were added but none buffered. // 2. it was seeked. @@ -1077,8 +1080,8 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str remove_entry_audio_buffer(entry); } // Slight optimization. A duplicate frame will still be sent but discarded in the video buffer. - bool ignore_video = rec->reconnect && VIDEO_IS_SINGLE_FRAME(entry); - if (!ignore_video && !VIDEO_ENDED_OR_EMPTY(entry)) { + bool skip_video = rec->reconnect && VIDEO_IS_SINGLE_FRAME(entry); + if (!skip_video && !VIDEO_ENDED_OR_EMPTY(entry)) { remove_entry_video_buffer(entry); remove_entry_video_buffer(entry); } @@ -1090,7 +1093,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str // Unlock to wait. nn_mutex_unlock(&sink->lock); - while (entry_audio_buffer_held(entry) || (!ignore_video && entry_video_buffer_held(entry))) { + while (entry_audio_buffer_held(entry) || (!skip_video && entry_video_buffer_held(entry))) { BLOCKING_SLEEP(NNWT_TS_FROM_USEC(2000)); } @@ -1107,7 +1110,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str if (VIDEO_STATE(entry) == BUFFER_ENDED) { al_assert(!VIDEO_EMPTY(entry)); VIDEO_STATE(entry) = BUFFER_CONFIGURED; - } else if (entry->ended && !ignore_video && VIDEO_EMPTY(entry)) { + } else if (entry->ended && VIDEO_EMPTY(entry)) { VIDEO_STATE(entry) = BUFFER_INIT; } entry->ended = false; @@ -1118,22 +1121,25 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } case LIANA_CLIENT_RESUME_AT: { struct lia_timing *time = (struct lia_timing *)opaque; - log_trace("resume_at("ENTRY_FMT"), paused_at: %f.", ENTRY_ARG(entry), entry->clock.paused_at); + log_trace("resume_at("ENTRY_FMT"), seek_pos: %f, paused_at: %f.", ENTRY_ARG(entry), time->seek_pos / 1000000.0, entry->clock.paused_at); // These buffers won't be re-added until after a CLIENT_RECONNECTED event. + bool ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry); if (!AUDIO_EMPTY(entry)) { camu_audio_buffer_reset(&entry->audio.buf); + // no_video is set in video BUFFER_EOF as a fail-safe. Reset it here. + camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video); } - if (!VIDEO_EMPTY(entry) && !VIDEO_IS_SINGLE_FRAME(entry)) { + if (!ignore_video) { camu_video_buffer_reset(&entry->video.buf, time->seek_pos); } -#ifdef CAMU_SINK_LOCAL - time->at = 0; -#endif nn_mutex_lock(&sink->lock); #ifdef LIANA_LIST_SCUFFED_LOOP if (time->seek_pos == 0) { camu_clock_loop(&entry->clock, camu_clock_get_last_pts(&entry->clock)); } else { +#endif +#ifdef CAMU_SINK_LOCAL + time->at = 0; #endif camu_clock_seek(&entry->clock, time->seek_pos / 1000000.0, time->at); #ifdef LIANA_LIST_SCUFFED_LOOP @@ -1198,7 +1204,6 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str while (entry_video_buffer_held(entry)) { BLOCKING_SLEEP(NNWT_TS_FROM_USEC(2000)); } } - // @TODO: Mark for removal here instead. bool removed = al_array_remove(sink->entries, entry); remove_from_queue_by_opaque(sink, entry); @@ -1525,7 +1530,6 @@ static struct nn_rpc_command commands[] = { { .op = CAMU_SINK_SEEK, .callback = seek_command_callback, .userdata = NULL } }; -// @TODO: Send sink id from list. static void identify_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet) { struct camu_sink *sink = (struct camu_sink *)userdata; @@ -1607,18 +1611,30 @@ void camu_sink_return_current(struct camu_sink *sink) nn_mutex_unlock(&sink->lock); } +static inline struct camu_sink_entry *get_entry_for_command(struct camu_sink *sink) +{ + if (sink->target && sink->target != (struct camu_sink_entry *)0xb00b) { + return sink->target; + } + return sink->current; +} + void camu_sink_skip(struct camu_sink *sink, s32 n) { + nn_mutex_lock(&sink->lock); + struct camu_sink_entry *current = get_entry_for_command(sink); + nn_mutex_unlock(&sink->lock); queue_cmd(sink, (struct camu_sink_cmd){ .op = SKIP, - .value.i = n + .value.i = n, + .opaque = current }); } void camu_sink_toggle_pause(struct camu_sink *sink) { nn_mutex_lock(&sink->lock); - struct camu_sink_entry *current = sink->current; + struct camu_sink_entry *current = get_entry_for_command(sink); nn_mutex_unlock(&sink->lock); if (!current) return; queue_cmd(sink, (struct camu_sink_cmd){ @@ -1631,7 +1647,7 @@ void camu_sink_toggle_pause(struct camu_sink *sink) void camu_sink_seek(struct camu_sink *sink, void *value, u8 mode) { nn_mutex_lock(&sink->lock); - struct camu_sink_entry *current = sink->current; + struct camu_sink_entry *current = get_entry_for_command(sink); u64 duration = 0; f64 pts = 0.0; if (current) { @@ -1713,7 +1729,7 @@ void camu_sink_close(struct camu_sink *sink) void camu_sink_free(struct camu_sink *sink) { - al_assert(sink->entries.count == 0); + al_assert(!sink->entries.count); al_array_free(sink->entries); nn_rpc_free(&sink->client); camu_queue_free(sink->queue); diff --git a/src/libsink/sink.h b/src/libsink/sink.h index c2ea6a1..40bd0a4 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -13,6 +13,7 @@ #ifndef CAMU_SINK_NO_VIDEO #include "../buffer/video.h" #else +// Stub video_buffer usage within libsink only. #include "../buffer/video_null.h" #endif diff --git a/src/portal/scripts/create_vendor.sh b/src/portal/scripts/create_vendor.sh index 3dd8bc3..6af8598 100755 --- a/src/portal/scripts/create_vendor.sh +++ b/src/portal/scripts/create_vendor.sh @@ -1,5 +1,7 @@ #! /usr/bin/env sh +set -e + python3 -m venv ./venv source ./venv/bin/activate diff --git a/src/portal/scripts/fetch_python_deps.sh b/src/portal/scripts/fetch_python_deps.sh new file mode 100755 index 0000000..2f8e098 --- /dev/null +++ b/src/portal/scripts/fetch_python_deps.sh @@ -0,0 +1,8 @@ +#! /usr/bin/env sh + +git clone https://github.com/yt-dlp/yt-dlp.git --depth=1 yt-dlp +git clone https://github.com/JustAnotherArchivist/snscrape.git --depth=1 snscrape +git clone https://github.com/trevorhobenshield/twitter-api-client.git --depth=1 twitter-api-client +git clone https://github.com/upbit/pixivpy.git --depth=1 pixivpy +git clone https://github.com/subzeroid/instagrapi.git --depth=1 instagrapi +git clone https://github.com/mikf/gallery-dl.git --depth=1 gallery-dl diff --git a/src/portal/src/search.c b/src/portal/src/search.c index 81d1e72..a1b7310 100644 --- a/src/portal/src/search.c +++ b/src/portal/src/search.c @@ -62,6 +62,7 @@ static nn_thread_result NNWT_THREADCALL queue_thread(void *userdata) { nn_thread_setcanceltype(NNWT_THREAD_CANCEL_ASYNCHRONOUS); struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; + nn_thread_set_name("portal_queue"); bool have_python = false; nn_mutex_lock(&bridge->mutex); diff --git a/src/render/meson.build b/src/render/meson.build index 557c459..b7927bb 100644 --- a/src/render/meson.build +++ b/src/render/meson.build @@ -19,7 +19,11 @@ elif get_option('renderer') == 'libplacebo' ] if get_option('renderer-api') == 'vulkan' or get_option('renderer-api') == 'dx11' - shaderc = dependency('shaderc', required: false, allow_fallback: false) + if is_windows + shaderc = dependency('shaderc_combined', required: false, allow_fallback: false) + else + shaderc = dependency('shaderc', required: false, allow_fallback: false) + endif if shaderc.found() render_deps += [shaderc] else diff --git a/src/screen/screen.c b/src/screen/screen.c index bb11757..09fbe64 100644 --- a/src/screen/screen.c +++ b/src/screen/screen.c @@ -9,6 +9,7 @@ #define SCREEN_MOD1(scr) ((scr)->flags & CAMU_SCREEN_MOD_CONTROL) #define SCREEN_MOD2(scr) ((scr)->flags & CAMU_SCREEN_MOD_SHIFT) #define SCREEN_IS_DRAGGING(scr) ((scr)->flags & CAMU_SCREEN_DRAGGING) +#define SCREEN_IS_ZOOMING(scr) ((scr)->flags & CAMU_SCREEN_ZOOMING) #define SCREEN_ZOOM_MODE(scr, mode) ((scr)->flags & mode) #define SCREEN_INVALID_TS ((u64)-1) #define SCREEN_LAST_CLICK_WITHIN(ns) \ @@ -29,15 +30,6 @@ static void render_callback(void *userdata) } } -static void refresh_callback(void *userdata) -{ - struct camu_screen *scr = (struct camu_screen *)userdata; - if (scr->renderer) { - // The thread-safety of this is questionable. - scr->renderer->render(scr->renderer, scr, true); - } -} - static void do_resize(struct camu_screen *scr, u32 width, u32 height) { if (scr->renderer) { @@ -76,7 +68,6 @@ static struct camu_view *get_view_from_mouse_pos(struct camu_screen *scr) return view; } - static bool pointer_pos_callback(void *userdata, f64 x, f64 y) { struct camu_screen *scr = (struct camu_screen *)userdata; @@ -84,8 +75,8 @@ static bool pointer_pos_callback(void *userdata, f64 x, f64 y) bool queue_refresh = false; if (SCREEN_IS_DRAGGING(scr)) { if (view && SCREEN_ZOOM_MODE(scr, CAMU_SCREEN_ZOOM_PAN_SIMPLE)) { - f64 dx = x - scr->last_mouse_x; - f64 dy = y - scr->last_mouse_y; + f64 dx = x - scr->last_pointer_x; + f64 dy = y - scr->last_pointer_y; if (camu_view_pan_simple(view, scr->width, scr->height, dx, dy)) { view->mode = CAMU_VIEW_DETACHED; queue_refresh = true; @@ -95,14 +86,67 @@ static bool pointer_pos_callback(void *userdata, f64 x, f64 y) } } } - scr->last_mouse_x = x; - scr->last_mouse_y = y; + scr->last_pointer_x = x; + scr->last_pointer_y = y; return queue_refresh; } +static bool scroll_callback(void *userdata, f64 y); + +static bool touch_callback(void *userdata, s32 index, u8 phase, f64 x, f64 y) +{ + struct camu_screen *scr = (struct camu_screen *)userdata; + // Rough testing stuff. + switch (phase) { + case STELA_TOUCH_BEGAN: { + if (index == 0) { + scr->flags |= CAMU_SCREEN_DRAGGING; + scr->last_click_ts = nn_get_timestamp(); + scr->last_pointer_x = x; + scr->last_pointer_y = y; + } else if (index == 1) { + scr->flags &= ~CAMU_SCREEN_DRAGGING; + scr->flags |= CAMU_SCREEN_ZOOMING; + } else if (index == 2) { + struct camu_view *view = get_view_from_mouse_pos(scr); + view->mode = CAMU_DEFAULT_VIEW; + camu_view_calculate(view, scr->width, scr->height); + } + break; + } + case STELA_TOUCH_MOVED: { + if (index == 0) { + if (SCREEN_IS_ZOOMING(scr)) { + f64 dy = (y - scr->last_pointer_y) / 50.0; + scr->last_pointer_x = x; + scr->last_pointer_y = y; + return scroll_callback(scr, dy); + } else if (SCREEN_IS_DRAGGING(scr)) { + return pointer_pos_callback(scr, x, y); + } + } + break; + } + case STELA_TOUCH_ENDED: { + if (index == 0) { + scr->flags &= ~(CAMU_SCREEN_DRAGGING | CAMU_SCREEN_ZOOMING); + if (SCREEN_LAST_CLICK_WITHIN(200000)) { + s32 n = (x >= scr->width / 2.0) ? 1 : -1; + scr->callback(scr->userdata, CAMU_SCREEN_SKIP, &n); + } + } else if (index == 1) { + scr->flags &= ~CAMU_SCREEN_ZOOMING; + scr->flags |= CAMU_SCREEN_DRAGGING; + } + break; + } + } + return false; +} + static void seek_to_percent_at_pointer(struct camu_screen *scr) { - f64 percent = scr->last_mouse_x / scr->width; + f64 percent = scr->last_pointer_x / scr->width; percent = CLAMP(percent, 0.0, 100.0); scr->callback(scr->userdata, CAMU_SCREEN_PERCENT_SEEK, &percent); } @@ -123,7 +167,7 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button) if (SCREEN_MOD1(scr)) { seek_to_percent_at_pointer(scr); } else { - s32 n = (scr->last_mouse_x >= scr->width / 2.0) ? 1 : -1; + s32 n = (scr->last_pointer_x >= scr->width / 2.0) ? 1 : -1; scr->callback(scr->userdata, CAMU_SCREEN_SKIP, &n); } } @@ -145,7 +189,7 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button) return false; } -static bool scroll_callback(void *userdata, f64 y) +bool scroll_callback(void *userdata, f64 y) { struct camu_screen *scr = (struct camu_screen *)userdata; if (SCREEN_MOD1(scr)) { @@ -162,7 +206,7 @@ static bool scroll_callback(void *userdata, f64 y) #endif struct camu_view *view = get_view_from_mouse_pos(scr); if (view && SCREEN_ZOOM_MODE(scr, CAMU_SCREEN_ZOOM_PAN_SIMPLE)) { - if (camu_view_zoom_simple(view, scr->width, scr->height, scr->last_mouse_x, scr->last_mouse_y, y)) { + if (camu_view_zoom_simple(view, scr->width, scr->height, scr->last_pointer_x, scr->last_pointer_y, y)) { view->mode = CAMU_VIEW_DETACHED; return true; } @@ -370,21 +414,21 @@ bool camu_screen_init(struct camu_screen *scr) al_atomic_store(s32)(&scr->state, CAMU_SCREEN_PAUSED, AL_ATOMIC_RELAXED); scr->window = stl_window_create(); scr->window->render_callback = render_callback; - scr->window->refresh_callback = refresh_callback; + scr->window->resize_callback = resize_callback; scr->window->pointer_pos_callback = pointer_pos_callback; + scr->window->touch_callback = touch_callback; scr->window->scroll_callback = scroll_callback; scr->window->mouse_button_callback = mouse_button_callback; scr->window->key_callback = key_callback; scr->window->key_immediate_callback = key_immediate_callback; - scr->window->resize_callback = resize_callback; scr->window->should_close_callback = should_close_callback; scr->window->userdata = scr; scr->renderer = NULL; al_atomic_store(u32)(&scr->force_refresh, 1, AL_ATOMIC_RELAXED); scr->flags = CAMU_SCREEN_ZOOM_PAN_SIMPLE; scr->last_click_ts = SCREEN_INVALID_TS; - scr->last_mouse_x = 0.0; - scr->last_mouse_y = 0.0; + scr->last_pointer_x = 0.0; + scr->last_pointer_y = 0.0; al_array_init(scr->videos); #ifdef CAMU_SCREEN_THREADED al_array_init(scr->add_queue); @@ -399,6 +443,7 @@ bool camu_screen_init(struct camu_screen *scr) static nn_thread_result NNWT_THREADCALL event_thread(void *userdata) { struct camu_screen *scr = (struct camu_screen *)userdata; + nn_thread_set_name("window_events"); while (scr->window->poll(scr->window, true) && al_atomic_load(s32)(&scr->state, AL_ATOMIC_RELAXED) != CAMU_SCREEN_STOPPED) { scr->window->process_events(scr->window); } @@ -613,7 +658,7 @@ void camu_screen_force_refresh(struct camu_screen *scr) camu_screen_wake(scr); } -#define CAMU_SCREEN_POLL_HZ 576 +#define CAMU_SCREEN_POLL_HZ 432 bool camu_screen_poll(struct camu_screen *scr, bool block) { diff --git a/src/screen/screen.h b/src/screen/screen.h index afd2dee..dacf756 100644 --- a/src/screen/screen.h +++ b/src/screen/screen.h @@ -21,7 +21,8 @@ enum { CAMU_SCREEN_MOD_SHIFT = 1, CAMU_SCREEN_MOD_CONTROL = 1 << 1, CAMU_SCREEN_DRAGGING = 1 << 2, - CAMU_SCREEN_ZOOM_PAN_SIMPLE = 1 << 3 + CAMU_SCREEN_ZOOMING = 1 << 3, + CAMU_SCREEN_ZOOM_PAN_SIMPLE = 1 << 4 }; enum { @@ -66,8 +67,8 @@ struct camu_screen { bool scaling_disabled; bool transparent_background; u64 last_click_ts; - f64 last_mouse_y; - f64 last_mouse_x; + f64 last_pointer_y; + f64 last_pointer_x; array(struct camu_screen_video) videos; #ifdef CAMU_SCREEN_THREADED array(struct camu_video_buffer *) add_queue; diff --git a/src/server/server.c b/src/server/server.c index ec99921..db55ae2 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -115,6 +115,8 @@ static bool identify_callback(void *userdata, struct nn_rpc_connection *conn, { struct camu_server *server = (struct camu_server *)userdata; + // @TODO: Return server-side IDs to the clients. + // This is important, for example, to identify which sink a list command is coming from. u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_NODE: { diff --git a/src/util/print_time.h b/src/util/print_time.h index d9dd4cc..89644b2 100644 --- a/src/util/print_time.h +++ b/src/util/print_time.h @@ -11,10 +11,9 @@ static inline s32 camu_print_time(char *buf, size_t size, u64 nanoseconds, bool seconds -= minutes * 60u; minutes -= hours * 60u; s32 offset = 0; - if (show_hour) { + if (show_hour && hour_padding >= 2 && hour_padding <= 4) { static const char *padding_table[] = { "%.2u:", "%.3u:", "%.4u:" }; - hour_padding = CLAMP(hour_padding, 2u, 4u) - 2u; - offset += al_snprintf(buf, size, padding_table[hour_padding], hours); + offset += al_snprintf(buf, size, padding_table[hour_padding - 2u], hours); } offset += al_snprintf(buf + offset, size - offset, "%.2u:%.2u", minutes, seconds); return offset; -- cgit v1.2.3-101-g0448