summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-09-08 18:20:00 -0400
committerAndrew Opalach <andrew@akon.city> 2025-09-08 18:20:00 -0400
commitac3fd1a688202375612e905da9ee63bd6cd36a6f (patch)
tree765f200b7a18a7483698751ea5227647656a4e56 /src/libsink
parent38c371d1fa82101d59eb202d2f4aa1b2b99c0e42 (diff)
downloadcamu-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.c143
-rw-r--r--src/libsink/sink.h3
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;