summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-10-23 18:15:07 -0400
committerAndrew Opalach <andrew@akon.city> 2024-10-23 18:15:07 -0400
commit831f260ba2f6bc89f0451f6cc628bd131913a363 (patch)
tree5db96c2782fd74f6799bd328742d21c8ffa3ca66
parent626f299dd3512d44df4bcf95fc7a21b9417c06b0 (diff)
downloadcamu-831f260ba2f6bc89f0451f6cc628bd131913a363.tar.gz
camu-831f260ba2f6bc89f0451f6cc628bd131913a363.tar.bz2
camu-831f260ba2f6bc89f0451f6cc628bd131913a363.zip
Sink-side list sync resilience, vcr fix
Signed-off-by: Andrew Opalach <andrew@akon.city>
-rw-r--r--doc/references/references.txt20
-rw-r--r--doc/style.css2
-rw-r--r--flake.nix2
-rw-r--r--meson.build2
-rw-r--r--src/buffer/audio.c18
-rw-r--r--src/buffer/audio.h1
-rw-r--r--src/liana/list.c13
-rw-r--r--src/liana/vcr.c16
-rw-r--r--src/libsink/common.h7
-rw-r--r--src/libsink/sink.c73
-rw-r--r--src/server/server.c5
11 files changed, 110 insertions, 49 deletions
diff --git a/doc/references/references.txt b/doc/references/references.txt
index fafba93..2ce8329 100644
--- a/doc/references/references.txt
+++ b/doc/references/references.txt
@@ -25,8 +25,8 @@ Latency/Present
~~~~~~~~~~~~~~~
* https://themaister.net/blog/2023/11/12/my-scuffed-game-streaming-adventure-pyrofling/
-Old Video Standards Stuff
-~~~~~~~~~~~~~~~~~~~~~~~~~
+Old Video Standards
+~~~~~~~~~~~~~~~~~~~
* https://en.wikipedia.org/wiki/24p
* https://en.wikipedia.org/wiki/Three-two_pull_down
* https://cinematography.com/index.php?/forums/topic/71346-why-23976-and-not-24-fps/&tab=comments#comment-455454
@@ -37,10 +37,6 @@ Old Video Standards Stuff
image:pal_vs_ntsc_aspect_ratio.png[PalVsNtsc,1100]
* https://en.wikipedia.org/wiki/Anamorphic_widescreen
-Python
-------
-* https://stackoverflow.com/questions/38243682/whats-the-standard-way-to-package-a-python-project-with-dependencies
-
Audio
-----
Formats
@@ -54,11 +50,19 @@ Drivers
* https://stackoverflow.com/questions/41263580/what-exactly-does-alsas-snd-pcm-delay-return
* https://stackoverflow.com/questions/24040672/the-meaning-of-period-in-alsa
-Interface
----------
+Latency
+^^^^^^^
+* https://blog.nirbheek.in/2018/04/a-simple-method-of-measuring-audio.html?m=1
+
+User Interface
+--------------
Terminal
~~~~~~~~
* https://www.amp-what.com/
* https://www.unicode.org/charts/beta/nameslist/n_2B00.html
+Python
+------
+* https://stackoverflow.com/questions/38243682/whats-the-standard-way-to-package-a-python-project-with-dependencies
+
// vim: set syntax=asciidoc:
diff --git a/doc/style.css b/doc/style.css
index e52223b..f99c06c 100644
--- a/doc/style.css
+++ b/doc/style.css
@@ -1,6 +1,6 @@
body {
font-family: serif;
- margin: 1em 45% 1em 5%;
+ margin: 1em 35% 1em 5%;
background-color: #fef9fa;
}
diff --git a/flake.nix b/flake.nix
index 2faaf91..28ccbe5 100644
--- a/flake.nix
+++ b/flake.nix
@@ -101,8 +101,6 @@
python-pkgs.python-creole
]))
perl
- sourceHighlight
- asciidoc
asciidoctor
pandoc
];
diff --git a/meson.build b/meson.build
index c088ff4..5e5d673 100644
--- a/meson.build
+++ b/meson.build
@@ -40,7 +40,7 @@ endif
if get_option('server').enabled()
subdir('src/server')
- if not meson.is_subproject()
+ if not meson.is_subproject() and not is_windows
subdir('src/fruits/cmsrv')
endif
endif
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index 1bf42bf..2af24b2 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -31,6 +31,7 @@ static void reset_buffer_state(struct camu_audio_buffer *buf)
{
buf->pts = -1.0;
buf->pause = PAUSE_PAUSED;
+ al_atomic_store(u8)(&buf->unpause, 0, AL_ATOMIC_RELAXED);
buf->volume = 1.f;
#ifdef CAMU_AUDIO_BUFFER_FADE
buf->fade_offset = 0;
@@ -185,13 +186,12 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_fra
al_free(frame);
}
-// Not thread-safe, must be called while the buffer is not being read from or written to.
void camu_audio_buffer_unpause(struct camu_audio_buffer *buf)
{
- buf->pause = PAUSE_PAUSED;
+ al_atomic_store(u8)(&buf->unpause, 1, AL_ATOMIC_RELAXED);
}
-// Not thread-safe.
+// Not thread-safe, must be called while the buffer is not being read from or written to.
void camu_audio_buffer_reset(struct camu_audio_buffer *buf)
{
reset_buffer_state(buf);
@@ -253,11 +253,19 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
return req;
#endif
}
+
+ if (al_atomic_load(u8)(&buf->unpause, AL_ATOMIC_ACQUIRE)) {
+ // If pause != PLAYING here, this will cause unexpected behavior.
+ buf->pause = PAUSE_PAUSED;
+ al_atomic_store(u8)(&buf->unpause, 0, AL_ATOMIC_RELEASE);
+ }
+
u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE);
if (flow == SIGNALED) {
buf->callback(buf->userdata, CAMU_BUFFER_EOF);
return 0;
}
+
f64 pts = camu_clock_get_pts(buf->clock, buf->latency);
size_t ret, signal = req;
size_t have = al_ring_buffer_occupied(&buf->rb);
@@ -298,6 +306,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
#ifdef CAMU_AUDIO_BUFFER_FADE
}
#endif
+
if (have < req) { // We don't have enough data to fulfill our request.
if (flow == FLUSHED) { // Stream is flushed.
// Check peak buffer for any remaining data.
@@ -335,6 +344,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
// The amount of available data could've increased in the flow = FLUSHED case.
req = AL_MIN(req, have);
}
+
if (req > 0) {
#ifdef CAMU_AUDIO_BUFFER_FADE
if (buf->pause == PAUSE_FADING) {
@@ -361,6 +371,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
}
#endif
}
+
// Check if we should request to uncork.
if (flow == FLOWING) {
ret = al_atomic_load(size_t)(&buf->uncork_at, AL_ATOMIC_RELAXED);
@@ -368,6 +379,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
buf->callback(buf->userdata, CAMU_BUFFER_UNCORK);
}
}
+
// To signal EOF, return less then req.
return signal;
}
diff --git a/src/buffer/audio.h b/src/buffer/audio.h
index 3adccc7..0f84d5a 100644
--- a/src/buffer/audio.h
+++ b/src/buffer/audio.h
@@ -17,6 +17,7 @@ struct camu_audio_buffer {
f64 pts;
u8 pause;
+ atomic(u8) unpause;
struct camu_clock *clock;
bool ignore_desync;
f64 latency;
diff --git a/src/liana/list.c b/src/liana/list.c
index 18bfd6c..a8d5f71 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -180,8 +180,12 @@ static bool assume_done(struct lia_list_entry *entry, u64 at)
void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index)
{
- if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
if (index == list->current) return;
+ if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
+ if (sequence == list->previous) {
+ // This can happen but almost certainly won't be expected behavior.
+ return;
+ }
struct lia_list_entry *current = get_entry_from_sequence(list, sequence);
struct lia_list_entry *target = get_entry_from_sequence(list, index);
@@ -208,7 +212,6 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index)
// This should only happen if `start` has never been set.
if (target->start == LIANA_TIMESTAMP_INVALID && target->paused_at == LIANA_TIMESTAMP_INVALID) {
- al_log_info("dd", "Hello");
target->start = at;
}
@@ -233,7 +236,7 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index)
target->start = at;
pause = !done ? LIANA_PAUSE_BOTH : LIANA_PAUSE_RESUME;
}
- al_log_info("dd", "%d: not paused, done: %d, pause: %d", sequence, done, pause);
+ al_log_debug("list", "%d: not paused, done: %d, pause: %d", sequence, done, pause);
} else {
if (assume_done(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID) {
pause = LIANA_PAUSE_NONE;
@@ -241,7 +244,7 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index)
target->start = at;
pause = LIANA_PAUSE_RESUME;
}
- al_log_info("dd", "%d: paused, pause: %d", sequence, pause);
+ al_log_debug("list", "%d: paused, pause: %d", sequence, pause);
}
list->current = index;
@@ -331,11 +334,11 @@ void lia_list_end(struct lia_list *list, s32 sequence)
if (sequence != list->current) return;
s32 size = (s32)list->entries.size;
s32 next = sequence + 1;
- list->previous = sequence;
struct lia_list_entry *current = al_array_at(list->entries, list->current);
current->offset = current->duration;
if (list->queued >= 0) {
list->current = list->queued;
+ list->previous = sequence;
list->queued = -1;
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 39fe0fe..5d7b851 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -5,7 +5,7 @@
#define VCR_BUFFER_INIT 64
#define VCR_BUFFER_BUFFERED (VCR_BUFFER_INIT - 8)
-#define VCR_BUFFER_LOW (VCR_BUFFER_INIT - 16)
+#define VCR_BUFFER_LOW (VCR_BUFFER_INIT - 32)
static void signal_callback(void *userdata)
{
@@ -91,7 +91,7 @@ void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track)
aki_cond_init(&track->cond);
aki_mutex_init(&track->mutex);
al_atomic_store(u8)(&track->buffered, 0, AL_ATOMIC_RELAXED);
- aki_packet_cache_init(&track->cache, VCR_BUFFER_INIT - 1);
+ aki_packet_cache_init(&track->cache, VCR_BUFFER_INIT);
al_array_push(vcr->tracks, track);
al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
}
@@ -112,13 +112,17 @@ static void cork_if_buffered(struct lia_vcr *vcr)
al_array_foreach(vcr->tracks, i, track) {
buffered &= al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED);
}
- if (buffered) aki_packet_stream_cork(vcr->data, true);
+ if (buffered) {
+ aki_packet_stream_cork(vcr->data, true);
+ al_array_foreach(vcr->tracks, i, track) {
+ aki_packet_cache_flush(&track->cache);
+ }
+ }
}
bool lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet)
{
struct lia_vcr_track *track;
- s32 count;
u8 op = aki_packet_read_u8(packet);
switch (op) {
case LIANA_PACKET_DATA:
@@ -134,9 +138,7 @@ bool lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet)
}
if (!aki_packet_cache_send_packet(&track->cache, packet)) {
aki_packet_free(packet);
- } else if ((count = al_atomic_add(s32)(&vcr->count, 1, AL_ATOMIC_RELAXED)) > vcr->mark.buffered) {
- vcr->mark.buffered = count;
- vcr->mark.low = vcr->mark.buffered - 16;
+ } else if (al_atomic_add(s32)(&vcr->count, 1, AL_ATOMIC_RELAXED) > vcr->mark.buffered) {
cork_if_buffered(vcr);
}
break;
diff --git a/src/libsink/common.h b/src/libsink/common.h
index f8dd9fa..0b1a6c9 100644
--- a/src/libsink/common.h
+++ b/src/libsink/common.h
@@ -3,10 +3,11 @@
#define CAMU_SINK_LOCAL 1
enum {
- CAMU_SINK_SET = 0,
+ CAMU_SINK_CLEAR = 0,
+ CAMU_SINK_SET,
CAMU_SINK_BUFFER,
CAMU_SINK_BUFFER_AND_QUEUE,
- CAMU_SINK_CLEAR,
CAMU_SINK_PAUSE,
- CAMU_SINK_SEEK
+ CAMU_SINK_SEEK,
+ CAMU_SINK_DURATION
};
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 990cd69..b9554c4 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -43,10 +43,10 @@ enum {
#define ENTRY_VIDEO_READY_OR_EMPTY(entry) true
#else
#define ENTRY_VIDEO_READY_OR_EMPTY(entry) \
- (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_QUEUED || entry->video.state == BUFFER_ADDED)
+ ((entry)->video.state == BUFFER_INIT || (entry)->video.state == BUFFER_QUEUED || (entry)->video.state == BUFFER_ADDED)
#endif
#define ENTRY_AUDIO_READY_OR_EMPTY(entry) \
- (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED || entry->audio.state == BUFFER_ADDED)
+ ((entry)->audio.state == BUFFER_INIT || (entry)->audio.state == BUFFER_QUEUED || (entry)->audio.state == BUFFER_ADDED)
#if defined CAMU_SCREEN_THREADED && defined CAMU_MIXER_THREADED
#define BLOCKING_SLEEP(delay) aki_thread_sleep(delay)
@@ -57,7 +57,7 @@ enum {
static bool entry_audio_buffer_held(struct camu_sink_entry *entry)
{
#ifdef CAMU_MIXER_THREADED
- return al_atomic_load(u8)(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED) == 1;
+ return al_atomic_load(u8)(&entry->audio.buf.ref, AL_ATOMIC_RELAXED) == 1;
#else
(void)entry;
return false;
@@ -68,7 +68,7 @@ static bool entry_audio_buffer_held(struct camu_sink_entry *entry)
static bool entry_video_buffer_held(struct camu_sink_entry *entry)
{
#ifdef CAMU_SCREEN_THREADED
- return al_atomic_load(u8)(&(entry)->video.buf.ref, AL_ATOMIC_RELAXED) == 1;
+ return al_atomic_load(u8)(&entry->video.buf.ref, AL_ATOMIC_RELAXED) == 1;
#else
(void)entry;
return false;
@@ -90,6 +90,8 @@ static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_e
if (entry->audio.state == BUFFER_ADDED) {
sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
entry->audio.state = BUFFER_SET_OR_BUFFERED;
+ } else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) {
+ entry->audio.state = BUFFER_CONFIGURED;
}
}
@@ -99,6 +101,8 @@ static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_e
if (entry->video.state == BUFFER_ADDED) {
sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
entry->video.state = BUFFER_SET_OR_BUFFERED;
+ } else if (entry->video.state == BUFFER_SET_OR_BUFFERED) {
+ entry->video.state = BUFFER_CONFIGURED;
}
}
#endif
@@ -212,8 +216,12 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
aki_packet_write_str(packet, &sink->default_list);
aki_packet_write_u8(packet, CAMU_LIST_SKIP);
- //s32 sequence = sink->current ? sink->current->sequence : LIANA_SEQUENCE_ANY;
s32 sequence = LIANA_SEQUENCE_ANY;
+ if (sink->target) {
+ sequence = sink->target->sequence;
+ } else if (sink->current) {
+ sequence = sink->current->sequence;
+ }
aki_packet_write_s32(packet, sequence);
aki_packet_write_s32(packet, cmd->value.i);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
@@ -305,17 +313,19 @@ static void maybe_remove_previous(struct camu_sink *sink)
void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
{
u8 state = entry->audio.state;
+ if (state == BUFFER_ADDED || state == BUFFER_SET_OR_BUFFERED) {
+ camu_audio_buffer_unpause(&entry->audio.buf);
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = START,
+ .value.i = CAMU_SINK_AUDIO
+ });
+ }
if (state == BUFFER_SET_OR_BUFFERED) {
bool can_resume = ENTRY_VIDEO_READY_OR_EMPTY(entry);
- camu_audio_buffer_unpause(&entry->audio.buf);
entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
if (can_resume) {
maybe_remove_previous(entry->sink);
}
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = START,
- .value.i = CAMU_SINK_AUDIO
- });
state = BUFFER_ADDED;
} else if (state == BUFFER_CONFIGURED) {
state = BUFFER_SET_OR_BUFFERED;
@@ -327,17 +337,20 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
{
u8 state = entry->video.state;
+ if (state == BUFFER_ADDED || state == BUFFER_SET_OR_BUFFERED) {
+ bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = single_frame ? STOP : START,
+ .value.i = CAMU_SINK_VIDEO
+ });
+ }
if (state == BUFFER_SET_OR_BUFFERED) {
bool can_resume = ENTRY_AUDIO_READY_OR_EMPTY(entry);
entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
if (can_resume) {
maybe_remove_previous(entry->sink);
}
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = single_frame ? STOP : START,
- .value.i = CAMU_SINK_VIDEO
- });
+
state = BUFFER_ADDED;
} else if (state == BUFFER_CONFIGURED) {
state = BUFFER_SET_OR_BUFFERED;
@@ -453,6 +466,10 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent
if (entry->audio.state != BUFFER_INIT && entry->audio.state != BUFFER_QUEUED &&
entry->video.state != BUFFER_INIT && entry->video.state != BUFFER_QUEUED) {
camu_video_buffer_set_latency(&entry->video.buf, -camu_mixer_get_latency(sink->audio.mixer));
+ } else if (entry->audio.state != BUFFER_INIT && entry->audio.state != BUFFER_QUEUED) {
+#if !CAMU_SINK_LOCAL
+ camu_audio_buffer_set_latency(&entry->audio.buf, -camu_mixer_get_latency(sink->audio.mixer));
+#endif
}
#else
(void)sink;
@@ -807,6 +824,10 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
#else
switch (pause) {
case LIANA_PAUSE_NONE:
+ if (sink->target) {
+ al_log_warn("sink", "Ignoring target on NONE.");
+ sink->target = NULL;
+ }
if (sink->current) {
al_array_push(sink->previous, sink->current);
}
@@ -814,22 +835,35 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
sink->current = entry;
break;
case LIANA_PAUSE_RESUME:
- camu_clock_resume(&entry->clock, at);
- if (sink->current) {
+ if (sink->target) {
+ al_log_warn("sink", "Ignoring target on RESUME.");
+ sink->target = NULL;
+ }
+ if (entry == sink->current) {
+ camu_audio_buffer_unpause(&entry->audio.buf);
+ } else if (sink->current) {
al_array_push(sink->previous, sink->current);
}
+ camu_clock_resume(&entry->clock, at);
set_or_queue_entry(entry);
sink->current = entry;
break;
case LIANA_PAUSE_PAUSE:
- if (sink->current) {
+ if (camu_clock_is_ended(&entry->clock)) {
+ // Server didn't know this entry was ended, but it is.
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = END,
+ .value.i = entry->sequence
+ });
+ camu_clock_pause(&sink->current->clock, at);
+ sink->current = entry;
+ } else if (sink->current) {
if (camu_clock_is_ended(&sink->current->clock)) {
// Server thought we weren't done, be we are.
al_array_push(sink->previous, sink->current);
set_or_queue_entry(entry);
sink->current = entry;
} else {
- al_printf("ay\n");
sink->target = entry;
camu_clock_pause(&sink->current->clock, at);
}
@@ -843,6 +877,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
sink->target = NULL;
} else {
sink->target = entry;
+ camu_audio_buffer_unpause(&entry->audio.buf);
}
camu_clock_pause(&prev_target->clock, at);
} else if (sink->current) {
diff --git a/src/server/server.c b/src/server/server.c
index 6c991eb..49ea234 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -61,6 +61,11 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
+ case CAMU_SINK_DURATION: {
+ struct camu_server_resource *resource = (struct camu_server_resource *)entry->opaque;
+ (void)resource;
+ break;
+ }
}
}