diff options
| author | 2024-10-21 19:22:50 -0400 | |
|---|---|---|
| committer | 2024-10-21 19:22:50 -0400 | |
| commit | 60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2 (patch) | |
| tree | 08ff2ce7975f523112e7ad2fe4f797b4fc7db5de /src/liana | |
| parent | 2f9a0945bfeee3296cec3d38d094e4c49f9cb65f (diff) | |
| download | camu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.tar.gz camu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.tar.bz2 camu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.zip | |
Everything before initial synced list
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 199 | ||||
| -rw-r--r-- | src/liana/client.h | 31 | ||||
| -rw-r--r-- | src/liana/common.h | 1 | ||||
| -rw-r--r-- | src/liana/handler.h | 58 | ||||
| -rw-r--r-- | src/liana/handlers.c | 37 | ||||
| -rw-r--r-- | src/liana/handlers.h | 13 | ||||
| -rw-r--r-- | src/liana/handlers/cdio.h | 18 | ||||
| -rw-r--r-- | src/liana/handlers/cdio_client.c | 49 | ||||
| -rw-r--r-- | src/liana/handlers/cdio_server.c | 106 | ||||
| -rw-r--r-- | src/liana/handlers/codec.h | 19 | ||||
| -rw-r--r-- | src/liana/handlers/codec_client.c | 116 | ||||
| -rw-r--r-- | src/liana/handlers/codec_server.c | 143 | ||||
| -rw-r--r-- | src/liana/handlers/dvd.h | 14 | ||||
| -rw-r--r-- | src/liana/handlers/dvd_server.c | 21 | ||||
| -rw-r--r-- | src/liana/list.c | 413 | ||||
| -rw-r--r-- | src/liana/list.h | 107 | ||||
| -rw-r--r-- | src/liana/list_cmp.h | 53 | ||||
| -rw-r--r-- | src/liana/meson.build | 36 | ||||
| -rw-r--r-- | src/liana/server.c | 345 | ||||
| -rw-r--r-- | src/liana/server.h | 42 | ||||
| -rw-r--r-- | src/liana/vcr.c | 232 | ||||
| -rw-r--r-- | src/liana/vcr.h | 46 |
22 files changed, 2099 insertions, 0 deletions
diff --git a/src/liana/client.c b/src/liana/client.c new file mode 100644 index 0000000..d556453 --- /dev/null +++ b/src/liana/client.c @@ -0,0 +1,199 @@ +#include "../server/common.h" +#ifdef CAMU_HAVE_FFMPEG +#include "../codec/ffmpeg/packet_ext.h" +#endif + +#include "client.h" +#include "handlers.h" + +static void data_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +{ + struct lia_client *client = (struct lia_client *)userdata; + (void)stream; + lia_vcr_push_packet(&client->vcr, packet); +} + +static void parse_info_packet(struct lia_client *client, struct aki_packet *packet) +{ + str liana; + aki_packet_read_str(packet, &liana); + client->duration = aki_packet_read_u64(packet); + u32 count = aki_packet_read_u32(packet); + for (u32 i = 0; i < count; i++) { + u8 mode = aki_packet_read_u8(packet); + struct lia_vcr_track *track = NULL; + switch (mode) { + case CAMU_NORMAL: { + client->mask |= 1 << 0; + track = al_alloc_object(struct lia_vcr_track); + track->stream.mode = mode; + u8 type = aki_packet_read_u8(packet); + track->stream.type = type; + if (type == CAMU_STREAM_AUDIO) { + struct camu_audio_format *fmt = &track->stream.audio.fmt; + fmt->format = aki_packet_read_s32(packet); + fmt->sample_rate = aki_packet_read_s32(packet); + fmt->channel_count = aki_packet_read_s32(packet); +#ifdef CAMU_HAVE_FFMPEG + av_channel_layout_default(&fmt->channel_layout, fmt->channel_count); +#endif + } else if (type == CAMU_STREAM_VIDEO) { + track->stream.video.width = aki_packet_read_s32(packet); + track->stream.video.height = aki_packet_read_s32(packet); + track->stream.video.format = aki_packet_read_s32(packet); + } + track->index = 0; + break; + } +#ifdef CAMU_HAVE_FFMPEG + case CAMU_FFMPEG_COMPAT: { + const AVCodec *codec = avcodec_find_decoder(aki_packet_read_av_codec_id(packet)); + AVFormatContext *format_context = avformat_alloc_context(); + AVStream *stream = aki_packet_read_av_stream(format_context, codec, packet); + s32 index = stream->index; + al_assert(index < 32); + switch (stream->codecpar->codec_type) { + case AVMEDIA_TYPE_AUDIO: + client->mask |= 1 << index; + break; + case AVMEDIA_TYPE_VIDEO: + client->mask |= 1 << index; + break; + case AVMEDIA_TYPE_SUBTITLE: + default: + continue; + } + track = al_alloc_object(struct lia_vcr_track); + track->stream.mode = CAMU_FFMPEG_COMPAT; + track->stream.type = stream->codecpar->codec_type; + track->stream.av.format_context = format_context; + track->stream.av.stream = stream; + if (track->stream.type == CAMU_STREAM_AUDIO) { + struct camu_audio_format *fmt = &track->stream.audio.fmt; + fmt->format = stream->codecpar->format; + fmt->sample_rate = stream->codecpar->sample_rate; + av_channel_layout_copy(&fmt->channel_layout, &stream->codecpar->ch_layout); + fmt->channel_count = stream->codecpar->ch_layout.nb_channels; + } + track->index = index; + break; + } +#endif + } + track->client = lia_handler_by_name(&liana)->create_client_handler(); + track->client->callback = client->callback; + track->client->userdata = client->userdata; + track->stream.mode = mode; + if (track->client->init(track->client, client->renderer, &track->stream)) { + client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, track->client->stream, track); + } + lia_vcr_add_track(&client->vcr, track); + } +} + +static void info_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +{ + struct lia_client *client = (struct lia_client *)userdata; + client->connection_id = aki_packet_read_u16(packet); + parse_info_packet(client, packet); + al_assert(client->mask != 0); + aki_packet_free(packet); + stream->packet_callback = data_packet_callback; + struct aki_packet *rpacket = aki_packet_create(); + aki_packet_write_s32(rpacket, client->mask); + aki_packet_stream_send_packet(stream, rpacket); +} + +static void packet_sent_callback(void *userdata, struct aki_packet *packet) +{ + (void)userdata; + aki_packet_free(packet); +} + +static void connection_callback(void *userdata, struct aki_packet_stream *stream) +{ + struct lia_client *client = (struct lia_client *)userdata; + stream->packet_sent_callback = packet_sent_callback; + struct aki_packet *packet = aki_packet_create(); + aki_packet_write_u16(packet, client->id); + aki_packet_write_u16(packet, 0); + aki_packet_write_s32(packet, client->mask); + aki_packet_write_u64(packet, client->pos); + if (client->mask == 0) { + stream->packet_callback = info_packet_callback; + } else { + stream->packet_callback = data_packet_callback; + } + aki_packet_stream_send_packet(stream, packet); +} + +static void connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +{ + struct lia_client *client = (struct lia_client *)userdata; + if (client->reconnect) { + lia_vcr_flush(&client->vcr); + } else { + lia_vcr_close_all(&client->vcr); + } + // REMOVE_BUFFERS can possibly run the event loop while waiting. + // This should be accounted for in the client code here to not cause + // any unexpected behavior. + client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect); + if (client->reconnect) { + client->reconnect = false; + aki_packet_stream_reconnect(stream, &client->addr, client->port); + } else { + client->callback(client->userdata, LIANA_CLIENT_CLOSED, NULL, NULL); + } +} + +void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop, u8 type, + str *addr, u16 port, u16 id, u64 pos, struct camu_renderer *renderer) +{ + client->loop = loop; + client->id = id; + client->pos = pos; + client->mask = 0; + client->reconnect = false; + lia_vcr_init(&client->vcr, &client->data); + // TODO: Starting and stopping of vcr could be more clear. + lia_vcr_start(&client->vcr, client->loop); + al_str_clone(&client->addr, addr); + client->port = port; + if (!aki_packet_stream_init(&client->data, type, connection_callback, connection_closed_callback, client)) { + connection_closed_callback(client, &client->data); + } + aki_packet_stream_set_multiplex(&client->data, CAMU_MULTIPLEX_LIANA); + client->renderer = renderer; + aki_packet_stream_connect(&client->data, client->loop, &client->addr, client->port); +} + +void lia_client_set_renderer(struct lia_client *client, struct camu_renderer *renderer) +{ + client->renderer = renderer; +} + +void lia_client_seek(struct lia_client *client, u64 pos) +{ + client->reconnect = true; + client->pos = pos; + aki_packet_stream_disconnect(&client->data); +} + +void lia_client_reseek(struct lia_client *client) +{ + (void)client; +} + +void lia_client_disconnect(struct lia_client *client) +{ + client->reconnect = false; + aki_packet_stream_disconnect(&client->data); +} + +void lia_client_free(struct lia_client *client) +{ + lia_vcr_free(&client->vcr); + aki_packet_stream_free(&client->data); + al_str_free(&client->addr); +} diff --git a/src/liana/client.h b/src/liana/client.h new file mode 100644 index 0000000..5c16cd8 --- /dev/null +++ b/src/liana/client.h @@ -0,0 +1,31 @@ +#pragma once + +#include <aki/packet_stream.h> + +#include "../codec/codec.h" + +#include "vcr.h" + +struct lia_client { + struct aki_event_loop *loop; + u16 id; + s32 mask; + u64 pos; + bool reconnect; + str addr; + u16 port; + u16 connection_id; + struct aki_packet_stream data; + u64 duration; + struct lia_vcr vcr; + struct camu_renderer *renderer; + void (*callback)(void *, u8, struct camu_codec_stream *stream, void *); + void *userdata; +}; + +void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop, u8 type, + str *addr, u16 port, u16 id, u64 pos, struct camu_renderer *renderer); +void lia_client_seek(struct lia_client *client, u64 pos); +void lia_client_reseek(struct lia_client *client); +void lia_client_disconnect(struct lia_client *client); +void lia_client_free(struct lia_client *client); diff --git a/src/liana/common.h b/src/liana/common.h new file mode 100644 index 0000000..6f70f09 --- /dev/null +++ b/src/liana/common.h @@ -0,0 +1 @@ +#pragma once diff --git a/src/liana/handler.h b/src/liana/handler.h new file mode 100644 index 0000000..6c23ac5 --- /dev/null +++ b/src/liana/handler.h @@ -0,0 +1,58 @@ +#pragma once + +#include <aki/packet_pool.h> + +#include "../codec/codec.h" +#include "../cache/handle.h" + +struct lia_server_handler { + bool (*init)(struct lia_server_handler *, struct cch_handle *); + void (*write_info)(struct lia_server_handler *, struct aki_packet *); + void (*subscribe)(struct lia_server_handler *, s32); + u64 (*get_duration)(struct lia_server_handler *); + bool (*seek)(struct lia_server_handler *, u64); + void (*step)(struct lia_server_handler *); + void (*write_packet)(struct lia_server_handler *, struct aki_packet *); + void (*free)(struct lia_server_handler **); + s32 status; +}; + +enum { + LIANA_PACKET_DATA = 0, + LIANA_PACKET_EOF, + LIANA_PACKET_ERROR +}; + +enum { + LIANA_CLIENT_CONFIGURE = 0, + LIANA_CLIENT_SET, + LIANA_CLIENT_PAUSE, + LIANA_CLIENT_RESUME, + LIANA_CLIENT_DATA, + LIANA_CLIENT_REMOVE_BUFFERS, + LIANA_CLIENT_EOF, + LIANA_CLIENT_CLOSED +}; + +struct lia_resume_req { + u64 start; + u64 offset; +}; + +struct lia_seek_req { + u64 base; + u64 ts; + u64 delay; + bool paused; + u64 paused_at; +}; + +struct lia_client_handler { + bool (*init)(struct lia_client_handler *, struct camu_renderer *, struct camu_codec_stream *); + bool (*handle_packet)(struct lia_client_handler *, struct aki_packet *); + void (*flush)(struct lia_client_handler *); + void (*free)(struct lia_client_handler **); + struct camu_codec_stream *stream; + void (*callback)(void *, u8, struct camu_codec_stream *stream, void *); + void *userdata; +}; diff --git a/src/liana/handlers.c b/src/liana/handlers.c new file mode 100644 index 0000000..e67250a --- /dev/null +++ b/src/liana/handlers.c @@ -0,0 +1,37 @@ +#include "handlers.h" + +#include "handlers/codec.h" +#include "handlers/cdio.h" + +struct lia_handler_entry liana_handlers[] = { + { + .name = al_str_c("codec"), +#ifdef LIANA_SERVER + .create_server_handler = lia_codec_server_create, +#endif +#ifdef LIANA_CLIENT + .create_client_handler = lia_codec_client_create +#endif + }, +#ifdef LIANA_HAVE_CDIO + { + .name = al_str_c("cdio"), +#ifdef LIANA_SERVER + .create_server_handler = lia_cdio_server_create, +#endif +#ifdef LIANA_CLIENT + .create_client_handler = lia_cdio_client_create +#endif + } +#endif +}; + +struct lia_handler_entry *lia_handler_by_name(str *name) +{ + for (u32 i = 0; i < AL_ARRAY_SIZE(liana_handlers); i++) { + if (al_str_eq(liana_handlers[i].name, name)) { + return &liana_handlers[i]; + } + } + return NULL; +} diff --git a/src/liana/handlers.h b/src/liana/handlers.h new file mode 100644 index 0000000..0193f18 --- /dev/null +++ b/src/liana/handlers.h @@ -0,0 +1,13 @@ +#pragma once + +#include "handler.h" + +struct lia_handler_entry { + str *name; + struct lia_server_handler *(*create_server_handler)(void); + struct lia_client_handler *(*create_client_handler)(void); +}; + +extern struct lia_handler_entry liana_handlers[]; + +struct lia_handler_entry *lia_handler_by_name(str *name); diff --git a/src/liana/handlers/cdio.h b/src/liana/handlers/cdio.h new file mode 100644 index 0000000..2727616 --- /dev/null +++ b/src/liana/handlers/cdio.h @@ -0,0 +1,18 @@ +#pragma once + +#include "../handler.h" + +struct lia_cdio_server { + struct lia_server_handler handler; + struct cch_handle *handle; + struct camu_audio_format fmt; + f64 pts; + struct aki_buffer buffer; +}; + +struct lia_cdio_client { + struct lia_client_handler handler; +}; + +struct lia_server_handler *lia_cdio_server_create(); +struct lia_client_handler *lia_cdio_client_create(void); diff --git a/src/liana/handlers/cdio_client.c b/src/liana/handlers/cdio_client.c new file mode 100644 index 0000000..93b1a27 --- /dev/null +++ b/src/liana/handlers/cdio_client.c @@ -0,0 +1,49 @@ +#include "cdio.h" + +static bool cdio_client_init(struct lia_client_handler *handler, struct camu_renderer *renderer, + struct camu_codec_stream *stream) +{ + struct lia_cdio_client *cdio = (struct lia_cdio_client *)handler; + (void)renderer; + cdio->handler.stream = stream; + return true; +} + +static bool cdio_client_handle_packet(struct lia_client_handler *handler, struct aki_packet *packet) +{ + struct lia_cdio_client *cdio = (struct lia_cdio_client *)handler; + if (!packet) { + cdio->handler.callback(cdio->handler.userdata, LIANA_CLIENT_EOF, cdio->handler.stream, NULL); + return true; + } + struct camu_codec_frame *frame = al_alloc_object(struct camu_codec_frame); + frame->mode = CAMU_NORMAL; + frame->pts = aki_packet_read_f64(packet); + frame->audio.sample_count = aki_packet_read_s32(packet); + struct aki_buffer buffer; + aki_packet_read_buffer(packet, &buffer); + frame->data = al_malloc(buffer.size); + al_memcpy(frame->data, aki_buffer_get_ptr(&buffer, 0), buffer.size); + cdio->handler.callback(cdio->handler.userdata, LIANA_CLIENT_DATA, cdio->handler.stream, frame); + return true; +} + +static void cdio_client_flush(struct lia_client_handler *handler) +{ + (void)handler; +} + +static void cdio_client_free(struct lia_client_handler **handler) +{ + (void)handler; +} + +struct lia_client_handler *lia_cdio_client_create(void) +{ + struct lia_cdio_client *cdio = al_alloc_object(struct lia_cdio_client); + cdio->handler.init = cdio_client_init; + cdio->handler.handle_packet = cdio_client_handle_packet; + cdio->handler.flush = cdio_client_flush; + cdio->handler.free = cdio_client_free; + return (struct lia_client_handler *)cdio; +} diff --git a/src/liana/handlers/cdio_server.c b/src/liana/handlers/cdio_server.c new file mode 100644 index 0000000..1eaf6f3 --- /dev/null +++ b/src/liana/handlers/cdio_server.c @@ -0,0 +1,106 @@ +#include <al/log.h> +#include <cdio/paranoia/cdda.h> + +#include "../../cache/entry.h" +#include "../../cache/handlers/cdio.h" + +#include "cdio.h" + +#define SECTORS_PER_PACKET 8 + +static bool cdio_server_init(struct lia_server_handler *handler, struct cch_handle *handle) +{ + struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; + cdio->handle = handle; + + cdio->fmt.format = CAMU_SAMPLE_FORMAT_S16; + cdio->fmt.sample_rate = 44100; + cdio->fmt.channel_count = 2; + cdio->pts = 0.0; + aki_buffer_init(&cdio->buffer); + + if (handle->entry->chapter) { + cch_handle_seek(handle, handle->entry->chapter->start * CDIO_CD_FRAMESIZE_RAW, SEEK_SET); + } + + return true; +} + +static void cdio_server_write_info(struct lia_server_handler *handler, struct aki_packet *packet) +{ + struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; + aki_packet_write_u64(packet, 0); + aki_packet_write_u32(packet, 1); + aki_packet_write_u8(packet, CAMU_NORMAL); + aki_packet_write_u8(packet, CAMU_STREAM_AUDIO); + aki_packet_write_s32(packet, cdio->fmt.format); + aki_packet_write_s32(packet, cdio->fmt.sample_rate); + aki_packet_write_s32(packet, cdio->fmt.channel_count); +} + +static void cdio_server_subscribe(struct lia_server_handler *handler, s32 mask) +{ + (void)handler; + (void)mask; +} + +static u64 cdio_server_get_duration(struct lia_server_handler *handler) +{ + (void)handler; + return 0; +} + +static bool cdio_server_seek(struct lia_server_handler *handler, u64 pos) +{ + struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; + return cch_handle_seek(cdio->handle, pos, SEEK_SET); +} + +static void cdio_server_step(struct lia_server_handler *handler) +{ + struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; + (void)cdio; +} + +static void cdio_server_write_packet(struct lia_server_handler *handler, struct aki_packet *packet) +{ + struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; + off_t size = CDIO_CD_FRAMESIZE_RAW * SECTORS_PER_PACKET; + aki_buffer_ensure_space(&cdio->buffer, size); + s32 ret = cch_handle_read(cdio->handle, aki_buffer_get_ptr(&cdio->buffer, 0), size); + if (ret == CAMU_ERR_EOF) { + cdio->handler.status = CAMU_ERR_EOF; + aki_packet_write_u8(packet, LIANA_PACKET_EOF); + return; + } + cdio->handler.status = CAMU_OK; + aki_packet_write_u8(packet, LIANA_PACKET_DATA); + aki_packet_write_s32(packet, 0); + aki_packet_write_f64(packet, cdio->pts); + s32 sample_count = ret / (camu_audio_format_bytes_per_sample(&cdio->fmt) * cdio->fmt.channel_count); + cdio->pts += camu_audio_format_samples_to_sec(&cdio->fmt, sample_count); + aki_packet_write_s32(packet, sample_count); + aki_buffer_set_size(&cdio->buffer, ret); + aki_packet_write_buffer(packet, &cdio->buffer); +} + +static void cdio_server_free(struct lia_server_handler **handler) +{ + struct lia_cdio_server *cdio = (struct lia_cdio_server *)*handler; + al_free(cdio); + *handler = NULL; +} + +struct lia_server_handler *lia_cdio_server_create() +{ + struct lia_cdio_server *cdio = al_alloc_object(struct lia_cdio_server); + cdio->handler.init = cdio_server_init; + cdio->handler.write_info = cdio_server_write_info; + cdio->handler.subscribe = cdio_server_subscribe; + cdio->handler.get_duration = cdio_server_get_duration; + cdio->handler.seek = cdio_server_seek; + cdio->handler.step = cdio_server_step; + cdio->handler.write_packet = cdio_server_write_packet; + cdio->handler.free = cdio_server_free; + return (struct lia_server_handler *)cdio; +} diff --git a/src/liana/handlers/codec.h b/src/liana/handlers/codec.h new file mode 100644 index 0000000..8c5528f --- /dev/null +++ b/src/liana/handlers/codec.h @@ -0,0 +1,19 @@ +#pragma once + +#include "../../codec/codec.h" + +#include "../handler.h" + +struct lia_codec_server { + struct lia_server_handler handler; + struct camu_demuxer *demux; + struct camu_codec_packet packet; +}; + +struct lia_codec_client { + struct lia_client_handler handler; + struct camu_decoder *dec; +}; + +struct lia_server_handler *lia_codec_server_create(void); +struct lia_client_handler *lia_codec_client_create(void); diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c new file mode 100644 index 0000000..050cc6f --- /dev/null +++ b/src/liana/handlers/codec_client.c @@ -0,0 +1,116 @@ +#include "../../codec/codec.h" +#include "../../codec/ffmpeg/decoder.h" +#include "../../codec/ffmpeg/packet_ext.h" +#include "../../codec/stb_image/decoder.h" +#include "../../codec/spng/decoder.h" +#include "../../codec/wuffs/decoder.h" + +#include "../vcr.h" + +#include "codec.h" + +static void data_callback(void *userdata, struct camu_codec_frame *frame) +{ + struct lia_codec_client *codec = (struct lia_codec_client *)userdata; + codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_DATA, codec->handler.stream, frame); +} + +static bool codec_client_init(struct lia_client_handler *handler, struct camu_renderer *renderer, + struct camu_codec_stream *stream) +{ + struct lia_codec_client *codec = (struct lia_codec_client *)handler; + codec->dec = camu_ff_decoder_create(); + //codec->dec = camu_stbi_decoder_create(); + //codec->dec = camu_spng_decoder_create(); + //codec->dec = camu_wuffs_decoder_create(); + codec->handler.stream = stream; + if (!codec->dec->init(codec->dec, renderer, stream, data_callback, codec)) { + return false; + } + return true; +} + +#ifdef CAMU_HAVE_FFMPEG +static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt) +{ + struct camu_codec_packet packet; + packet.av.pkt = pkt; + s32 ret = codec->dec->push(codec->dec, &packet); + av_packet_unref(pkt); + av_packet_free(&pkt); + return ret == CAMU_OK; +} +#endif + +static bool push_packet(struct lia_codec_client *codec, struct aki_buffer *buffer) +{ + struct camu_codec_packet packet; + packet.buffer = buffer; + s32 ret = codec->dec->push(codec->dec, &packet); + return ret == CAMU_OK; +} + +static bool codec_client_handle_packet(struct lia_client_handler *handler, struct aki_packet *packet) +{ + struct lia_codec_client *codec = (struct lia_codec_client *)handler; + if (!packet) { + struct lia_codec_client *codec = (struct lia_codec_client *)handler; + s32 ret = codec->dec->push(codec->dec, NULL); + // Flush returns success. + ret = codec->dec->process(codec->dec); + // process() could still error. + codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); + return ret == CAMU_ERR_EOF; + } + bool success = false; + u8 type = aki_packet_read_u8(packet); + switch (type) { + case CAMU_NORMAL: { + struct aki_buffer buffer; + aki_packet_read_buffer(packet, &buffer); + success = push_packet(codec, &buffer); + break; + } +#ifdef CAMU_HAVE_FFMPEG + case CAMU_FFMPEG_COMPAT: { + success = push_av_packet(codec, aki_packet_read_av_packet(packet)); + break; + } +#endif + } + // Forcing in EOF on errors is not necessary but should be a better experience client-side. + if (!success) { + codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); + return false; + } + s32 ret = codec->dec->process(codec->dec); + if (!(ret == CAMU_ERR_AGAIN || ret == CAMU_ERR_EOF)) { + codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); + return false; + } + return true; +} + +static void codec_client_flush(struct lia_client_handler *handler) +{ + struct lia_codec_client *codec = (struct lia_codec_client *)handler; + codec->dec->flush(codec->dec); +} + +static void codec_client_free(struct lia_client_handler **handler) +{ + struct lia_codec_client *codec = (struct lia_codec_client *)*handler; + codec->dec->free(&codec->dec); + al_free(codec); + *handler = NULL; +} + +struct lia_client_handler *lia_codec_client_create(void) +{ + struct lia_codec_client *codec = al_alloc_object(struct lia_codec_client); + codec->handler.init = codec_client_init; + codec->handler.handle_packet = codec_client_handle_packet; + codec->handler.flush = codec_client_flush; + codec->handler.free = codec_client_free; + return (struct lia_client_handler *)codec; +} diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c new file mode 100644 index 0000000..07fafdf --- /dev/null +++ b/src/liana/handlers/codec_server.c @@ -0,0 +1,143 @@ +#include "../../codec/codec.h" +#include "../../codec/ffmpeg/demuxer.h" +#include "../../codec/ffmpeg/packet_ext.h" +#include "../../codec/stb_image/demuxer.h" +#include "../../codec/spng/demuxer.h" +#include "../../codec/wuffs/demuxer.h" + +#include "../handler.h" + +#include "codec.h" + +static bool codec_server_init(struct lia_server_handler *handler, struct cch_handle *handle) +{ + struct lia_codec_server *codec = (struct lia_codec_server *)handler; + codec->demux = camu_ff_demuxer_create(); + //codec->demux = camu_stbi_demuxer_create(); + //codec->demux = camu_spng_demuxer_create(); + //codec->demux = camu_wuffs_demuxer_create(); + if (!codec->demux->init(codec->demux, handle)) { + return false; + } + switch (codec->demux->mode) { + case CAMU_NORMAL: + break; +#ifdef CAMU_HAVE_FFMPEG + case CAMU_FFMPEG_COMPAT: + codec->packet.av.pkt = av_packet_alloc(); + break; +#endif + } + return true; +} + +static void codec_server_write_info(struct lia_server_handler *handler, struct aki_packet *packet) +{ + struct lia_codec_server *codec = (struct lia_codec_server *)handler; + u64 duration = codec->demux->get_duration(codec->demux); + aki_packet_write_u64(packet, duration); + aki_packet_write_u32(packet, codec->demux->streams.size); + struct camu_codec_stream *stream; + al_array_foreach_ptr(codec->demux->streams, i, stream) { + aki_packet_write_u8(packet, stream->mode); + switch (stream->mode) { + case CAMU_NORMAL: + aki_packet_write_u8(packet, stream->type); + aki_packet_write_s32(packet, stream->video.width); + aki_packet_write_s32(packet, stream->video.height); + aki_packet_write_s32(packet, stream->video.format); + break; +#ifdef CAMU_HAVE_FFMPEG + case CAMU_FFMPEG_COMPAT: + aki_packet_write_av_codec_id(packet, stream->av.stream->codecpar->codec_id); + aki_packet_write_av_stream(packet, stream->av.stream); + break; +#endif + } + } +} + +static void codec_server_subscribe(struct lia_server_handler *handler, s32 mask) +{ + struct lia_codec_server *codec = (struct lia_codec_server *)handler; + codec->demux->subscribed = mask; +} + +static u64 codec_server_get_duration(struct lia_server_handler *handler) +{ + struct lia_codec_server *codec = (struct lia_codec_server *)handler; + return codec->demux->get_duration(codec->demux); +} + +static bool codec_server_seek(struct lia_server_handler *handler, u64 pos) +{ + struct lia_codec_server *codec = (struct lia_codec_server *)handler; + return codec->demux->seek(codec->demux, pos); +} + +static void codec_server_step(struct lia_server_handler *handler) +{ + struct lia_codec_server *codec = (struct lia_codec_server *)handler; + codec->handler.status = codec->demux->get_packet(codec->demux, &codec->packet); +} + +static void codec_server_write_packet(struct lia_server_handler *handler, struct aki_packet *packet) +{ + struct lia_codec_server *codec = (struct lia_codec_server *)handler; + if (codec->handler.status == CAMU_OK) { + aki_packet_write_u8(packet, LIANA_PACKET_DATA); + switch (codec->packet.mode) { + case CAMU_NORMAL: { + aki_packet_write_s32(packet, 0); + aki_packet_write_u8(packet, codec->packet.mode); + aki_packet_write_buffer(packet, codec->packet.buffer); + break; + } +#ifdef CAMU_HAVE_FFMPEG + case CAMU_FFMPEG_COMPAT: { + AVPacket *pkt = codec->packet.av.pkt; + aki_packet_write_s32(packet, pkt->stream_index); + aki_packet_write_u8(packet, codec->packet.mode); + aki_packet_write_av_packet(packet, pkt); + av_packet_unref(pkt); + break; + } +#endif + } + } else if (codec->handler.status == CAMU_ERR_EOF) { + aki_packet_write_u8(packet, LIANA_PACKET_EOF); + } else { + aki_packet_write_u8(packet, LIANA_PACKET_ERROR); + } +} + +static void codec_server_free(struct lia_server_handler **handler) +{ + struct lia_codec_server *codec = (struct lia_codec_server *)*handler; + switch (codec->demux->mode) { + case CAMU_NORMAL: + break; +#ifdef CAMU_HAVE_FFMPEG + case CAMU_FFMPEG_COMPAT: + av_packet_free(&codec->packet.av.pkt); + break; +#endif + } + codec->demux->free(&codec->demux); + al_free(codec); + *handler = NULL; +} + +struct lia_server_handler *lia_codec_server_create(void) +{ + struct lia_codec_server *codec = al_alloc_object(struct lia_codec_server); + codec->handler.init = codec_server_init; + codec->handler.write_info = codec_server_write_info; + codec->handler.subscribe = codec_server_subscribe; + codec->handler.get_duration = codec_server_get_duration; + codec->handler.seek = codec_server_seek; + codec->handler.step = codec_server_step; + codec->handler.write_packet = codec_server_write_packet; + codec->handler.free = codec_server_free; + return (struct lia_server_handler *)codec; +} diff --git a/src/liana/handlers/dvd.h b/src/liana/handlers/dvd.h new file mode 100644 index 0000000..788c30e --- /dev/null +++ b/src/liana/handlers/dvd.h @@ -0,0 +1,14 @@ +#pragma once + +#include "../handler.h" + +struct lia_dvd_server { + struct lia_server_handler handler; +}; + +struct lia_dvd_client { + struct lia_client_handler handler; +}; + +struct lia_server_handler *lia_dvd_server_create(void); +struct lia_client_handler *lia_dvd_client_create(void); diff --git a/src/liana/handlers/dvd_server.c b/src/liana/handlers/dvd_server.c new file mode 100644 index 0000000..a487de4 --- /dev/null +++ b/src/liana/handlers/dvd_server.c @@ -0,0 +1,21 @@ +#include "dvd.h" + +static bool dvd_server_init(struct lia_server_handler *handler, struct cch_handle *handle) +{ + struct lia_dvd_server *dvd = (struct lia_dvd_server *)handler; + return true; +} + +struct lia_server_handler *lia_dvd_server_create(void) +{ + struct lia_dvd_server *dvd = al_alloc_object(struct lia_dvd_server); + dvd->handler.init = dvd_server_init; +// dvd->handler.write_info = dvd_server_write_info; +// dvd->handler.subscribe = dvd_server_subscribe; +// dvd->handler.get_duration = dvd_server_get_duration; +// dvd->handler.seek = dvd_server_seek; +// dvd->handler.step = dvd_server_step; +// dvd->handler.write_packet = dvd_server_write_packet; +// dvd->handler.free = dvd_server_free; + return (struct lia_server_handler *)dvd; +} diff --git a/src/liana/list.c b/src/liana/list.c new file mode 100644 index 0000000..37a144f --- /dev/null +++ b/src/liana/list.c @@ -0,0 +1,413 @@ +#include <aki/thread.h> +#include <al/random.h> +#include <al/lib.h> +#include <al/log.h> + +#include "../libsink/common.h" + +#include "list.h" +#include "list_cmp.h" + +void lia_list_init(struct lia_list *list) +{ + list->current = -1; + list->previous = -1; + list->queued = -1; + list->idle = true; + al_array_init(list->entries); + al_array_init(list->sinks); +} + +/* +static void buffer_ahead(struct lia_list *list) +{ + s32 size = (s32)list->entries.size; + if (list->queued >= 0 && list->queued + 1 < size) { + s32 ahead = list->queued + 1; + for (s32 i = ahead; i < AL_MIN(ahead + LIANA_BUFFER_AHEAD, size); i++) { + struct lia_list_entry *buffered = al_array_at(list->entries, i); + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->callback(sink->userdata, CAMU_SINK_BUFFER, buffered, i, LIANA_TIMESTAMP_INVALID); + } + } + } +} +*/ + +static void unset_all(struct lia_list *list) +{ + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->set = -1; + sink->queued = -1; + } +} + +void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *), void *userdata) +{ + struct lia_list_sink *sink = al_alloc_object(struct lia_list_sink); + sink->callback = callback; + sink->userdata = userdata; + al_array_push(list->sinks, sink); + if (list->current >= 0) { + sink->set = list->current; + struct lia_list_entry *current = al_array_at(list->entries, list->current); + u64 now = aki_get_timestamp(); + u8 pause; + u64 at = LIANA_TIMESTAMP_INVALID; + u64 seek_pos = current->offset; + if (current->paused_at == LIANA_TIMESTAMP_INVALID) { + now += LIANA_BASE_DELAY; + if (now > current->start) { + at = now; + seek_pos += now - current->start; + } else { + at = current->start; + } + pause = LIANA_PAUSE_RESUME; + } else { + pause = LIANA_PAUSE_NONE; + } + struct lia_timing time = { + .at = at, + .seek_pos = seek_pos, + .pause = pause + }; + sink->callback(sink->userdata, CAMU_SINK_SET, current, list->current, &time); + } else { + sink->set = -1; + } + sink->queued = -1; +} + +void lia_list_unset(struct lia_list *list) +{ + unset_all(list); + list->current = list->entries.size - 1; + list->previous = -1; + list->queued = -1; + list->idle = true; + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->callback(sink->userdata, CAMU_SINK_CLEAR, NULL, -1, NULL); + } +} + +void lia_list_remove_sink(struct lia_list *list, void *userdata) +{ + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + if (sink->userdata == userdata) { + al_array_remove_at(list->sinks, i); + al_free(sink); + break; + } + } +} + +void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) +{ + struct lia_list_entry *entry = al_alloc_object(struct lia_list_entry); + entry->opaque = opaque; + entry->paused_at = LIANA_TIMESTAMP_INVALID; + entry->held = false; + entry->offset = 0; + entry->duration = duration; + al_str_clone(&entry->name, name); + al_array_push(list->entries, entry); + if (list->idle) { + list->current++; + list->idle = false; + entry->start = aki_get_timestamp() + LIANA_BASE_DELAY; + struct lia_timing time = { + .at = entry->start, + .seek_pos = entry->offset, + .pause = LIANA_PAUSE_RESUME + }; + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + al_assert(sink->set == -1); + sink->set = list->current; + sink->callback(sink->userdata, CAMU_SINK_SET, entry, list->current, &time); + } + if (list->callback) list->callback(list->userdata, LIANA_META_PLAYING, entry); + } else { + /* + if (list->queued == -1) { + struct lia_list_entry *current = al_array_at(list->entries, list->current); + list->queued = list->current + 1; + entry->start = current->start + (current->duration - current->offset); + struct lia_timing time = { + .at = entry->start, + .seek_pos = entry->offset, + .pause = LIANA_PAUSE_RESUME + }; + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->queued = list->queued; + sink->callback(sink->userdata, CAMU_SINK_BUFFER_AND_QUEUE, entry, list->queued, &time); + } + } else { + */ + entry->start = LIANA_TIMESTAMP_INVALID; + //} + if (list->callback) list->callback(list->userdata, LIANA_META_QUEUED, entry); + } +} + +static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 sequence) +{ + s32 size = (s32)list->entries.size; + if (sequence < 0 || sequence >= size) return NULL; + return al_array_at(list->entries, sequence); +} + +static bool assume_done(struct lia_list_entry *entry, u64 at) +{ + if (entry->duration == 0) return true; + if (entry->paused_at == LIANA_TIMESTAMP_INVALID && entry->start != LIANA_TIMESTAMP_INVALID) { + if (entry->offset > entry->duration) { + return true; + } + if (at < entry->start) return false; + if (at - entry->start >= entry->duration - entry->offset) { + return true; + } + } + return false; +} + +void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) +{ + if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; + if (index == list->current) return; + + struct lia_list_entry *current = get_entry_from_sequence(list, sequence); + struct lia_list_entry *target = get_entry_from_sequence(list, index); + if (!current || !target) return; + + u64 now = aki_get_timestamp(); + u64 at = now + LIANA_BASE_DELAY; + u8 pause; + + // TODO: + // - Think about what is means for an entry to be done. never put into a pause state? + // - clock_end()?? + // - Sink needs to handle case where entry gets queued but the list already expects it to be playing + // - It's possible to know if sink->current is done during a set command, synchronously. + // So, check that when queueing an entry. + // - Can clock be ended during a queue command in any other case? + // - In the simplest case of our only operation being skip, how could client's become desynced? + // - Then with toggle pause + // - Is it safe to assert paused state on the client. + // - Do queued + // - Do seek + + //if (current->start != LIANA_TIMESTAMP_INVALID && current->start > ts - LIANA_BASE_PING) { + + // This should only happen if `start` has never been set. + if (target->start == LIANA_TIMESTAMP_INVALID && target->paused_at == LIANA_TIMESTAMP_INVALID) { + al_log_info("dd", "Hello"); + target->start = at; + } + + if (target->held) { + target->paused_at = LIANA_TIMESTAMP_INVALID; + target->held = false; + } + + al_assert(!current->held); + + bool done = assume_done(current, now); + if (done || current->paused_at == LIANA_TIMESTAMP_INVALID) { + if (!done) { + current->paused_at = at; + current->offset += current->paused_at - current->start; + current->start = LIANA_TIMESTAMP_INVALID; + current->held = true; + } + if (assume_done(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID) { + pause = !done ? LIANA_PAUSE_PAUSE : LIANA_PAUSE_NONE; + } else { //if (target->paused_at == LIANA_TIMESTAMP_INVALID) { + target->start = at; + pause = !done ? LIANA_PAUSE_BOTH : LIANA_PAUSE_RESUME; + } + al_log_info("dd", "%d: not paused, done: %d, pause: %d", sequence, done, pause); + } else { + if (assume_done(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID) { + pause = LIANA_PAUSE_NONE; + } else { //if (target->paused_at == LIANA_TIMESTAMP_INVALID) { + target->start = at; + pause = LIANA_PAUSE_RESUME; + } + al_log_info("dd", "%d: paused, pause: %d", sequence, pause); + } + + list->current = index; + list->previous = sequence; + list->idle = false; + + struct lia_timing time = { + .at = at, + .seek_pos = target->offset, + .pause = pause + }; + + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->set = index; + sink->callback(sink->userdata, CAMU_SINK_SET, target, index, &time); + } + + list->callback(list->userdata, LIANA_META_PLAYING, target); +} + +void lia_list_skip(struct lia_list *list, s32 sequence, s32 n) +{ + if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; + lia_list_skipto(list, sequence, sequence + n); +} + +void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) +{ + if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; + struct lia_list_entry *current = get_entry_from_sequence(list, sequence); + if (!current) return; + + u64 now = aki_get_timestamp(); + u8 pause = current->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_PAUSE : LIANA_PAUSE_RESUME; + u64 at; + switch (pause) { + case LIANA_PAUSE_PAUSE: + al_assert(pts != -1.0); + current->paused_at = now + LIANA_PAUSE_DELAY; + current->offset += current->paused_at - current->start; + current->start = LIANA_TIMESTAMP_INVALID; + at = current->paused_at; + break; + case LIANA_PAUSE_RESUME: + current->paused_at = LIANA_TIMESTAMP_INVALID; + current->start = now + LIANA_PAUSE_DELAY; + at = current->start; + break; + } + + struct lia_timing time = { + .at = at, + .seek_pos = LIANA_TIMESTAMP_INVALID, + .pause = pause + }; + + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->callback(sink->userdata, CAMU_SINK_PAUSE, current, sequence, &time); + } +} + +void lia_list_seek(struct lia_list *list, s32 sequence, f64 percent) +{ + if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; + list->idle = false; + struct lia_list_entry *current = get_entry_from_sequence(list, sequence); + u64 pos = (u64)(current->duration * percent); + u64 now = aki_get_timestamp(); + current->offset = pos; + current->start = now + LIANA_BASE_DELAY; + struct lia_timing time = { + .at = current->start, + .seek_pos = current->offset, + .pause = LIANA_PAUSE_NONE + }; + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->callback(sink->userdata, CAMU_SINK_SEEK, current, sequence, &time); + } +} + +void lia_list_end(struct lia_list *list, s32 sequence) +{ + al_assert(sequence != LIANA_SEQUENCE_ANY); + if (sequence != list->current) return; + s32 size = (s32)list->entries.size; + s32 next = sequence + 1; + list->previous = sequence; + struct lia_list_entry *current = al_array_at(list->entries, list->current); + //current->offset = current->duration; + if (list->queued >= 0) { + list->current = list->queued; + list->queued = -1; + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->queued = -1; + } + if (list->callback) { + list->callback(list->userdata, LIANA_META_PLAYING, al_array_at(list->entries, list->current)); + } + } else if (next < size) { + lia_list_skipto(list, sequence, next); + } else { + list->idle = true; + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->set = -1; + } + } +} + +void lia_list_reverse(struct lia_list *list) +{ + u32 size = list->entries.size; + for (u32 i = 0; i < size; i++) { + u32 tail = size - (i + 1); + if (tail <= i) break; + AL_SWAP(al_array_at(list->entries, i), al_array_at(list->entries, tail), struct lia_list_entry *); + } + unset_all(list); +} + +void lia_list_sort(struct lia_list *list) +{ + al_array_sort(list->entries, struct lia_list_entry *, camu_db_compare); + unset_all(list); +} + +void lia_list_shuffle(struct lia_list *list) +{ + u32 size = list->entries.size; + if (size == 0) return; + for (u32 i = 0; i < size - 1; i++) { + u32 j = i + al_rand() / (AL_RAND_MAX / (size - i) + 1); + AL_SWAP(al_array_at(list->entries, i), al_array_at(list->entries, j), struct lia_list_entry *); + } + unset_all(list); +} + +void lia_list_clear(struct lia_list *list) +{ + struct lia_list_entry *entry; + al_array_foreach(list->entries, i, entry) { + al_str_free(&entry->name); + al_free(entry); + } + list->entries.size = 0; + unset_all(list); + list->current = -1; + list->previous = -1; + list->idle = true; +} + +void lia_list_free(struct lia_list *list) +{ + struct lia_list_entry *entry; + al_array_foreach(list->entries, i, entry) { + al_str_free(&entry->name); + al_free(entry); + } + al_array_free(list->entries); + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + al_free(sink); + } + al_array_free(list->sinks); +} diff --git a/src/liana/list.h b/src/liana/list.h new file mode 100644 index 0000000..8adea17 --- /dev/null +++ b/src/liana/list.h @@ -0,0 +1,107 @@ +#pragma once + +#include <al/types.h> +#include <al/str.h> +#include <al/array.h> + +#define LIANA_SEQUENCE_ANY -1 +#define LIANA_TIMESTAMP_INVALID UINT64_MAX + +#define LIANA_BASE_PING 200000Lu // 100ms +#define LIANA_BASE_DELAY 450000Lu // 250ms +#define LIANA_PAUSE_DELAY LIANA_BASE_PING +#define LIANA_DELAY_IGNORE 0Lu + +#define LIANA_BUFFER_AHEAD 2 + +enum { + LIANA_ENTRY_DURATION = 0, + LIANA_ENTRY_ID, + LIANA_ENTRY_SKIP +}; + +enum { + LIANA_META_PLAYING = 0, + LIANA_META_QUEUED +}; + +// Pause: +// - Set `paused_at` to now() + PAUSE_DELAY. +// - Increment `offset` by how long the entry will have been +// playing when paused at the requested timestamp (`paused_at` - `start`). +// Resume: +// - Unset `paused_at` +// - Set `start` to now() + PAUSE_DELAY. +// Next/Prev: +// Current Playing, Target Playing: +// - Set `current->held_at` to now() + PAUSE_DELAY. +// - Increment `current->offset` by the same logic as pause. +// - +// Current Playing, Target Paused: +// Current Paused, Target Playing: +// Current Paused, Target Paused: +// Seek: + +// NOTE: keep global per list max time until all sinks _should_ by buffered +// possibly use that instead of LIANA_PAUSE_DELAY. + +enum { + LIANA_PAUSE_NONE = 0, + LIANA_PAUSE_RESUME, + LIANA_PAUSE_PAUSE, + LIANA_PAUSE_BOTH +}; + +struct lia_timing { + u64 at; + u64 seek_pos; + u8 pause; +}; + +struct lia_list_entry { + void *opaque; + u64 start; + u64 paused_at; + bool held; + u64 offset; + u64 duration; + str name; +}; + +struct lia_list_sink { + s32 set; + s32 queued; + void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *); + void *userdata; +}; + +struct lia_list { + s32 current; + s32 previous; + s32 queued; + bool idle; + array(struct lia_list_entry *) entries; + array(struct lia_list_sink *) sinks; + void (*callback)(void *, u8, struct lia_list_entry *); + void *userdata; +}; + +void lia_list_init(struct lia_list *list); + +void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *), void *userdata); +void lia_list_remove_sink(struct lia_list *list, void *userdata); + +void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name); +void lia_list_unset(struct lia_list *list); +void lia_list_skipto(struct lia_list *list, s32 sequence, s32 i); +void lia_list_skip(struct lia_list *list, s32 sequence, s32 n); +void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts); +void lia_list_seek(struct lia_list *list, s32 sequence, f64 percent); +void lia_list_end(struct lia_list *list, s32 sequence); + +void lia_list_reverse(struct lia_list *list); +void lia_list_sort(struct lia_list *list); +void lia_list_shuffle(struct lia_list *list); +void lia_list_clear(struct lia_list *list); + +void lia_list_free(struct lia_list *list); diff --git a/src/liana/list_cmp.h b/src/liana/list_cmp.h new file mode 100644 index 0000000..1f77b46 --- /dev/null +++ b/src/liana/list_cmp.h @@ -0,0 +1,53 @@ +#pragma once + +#include <al/str.h> + +#include "list.h" + +static void camu_db_num_from_path(str *path, s64 *id, s64 *index) +{ + s32 index0 = al_str_rfind(path, '/'); + if (index0 < 0) return; + + str s = *al_str_substr(path, index0 + 1, path->len); + index0 = al_str_find(&s, '_'); + if (index0 < 0) return; + s = *al_str_substr(&s, index0 + 1, s.len); + index0 = al_str_find(&s, '_'); + if (index0 < 0) return; + s = *al_str_substr(&s, index0 + 1, s.len); + index0 = al_str_find(&s, '_'); + if (index0 < 0) return; + s = *al_str_substr(&s, 0, index0); + + s64 num = al_str_to_long(&s, 10); + if (num == INT64_MIN || num == INT64_MAX) return; + *id = num; + + index0 = al_str_rfind(path, '.'); + s32 index1 = al_str_rfind(path, 'a'); + if (index0 < 0 || index1 < 0) return; + s = *al_str_substr(path, index1 + 1, index0); + + num = al_str_to_long(&s, 10); + if (num == INT64_MIN || num == INT64_MAX) return; + *index = num; +} + +static s32 camu_db_compare(const void *a, const void *b) +{ + struct lia_list_entry *aa = *((struct lia_list_entry **)a); + struct lia_list_entry *bb = *((struct lia_list_entry **)b); + s64 anum = -1, aindex = -1; + s64 bnum = -1, bindex = -1; + camu_db_num_from_path(&aa->name, &anum, &aindex); + camu_db_num_from_path(&bb->name, &bnum, &bindex); + if (anum == bnum) { + if (aindex > bindex) return 1; + else if (aindex < bindex) return -1; + } else { + if (anum > bnum) return -1; + else if (anum < bnum) return 1; + } + return 0; +} diff --git a/src/liana/meson.build b/src/liana/meson.build new file mode 100644 index 0000000..0859e99 --- /dev/null +++ b/src/liana/meson.build @@ -0,0 +1,36 @@ +liana_server_src = [ + 'list.c', + 'server.c', + 'handlers/codec_server.c', + 'handlers.c' +] +liana_client_src = [ + 'client.c', + 'vcr.c', + 'handlers/codec_client.c', + 'handlers.c' +] +liana_deps = [codecs] +liana_args = [] + +libcdio_paranoia = dependency('libcdio_paranoia', required: false, allow_fallback: true) +libcdio_cdda = dependency('libcdio_cdda', required: false, allow_fallback: true) +if libcdio_paranoia.found() and libcdio_cdda.found() + liana_server_src += ['handlers/cdio_server.c'] + liana_client_src += ['handlers/cdio_client.c'] + liana_deps += [libcdio_paranoia, libcdio_cdda] + liana_args += ['-DLIANA_HAVE_CDIO'] +endif + +#libdvdcss = dependency('libdvdcss', required: false, allow_fallback: true) +#libdvdread = dependency('libdvdread', required: false, allow_fallback: true) +#libdvdnav = dependency('libdvdnav', required: false, allow_fallback: true) +#if libdvdcss.found() and libdvdread.found() and libdvdnav.found() +# liana_server_src += ['handlers/dvd_server.c'] +# liana_server_deps += [libdvdcss] +#endif + +liana_server = declare_dependency(sources: liana_server_src, dependencies: liana_deps, + compile_args: [liana_args, '-DLIANA_SERVER']) +liana_client = declare_dependency(sources: liana_client_src, dependencies: liana_deps, + compile_args: [liana_args, '-DLIANA_CLIENT']) diff --git a/src/liana/server.c b/src/liana/server.c new file mode 100644 index 0000000..1fc96fb --- /dev/null +++ b/src/liana/server.c @@ -0,0 +1,345 @@ +#include <al/random.h> + +#include "server.h" +#include "handler.h" +#include "handlers.h" +#include "list.h" + +bool lia_server_init(struct lia_server *server, struct aki_event_loop *loop) +{ + server->loop = loop; + al_array_init(server->nodes); + al_array_init(server->zombies); + return true; +} + +static void remove_zombie(struct lia_server *server, struct aki_packet_stream *stream) +{ + struct aki_packet_stream *zombie; + al_array_foreach(server->zombies, i, zombie) { + if (zombie == stream) { + al_array_remove_at(server->zombies, i); + break; + } + } +} + +static void data_packet_sent_callback(void *userdata, struct aki_packet *packet) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + aki_packet_pool_return(&conn->pool, packet); +} + +static u8 packet_pool_callback(void *userdata, struct aki_packet *packet) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + if (!aki_packet_stream_send_packet(conn->stream, packet)) { + return AKI_PACKET_POOL_RETURN; + } + return AKI_PACKET_POOL_KEEP; +} + +static aki_thread_result AKI_THREADCALL handler_thread(void *userdata) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + do { + struct aki_packet *packet = aki_packet_pool_get(&conn->pool); + if (!packet) break; + conn->handler->step(conn->handler); + conn->handler->write_packet(conn->handler, packet); + aki_packet_pool_submit(&conn->pool, packet); + if (conn->handler->status != CAMU_OK) break; + } while (1); + aki_packet_pool_flush(&conn->pool); + return 0; +} + +static void discard_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +{ + (void)userdata; + (void)stream; + aki_packet_free(packet); + // We should never be here. Although, we also shouldn't assert because + // any erroneous connection can bring us here. + al_assert(false); +} + +static void close_connection_internal(struct lia_node_connection *conn) +{ + al_assert(conn->handler); + conn->handler->free(&conn->handler); + cch_entry_return_handle(conn->node->entry, &conn->handle); + aki_packet_stream_free(conn->stream); + al_free(conn->stream); + conn->stream = NULL; +} + +static void data_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + (void)stream; + aki_packet_pool_disable(&conn->pool); + cch_handle_disable(&conn->handle); + // This is joining handler_thread(), we will never be here if init_thread() blocks or fails. + aki_thread_join(&conn->thread); + close_connection_internal(conn); + aki_packet_pool_free(&conn->pool); +} + +static void subscribe_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + s32 mask = aki_packet_read_s32(packet); + conn->handler->subscribe(conn->handler, mask); + stream->packet_callback = discard_packet_callback; + stream->packet_sent_callback = data_packet_sent_callback; + stream->connection_closed_callback = data_connection_closed_callback; + aki_thread_create(&conn->thread, handler_thread, conn); + aki_packet_free(packet); +} + +static void subscribe_packet_sent_callback(void *userdata, struct aki_packet *packet) +{ + (void)userdata; + aki_packet_free(packet); +} + +static void subscribe_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + (void)stream; + close_connection_internal(conn); + aki_packet_pool_free(&conn->pool); +} + +static void handle_connection(struct lia_node_connection *conn, struct aki_packet *packet) +{ + struct aki_packet_stream *stream = conn->stream; + + s32 mask = aki_packet_read_s32(packet); + u64 seek_pos = aki_packet_read_u64(packet); + + // Besides being wasteful, seeking to 0 on a new stream can skip data. + if (seek_pos > 0) conn->handler->seek(conn->handler, seek_pos); + + if (mask == 0) { + struct aki_packet *rpacket = aki_packet_create(); + aki_packet_write_u16(rpacket, conn->id); + aki_packet_write_str(rpacket, cch_entry_get_liana(conn->node->entry)); + conn->handler->write_info(conn->handler, rpacket); + stream->packet_callback = subscribe_packet_callback; + stream->packet_sent_callback = subscribe_packet_sent_callback; + stream->connection_closed_callback = subscribe_connection_closed_callback; + aki_packet_stream_send_packet(stream, rpacket); + } else { + conn->handler->subscribe(conn->handler, mask); + stream->packet_callback = discard_packet_callback; + stream->packet_sent_callback = data_packet_sent_callback; + stream->connection_closed_callback = data_connection_closed_callback; + aki_thread_create(&conn->thread, handler_thread, conn); + } +} + +static void connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +{ + struct lia_server *server = (struct lia_server *)userdata; + // This connection might no longer be a zombie, but that's fine. + remove_zombie(server, stream); + aki_packet_stream_free(stream); + al_free(stream); +} + +static void signal_callback(void *userdata) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + struct lia_node *node = conn->node; + struct lia_server *server = node->server; + aki_signal_stop(&conn->signal); + aki_thread_join(&conn->thread); + struct aki_packet *packet = conn->packet; + if (!packet) { + // Connection was closed before init was done. + close_connection_internal(conn); + return; + } + if (!conn->errored) { + conn->id = al_inc_u16(); + aki_packet_pool_init(&conn->pool, 48, server->loop, packet_pool_callback, conn); + al_array_push(node->connections, conn); + handle_connection(conn, packet); + } else { + conn->handler->free(&conn->handler); + cch_entry_return_handle(node->entry, &conn->handle); + struct aki_packet_stream *stream = conn->stream; + al_free(conn); + conn = NULL; + // This connection is now nothing but a packet stream. + stream->userdata = server; + stream->connection_closed_callback = connection_closed_callback; + aki_packet_stream_disconnect(stream); + } + aki_packet_free(packet); +} + +static aki_thread_result AKI_THREADCALL init_thread(void *userdata) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + if (!conn->handler->init(conn->handler, &conn->handle)) { + conn->errored = true; + } + aki_signal_send(&conn->signal); + return 0; +} + +static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u16 id) +{ + struct lia_node_connection *conn = NULL; + al_array_foreach(node->connections, i, conn) { + if (conn->id == id) return conn; + } + return NULL; +} + +static struct lia_node *get_node_from_id(struct lia_server *server, u16 id) +{ + struct lia_node *node; + al_array_foreach(server->nodes, i, node) { + if (node->id == id) return node; + } + return NULL; +} + +static void pre_init_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + (void)stream; + aki_packet_free(conn->packet); + // Checked in signal_callback and will signal to cleanup the connection. + conn->packet = NULL; +} + +static void packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +{ + struct lia_server *server = (struct lia_server *)userdata; + + // We got a packet, this connection is no longer a zombie. + remove_zombie(server, stream); + + u16 node_id = aki_packet_read_u16(packet); + u16 connection_id = aki_packet_read_u16(packet); + + struct lia_node *node = get_node_from_id(server, node_id); + struct lia_node_connection *conn = NULL; + if (connection_id == 0) { + conn = al_alloc_object(struct lia_node_connection); + stream->userdata = conn; + conn->node = node; + conn->stream = stream; + conn->packet = packet; + aki_signal_init(&conn->signal, signal_callback, conn); + aki_signal_start(&conn->signal, server->loop); + cch_entry_get_handle(node->entry, &conn->handle); + conn->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler(); + conn->errored = false; + stream->packet_callback = discard_packet_callback; + stream->connection_closed_callback = pre_init_connection_closed_callback; + aki_thread_create(&conn->thread, init_thread, conn); + } else { + // This is completely unused and connections never get removed from node->connection. + if ((conn = get_connection_from_id(node, connection_id))) { + handle_connection(conn, packet); + } else { + aki_packet_stream_disconnect(stream); + } + aki_packet_free(packet); + } +} + +static void packet_sent_callback(void *userdata, struct aki_packet *packet) +{ + (void)userdata; + aki_packet_free(packet); +} + +static void connection_callback(void *userdata, struct aki_packet_stream *stream) +{ + struct lia_server *server = (struct lia_server *)userdata; + stream->packet_callback = packet_callback; + stream->packet_sent_callback = packet_sent_callback; + al_array_push(server->zombies, stream); +} + +void lia_server_add_socket(struct lia_server *server, struct aki_socket *sock) +{ + struct aki_packet_stream *stream = al_alloc_object(struct aki_packet_stream); + stream->connection_callback = connection_callback; + stream->connection_closed_callback = connection_closed_callback; + stream->userdata = server; + aki_packet_stream_from_socket(stream, server->loop, sock); +} + +struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry) +{ + struct lia_node *node = al_alloc_object(struct lia_node); + node->id = al_inc_u16(); + node->entry = entry; + al_array_init(node->connections); + node->server = server; + al_array_push(server->nodes, node); + return node; +} + +u64 lia_node_get_duration(struct lia_node *node) +{ + (void)node; +#if 0 + return 0; +#else + struct cch_handle handle; + cch_entry_get_handle(node->entry, &handle); + struct lia_server_handler *handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler(); + if (!handler->init(handler, &handle)) { + return LIANA_TIMESTAMP_INVALID; + } + u64 duration = handler->get_duration(handler); + handler->free(&handler); + cch_entry_return_handle(node->entry, &handle); + return duration; +#endif +} + +void lia_server_close(struct lia_server *server) +{ + struct lia_node *node; + al_array_foreach(server->nodes, i, node) { + struct lia_node_connection *connection; + al_array_foreach(node->connections, j, connection) { + if (connection->stream) { + aki_packet_stream_disconnect(connection->stream); + } + } + } + struct aki_packet_stream *zombie; + al_array_foreach_rev(server->zombies, i, zombie) { + aki_packet_stream_disconnect(zombie); + } +} + +void lia_server_free(struct lia_server *server) +{ + struct lia_node *node; + al_array_foreach(server->nodes, i, node) { + struct lia_node_connection *connection; + al_array_foreach(node->connections, j, connection) { + al_free(connection); + } + al_array_free(node->connections); + al_free(node); + } + al_array_free(server->nodes); + struct aki_packet_stream *zombie; + al_array_foreach(server->zombies, i, zombie) { + al_free(zombie); + } + al_array_free(server->zombies); +} diff --git a/src/liana/server.h b/src/liana/server.h new file mode 100644 index 0000000..3c4c294 --- /dev/null +++ b/src/liana/server.h @@ -0,0 +1,42 @@ +#pragma once + +#include <aki/event_loop.h> +#include <aki/packet_stream.h> +#include <aki/packet_pool.h> +#include <aki/socket.h> +#include <aki/signal.h> + +#include "../cache/entry.h" + +struct lia_node_connection { + u16 id; + struct aki_packet *packet; + struct aki_packet_stream *stream; + struct lia_server_handler *handler; + bool errored; + struct cch_handle handle; + struct aki_thread thread; + struct aki_signal signal; + struct aki_packet_pool pool; + struct lia_node *node; +}; + +struct lia_node { + u16 id; + struct cch_entry *entry; + array(struct lia_node_connection *) connections; + struct lia_server *server; +}; + +struct lia_server { + struct aki_event_loop *loop; + array(struct lia_node *) nodes; + array(struct aki_packet_stream *) zombies; +}; + +bool lia_server_init(struct lia_server *server, struct aki_event_loop *loop); +void lia_server_add_socket(struct lia_server *server, struct aki_socket *sock); +struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry); +u64 lia_node_get_duration(struct lia_node *node); +void lia_server_close(struct lia_server *server); +void lia_server_free(struct lia_server *server); diff --git a/src/liana/vcr.c b/src/liana/vcr.c new file mode 100644 index 0000000..fad8b78 --- /dev/null +++ b/src/liana/vcr.c @@ -0,0 +1,232 @@ +#include <al/log.h> + +#include "vcr.h" +#include "handler.h" + +#define VCR_BUFFER_INIT 64 +#define VCR_BUFFER_BUFFERED (VCR_BUFFER_INIT - 8) +#define VCR_BUFFER_LOW (VCR_BUFFER_INIT - 16) + +static void signal_callback(void *userdata) +{ + struct lia_vcr *vcr = (struct lia_vcr *)userdata; + aki_packet_stream_cork(vcr->data, false); +} + +void lia_vcr_init(struct lia_vcr *vcr, struct aki_packet_stream *data) +{ + al_array_init(vcr->tracks); + al_atomic_store(s32)(&vcr->count, 0, AL_ATOMIC_RELAXED); + vcr->mark.low = VCR_BUFFER_LOW; + vcr->mark.buffered = VCR_BUFFER_BUFFERED; + vcr->data = data; + aki_signal_init(&vcr->signal, signal_callback, vcr); +} + +void lia_vcr_start(struct lia_vcr *vcr, struct aki_event_loop *loop) +{ + aki_signal_start(&vcr->signal, loop); +} + +static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata) +{ + struct lia_vcr_track *track = (struct lia_vcr_track *)userdata; + struct lia_vcr *vcr = track->vcr; + u32 packets; + struct aki_packet *packet = NULL; + while (aki_packet_cache_wait(&track->cache, &packets)) { + s32 state = 0; + for (u32 i = 0; i < packets; i++) { + packet = aki_packet_cache_pop(&track->cache); + state = 0; + if (packet) { + state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED); + if (state == LIANA_STREAM_CLOSED) { + // We were signaled to close, exit thread. + goto out; + } + u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED); + s32 count = al_atomic_sub(s32)(&vcr->count, 1, AL_ATOMIC_RELAXED); + if (count < vcr->mark.low && buffered) { + aki_signal_send(&vcr->signal); + } + } + if (!track->client->handle_packet(track->client, packet)) { + // Handler error, exit. + goto out; + } + if (packet) { + aki_packet_free(packet); + packet = NULL; + if (state == LIANA_STREAM_STOPPED) { + // Unlock here to accumulate packets while waiting. + aki_packet_cache_unlock(&track->cache); + aki_mutex_lock(&track->mutex); + aki_cond_wait(&track->cond, &track->mutex); + aki_mutex_unlock(&track->mutex); + break; + } + } else { + // NULL packet, exit thread. + goto out; + } + } + if (state != LIANA_STREAM_STOPPED) { + // If state = STOPPED, we already unlocked. + aki_packet_cache_unlock(&track->cache); + } + } + // packet_cache_wait() returned false. + return 0; +out: + // packet_cache_wait() returned true and we are jumping out of the loop. + if (packet) aki_packet_free(packet); + aki_packet_cache_unlock(&track->cache); + return 0; +} + +void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track) +{ + track->vcr = vcr; + aki_cond_init(&track->cond); + aki_mutex_init(&track->mutex); + al_atomic_store(u8)(&track->buffered, 0, AL_ATOMIC_RELAXED); + aki_packet_cache_init(&track->cache, VCR_BUFFER_INIT - 1); + al_array_push(vcr->tracks, track); + al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); +} + +static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index) +{ + struct lia_vcr_track *track; + al_array_foreach(vcr->tracks, i, track) { + if (track->index == index) return track; + } + return NULL; +} + +static void cork_if_buffered(struct lia_vcr *vcr) +{ + u8 buffered = 1; + struct lia_vcr_track *track; + al_array_foreach(vcr->tracks, i, track) { + buffered &= al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED); + } + if (buffered) aki_packet_stream_cork(vcr->data, true); +} + +bool lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet) +{ + struct lia_vcr_track *track; + s32 count; + u8 op = aki_packet_read_u8(packet); + switch (op) { + case LIANA_PACKET_DATA: + track = get_track_from_index(vcr, aki_packet_read_s32(packet)); + al_assert(track); + if (!track->running) { + aki_thread_create(&track->thread, vcr_track_thread, track); + track->running = true; + } + if (!aki_packet_cache_send_packet(&track->cache, packet)) { + aki_packet_free(packet); + } else if ((count = al_atomic_add(s32)(&vcr->count, 1, AL_ATOMIC_RELAXED)) > vcr->mark.buffered) { + vcr->mark.buffered = count; + vcr->mark.low = vcr->mark.buffered - 16; + cork_if_buffered(vcr); + } + break; + case LIANA_PACKET_EOF: + al_array_foreach(vcr->tracks, i, track) { + aki_packet_cache_send_packet(&track->cache, NULL); + } + aki_packet_free(packet); + break; + case LIANA_PACKET_ERROR: + al_log_warn("liana", "Unhandled error packet."); + aki_packet_free(packet); + break; + default: + al_assert_and_return(false); + } + return true; +} + +void lia_vcr_cork(struct lia_vcr_track *track) +{ + al_atomic_store(s32)(&track->state, LIANA_STREAM_STOPPED, AL_ATOMIC_RELAXED); + al_atomic_store(u8)(&track->buffered, 1, AL_ATOMIC_RELAXED); +} + +void lia_vcr_uncork(struct lia_vcr_track *track) +{ + if (al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED) != LIANA_STREAM_STOPPED) { + // This can be reached during normal operation. Whether or not that should + // be allowed is up for consideration. + return; + } + al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); + aki_mutex_lock(&track->mutex); + if (aki_cond_is_waiting(&track->cond)) { + aki_cond_signal(&track->cond); + } + aki_mutex_unlock(&track->mutex); +} + +static void vcr_track_close_internal(struct lia_vcr_track *track) +{ + al_atomic_store(s32)(&track->state, LIANA_STREAM_CLOSED, AL_ATOMIC_RELAXED); + aki_packet_cache_disable(&track->cache); + aki_mutex_lock(&track->mutex); + if (aki_cond_is_waiting(&track->cond)) { + aki_cond_signal(&track->cond); + } + aki_mutex_unlock(&track->mutex); + if (track->running) { + aki_thread_join(&track->thread); + track->running = false; + } + struct aki_packet *packet; + while ((packet = aki_packet_cache_pop(&track->cache))) { + aki_packet_free(packet); + } +} + +void lia_vcr_flush(struct lia_vcr *vcr) +{ + struct lia_vcr_track *track; + al_array_foreach(vcr->tracks, i, track) { + vcr_track_close_internal(track); + track->client->flush(track->client); + aki_packet_cache_enable(&track->cache); + track->running = false; + al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); + } + aki_signal_send(&vcr->signal); + al_atomic_store(s32)(&vcr->count, 0, AL_ATOMIC_RELAXED); +} + +void lia_vcr_close_all(struct lia_vcr *vcr) +{ + struct lia_vcr_track *track; + al_array_foreach(vcr->tracks, i, track) { + vcr_track_close_internal(track); + } + aki_signal_stop(&vcr->signal); +} + +void lia_vcr_free(struct lia_vcr *vcr) +{ + struct lia_vcr_track *track; + al_array_foreach(vcr->tracks, i, track) { + aki_packet_cache_free(&track->cache); + track->client->free(&track->client); +#if CAMU_HAVE_FFMPEG + if (track->stream.mode == CAMU_FFMPEG_COMPAT) { + avformat_free_context(track->stream.av.format_context); + } +#endif + al_free(track); + } + al_array_free(vcr->tracks); +} diff --git a/src/liana/vcr.h b/src/liana/vcr.h new file mode 100644 index 0000000..f471c4e --- /dev/null +++ b/src/liana/vcr.h @@ -0,0 +1,46 @@ +#pragma once + +#include <al/atomic.h> +#include <aki/packet_cache.h> +#include <aki/packet_stream.h> +#include <aki/signal.h> + +#include "../codec/codec.h" + +enum { + LIANA_STREAM_RUNNING = 0, + LIANA_STREAM_STOPPED, + LIANA_STREAM_CLOSED +}; + +struct lia_vcr_track { + s32 index; + struct camu_codec_stream stream; + struct lia_client_handler *client; + atomic(s32) state; + atomic(u8) buffered; + struct aki_packet_cache cache; + struct aki_cond cond; + struct aki_mutex mutex; + bool running; + struct aki_thread thread; + struct lia_vcr *vcr; +}; + +struct lia_vcr { + array(struct lia_vcr_track *) tracks; + atomic(s32) count; + struct { s32 low, buffered; } mark; + struct aki_packet_stream *data; + struct aki_signal signal; +}; + +void lia_vcr_init(struct lia_vcr *vcr, struct aki_packet_stream *data); +void lia_vcr_start(struct lia_vcr *vcr, struct aki_event_loop *loop); +void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track); +bool lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet); +void lia_vcr_cork(struct lia_vcr_track *track); +void lia_vcr_uncork(struct lia_vcr_track *track); +void lia_vcr_flush(struct lia_vcr *vcr); +void lia_vcr_close_all(struct lia_vcr *vcr); +void lia_vcr_free(struct lia_vcr *vcr); |