summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/sink.c110
-rw-r--r--src/libsink/sink.h2
2 files changed, 69 insertions, 43 deletions
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 65a082e..c6727ab 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -129,8 +129,8 @@ 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);
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_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);
@@ -153,7 +153,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
#endif
} else {
- camu_clock_pause(&entry->clock);
+ camu_clock_pause(&entry->clock, -1.0);
#ifndef CAMU_SINK_NO_VIDEO
if (sink->video.state == SINK_PLAYING) {
sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL);
@@ -167,8 +167,8 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
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);
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_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);
@@ -177,8 +177,8 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
}
case SEEK: {
- struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SRV_LIST_ACTION);
- aki_packet_write_u8(packet, CAMU_LIST_SEEK);
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_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);
@@ -186,8 +186,8 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
}
case FINISHED: {
- struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SRV_LIST_ACTION);
- aki_packet_write_u8(packet, CAMU_LIST_FINISHED);
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_FINISHED);
aki_packet_write_s32(packet, cmd->value.i);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
@@ -257,7 +257,7 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *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) {
- camu_clock_resume(&entry->clock);
+ camu_clock_ready(&entry->clock);
maybe_remove_previous(entry->sink);
}
//queue_cmd(entry->sink, (struct camu_sink_cmd){
@@ -283,7 +283,7 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
});
entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
if (can_resume) {
- camu_clock_resume(&entry->clock);
+ camu_clock_ready(&entry->clock);
maybe_remove_previous(entry->sink);
}
//queue_cmd(entry->sink, (struct camu_sink_cmd){
@@ -301,11 +301,12 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
static void audio_buffer_callback(void *userdata, u8 op)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
+ struct camu_sink *sink = entry->sink;
switch (op) {
case CAMU_BUFFER_BUFFERED:
- aki_mutex_lock(&entry->sink->mutex);
+ aki_mutex_lock(&sink->mutex);
add_audio_if_set_and_buffered(entry);
- aki_mutex_unlock(&entry->sink->mutex);
+ aki_mutex_unlock(&sink->mutex);
break;
case CAMU_BUFFER_CORK:
shrb_vcr_stream_cork(entry->audio.stream);
@@ -314,25 +315,22 @@ static void audio_buffer_callback(void *userdata, u8 op)
shrb_vcr_stream_uncork(entry->audio.stream);
break;
case CAMU_BUFFER_PAUSED:
- aki_mutex_lock(&entry->sink->mutex);
+ aki_mutex_lock(&sink->mutex);
if (!entry->audio.armed) {
- queue_cmd(entry->sink, (struct camu_sink_cmd){
+ queue_cmd(sink, (struct camu_sink_cmd){
.op = STOP,
.value.i = CAMU_SINK_AUDIO
});
}
- aki_mutex_unlock(&entry->sink->mutex);
+ aki_mutex_unlock(&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) {
+ al_log_info("sink", "Swapping audio buffers (gapless).");
swap_buffers_internal(sink);
swapped = true;
- al_log_info("sink", "Audio buffers swapped (gapless).");
}
aki_mutex_unlock(&sink->mutex);
if (!swapped) {
@@ -347,9 +345,9 @@ static void audio_buffer_callback(void *userdata, u8 op)
//});
#endif
}
- queue_cmd(entry->sink, (struct camu_sink_cmd){
+ queue_cmd(sink, (struct camu_sink_cmd){
.op = FINISHED,
- .value.i = sequence
+ .value.i = entry->sequence
});
break;
}
@@ -360,11 +358,12 @@ static void audio_buffer_callback(void *userdata, u8 op)
static void video_buffer_callback(void *userdata, u8 op)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
+ struct camu_sink *sink = entry->sink;
switch (op) {
case CAMU_BUFFER_BUFFERED:
- aki_mutex_lock(&entry->sink->mutex);
+ aki_mutex_lock(&sink->mutex);
add_video_if_set_and_buffered(entry);
- aki_mutex_unlock(&entry->sink->mutex);
+ aki_mutex_unlock(&sink->mutex);
break;
case CAMU_BUFFER_CORK:
shrb_vcr_stream_cork(entry->video.stream);
@@ -372,13 +371,23 @@ static void video_buffer_callback(void *userdata, u8 op)
case CAMU_BUFFER_UNCORK:
shrb_vcr_stream_uncork(entry->video.stream);
break;
- case CAMU_BUFFER_EOF:
- queue_cmd(entry->sink, (struct camu_sink_cmd){
+ case CAMU_BUFFER_EOF: {
+ queue_cmd(sink, (struct camu_sink_cmd){
.op = STOP,
.value.i = CAMU_SINK_VIDEO
});
+ aki_mutex_lock(&sink->mutex);
+ bool no_audio = entry->audio.state == BUFFER_INIT;
+ aki_mutex_unlock(&sink->mutex);
+ if (no_audio) {
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = FINISHED,
+ .value.i = entry->sequence
+ });
+ }
break;
}
+ }
}
#endif
@@ -419,8 +428,9 @@ static void client_callback(void *userdata, u8 op, struct shrb_vcr_stream *strea
}
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);
+ req->delay = SHRUB_DELAY_IGNORE;
+ camu_clock_set(&entry->clock, req);
+ if (req->paused) req->base = req->paused_at;
break;
}
case SHRUB_CLIENT_PAUSE: {
@@ -432,9 +442,12 @@ static void client_callback(void *userdata, u8 op, struct shrb_vcr_stream *strea
break;
}
case SHRUB_CLIENT_RESUME: {
- u64 ts = *(u64 *)opaque;
aki_mutex_lock(&entry->sink->mutex);
- camu_clock_arm_resume(&entry->clock, ts);
+ struct shrb_resume_req *req = (struct shrb_resume_req *)opaque;
+ if (sink->queued) {
+ camu_clock_offset(&sink->queued->clock, req->offset);
+ }
+ camu_clock_arm_resume(&entry->clock, req->start);
if (sink->audio.state == SINK_PAUSED) {
sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
sink->audio.state = SINK_PLAYING;
@@ -569,17 +582,16 @@ static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink *
entry->audio.buf.callback = audio_buffer_callback;
entry->audio.buf.userdata = entry;
- struct camu_renderer *renderer = NULL;
#ifndef CAMU_SINK_NO_VIDEO
- renderer = sink->video.renderer;
entry->video.state = BUFFER_INIT;
- camu_video_buffer_init(&entry->video.buf, &entry->clock, 0.0, renderer);
+ camu_video_buffer_init(&entry->video.buf, &entry->clock, 0.0, sink->video.renderer);
entry->video.buf.callback = video_buffer_callback;
entry->video.buf.userdata = entry;
#endif
entry->client.callback = client_callback;
entry->client.userdata = entry;
+ // Assume we can error out at any point after connect is called (even from within connect() itself).
shrb_client_connect(&entry->client, sink->loop, addr, port, node_id);
return entry;
@@ -598,16 +610,19 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
u16 node_id = aki_packet_read_u16(packet);
s32 sequence = aki_packet_read_s32(packet);
u64 start = aki_packet_read_u64(packet);
+ u8 paused = aki_packet_read_u8(packet);
+ u64 paused_at = 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;
- // This will only have an effect if the clock is already set.
- camu_clock_arm_resume(&entry->clock, start + SHRUB_BASE_DELAY);
+ if (!paused) {
+ // 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);
+ camu_clock_pause(&sink->current->clock, paused_at / 1000000.0);
al_array_push(sink->previous, sink->current);
}
sink->current = entry;
@@ -662,7 +677,7 @@ 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, CAMU_SRV_IDENTIFY);
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_IDENTIFY);
aki_packet_write_u8(packet, CAMU_SINK);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
}
@@ -676,14 +691,17 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection
bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port)
{
- aki_rpc_init(&sink->client, AKI_SOCKET_TCP, connection_callback,
- connection_closed_callback, sink);
+ aki_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink);
for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) {
commands[i].userdata = sink;
al_assert(sink->callback);
aki_rpc_add_command(&sink->client, &commands[i]);
}
- return aki_rpc_connect(&sink->client, sink->loop, addr, port);
+ if (!aki_rpc_prepare_client(&sink->client, AKI_SOCKET_TCP, CAMU_MULTIPLEX_RPC)) {
+ return false;
+ }
+ aki_rpc_connect(&sink->client, addr, port);
+ return true;
}
struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink)
@@ -697,6 +715,14 @@ void camu_sink_return_current(struct camu_sink *sink)
aki_mutex_unlock(&sink->mutex);
}
+void camu_sink_add(struct camu_sink *sink, str *line)
+{
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
+ aki_packet_write_u8(packet, CAMU_ADD);
+ aki_packet_write_str(packet, line);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+}
+
void camu_sink_skip(struct camu_sink *sink, s32 n)
{
queue_cmd(sink, (struct camu_sink_cmd){
@@ -763,7 +789,7 @@ void camu_sink_stop(struct camu_sink *sink)
void camu_sink_close(struct camu_sink *sink)
{
// TODO: Make sure no commands can come in after this.
- aki_rpc_disconnect(&sink->client);
+ aki_rpc_conn_disconnect(sink->client.conn);
aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *entry;
al_array_foreach(sink->entries, i, entry) {
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index 4d7dc06..b789c9a 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -1,6 +1,5 @@
#pragma once
-//#define CAMU_SINK_NO_VIDEO
//#define CAMU_SINK_LOCAL_PAUSE
#include <al/str.h>
@@ -100,6 +99,7 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port);
struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink);
void camu_sink_return_current(struct camu_sink *sink);
+void camu_sink_add(struct camu_sink *sink, str *line);
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);