diff options
| author | 2025-09-08 18:20:00 -0400 | |
|---|---|---|
| committer | 2025-09-08 18:20:00 -0400 | |
| commit | ac3fd1a688202375612e905da9ee63bd6cd36a6f (patch) | |
| tree | 765f200b7a18a7483698751ea5227647656a4e56 | |
| parent | 38c371d1fa82101d59eb202d2f4aa1b2b99c0e42 (diff) | |
| download | camu-ac3fd1a688202375612e905da9ee63bd6cd36a6f.tar.gz camu-ac3fd1a688202375612e905da9ee63bd6cd36a6f.tar.bz2 camu-ac3fd1a688202375612e905da9ee63bd6cd36a6f.zip | |
Sink fixes
- Invalidate sink entry IDs on reconnect.
- Fix A/V latency calculations.
- Fix reconnect_timer behavior and log reconnect attempts.
Signed-off-by: Andrew Opalach <andrew@akon.city>
| -rw-r--r-- | src/buffer/video.c | 8 | ||||
| -rw-r--r-- | src/buffer/video.h | 2 | ||||
| -rw-r--r-- | src/fruits/cmv/cmv.c | 9 | ||||
| -rw-r--r-- | src/liana/vcr.c | 10 | ||||
| -rw-r--r-- | src/libsink/sink.c | 143 | ||||
| -rw-r--r-- | src/libsink/sink.h | 3 | ||||
| -rw-r--r-- | subprojects/libnaunet.wrap | 4 |
7 files changed, 103 insertions, 76 deletions
diff --git a/src/buffer/video.c b/src/buffer/video.c index 7874615..0e9ce28 100644 --- a/src/buffer/video.c +++ b/src/buffer/video.c @@ -15,8 +15,8 @@ // quite incompatible with very low frame rates. // MARK_LOW is considered directly after reading a frame, so in other words, // it will trigger at the point where there is about (LOW+1) frames left. Even at -// something like 240fps that's still ~(4*(LOW+1))ms to uncork and produce a new frame. -#define BUFFER_MARK_LOW 3 +// something like 240fps that's still ~(4.16*(LOW+1))ms to uncork and produce a new frame. +#define BUFFER_MARK_LOW 4 #define BUFFER_MARK_BUFFERED 6 // Must be >1. #define BUFFER_MARK_HIGH 8 #define BUFFER_MARK_RESET (BUFFER_MARK_HIGH * 2) @@ -112,9 +112,9 @@ bool camu_video_buffer_configure_subtitles(struct camu_video_buffer *buf, struct return buf->queue->configure_subtitles(buf->queue, fmt->width, fmt->height, stream); } -void camu_video_buffer_set_latency(struct camu_video_buffer *buf, s32 frames) +void camu_video_buffer_set_latency(struct camu_video_buffer *buf, f64 latency) { - buf->latency = frames * buf->avg_frame_duration; + buf->latency = latency; } static void after_push_internal(struct camu_video_buffer *buf) diff --git a/src/buffer/video.h b/src/buffer/video.h index d844d80..380943e 100644 --- a/src/buffer/video.h +++ b/src/buffer/video.h @@ -49,7 +49,7 @@ bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *cl bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_codec_stream *stream, struct camu_renderer *renderer); bool camu_video_buffer_configure_subtitles(struct camu_video_buffer *buf, struct camu_codec_stream *stream); -void camu_video_buffer_set_latency(struct camu_video_buffer *buf, s32 frames); +void camu_video_buffer_set_latency(struct camu_video_buffer *buf, f64 latency); void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_frame *frame); void camu_video_buffer_push_subtitle(struct camu_video_buffer *buf, struct camu_codec_packet *packet); void camu_video_buffer_flush(struct camu_video_buffer *buf); diff --git a/src/fruits/cmv/cmv.c b/src/fruits/cmv/cmv.c index 63e035d..840eb6a 100644 --- a/src/fruits/cmv/cmv.c +++ b/src/fruits/cmv/cmv.c @@ -220,10 +220,17 @@ s32 window_system_main(u32 argc, str *argv, void *extra) return EXIT_SUCCESS; } +static bool sigint_force = false; static void sigint_handler(int signum) { (void)signum; - c.desktop.should_quit = 1; + if (sigint_force) { + exit(128+SIGINT); + } else { + c.desktop.should_quit = 1; + log_warn("Attempting graceful exit. ^C again to force exit."); + sigint_force = true; + } } #ifdef NAUNET_ON_WINDOWS diff --git a/src/liana/vcr.c b/src/liana/vcr.c index d2c6b33..d8a1fc9 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -5,7 +5,7 @@ #include "vcr.h" -#define VCR_BUFFER_BUFFERED MB(2) +#define VCR_BUFFER_BUFFERED MB(4) enum { VCR_EXPAND_UNTOUCHED = 0, @@ -68,7 +68,7 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) struct lia_vcr_track *track = (struct lia_vcr_track *)userdata; struct lia_vcr *vcr = track->vcr; const char thread_name[16] = "\0"; // 16 = limit. - al_snprintf((char *)thread_name, sizeof(thread_name), "vcr:%hu_%hu", vcr->node_id, track->stream->index); + al_snprintf((char *)thread_name, sizeof(thread_name), "vcr:%hu_%d", vcr->node_id, track->stream->index); nn_thread_set_name(thread_name); s32 state; @@ -274,9 +274,9 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) case LIANA_PACKET_DATA: { u32 size = nn_packet_get_size(packet); update_metrics(vcr, size); - track = get_track_from_index(vcr, nn_packet_read_s32(packet)); - if (!track) { - log_debug("Received data from errored or unknown track."); + s32 index = nn_packet_read_s32(packet); + if (!(track = get_track_from_index(vcr, index))) { + log_debug("Received data from errored or unknown track (index: %d).", index); break; } if (VCR_TRACK_THREADED(track)) { diff --git a/src/libsink/sink.c b/src/libsink/sink.c index ff91a51..f410636 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -59,15 +59,21 @@ enum { // Number of entries to keep buffered at one time. #define ENTRY_MAX_AGE 4 +#define ENTRY_IS_VALID(entry) ((entry) && (entry) != (struct camu_sink_entry *)0xb00b) + +// Store connection number in the upper 16 bits so IDs don't conflict after a server restart. +// If the sink disconnects but the server didn't restart, this will invalidate IDs that _do_ map +// to the same resource. So it's a trade off. +#define LOCAL_ENTRY_ID(sink, id) (((u64)(sink)->connection_number) << 47 | (u64)(id)) +#define REMOTE_ENTRY_ID(id) ((u32)(id & 0x7fffffff)) + // printf format for entries. #ifdef AL_DEBUG #define ENTRY_FMT "#%u(%p)" -#define ENTRY_ARG(entry) \ - ((entry) && (entry) != (struct camu_sink_entry *)0xb00b) ? (entry)->id : 0, (entry) ? (entry) : 0x0 +#define ENTRY_ARG(entry) (ENTRY_IS_VALID(entry) ? REMOTE_ENTRY_ID((entry)->id) : 0), ((entry) ? (entry) : 0x0) #else #define ENTRY_FMT "#%u" -#define ENTRY_ARG(entry) \ - ((entry) && (entry) != (struct camu_sink_entry *)0xb00b) ? (entry)->id : 0 +#define ENTRY_ARG(entry) (ENTRY_IS_VALID(entry) ? REMOTE_ENTRY_ID((entry)->id) : 0) #endif #define AUDIO_STATE(entry) ((entry)->audio.state) @@ -418,7 +424,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) nn_packet_write_u8(packet, CAMU_LIST_SEEK); struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; nn_packet_write_s32(packet, entry->sequence); - nn_packet_write_u32(packet, entry->id); + nn_packet_write_u32(packet, REMOTE_ENTRY_ID(entry->id)); nn_packet_write_u64(packet, cmd->value.u); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; @@ -442,7 +448,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) nn_packet_write_str(packet, &sink->default_list); nn_packet_write_u8(packet, CAMU_LIST_END); struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - nn_packet_write_u32(packet, entry->id); + nn_packet_write_u32(packet, REMOTE_ENTRY_ID(entry->id)); nn_packet_write_u32(packet, (u32)cmd->value.u); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; @@ -481,30 +487,6 @@ static void mixer_callback(void *userdata, u8 op) } } -bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop, - struct camu_mixer *mixer, struct camu_renderer *renderer) -{ - sink->loop = loop; - nn_mutex_init(&sink->lock); - nn_signal_init(&sink->queue_signal, sink->loop, queue_signal_callback, sink); - nn_signal_start(&sink->queue_signal); - camu_queue_init(sink->queue); - sink->current = NULL; - sink->queued = NULL; - sink->target = NULL; - sink->detached = NULL; - al_array_init(sink->previous); - al_array_init(sink->entries); - sink->lru = 0; - mixer->callback = mixer_callback; - mixer->userdata = sink; - sink->audio.state = SINK_PAUSED; - sink->audio.mixer = mixer; - sink->video.state = SINK_PAUSED; - sink->video.renderer = renderer; - return true; -} - static s32 entry_lru_compare(const void *a, const void *b) { struct camu_sink_entry *aa = *((struct camu_sink_entry **)a); @@ -908,26 +890,31 @@ static void clock_callback(void *userdata, u8 op) static void evaluate_and_set_buffer_params(struct camu_sink *sink, struct camu_sink_entry *entry) { bool ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry); + f64 avg_frame_duration = entry->video.buf.avg_frame_duration; #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; struct camu_renderer *renderer = sink->video.renderer; - if (renderer) frames -= renderer->get_latency(renderer); - camu_video_buffer_set_latency(&entry->video.buf, -frames); + f64 video = renderer->get_latency(renderer) * avg_frame_duration; + // Start either the audio or video early so we can start the clock + // as soon as possible while keeping A/V sync. + if (audio > video) { + camu_video_buffer_set_latency(&entry->video.buf, video - audio); + } else if (video > audio) { + camu_audio_buffer_set_latency(&entry->audio.buf, audio - video); + } } // If we're local we don't have to worry about syncing audio-only entries. camu_audio_buffer_set_ignore_desync(&entry->audio.buf, ignore_video); camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video); #else - // To sync clients with differing audio latencies our only option is to factor the mixer - // latency directly into the audio buffer. + // When trying to sync clients with different audio/video latencies (common case), + // our only option is to factor the latency directly into the buffers. f64 audio = camu_mixer_get_latency(sink->audio.mixer); 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); + f64 video = renderer->get_latency(renderer) * avg_frame_duration; + camu_video_buffer_set_latency(&entry->video.buf, video); } camu_audio_buffer_set_latency(&entry->audio.buf, audio); camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video); @@ -1253,7 +1240,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } } -static struct camu_sink_entry *create_entry(struct camu_sink *sink, u32 id) +static struct camu_sink_entry *create_entry(struct camu_sink *sink, u64 id) { struct camu_sink_entry *entry = al_alloc_object(struct camu_sink_entry); entry->sink = sink; @@ -1279,13 +1266,15 @@ static struct camu_sink_entry *create_entry(struct camu_sink *sink, u32 id) camu_video_buffer_init(&entry->video.buf, &entry->clock); entry->video.buf.callback = video_buffer_callback; entry->video.buf.userdata = entry; + union { f64 f; u64 u; } fv = { .u = id }; + entry->video.buf.reset_pts = fv.f; al_array_push(sink->entries, entry); return entry; } -static struct camu_sink_entry *get_entry_from_id(struct camu_sink *sink, u32 id) +static struct camu_sink_entry *get_entry_from_id(struct camu_sink *sink, u64 id) { struct camu_sink_entry *entry; al_array_foreach(sink->entries, i, entry) { @@ -1319,7 +1308,7 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn, u32 node_id = nn_packet_read_u32(packet); // List entry info. - u32 id = nn_packet_read_u32(packet); + u64 id = LOCAL_ENTRY_ID(sink, nn_packet_read_u32(packet)); s32 sequence = nn_packet_read_s32(packet); u64 at = nn_packet_read_u64(packet); u64 seek_pos = nn_packet_read_u64(packet); @@ -1457,7 +1446,7 @@ static bool pause_command_callback(void *userdata, struct nn_rpc_connection *con struct camu_sink *sink = (struct camu_sink *)userdata; (void)rpacket; - u32 id = nn_packet_read_u32(packet); + u64 id = LOCAL_ENTRY_ID(sink, nn_packet_read_u32(packet)); s32 sequence = nn_packet_read_s32(packet); u64 at = nn_packet_read_u64(packet); u8 pause = nn_packet_read_u8(packet); @@ -1511,7 +1500,7 @@ static bool seek_command_callback(void *userdata, struct nn_rpc_connection *conn struct camu_sink *sink = (struct camu_sink *)userdata; (void)rpacket; - u32 id = nn_packet_read_u32(packet); + u64 id = LOCAL_ENTRY_ID(sink, nn_packet_read_u32(packet)); s32 sequence = nn_packet_read_s32(packet); (void)sequence; u64 at = nn_packet_read_u64(packet); @@ -1553,9 +1542,8 @@ static void identify_callback(void *userdata, struct nn_rpc_connection *conn, st nn_packet_stream_return_packet(conn->stream, packet); } -static void identify_on_connection(struct camu_sink *sink, struct nn_rpc_connection *conn) +static void identify_on_connection(struct camu_sink *sink) { - sink->conn = conn; struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_IDENTIFY); nn_packet_write_u8(packet, CAMU_SINK); nn_packet_write_str(packet, &sink->name); @@ -1565,25 +1553,64 @@ static void identify_on_connection(struct camu_sink *sink, struct nn_rpc_connect static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; - identify_on_connection(sink, conn); - nn_timer_stop(&sink->reconnect_timer); + sink->conn = conn; + sink->connection_number++; + identify_on_connection(sink); } static void reconnect_timer_callback(void *userdata, struct nn_timer *timer) { struct camu_sink *sink = (struct camu_sink *)userdata; (void)timer; + nn_timer_stop(&sink->reconnect_timer); nn_rpc_reconnect(&sink->client, &sink->addr, sink->port); } static void connection_closed_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; + bool reconnect = !sink->reconnect_timer.disabled; + if (reconnect) { + if (sink->conn && sink->connection_number > 0) { + log_info("Connection to server closed, attempting reconnect..."); + } else { + log_info("Failed to connect to server, trying again..."); + } + } if (sink->conn) { al_assert(sink->conn == conn); sink->conn = NULL; } - nn_timer_again(&sink->reconnect_timer); + if (reconnect) nn_timer_again(&sink->reconnect_timer); +} + +bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop, + struct camu_mixer *mixer, struct camu_renderer *renderer) +{ + sink->loop = loop; + nn_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink); + sink->conn = NULL; + sink->connection_number = 0; + nn_mutex_init(&sink->lock); + nn_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink); + nn_timer_set_repeat(&sink->reconnect_timer, NNWT_TS_FROM_USEC(2000000)); + nn_signal_init(&sink->queue_signal, sink->loop, queue_signal_callback, sink); + nn_signal_start(&sink->queue_signal); + camu_queue_init(sink->queue); + sink->current = NULL; + sink->queued = NULL; + sink->target = NULL; + sink->detached = NULL; + al_array_init(sink->previous); + al_array_init(sink->entries); + sink->lru = 0; + mixer->callback = mixer_callback; + mixer->userdata = sink; + sink->audio.state = SINK_PAUSED; + sink->audio.mixer = mixer; + sink->video.state = SINK_PAUSED; + sink->video.renderer = renderer; + return true; } bool camu_sink_connect(struct camu_sink *sink, str *name, u8 type, str *addr, u16 port) @@ -1592,20 +1619,14 @@ bool camu_sink_connect(struct camu_sink *sink, str *name, u8 type, str *addr, u1 sink->type = type; al_str_clone(&sink->addr, addr); sink->port = port; - nn_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink); + sink->conn = nn_rpc_prepare_client(&sink->client); + al_assert(sink->callback); for (u32 i = 0; i < ARRAY_SIZE(commands); i++) { commands[i].userdata = sink; - al_assert(sink->callback); nn_rpc_add_command(&sink->client, &commands[i]); } - nn_rpc_prepare_client(&sink->client); - // This client may never connect, make sure conn is still set in that case. - sink->conn = sink->client.conn; - nn_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink); - nn_timer_set_repeat(&sink->reconnect_timer, NNWT_TS_FROM_USEC(1000000)); #ifdef CAMU_DIRECT_MODE - struct nn_rpc_connection *conn = sink->client.conn; - nn_multiplex_direct_connect(conn->stream, CAMU_MULTIPLEX_RPC); + nn_multiplex_direct_connect(sink->conn->stream, CAMU_MULTIPLEX_RPC); #else nn_rpc_connect(&sink->client, CAMU_MULTIPLEX_RPC, sink->type, &sink->addr, sink->port); #endif @@ -1625,10 +1646,7 @@ void camu_sink_return_current(struct camu_sink *sink) 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; + return ENTRY_IS_VALID(sink->target) ? sink->target : sink->current; } void camu_sink_skip(struct camu_sink *sink, s32 n) @@ -1685,6 +1703,7 @@ void camu_sink_seek(struct camu_sink *sink, void *value, u8 mode) break; } case CAMU_SEEK_PERCENT: { + // There's probably a way to lose less precision here. f64 percent = *(f64 *)value; cmd.value.u = (u64)(duration * percent); break; @@ -1746,6 +1765,6 @@ void camu_sink_free(struct camu_sink *sink) nn_rpc_free(&sink->client); camu_queue_free(sink->queue); nn_mutex_destroy(&sink->lock); - al_str_free(&sink->addr); al_str_free(&sink->name); + al_str_free(&sink->addr); } diff --git a/src/libsink/sink.h b/src/libsink/sink.h index 40bd0a4..3face58 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -44,8 +44,8 @@ enum { }; struct camu_sink_entry { - u32 id; struct lia_client client; + u64 id; s32 sequence; u16 lru; bool disconnected; @@ -83,6 +83,7 @@ struct camu_sink { u16 port; struct nn_rpc client; struct nn_rpc_connection *conn; + u16 connection_number; struct nn_mutex lock; struct nn_timer reconnect_timer; struct nn_signal queue_signal; diff --git a/subprojects/libnaunet.wrap b/subprojects/libnaunet.wrap index e4f9e1d..6619f18 100644 --- a/subprojects/libnaunet.wrap +++ b/subprojects/libnaunet.wrap @@ -1,6 +1,6 @@ [wrap-git] -directory = libnaunet-4aea807 +directory = libnaunet-ed40327 url = https://git.akon.city/libnaunet.git push-url = git@git.akon.city:libnaunet.git -revision = 4aea807a167a894fd264718efc8c9fa8bced7dcb +revision = ed403279224444f9fb3bac68a1ab008f8e7dfd8c depth = 1 |