summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/sink.c170
-rw-r--r--src/libsink/sink.h7
2 files changed, 126 insertions, 51 deletions
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 58e3db3..343275b 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -1,8 +1,11 @@
#include <al/log.h>
#include "../tree/common.h"
+#include "../bimu/common.h"
+#include "../bimu/handler.h"
#include "sink.h"
+#include "common.h"
enum {
SINK_EMPTY = 0,
@@ -12,6 +15,7 @@ enum {
enum {
ENTRY_LOADED = 0,
+ ENTRY_BUFFERED,
ENTRY_DISREGUARDED
};
@@ -27,8 +31,9 @@ enum {
START,
STOP,
TOGGLE_PAUSE,
+ RESEEK,
SEEK,
- ENTRY_BUFFERED, // Currently set entry is buffered.
+ SET_BUFFERED, // Currently set entry is buffered.
CLOSE,
// Internal.
CORK,
@@ -94,6 +99,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
case TOGGLE_PAUSE: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
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);
@@ -116,36 +122,59 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
break;
}
+ case RESEEK: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ bmu_client_reseek(&entry->client);
+ break;
+ }
case SEEK: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
bmu_client_seek(&entry->client, cmd->value.f);
break;
}
- case ENTRY_BUFFERED: {
+// case SET_BUFFERED: {
+// switch (cmd->value.i) {
+// case CAMU_SINK_AUDIO:
+// sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL);
+// break;
+//#ifndef CAMU_SINK_NO_VIDEO
+// case CAMU_SINK_VIDEO:
+// sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL);
+// break;
+//#endif
+// }
+// break;
+// }
+ case CLOSE: {
+ 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:
- sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL);
+ bmu_vcr_stream_cork(entry->audio.stream);
break;
#ifndef CAMU_SINK_NO_VIDEO
case CAMU_SINK_VIDEO:
- sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL);
+ bmu_vcr_stream_cork(entry->video.stream);
break;
#endif
}
break;
}
- case CLOSE: {
- sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
- return;
- }
- case CORK: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- bmu_client_cork(&entry->client, true);
- break;
- }
case UNCORK: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- bmu_client_cork(&entry->client, false);
+ 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;
}
}
@@ -173,18 +202,27 @@ static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
{
u8 state = entry->audio.state;
if (state == BUFFER_SET_OR_BUFFERED) {
- 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 = ENTRY_BUFFERED,
- .value.i = CAMU_SINK_AUDIO
- });
+ bool should_resume = ENTRY_VIDEO_EMPTY(entry) || entry->video.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;
+ }
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
.value.i = CAMU_SINK_AUDIO
});
- if (ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED) {
+ if (should_resume) {
camu_clock_resume(&entry->clock);
+ entry->state = ENTRY_BUFFERED;
}
+ 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
+ //});
state = BUFFER_ADDED;
} else if (state == BUFFER_CONFIGURED) {
state = BUFFER_SET_OR_BUFFERED;
@@ -204,12 +242,14 @@ static void audio_buffer_callback(void *userdata, u8 op)
case CAMU_BUFFER_CORK:
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = CORK,
+ .value.i = CAMU_SINK_AUDIO,
.opaque = entry
});
break;
case CAMU_BUFFER_UNCORK:
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = UNCORK,
+ .value.i = CAMU_SINK_AUDIO,
.opaque = entry
});
break;
@@ -234,10 +274,10 @@ static void audio_buffer_callback(void *userdata, u8 op)
.value.i = CAMU_SINK_AUDIO
});
#ifndef CAMU_SINK_NO_VIDEO
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = ENTRY_BUFFERED,
- .value.i = CAMU_SINK_VIDEO
- });
+ //queue_cmd(entry->sink, (struct camu_sink_cmd){
+ // .op = SET_BUFFERED,
+ // .value.i = CAMU_SINK_VIDEO
+ //});
#endif
}
break;
@@ -250,18 +290,27 @@ static void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
{
u8 state = entry->video.state;
if (state == BUFFER_SET_OR_BUFFERED) {
- 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 = ENTRY_BUFFERED,
- .value.i = CAMU_SINK_VIDEO
- });
+ 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;
+ }
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
.value.i = CAMU_SINK_VIDEO
});
- if (ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED) {
+ if (should_resume) {
camu_clock_resume(&entry->clock);
+ entry->state = ENTRY_BUFFERED;
}
+ 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
+ //});
state = BUFFER_ADDED;
} else if (state == BUFFER_CONFIGURED) {
state = BUFFER_SET_OR_BUFFERED;
@@ -281,12 +330,14 @@ static void video_buffer_callback(void *userdata, u8 op)
case CAMU_BUFFER_CORK:
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = CORK,
+ .value.i = CAMU_SINK_VIDEO,
.opaque = entry
});
break;
case CAMU_BUFFER_UNCORK:
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = UNCORK,
+ .value.i = CAMU_SINK_VIDEO,
.opaque = entry
});
break;
@@ -300,22 +351,10 @@ static void video_buffer_callback(void *userdata, u8 op)
}
#endif
-static void client_callback(void *userdata, u8 op, struct bmu_client_stream *stream, void *opaque)
+static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream, void *opaque)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
switch (op) {
- case BIMU_CLIENT_SET: {
- struct bmu_seek_req *req = (struct bmu_seek_req *)opaque;
- camu_clock_set(&entry->clock, req->base, req->start);
- camu_audio_buffer_reset(&entry->audio.buf);
-#ifndef CAMU_SINK_NO_VIDEO
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- if (!single_frame) {
- camu_video_buffer_reset(&entry->video.buf);
- }
-#endif
- break;
- }
case BIMU_CLIENT_CONFIGURE: {
// No data can be sent until all active streams are configured.
switch (stream->type) {
@@ -346,6 +385,23 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
}
break;
}
+ case BIMU_CLIENT_SET: {
+ u64 ts = aki_get_timestamp();
+ struct bmu_seek_req *req = (struct bmu_seek_req *)opaque;
+ //req->delay = 0; // Ignore delay.
+ req->ts -= req->delay;
+ bool late = ts <= req->ts;
+ if (late) {}
+ 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 BIMU_CLIENT_DATA: {
struct camu_frame *frame = (struct camu_frame *)opaque;
switch (stream->type) {
@@ -395,6 +451,11 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
#ifndef CAMU_SINK_NO_VIDEO
}
#endif
+#ifndef CAMU_SINK_NO_VIDEO
+ if (!single_frame) {
+ camu_video_buffer_reset(&entry->video.buf);
+ }
+#endif
break;
}
case BIMU_CLIENT_EOF: {
@@ -417,6 +478,7 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
}
case BIMU_CLIENT_CLOSED: {
bmu_client_free(&entry->client);
+ aki_timer_stop(&entry->timer);
camu_audio_buffer_free(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
camu_video_buffer_free(&entry->video.buf);
@@ -460,6 +522,15 @@ 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 bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn,
struct aki_packet *packet, struct aki_packet *rpacket)
{
@@ -492,11 +563,14 @@ static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *c
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);
aki_packet_free(packet);
+
return false;
}
@@ -524,6 +598,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
u16 node_id = aki_packet_read_u16(packet);
aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *entry = entry_from_node_id(sink, node_id);
+ if (!entry) goto out;
struct camu_sink_entry *previous = sink->current;
sink->current = entry;
if (entry->audio.state == BUFFER_INIT) {
@@ -542,8 +617,9 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
camu_clock_pause(&previous->clock);
remove_entry_buffers(sink, previous);
}
- aki_mutex_unlock(&sink->mutex);
+out:
+ aki_mutex_unlock(&sink->mutex);
aki_packet_free(packet);
return false;
}
@@ -644,13 +720,13 @@ void camu_sink_stop(struct camu_sink *sink)
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);
}
- */
- al_array_free(sink->entries);
aki_mutex_unlock(&sink->mutex);
+ al_array_free(sink->entries);
+ aki_rpc_disconnect(&sink->client);
aki_signal_stop(&sink->signal);
camu_queue_free(sink->queue);
aki_mutex_destroy(&sink->mutex);
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index 95c280f..e3550bb 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -18,8 +18,6 @@
#include "../bimu/client.h"
-#include "common.h"
-
enum {
CAMU_SINK_AUDIO = 0,
#ifndef CAMU_SINK_NO_VIDEO
@@ -41,16 +39,17 @@ struct camu_sink_entry {
u8 state;
struct bmu_client client;
struct camu_clock clock;
+ struct aki_timer timer;
struct {
u8 state;
struct camu_audio_buffer buf;
- struct bmu_client_stream *stream;
+ struct bmu_vcr_stream *stream;
} audio;
#ifndef CAMU_SINK_NO_VIDEO
struct {
u8 state;
struct camu_video_buffer buf;
- struct bmu_client_stream *stream;
+ struct bmu_vcr_stream *stream;
} video;
#endif
struct camu_sink *sink;