diff options
| author | 2025-09-08 18:20:00 -0400 | |
|---|---|---|
| committer | 2025-09-08 18:20:00 -0400 | |
| commit | ac3fd1a688202375612e905da9ee63bd6cd36a6f (patch) | |
| tree | 765f200b7a18a7483698751ea5227647656a4e56 /src/libsink | |
| 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>
Diffstat (limited to 'src/libsink')
| -rw-r--r-- | src/libsink/sink.c | 143 | ||||
| -rw-r--r-- | src/libsink/sink.h | 3 |
2 files changed, 83 insertions, 63 deletions
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; |