diff options
| author | 2024-04-09 11:24:01 -0400 | |
|---|---|---|
| committer | 2024-04-09 11:24:01 -0400 | |
| commit | 02f3d3565602146bbbfce85b2719246f24036cb9 (patch) | |
| tree | c6588ffe297b777e36260effa4fa42b958ca6ba3 /src/libsink | |
| parent | bbf3314165182e402ff25acccddc004a87f81ef0 (diff) | |
| download | camu-02f3d3565602146bbbfce85b2719246f24036cb9.tar.gz camu-02f3d3565602146bbbfce85b2719246f24036cb9.tar.bz2 camu-02f3d3565602146bbbfce85b2719246f24036cb9.zip | |
Massive restructure and many changes
- The server-side list concept is still a wip
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/libsink')
| -rw-r--r-- | src/libsink/common.h | 5 | ||||
| -rw-r--r-- | src/libsink/meson.build | 2 | ||||
| -rw-r--r-- | src/libsink/sink.c | 405 | ||||
| -rw-r--r-- | src/libsink/sink.h | 22 |
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); |