summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/common.h5
-rw-r--r--src/libsink/meson.build2
-rw-r--r--src/libsink/sink.c405
-rw-r--r--src/libsink/sink.h22
4 files changed, 194 insertions, 240 deletions
diff --git a/src/libsink/common.h b/src/libsink/common.h
index 48e43e8..a5ad25c 100644
--- a/src/libsink/common.h
+++ b/src/libsink/common.h
@@ -1,7 +1,6 @@
#pragma once
enum {
- CAMU_SINK_CMD_BUFFER = 0,
- CAMU_SINK_CMD_SET,
- CAMU_SINK_CMD_QUEUE,
+ CAMU_SINK_SET = 0,
+ CAMU_SINK_QUEUE
};
diff --git a/src/libsink/meson.build b/src/libsink/meson.build
index 0e6dfb3..83cfdd7 100644
--- a/src/libsink/meson.build
+++ b/src/libsink/meson.build
@@ -1,3 +1,3 @@
libsink_src = ['sink.c']
-libsink_deps = [bimu_client]
+libsink_deps = [shrub_client]
libsink = declare_dependency(sources: libsink_src, dependencies: libsink_deps)
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 3468a48..65a082e 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -1,12 +1,14 @@
#include <al/log.h>
-#include "../tree/common.h"
+#include "../server/common.h"
-#include "../bimu/common.h"
-#include "../bimu/handler.h"
+#include "../shrub/common.h"
+#include "../shrub/handler.h"
#include "../buffer/common.h"
+#include "../list/list.h"
+
#include "sink.h"
#include "common.h"
@@ -17,35 +19,25 @@ enum {
};
enum {
- ENTRY_LOADED = 0,
- ENTRY_BUFFERED,
- ENTRY_DISREGUARDED
-};
-
-enum {
BUFFER_INIT = 0,
BUFFER_QUEUED,
BUFFER_CONFIGURED,
BUFFER_SET_OR_BUFFERED,
- BUFFER_ADDED,
+ BUFFER_ADDED
};
enum {
START,
STOP,
TOGGLE_PAUSE,
+ SEEK,
SKIP,
+ FINISHED,
RESEEK,
- SEEK,
SET_BUFFERED, // Currently set entry is buffered.
- CLOSE,
- // Internal.
- CORK,
- UNCORK,
+ CLOSE
};
-#define BIMU_DELAY_IGNORE 0
-
#define ENTRY_AUDIO_BUFFER_HELD(entry) \
(al_atomic_bool_load(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED))
#ifdef CAMU_SINK_NO_VIDEO
@@ -56,10 +48,14 @@ enum {
#define ENTRY_BUFFERS_HELD(entry) (ENTRY_AUDIO_BUFFER_HELD(entry) || ENTRY_VIDEO_BUFFER_HELD(entry))
#endif
-#define ENTRY_VIDEO_EMPTY(entry) \
- (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_QUEUED)
-#define ENTRY_AUDIO_EMPTY(entry) \
- (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED)
+#ifdef CAMU_SINK_NO_VIDEO
+#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)
+#endif
+#define ENTRY_AUDIO_READY_OR_EMPTY(entry) \
+ (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED || entry->audio.state == BUFFER_ADDED)
static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry)
{
@@ -132,10 +128,19 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
break;
}
+ case SKIP: {
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SRV_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_LIST_SKIP);
+ s32 sequence = sink->current ? sink->current->sequence : CAMU_SEQUENCE_INVALID;
+ aki_packet_write_s32(packet, sequence);
+ aki_packet_write_s32(packet, cmd->value.i);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+ break;
+ }
case TOGGLE_PAUSE: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+#ifdef CAMU_SINK_LOCAL_PAUSE
if (camu_clock_is_paused(&entry->clock)) {
- camu_clock_calc_tick(&entry->clock);
camu_clock_resume(&entry->clock);
if (sink->audio.state == SINK_PAUSED) {
sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
@@ -156,22 +161,40 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
#endif
}
+#else
+ f64 base = camu_clock_get_base_pts(&entry->clock);
+ if (!camu_clock_is_paused(&entry->clock)) {
+ f64 pts = camu_clock_get_pts(&entry->clock, 0.0);
+ if (pts > base) base = pts;
+ }
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SRV_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_LIST_TOGGLE_PAUSE);
+ s32 sequence = sink->current ? sink->current->sequence : CAMU_SEQUENCE_INVALID;
+ aki_packet_write_s32(packet, sequence);
+ aki_packet_write_u64(packet, base * 1000000.0);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+#endif
break;
}
- case SKIP: {
- struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_SKIP);
- aki_packet_write_s32(packet, cmd->value.i);
+ case SEEK: {
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SRV_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_LIST_SEEK);
+ s32 sequence = sink->current ? sink->current->sequence : CAMU_SEQUENCE_INVALID;
+ aki_packet_write_s32(packet, sequence);
+ aki_packet_write_f64(packet, cmd->value.f);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
- case RESEEK: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- bmu_client_reseek(&entry->client);
+ case FINISHED: {
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SRV_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_LIST_FINISHED);
+ aki_packet_write_s32(packet, cmd->value.i);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
- case SEEK: {
+ case RESEEK: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- bmu_client_seek(&entry->client, cmd->value.f);
+ shrb_client_reseek(&entry->client);
break;
}
// case SET_BUFFERED: {
@@ -192,34 +215,6 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
return;
}
- case CORK: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO:
- bmu_vcr_stream_cork(entry->audio.stream);
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO:
- bmu_vcr_stream_cork(entry->video.stream);
- break;
-#endif
- }
- break;
- }
- case UNCORK: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO:
- bmu_vcr_stream_uncork(entry->audio.stream);
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO:
- bmu_vcr_stream_uncork(entry->video.stream);
- break;
-#endif
- }
- break;
- }
}
}
@@ -241,31 +236,30 @@ static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
aki_signal_send(&sink->signal);
}
+static void maybe_remove_previous(struct camu_sink *sink)
+{
+ struct camu_sink_entry *previous;
+ al_array_foreach(sink->previous, i, previous) {
+ remove_entry_buffers(sink, previous);
+ }
+ sink->previous.size = 0;
+}
+
void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
{
u8 state = entry->audio.state;
if (state == BUFFER_SET_OR_BUFFERED) {
- // TODO: should_resume is not robust. CAMU_SINK_NO_VIDEO doesn't work but that's
- // the least of our problems.
-#ifndef CAMU_SINK_NO_VIDEO
- bool should_resume = ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED;
-#endif
- if (should_resume && !camu_clock_calc_tick(&entry->clock)) {
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = RESEEK,
- .opaque = entry
- });
- return;
- }
+ bool can_resume = ENTRY_VIDEO_READY_OR_EMPTY(entry);
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
.value.i = CAMU_SINK_AUDIO
});
- if (should_resume) {
+ 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) {
camu_clock_resume(&entry->clock);
- entry->state = ENTRY_BUFFERED;
+ maybe_remove_previous(entry->sink);
}
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
//queue_cmd(entry->sink, (struct camu_sink_cmd){
// .op = SET_BUFFERED,
// .value.i = CAMU_SINK_AUDIO
@@ -282,23 +276,16 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
{
u8 state = entry->video.state;
if (state == BUFFER_SET_OR_BUFFERED) {
- bool should_resume = ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED;
- if (should_resume && !camu_clock_calc_tick(&entry->clock)) {
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = RESEEK,
- .opaque = entry
- });
- return;
- }
+ bool can_resume = ENTRY_AUDIO_READY_OR_EMPTY(entry);
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
.value.i = CAMU_SINK_VIDEO
});
- if (should_resume) {
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ if (can_resume) {
camu_clock_resume(&entry->clock);
- entry->state = ENTRY_BUFFERED;
+ maybe_remove_previous(entry->sink);
}
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
//queue_cmd(entry->sink, (struct camu_sink_cmd){
// .op = SET_BUFFERED,
// .value.i = CAMU_SINK_VIDEO
@@ -321,35 +308,27 @@ static void audio_buffer_callback(void *userdata, u8 op)
aki_mutex_unlock(&entry->sink->mutex);
break;
case CAMU_BUFFER_CORK:
- bmu_vcr_stream_cork(entry->audio.stream);
- /*
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = CORK,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = entry
- });
- */
+ shrb_vcr_stream_cork(entry->audio.stream);
break;
case CAMU_BUFFER_UNCORK:
- bmu_vcr_stream_uncork(entry->audio.stream);
- /*
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = UNCORK,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = entry
- });
- */
+ shrb_vcr_stream_uncork(entry->audio.stream);
break;
case CAMU_BUFFER_PAUSED:
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_AUDIO
- });
+ aki_mutex_lock(&entry->sink->mutex);
+ if (!entry->audio.armed) {
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = STOP,
+ .value.i = CAMU_SINK_AUDIO
+ });
+ }
+ aki_mutex_unlock(&entry->sink->mutex);
break;
case CAMU_BUFFER_EOF: {
+ // TODO: Video only case not handled (video not image).
struct camu_sink *sink = entry->sink;
bool swapped = false;
aki_mutex_lock(&sink->mutex);
+ s32 sequence = sink->current->sequence;
if (sink->queued) {
swap_buffers_internal(sink);
swapped = true;
@@ -368,6 +347,10 @@ static void audio_buffer_callback(void *userdata, u8 op)
//});
#endif
}
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = FINISHED,
+ .value.i = sequence
+ });
break;
}
}
@@ -384,24 +367,10 @@ static void video_buffer_callback(void *userdata, u8 op)
aki_mutex_unlock(&entry->sink->mutex);
break;
case CAMU_BUFFER_CORK:
- bmu_vcr_stream_cork(entry->video.stream);
- /*
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = CORK,
- .value.i = CAMU_SINK_VIDEO,
- .opaque = entry
- });
- */
+ shrb_vcr_stream_cork(entry->video.stream);
break;
case CAMU_BUFFER_UNCORK:
- bmu_vcr_stream_uncork(entry->video.stream);
- /*
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = UNCORK,
- .value.i = CAMU_SINK_VIDEO,
- .opaque = entry
- });
- */
+ shrb_vcr_stream_uncork(entry->video.stream);
break;
case CAMU_BUFFER_EOF:
queue_cmd(entry->sink, (struct camu_sink_cmd){
@@ -413,14 +382,15 @@ static void video_buffer_callback(void *userdata, u8 op)
}
#endif
-static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream, void *opaque)
+static void client_callback(void *userdata, u8 op, struct shrb_vcr_stream *stream, void *opaque)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
+ struct camu_sink *sink = entry->sink;
switch (op) {
- case BIMU_CLIENT_CONFIGURE: {
- // No data can be sent until all active streams are configured.
+ case SHRUB_CLIENT_CONFIGURE: {
+ // No data can be sent until all selected streams are configured.
switch (stream->type) {
- case BIMU_STREAM_AUDIO:
+ case SHRUB_STREAM_AUDIO:
entry->audio.stream = stream;
camu_audio_buffer_configure(&entry->audio.buf, &stream->stream);
aki_mutex_lock(&entry->sink->mutex);
@@ -432,7 +402,7 @@ static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream
aki_mutex_unlock(&entry->sink->mutex);
break;
#ifndef CAMU_SINK_NO_VIDEO
- case BIMU_STREAM_VIDEO:
+ case SHRUB_STREAM_VIDEO:
entry->video.stream = stream;
camu_video_buffer_configure(&entry->video.buf, &stream->stream);
aki_mutex_lock(&entry->sink->mutex);
@@ -447,46 +417,51 @@ static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream
}
break;
}
- case BIMU_CLIENT_SET: {
- u64 ts = aki_get_timestamp();
- struct bmu_seek_req *req = (struct bmu_seek_req *)opaque;
- req->delay = BIMU_DELAY_IGNORE;
- req->ts -= req->delay;
- bool late = ts <= req->ts;
- if (late) {}
+ case SHRUB_CLIENT_SET: {
+ struct shrb_seek_req *req = (struct shrb_seek_req *)opaque;
+ //req->delay = SHRUB_DELAY_IGNORE;
camu_clock_set(&entry->clock, req->base, req->ts, req->delay);
- camu_audio_buffer_reset(&entry->audio.buf);
- if (!late) {
- // This is a fail-safe for the case where receiving data is extremely slow.
- // It should not be depended on for checking if an entry buffers in time.
- aki_timer_set_repeat(&entry->timer, AKI_TS_FROM_USEC(ts - req->ts));
- aki_timer_again(&entry->timer);
+ break;
+ }
+ case SHRUB_CLIENT_PAUSE: {
+ f64 pts = *(f64 *)opaque;
+ aki_mutex_lock(&entry->sink->mutex);
+ camu_clock_arm_pause(&entry->clock, pts);
+ entry->audio.armed = false;
+ aki_mutex_unlock(&entry->sink->mutex);
+ break;
+ }
+ case SHRUB_CLIENT_RESUME: {
+ u64 ts = *(u64 *)opaque;
+ aki_mutex_lock(&entry->sink->mutex);
+ camu_clock_arm_resume(&entry->clock, ts);
+ if (sink->audio.state == SINK_PAUSED) {
+ sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
+ sink->audio.state = SINK_PLAYING;
+ } else {
+ entry->audio.armed = true;
}
+ aki_mutex_unlock(&entry->sink->mutex);
break;
}
- case BIMU_CLIENT_DATA: {
+ case SHRUB_CLIENT_DATA: {
struct camu_frame *frame = (struct camu_frame *)opaque;
switch (stream->type) {
- case BIMU_STREAM_AUDIO:
+ case SHRUB_STREAM_AUDIO:
camu_audio_buffer_push(&entry->audio.buf, frame);
break;
#ifndef CAMU_SINK_NO_VIDEO
- case BIMU_STREAM_VIDEO:
+ case SHRUB_STREAM_VIDEO:
camu_video_buffer_push(&entry->video.buf, frame);
break;
#endif
default:
-#ifdef HAVE_FFMPEG
- if (frame->type == CAMU_FFMPEG_COMPAT) {
- av_frame_free(&frame->av.frame);
- }
-#endif
- al_free(frame);
+ camu_frame_discard(frame);
break;
}
break;
}
- case BIMU_CLIENT_REMOVE_BUFFERS: {
+ case SHRUB_CLIENT_REMOVE_BUFFERS: {
aki_mutex_lock(&entry->sink->mutex);
if (entry->audio.state == BUFFER_ADDED) {
entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
@@ -502,32 +477,25 @@ static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream
aki_mutex_unlock(&entry->sink->mutex);
#ifndef CAMU_SINK_NO_VIDEO
if (single_frame) {
- while (ENTRY_AUDIO_BUFFER_HELD(entry)) {
- aki_thread_sleep(AKI_TS_FROM_USEC(200));
- }
+ while (ENTRY_AUDIO_BUFFER_HELD(entry)) { aki_thread_sleep(AKI_TS_FROM_USEC(200)); }
} else {
#endif
- while (ENTRY_BUFFERS_HELD(entry)) {
- aki_thread_sleep(AKI_TS_FROM_USEC(200));
- }
+ while (ENTRY_BUFFERS_HELD(entry)) { aki_thread_sleep(AKI_TS_FROM_USEC(200)); }
#ifndef CAMU_SINK_NO_VIDEO
- }
-#endif
-#ifndef CAMU_SINK_NO_VIDEO
- if (!single_frame) {
camu_video_buffer_reset(&entry->video.buf);
}
#endif
+ camu_audio_buffer_reset(&entry->audio.buf);
break;
}
- case BIMU_CLIENT_EOF: {
+ case SHRUB_CLIENT_EOF: {
switch (stream->type) {
- case BIMU_STREAM_AUDIO: {
+ case SHRUB_STREAM_AUDIO: {
camu_audio_buffer_flush(&entry->audio.buf);
break;
}
#ifndef CAMU_SINK_NO_VIDEO
- case BIMU_STREAM_VIDEO: {
+ case SHRUB_STREAM_VIDEO: {
bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
if (!single_frame) {
camu_video_buffer_flush(&entry->video.buf);
@@ -538,9 +506,8 @@ static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream
}
break;
}
- case BIMU_CLIENT_CLOSED: {
- bmu_client_free(&entry->client);
- aki_timer_stop(&entry->timer);
+ case SHRUB_CLIENT_CLOSED: {
+ shrb_client_free(&entry->client);
camu_audio_buffer_free(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
camu_video_buffer_free(&entry->video.buf);
@@ -566,6 +533,7 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
aki_mutex_init(&sink->mutex);
sink->queued = NULL;
sink->current = NULL;
+ al_array_init(sink->previous);
al_array_init(sink->entries);
sink->audio.mixer = mixer;
sink->audio.state = SINK_PAUSED;
@@ -585,16 +553,8 @@ static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 no
return NULL;
}
-static void buffer_timer_callback(void *userdata, struct aki_timer *timer)
-{
- struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
- (void)timer;
- //aki_mutex_lock(&entry->sink->mutex);
- //aki_mutex_unlock(&entry->sink->mutex);
- aki_timer_stop(&entry->timer);
-}
-
-static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink *sink, str *addr, s32 port, u32 node_id)
+static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink *sink,
+ str *addr, s32 port, u32 node_id)
{
struct camu_sink_entry *entry = entry_from_node_id(sink, node_id);
if (entry) return entry;
@@ -602,9 +562,9 @@ static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink *
entry = al_alloc_object(struct camu_sink_entry);
al_array_push(sink->entries, entry);
entry->sink = sink;
- entry->state = ENTRY_LOADED;
entry->audio.state = BUFFER_INIT;
+ entry->audio.armed = false;
camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer);
entry->audio.buf.callback = audio_buffer_callback;
entry->audio.buf.userdata = entry;
@@ -613,41 +573,18 @@ static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink *
#ifndef CAMU_SINK_NO_VIDEO
renderer = sink->video.renderer;
entry->video.state = BUFFER_INIT;
- camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer);
+ camu_video_buffer_init(&entry->video.buf, &entry->clock, 0.0, renderer);
entry->video.buf.callback = video_buffer_callback;
entry->video.buf.userdata = entry;
#endif
- aki_timer_init(&entry->timer, sink->loop, buffer_timer_callback, entry);
-
entry->client.callback = client_callback;
entry->client.userdata = entry;
- bmu_client_connect(&entry->client, sink->loop, addr, port, node_id);
+ shrb_client_connect(&entry->client, sink->loop, addr, port, node_id);
return entry;
}
-static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn,
- struct aki_packet *packet, struct aki_packet *rpacket)
-{
- struct camu_sink *sink = (struct camu_sink *)userdata;
- (void)conn;
- (void)rpacket;
-
- str addr;
- aki_packet_read_str(packet, &addr);
- s32 port = aki_packet_read_s32(packet);
- u16 node_id = aki_packet_read_u16(packet);
-
- aki_mutex_lock(&sink->mutex);
- ensure_entry_buffered_internal(sink, &addr, port, node_id);
- aki_mutex_unlock(&sink->mutex);
-
- aki_packet_free(packet);
-
- return false;
-}
-
static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn,
struct aki_packet *packet, struct aki_packet *rpacket)
{
@@ -659,11 +596,20 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
aki_packet_read_str(packet, &addr);
s32 port = aki_packet_read_s32(packet);
u16 node_id = aki_packet_read_u16(packet);
+ s32 sequence = aki_packet_read_s32(packet);
+ u64 start = aki_packet_read_u64(packet);
aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *entry = ensure_entry_buffered_internal(sink, &addr, port, node_id);
+ entry->sequence = sequence;
if (entry == sink->current) goto out;
- struct camu_sink_entry *previous = sink->current;
+ // This will only have an effect if the clock is already set.
+ camu_clock_arm_resume(&entry->clock, start + SHRUB_BASE_DELAY);
+ if (sink->current) {
+ // TODOODO, pause_at sent from list calc difference in clock and set.
+ camu_clock_pause(&sink->current->clock);
+ al_array_push(sink->previous, sink->current);
+ }
sink->current = entry;
if (sink->queued) sink->queued = NULL;
if (entry->audio.state == BUFFER_INIT) {
@@ -678,11 +624,6 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
add_video_if_set_and_buffered(entry);
}
#endif
- if (previous) {
- camu_clock_pause(&previous->clock);
- // TODO: sink->previous that gets swapped in add_x_if_set_and_buffered?
- remove_entry_buffers(sink, previous);
- }
out:
aki_mutex_unlock(&sink->mutex);
aki_packet_free(packet);
@@ -700,9 +641,11 @@ static bool queue_command_callback(void *userdata, struct aki_rpc_connection *co
aki_packet_read_str(packet, &addr);
s32 port = aki_packet_read_s32(packet);
u16 node_id = aki_packet_read_u16(packet);
+ s32 sequence = aki_packet_read_s32(packet);
aki_mutex_lock(&sink->mutex);
sink->queued = ensure_entry_buffered_internal(sink, &addr, port, node_id);
+ sink->queued->sequence = sequence;
aki_mutex_unlock(&sink->mutex);
aki_packet_free(packet);
@@ -711,17 +654,16 @@ static bool queue_command_callback(void *userdata, struct aki_rpc_connection *co
}
static struct aki_rpc_command commands[] = {
- { .op = CAMU_SINK_CMD_BUFFER, .callback = buffer_command_callback, .userdata = NULL },
- { .op = CAMU_SINK_CMD_SET, .callback = set_command_callback, .userdata = NULL },
- { .op = CAMU_SINK_CMD_QUEUE, .callback = queue_command_callback, .userdata = NULL }
+ { .op = CAMU_SINK_SET, .callback = set_command_callback, .userdata = NULL },
+ { .op = CAMU_SINK_QUEUE, .callback = queue_command_callback, .userdata = NULL }
};
static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
{
struct camu_sink *sink = (struct camu_sink *)userdata;
sink->conn = conn;
- struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_IDENTIFY);
- aki_packet_write_u8(packet, TREE_SINK);
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SRV_IDENTIFY);
+ aki_packet_write_u8(packet, CAMU_SINK);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
}
@@ -744,6 +686,25 @@ bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port)
return aki_rpc_connect(&sink->client, sink->loop, addr, port);
}
+struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink)
+{
+ aki_mutex_lock(&sink->mutex);
+ return sink->current;
+}
+
+void camu_sink_return_current(struct camu_sink *sink)
+{
+ aki_mutex_unlock(&sink->mutex);
+}
+
+void camu_sink_skip(struct camu_sink *sink, s32 n)
+{
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = SKIP,
+ .value.i = n
+ });
+}
+
void camu_sink_toggle_pause(struct camu_sink *sink)
{
aki_mutex_lock(&sink->mutex);
@@ -771,23 +732,15 @@ void camu_sink_seek(struct camu_sink *sink, f64 precent)
}
}
-void camu_sink_skip(struct camu_sink *sink, s32 n)
-{
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = SKIP,
- .value.i = n
- });
-}
-
-struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink)
+void camu_sink_reseek(struct camu_sink *sink)
{
aki_mutex_lock(&sink->mutex);
- return sink->current;
-}
-
-void camu_sink_return_current(struct camu_sink *sink)
-{
+ struct camu_sink_entry *current = sink->current;
aki_mutex_unlock(&sink->mutex);
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = RESEEK,
+ .opaque = current
+ });
}
void camu_sink_stop(struct camu_sink *sink)
@@ -814,7 +767,7 @@ void camu_sink_close(struct camu_sink *sink)
aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *entry;
al_array_foreach(sink->entries, i, entry) {
- bmu_client_close(&entry->client);
+ shrb_client_close(&entry->client);
}
sink->entries.size = 0;
aki_mutex_unlock(&sink->mutex);
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index 40d20d4..4d7dc06 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -1,12 +1,12 @@
#pragma once
//#define CAMU_SINK_NO_VIDEO
+//#define CAMU_SINK_LOCAL_PAUSE
#include <al/str.h>
#include <al/array.h>
#include <aki/rpc2.h>
#include <aki/signal.h>
-#include <aki/timer.h>
#include "../util/queue.h"
@@ -16,7 +16,7 @@
#include "../buffer/video.h"
#endif
-#include "../bimu/client.h"
+#include "../shrub/client.h"
enum {
CAMU_SINK_AUDIO = 0,
@@ -36,20 +36,20 @@ enum {
};
struct camu_sink_entry {
- u8 state;
- struct bmu_client client;
+ s32 sequence;
struct camu_clock clock;
- struct aki_timer timer;
+ struct shrb_client client;
struct {
u8 state;
+ bool armed;
struct camu_audio_buffer buf;
- struct bmu_vcr_stream *stream;
+ struct shrb_vcr_stream *stream;
} audio;
#ifndef CAMU_SINK_NO_VIDEO
struct {
u8 state;
struct camu_video_buffer buf;
- struct bmu_vcr_stream *stream;
+ struct shrb_vcr_stream *stream;
} video;
#endif
struct camu_sink *sink;
@@ -75,6 +75,7 @@ struct camu_sink {
struct aki_mutex mutex;
struct camu_sink_entry *queued;
struct camu_sink_entry *current;
+ array(struct camu_sink_entry *) previous;
array(struct camu_sink_entry *) entries;
struct {
u8 state;
@@ -97,11 +98,12 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
#endif
);
bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port);
-void camu_sink_toggle_pause(struct camu_sink *sink);
-void camu_sink_seek(struct camu_sink *sink, f64 pos);
-void camu_sink_skip(struct camu_sink *sink, s32 n);
struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink);
void camu_sink_return_current(struct camu_sink *sink);
+void camu_sink_skip(struct camu_sink *sink, s32 n);
+void camu_sink_toggle_pause(struct camu_sink *sink);
+void camu_sink_seek(struct camu_sink *sink, f64 pos);
+void camu_sink_reseek(struct camu_sink *sink);
void camu_sink_stop(struct camu_sink *sink);
void camu_sink_close(struct camu_sink *sink);
void camu_sink_free(struct camu_sink *sink);