summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--src/buffer/video.c8
-rw-r--r--src/buffer/video.h2
-rw-r--r--src/fruits/cmv/cmv.c9
-rw-r--r--src/liana/vcr.c10
-rw-r--r--src/libsink/sink.c143
-rw-r--r--src/libsink/sink.h3
-rw-r--r--subprojects/libnaunet.wrap4
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