diff options
| author | 2024-05-25 20:22:59 -0400 | |
|---|---|---|
| committer | 2024-05-25 20:22:59 -0400 | |
| commit | 2f9a0945bfeee3296cec3d38d094e4c49f9cb65f (patch) | |
| tree | ce80b220accbb44920932b1a2a6ca49ed31752a4 /src/libsink | |
| parent | 90e0930d820972802d1bb2fb27d7dded304483d6 (diff) | |
| download | camu-2f9a0945bfeee3296cec3d38d094e4c49f9cb65f.tar.gz camu-2f9a0945bfeee3296cec3d38d094e4c49f9cb65f.tar.bz2 camu-2f9a0945bfeee3296cec3d38d094e4c49f9cb65f.zip | |
wip
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/libsink')
| -rw-r--r-- | src/libsink/sink.c | 110 | ||||
| -rw-r--r-- | src/libsink/sink.h | 2 |
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); |