summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c199
-rw-r--r--src/liana/client.h31
-rw-r--r--src/liana/common.h1
-rw-r--r--src/liana/handler.h58
-rw-r--r--src/liana/handlers.c37
-rw-r--r--src/liana/handlers.h13
-rw-r--r--src/liana/handlers/cdio.h18
-rw-r--r--src/liana/handlers/cdio_client.c49
-rw-r--r--src/liana/handlers/cdio_server.c106
-rw-r--r--src/liana/handlers/codec.h19
-rw-r--r--src/liana/handlers/codec_client.c116
-rw-r--r--src/liana/handlers/codec_server.c143
-rw-r--r--src/liana/handlers/dvd.h14
-rw-r--r--src/liana/handlers/dvd_server.c21
-rw-r--r--src/liana/list.c413
-rw-r--r--src/liana/list.h107
-rw-r--r--src/liana/list_cmp.h53
-rw-r--r--src/liana/meson.build36
-rw-r--r--src/liana/server.c345
-rw-r--r--src/liana/server.h42
-rw-r--r--src/liana/vcr.c232
-rw-r--r--src/liana/vcr.h46
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);