summaryrefslogtreecommitdiff
path: root/src/bimu
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2023-11-06 11:54:08 -0500
committerAndrew Opalach <andrew@akon.city> 2023-11-06 11:54:08 -0500
commita48a68cfb04b2020737c0adfb0e7667451a22b5c (patch)
treea92121897f84298c3798a4f2ba56de45cd100eab /src/bimu
downloadcamu-a48a68cfb04b2020737c0adfb0e7667451a22b5c.tar.gz
camu-a48a68cfb04b2020737c0adfb0e7667451a22b5c.tar.bz2
camu-a48a68cfb04b2020737c0adfb0e7667451a22b5c.zip
Add camu
- Subprojects temporarily omitted Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/bimu')
-rw-r--r--src/bimu/client.c1
-rw-r--r--src/bimu/client.h11
-rw-r--r--src/bimu/common.h7
-rw-r--r--src/bimu/handler.h35
-rw-r--r--src/bimu/handlers/codec.h20
-rw-r--r--src/bimu/handlers/codec_client.c109
-rw-r--r--src/bimu/handlers/codec_server.c159
-rw-r--r--src/bimu/local.c281
-rw-r--r--src/bimu/local.h53
-rw-r--r--src/bimu/meson.build8
-rw-r--r--src/bimu/server.c46
-rw-r--r--src/bimu/server.h21
12 files changed, 751 insertions, 0 deletions
diff --git a/src/bimu/client.c b/src/bimu/client.c
new file mode 100644
index 0000000..f679c0d
--- /dev/null
+++ b/src/bimu/client.c
@@ -0,0 +1 @@
+#include "client.h"
diff --git a/src/bimu/client.h b/src/bimu/client.h
new file mode 100644
index 0000000..ec48275
--- /dev/null
+++ b/src/bimu/client.h
@@ -0,0 +1,11 @@
+#pragma once
+
+#include <al/types.h>
+
+#include "handlers/codec.h"
+
+struct bmu_client_stream {
+ u8 type;
+ s32 index;
+ struct camu_stream stream;
+};
diff --git a/src/bimu/common.h b/src/bimu/common.h
new file mode 100644
index 0000000..1c4897a
--- /dev/null
+++ b/src/bimu/common.h
@@ -0,0 +1,7 @@
+#pragma once
+
+enum {
+ BIMU_STREAM_AUDIO = 0,
+ BIMU_STREAM_VIDEO,
+ BIMU_STREAM_UNKNOWN,
+};
diff --git a/src/bimu/handler.h b/src/bimu/handler.h
new file mode 100644
index 0000000..e51bc2d
--- /dev/null
+++ b/src/bimu/handler.h
@@ -0,0 +1,35 @@
+#pragma once
+
+#include <aki/packet_pool.h>
+
+#include "../codec/codec.h"
+#include "../cache/handle.h"
+
+struct bmu_server_handler {
+ bool (*init)(struct bmu_server_handler *, struct cch_handle *, struct aki_packet_pool *);
+ void (*write_info)(struct bmu_server_handler *, struct aki_packet *);
+ s64 (*get_duration)(struct bmu_server_handler *);
+ bool (*seek)(struct bmu_server_handler *, s64);
+ bool (*step)(struct bmu_server_handler *);
+ void (*free)(struct bmu_server_handler **);
+};
+
+enum {
+ BIMU_CLIENT_CONFIGURE = 0,
+ BIMU_CLIENT_DATA,
+ BIMU_CLIENT_SEEK,
+ BIMU_CLIENT_EOF,
+ BIMU_CLIENT_CLOSED
+};
+
+struct bmu_client_stream;
+struct bmu_client_handler {
+ bool (*init)(struct bmu_client_handler *, struct camu_renderer *, struct bmu_client_stream *);
+ void (*handle_data_packet)(struct bmu_client_handler *, struct aki_packet *);
+ void (*handle_eof)(struct bmu_client_handler *);
+ void (*flush)(struct bmu_client_handler *);
+ void (*free)(struct bmu_client_handler **);
+ struct bmu_client_stream *stream;
+ void (*callback)(void *, u8, struct bmu_client_stream *stream, void *);
+ void *userdata;
+};
diff --git a/src/bimu/handlers/codec.h b/src/bimu/handlers/codec.h
new file mode 100644
index 0000000..a8031b1
--- /dev/null
+++ b/src/bimu/handlers/codec.h
@@ -0,0 +1,20 @@
+#pragma once
+
+#include "../handler.h"
+
+#include "../../codec/codec.h"
+
+struct bmu_codec_server {
+ struct bmu_server_handler handler;
+ struct camu_demuxer *demux;
+ struct camu_packet packet;
+ struct aki_packet_pool *pool;
+};
+
+struct bmu_codec_client {
+ struct bmu_client_handler handler;
+ struct camu_decoder *dec;
+};
+
+struct bmu_server_handler *bmu_codec_server_create(void);
+struct bmu_client_handler *bmu_codec_client_create(void);
diff --git a/src/bimu/handlers/codec_client.c b/src/bimu/handlers/codec_client.c
new file mode 100644
index 0000000..ffe8c16
--- /dev/null
+++ b/src/bimu/handlers/codec_client.c
@@ -0,0 +1,109 @@
+#include "../../codec/libav/decoder.h"
+#include "../../codec/libav/packet_ext.h"
+#include "../../codec/stb_image/decoder.h"
+#include "../../codec/wuffs/decoder.h"
+#include "../../codec/spng/decoder.h"
+
+#include "../../bimu/local.h"
+
+#include "codec.h"
+
+static void data_callback(void *userdata, struct camu_frame *frame)
+{
+ struct bmu_codec_client *codec = (struct bmu_codec_client *)userdata;
+ codec->handler.callback(codec->handler.userdata, BIMU_CLIENT_DATA, codec->handler.stream, frame);
+}
+
+static bool codec_client_init(struct bmu_client_handler *handler, struct camu_renderer *renderer,
+ struct bmu_client_stream *stream)
+{
+ struct bmu_codec_client *codec = (struct bmu_codec_client *)handler;
+ codec->dec = camu_lav_decoder_create();
+ //codec->dec = camu_spng_decoder_create();
+ //codec->dec = camu_stbi_decoder_create();
+ //codec->dec = camu_wuffs_decoder_create();
+ codec->handler.stream = stream;
+ if (!codec->dec->init(codec->dec, renderer, &stream->stream)) {
+ return false;
+ }
+ codec->dec->stream = &stream->stream;
+ codec->handler.callback(codec->handler.userdata, BIMU_CLIENT_CONFIGURE, codec->handler.stream, NULL);
+ codec->dec->set_callback(codec->dec, data_callback, codec);
+ return true;
+}
+
+#ifdef HAVE_FFMPEG
+static void push_av_packet(struct bmu_codec_client *codec, AVPacket *pkt)
+{
+ struct camu_packet packet;
+ packet.av.pkt = pkt;
+ s32 ret = codec->dec->push(codec->dec, &packet);
+ av_packet_unref(pkt);
+ av_packet_free(&pkt);
+ al_assert(ret == CAMU_OK);
+}
+#endif
+
+static void push_packet(struct bmu_codec_client *codec, struct aki_buffer *buffer)
+{
+ struct camu_packet packet;
+ packet.buffer = buffer;
+ s32 ret = codec->dec->push(codec->dec, &packet);
+ al_assert(ret == CAMU_OK);
+}
+
+static void codec_client_handle_data_packet(struct bmu_client_handler *handler, struct aki_packet *packet)
+{
+ struct bmu_codec_client *codec = (struct bmu_codec_client *)handler;
+ switch (aki_packet_read_u8(packet)) {
+ case CAMU_NORMAL: {
+ struct aki_buffer buffer;
+ aki_packet_read_buffer(packet, &buffer);
+ push_packet(codec, &buffer);
+ break;
+ }
+#ifdef HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT: {
+ push_av_packet(codec, aki_packet_read_av_packet(packet));
+ break;
+ }
+#endif
+ }
+ s32 ret = codec->dec->process(codec->dec);
+ al_assert(ret == CAMU_ERR_AGAIN || ret == CAMU_ERR_EOF);
+}
+
+static void codec_client_handle_eof(struct bmu_client_handler *handler)
+{
+ struct bmu_codec_client *codec = (struct bmu_codec_client *)handler;
+ s32 ret = codec->dec->push(codec->dec, NULL);
+ // Flush returns success.
+ ret = codec->dec->process(codec->dec);
+ al_assert(ret == CAMU_ERR_EOF);
+ codec->handler.callback(codec->handler.userdata, BIMU_CLIENT_EOF, codec->handler.stream, NULL);
+}
+
+static void codec_client_flush(struct bmu_client_handler *handler)
+{
+ struct bmu_codec_client *codec = (struct bmu_codec_client *)handler;
+ codec->dec->flush(codec->dec);
+}
+
+static void codec_client_free(struct bmu_client_handler **handler)
+{
+ struct bmu_codec_client *codec = (struct bmu_codec_client *)*handler;
+ codec->dec->free(&codec->dec);
+ al_free(codec);
+ *handler = NULL;
+}
+
+struct bmu_client_handler *bmu_codec_client_create(void)
+{
+ struct bmu_codec_client *codec = al_alloc_object(struct bmu_codec_client);
+ codec->handler.init = codec_client_init;
+ codec->handler.handle_data_packet = codec_client_handle_data_packet;
+ codec->handler.handle_eof = codec_client_handle_eof;
+ codec->handler.flush = codec_client_flush;
+ codec->handler.free = codec_client_free;
+ return (struct bmu_client_handler *)codec;
+}
diff --git a/src/bimu/handlers/codec_server.c b/src/bimu/handlers/codec_server.c
new file mode 100644
index 0000000..b555370
--- /dev/null
+++ b/src/bimu/handlers/codec_server.c
@@ -0,0 +1,159 @@
+#include "../../codec/codec.h"
+#include "../../codec/libav/demuxer.h"
+#include "../../codec/libav/packet_ext.h"
+#include "../../codec/stb_image/demuxer.h"
+#include "../../codec/wuffs/demuxer.h"
+#include "../../codec/spng/demuxer.h"
+
+#include "../handler.h"
+
+#include "codec.h"
+
+static bool codec_server_init(struct bmu_server_handler *handler, struct cch_handle *handle, struct aki_packet_pool *pool)
+{
+ struct bmu_codec_server *codec = (struct bmu_codec_server *)handler;
+ codec->demux = camu_lav_demuxer_create();
+ //codec->demux = camu_spng_demuxer_create();
+ //codec->demux = camu_stbi_demuxer_create();
+ //codec->demux = camu_wuffs_demuxer_create();
+ if (!codec->demux->init(codec->demux, handle)) {
+ return false;
+ }
+ switch (codec->demux->type) {
+ case CAMU_NORMAL:
+ break;
+#ifdef HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT:
+ codec->packet.av.pkt = av_packet_alloc();
+ break;
+#endif
+ }
+ codec->pool = pool;
+ return true;
+}
+
+static void codec_server_write_info(struct bmu_server_handler *handler, struct aki_packet *packet)
+{
+ struct bmu_codec_server *codec = (struct bmu_codec_server *)handler;
+ aki_packet_write_u32(packet, codec->demux->streams.size);
+ struct camu_stream *stream;
+ al_array_foreach_ptr(codec->demux->streams, i, stream) {
+ aki_packet_write_u8(packet, stream->type);
+ switch (stream->type) {
+ case CAMU_NORMAL:
+ aki_packet_write_s32(packet, stream->video.width);
+ aki_packet_write_s32(packet, stream->video.height);
+ break;
+#ifdef 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 s64 codec_server_get_duration(struct bmu_server_handler *handler)
+{
+ struct bmu_codec_server *codec = (struct bmu_codec_server *)handler;
+ return codec->demux->get_duration(codec->demux);
+}
+
+static bool codec_server_seek(struct bmu_server_handler *handler, s64 pos)
+{
+ struct bmu_codec_server *codec = (struct bmu_codec_server *)handler;
+ return codec->demux->seek(codec->demux, pos);
+}
+
+static bool codec_server_step(struct bmu_server_handler *handler)
+{
+ struct bmu_codec_server *codec = (struct bmu_codec_server *)handler;
+ struct aki_packet *packet = aki_packet_pool_get(codec->pool);
+ if (!packet) return false;
+ aki_packet_write_u8(packet, 0);
+ s32 ret = codec->demux->get_packet(codec->demux, &codec->packet);
+ do {
+ if (ret == CAMU_OK) {
+ switch (codec->packet.type) {
+ case CAMU_NORMAL: {
+ aki_packet_write_s32(packet, 0);
+ aki_packet_write_u8(packet, 0);
+ aki_packet_write_u8(packet, codec->packet.type);
+ aki_packet_write_buffer(packet, codec->packet.buffer);
+ break;
+ }
+#ifdef 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, 0);
+ aki_packet_write_u8(packet, codec->packet.type);
+ aki_packet_write_av_packet(packet, pkt);
+ av_packet_unref(pkt);
+ break;
+ }
+#endif
+ }
+ aki_packet_pool_submit(codec->pool, packet);
+ } else if (ret == CAMU_ERR_EOF) {
+ struct camu_stream *stream;
+ al_array_foreach_ptr(codec->demux->streams, i, stream) {
+ if (i != 0) { // HACK
+ packet = aki_packet_pool_get(codec->pool);
+ if (!packet) return false;
+ aki_packet_write_u8(packet, 0);
+ }
+ switch (stream->type) {
+ case CAMU_NORMAL:
+ aki_packet_write_s32(packet, 0);
+ break;
+#ifdef HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT:
+ aki_packet_write_s32(packet, stream->av.stream->index);
+ break;
+#endif
+ }
+ aki_packet_write_u8(packet, 1);
+ aki_packet_pool_submit(codec->pool, packet);
+ }
+ return false;
+ } else {
+ aki_packet_write_s32(packet, -2);
+ aki_packet_write_u8(packet, 2);
+ aki_packet_pool_submit(codec->pool, packet);
+ return false;
+ }
+ break;
+ } while (1);
+ return true;
+}
+
+static void codec_server_free(struct bmu_server_handler **handler)
+{
+ struct bmu_codec_server *codec = (struct bmu_codec_server *)*handler;
+ switch (codec->demux->type) {
+ case CAMU_NORMAL:
+ break;
+#ifdef 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 bmu_server_handler *bmu_codec_server_create(void)
+{
+ struct bmu_codec_server *codec = al_alloc_object(struct bmu_codec_server);
+ codec->handler.init = codec_server_init;
+ codec->handler.write_info = codec_server_write_info;
+ codec->handler.get_duration = codec_server_get_duration;
+ codec->handler.seek = codec_server_seek;
+ codec->handler.step = codec_server_step;
+ codec->handler.free = codec_server_free;
+ return (struct bmu_server_handler *)codec;
+}
diff --git a/src/bimu/local.c b/src/bimu/local.c
new file mode 100644
index 0000000..3f62f75
--- /dev/null
+++ b/src/bimu/local.c
@@ -0,0 +1,281 @@
+#include "../codec/libav/packet_ext.h"
+
+#include "local.h"
+
+#define PACKETS 1400
+
+static struct bmu_local_stream *stream_at_index(struct bmu_local *runner, s32 index)
+{
+ struct bmu_local_stream *stream = NULL;
+ al_array_foreach(runner->streams, i, stream) {
+ if (stream->s.index == index) {
+ return stream;
+ }
+ }
+ return NULL;
+}
+
+static aki_thread_result AKI_THREADCALL stream_thread(void *userdata)
+{
+ struct bmu_local_stream *stream = (struct bmu_local_stream *)userdata;
+ do {
+ if (!aki_packet_cache_wait(&stream->cache)) break;
+ struct aki_packet *packet = aki_packet_cache_pop(&stream->cache);
+ if (!packet) break;
+ aki_mutex_lock(&stream->mutex);
+ if (al_atomic_s32_load(&stream->state, AL_ATOMIC_RELAXED) == BIMU_LOCAL_STOPPED) {
+ aki_cond_wait(&stream->cond, &stream->mutex);
+ }
+ s32 state = al_atomic_s32_load(&stream->state, AL_ATOMIC_RELAXED);
+ aki_mutex_unlock(&stream->mutex);
+ if (state == BIMU_LOCAL_CLOSED) {
+ aki_packet_pool_return(&stream->runner->pool, packet);
+ break;
+ }
+ switch (aki_packet_read_u8(packet)) {
+ case 0:
+ stream->client->handle_data_packet(stream->client, packet);
+ break;
+ case 1:
+ stream->client->handle_eof(stream->client);
+ break;
+ }
+ aki_packet_pool_return(&stream->runner->pool, packet);
+ } while (1);
+ return 0;
+}
+
+static u8 packet_pool_callback(void *userdata, struct aki_packet *packet)
+{
+ struct bmu_local *runner = (struct bmu_local *)userdata;
+ if (!packet) {
+ aki_event_loop_break(&runner->loop);
+ return AKI_PACKET_POOL_CLOSED;
+ }
+ aki_packet_write_size(packet);
+ u8 close = aki_packet_read_u8(packet);
+ if (close) return AKI_PACKET_POOL_DISABLE;
+ struct bmu_local_stream *stream = NULL;
+ s32 index = aki_packet_read_s32(packet);
+ if (index >= 0) {
+ stream = stream_at_index(runner, index);
+ if (!stream) return AKI_PACKET_POOL_RETURN;
+ if (al_atomic_s32_load(&stream->state, AL_ATOMIC_RELAXED) != BIMU_LOCAL_CLOSED) {
+ if (aki_packet_cache_send_packet(&stream->cache, packet)) {
+ return AKI_PACKET_POOL_KEEP;
+ }
+ }
+ }
+ return AKI_PACKET_POOL_RETURN;
+}
+
+bool bmu_local_init(struct bmu_local *runner, struct cch_entry *entry)
+{
+ runner->entry = entry;
+ runner->server = bmu_codec_server_create();
+ cch_entry_get_handle(runner->entry, &runner->handle);
+ aki_event_loop_init(&runner->loop);
+ aki_packet_pool_init(&runner->pool, PACKETS, &runner->loop, packet_pool_callback, runner);
+ if (!runner->server->init(runner->server, &runner->handle, &runner->pool)) {
+ return false;
+ }
+ al_array_init(runner->streams);
+ al_atomic_bool_store(&runner->closed, false, AL_ATOMIC_RELAXED);
+ return true;
+}
+
+bool bmu_local_prepare_clients(struct bmu_local *runner, struct camu_renderer *renderer)
+{
+ struct aki_packet *packet = aki_packet_create();
+ runner->server->write_info(runner->server, packet);
+ aki_packet_write_size(packet);
+ u32 count = aki_packet_read_u32(packet);
+ for (u32 i = 0; i < count; i++) {
+ struct bmu_local_stream *local = al_alloc_object(struct bmu_local_stream);
+ struct bmu_client_stream *stream = &local->s;
+ local->runner = runner;
+ local->client = bmu_codec_client_create();
+ local->client->callback = runner->callback;
+ local->client->userdata = runner->userdata;
+ aki_cond_init(&local->cond);
+ aki_mutex_init(&local->mutex);
+ aki_packet_cache_init(&local->cache, PACKETS);
+ al_atomic_s32_store(&local->state, BIMU_LOCAL_RUNNING, AL_ATOMIC_RELAXED);
+ stream->stream.type = aki_packet_read_u8(packet);
+ switch (stream->stream.type) {
+ case CAMU_NORMAL:
+ stream->index = 0;
+ stream->type = BIMU_STREAM_VIDEO;
+ stream->stream.video.width = aki_packet_read_s32(packet);
+ stream->stream.video.height = aki_packet_read_s32(packet);
+ break;
+#ifdef HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT: {
+ stream->stream.type = CAMU_FFMPEG_COMPAT;
+ stream->stream.av.format_context = avformat_alloc_context();
+ const AVCodec *codec = avcodec_find_decoder(aki_packet_read_av_codec_id(packet));
+ stream->stream.av.stream = aki_packet_read_av_stream(stream->stream.av.format_context, codec, packet);
+ stream->index = stream->stream.av.stream->index;
+ switch (stream->stream.av.stream->codecpar->codec_type) {
+ case AVMEDIA_TYPE_AUDIO:
+ stream->type = BIMU_STREAM_AUDIO;
+ break;
+ case AVMEDIA_TYPE_VIDEO:
+ stream->type = BIMU_STREAM_VIDEO;
+ break;
+ case AVMEDIA_TYPE_SUBTITLE:
+ default:
+ stream->type = BIMU_STREAM_UNKNOWN;
+ break;
+ }
+ }
+#endif
+ }
+ if (!local->client->init(local->client, renderer, stream)) {
+ aki_packet_free(packet);
+ return false;
+ }
+ aki_thread_create(&local->thread, stream_thread, stream);
+ al_array_push(runner->streams, local);
+ }
+ aki_packet_free(packet);
+ return true;
+}
+
+void bmu_local_stream_stop(struct bmu_local_stream *stream)
+{
+ aki_mutex_lock(&stream->mutex);
+ if (al_atomic_s32_load(&stream->state, AL_ATOMIC_RELAXED) == BIMU_LOCAL_RUNNING) {
+ al_atomic_s32_store(&stream->state, BIMU_LOCAL_STOPPED, AL_ATOMIC_RELAXED);
+ }
+ aki_mutex_unlock(&stream->mutex);
+}
+
+void bmu_local_stream_continue(struct bmu_local_stream *stream)
+{
+ aki_mutex_lock(&stream->mutex);
+ al_atomic_s32_store(&stream->state, BIMU_LOCAL_RUNNING, AL_ATOMIC_RELAXED);
+ if (aki_cond_is_waiting(&stream->cond)) {
+ aki_cond_signal(&stream->cond);
+ }
+ aki_mutex_unlock(&stream->mutex);
+}
+
+void bmu_local_stream_close(struct bmu_local_stream *stream)
+{
+ aki_mutex_lock(&stream->mutex);
+ u8 state = al_atomic_s32_load(&stream->state, AL_ATOMIC_RELAXED);
+ if (state == BIMU_LOCAL_CLOSED) {
+ aki_mutex_unlock(&stream->mutex);
+ return;
+ }
+ al_atomic_s32_store(&stream->state, BIMU_LOCAL_CLOSED, AL_ATOMIC_RELAXED);
+ aki_packet_cache_disable(&stream->cache);
+ if (aki_cond_is_waiting(&stream->cond)) {
+ aki_cond_signal(&stream->cond);
+ }
+ aki_mutex_unlock(&stream->mutex);
+ aki_thread_join(&stream->thread);
+}
+
+static aki_thread_result AKI_THREADCALL decode_thread(void *userdata)
+{
+ struct bmu_local *runner = (struct bmu_local *)userdata;
+ while (!al_atomic_bool_load(&runner->closed, AL_ATOMIC_RELAXED)
+ && runner->server->step(runner->server)) {}
+ return 0;
+}
+
+static aki_thread_result AKI_THREADCALL event_loop_thread(void *userdata)
+{
+ struct bmu_local *runner = (struct bmu_local *)userdata;
+ aki_event_loop_run(&runner->loop);
+ return 0;
+}
+
+void bmu_local_run(struct bmu_local *runner)
+{
+ aki_thread_create(&runner->decode_thread, decode_thread, runner);
+ aki_thread_create(&runner->event_loop_thread, event_loop_thread, runner);
+}
+
+f64 bmu_local_get_duration(struct bmu_local *runner)
+{
+ return runner->server->get_duration(runner->server) / 1000000.0;
+}
+
+void bmu_local_seek(struct bmu_local *runner, f64 percent)
+{
+ s64 pos = runner->server->get_duration(runner->server) * percent;
+ f64 fpos = pos / 1000000.0;
+ al_atomic_bool_store(&runner->closed, true, AL_ATOMIC_RELAXED);
+ aki_packet_pool_disable(&runner->pool);
+ struct bmu_local_stream *stream = NULL;
+ al_array_foreach(runner->streams, i, stream) {
+ bmu_local_stream_close(stream);
+ struct aki_packet *packet;
+ while ((packet = aki_packet_cache_pop(&stream->cache))) {
+ aki_packet_pool_return(&stream->runner->pool, packet);
+ }
+ }
+ aki_thread_join(&runner->decode_thread);
+ runner->callback(runner->userdata, BIMU_CLIENT_SEEK, NULL, &fpos);
+ runner->server->seek(runner->server, pos);
+ al_atomic_bool_store(&runner->closed, false, AL_ATOMIC_RELAXED);
+ aki_packet_pool_enable(&runner->pool);
+ al_array_foreach(runner->streams, i, stream) {
+ stream->client->flush(stream->client);
+ al_atomic_s32_store(&stream->state, BIMU_LOCAL_RUNNING, AL_ATOMIC_RELAXED);
+ aki_packet_cache_enable(&stream->cache);
+ aki_thread_create(&stream->thread, stream_thread, stream);
+ }
+ aki_thread_create(&runner->decode_thread, decode_thread, runner);
+}
+
+void bmu_local_stop(struct bmu_local *runner)
+{
+ al_atomic_bool_store(&runner->closed, true, AL_ATOMIC_RELAXED);
+ cch_handle_disable(&runner->handle);
+ struct bmu_local_stream *stream = NULL;
+ al_array_foreach(runner->streams, i, stream) {
+ bmu_local_stream_close(stream);
+ }
+ aki_thread_join(&runner->decode_thread);
+ cch_entry_return_handle(runner->entry, &runner->handle);
+ struct aki_packet *packet = aki_packet_pool_get(&runner->pool);
+ if (!packet) al_assert(false);
+ aki_packet_write_u8(packet, 1);
+ aki_packet_pool_submit(&runner->pool, packet);
+ aki_thread_join(&runner->event_loop_thread);
+ runner->callback(runner->userdata, BIMU_CLIENT_CLOSED, NULL, NULL);
+}
+
+void bmu_local_close(struct bmu_local *runner)
+{
+ struct bmu_local_stream *stream = NULL;
+ al_array_foreach(runner->streams, i, stream) {
+ stream->client->free(&stream->client);
+ switch (stream->s.stream.type) {
+ case CAMU_NORMAL:
+ break;
+#ifdef HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT:
+ avformat_free_context(stream->s.stream.av.format_context);
+ break;
+#endif
+ }
+ struct aki_packet *packet;
+ while ((packet = aki_packet_cache_pop(&stream->cache))) {
+ aki_packet_pool_return(&stream->runner->pool, packet);
+ }
+ aki_mutex_destroy(&stream->mutex);
+ aki_cond_destroy(&stream->cond);
+ aki_packet_cache_free(&stream->cache);
+ al_free(stream);
+ al_array_remove_at_iter(runner->streams, i);
+ }
+ al_array_free(runner->streams);
+ runner->server->free(&runner->server);
+ aki_packet_pool_free(&runner->pool);
+ aki_event_loop_destroy(&runner->loop);
+}
diff --git a/src/bimu/local.h b/src/bimu/local.h
new file mode 100644
index 0000000..1928249
--- /dev/null
+++ b/src/bimu/local.h
@@ -0,0 +1,53 @@
+#pragma once
+
+#include <al/atomic.h>
+#include <aki/packet_pool.h>
+#include <aki/packet_cache.h>
+#include <aki/signal.h>
+
+#include "../cache/entry.h"
+
+#include "client.h"
+#include "common.h"
+
+enum {
+ BIMU_LOCAL_RUNNING = 0,
+ BIMU_LOCAL_STOPPED,
+ BIMU_LOCAL_CLOSED
+};
+
+struct bmu_local_stream {
+ struct bmu_client_stream s;
+ struct bmu_client_handler *client;
+ atomic_s32 state;
+ struct aki_cond cond;
+ struct aki_mutex mutex;
+ struct aki_packet_cache cache;
+ struct aki_thread thread;
+ struct bmu_local *runner;
+};
+
+struct bmu_local {
+ struct cch_entry *entry;
+ struct cch_handle handle;
+ struct bmu_server_handler *server;
+ struct aki_thread decode_thread;
+ struct aki_event_loop loop;
+ struct aki_packet_pool pool;
+ struct aki_thread event_loop_thread;
+ atomic_bool closed;
+ array(struct bmu_local_stream *) streams;
+ void (*callback)(void *, u8, struct bmu_client_stream *stream, void *);
+ void *userdata;
+};
+
+bool bmu_local_init(struct bmu_local *runner, struct cch_entry *entry);
+bool bmu_local_prepare_clients(struct bmu_local *runner, struct camu_renderer *renderer);
+void bmu_local_stream_stop(struct bmu_local_stream *stream);
+void bmu_local_stream_continue(struct bmu_local_stream *stream);
+void bmu_local_stream_close(struct bmu_local_stream *stream);
+void bmu_local_run(struct bmu_local *runner);
+f64 bmu_local_get_duration(struct bmu_local *runner);
+void bmu_local_seek(struct bmu_local *runner, f64 percent);
+void bmu_local_stop(struct bmu_local *runner);
+void bmu_local_close(struct bmu_local *runner);
diff --git a/src/bimu/meson.build b/src/bimu/meson.build
new file mode 100644
index 0000000..086114d
--- /dev/null
+++ b/src/bimu/meson.build
@@ -0,0 +1,8 @@
+bimu_src = [
+ 'local.c',
+ 'server.c',
+ 'client.c',
+ 'handlers/codec_server.c',
+ 'handlers/codec_client.c'
+]
+bimu = declare_dependency(sources: bimu_src)
diff --git a/src/bimu/server.c b/src/bimu/server.c
new file mode 100644
index 0000000..c1b0c38
--- /dev/null
+++ b/src/bimu/server.c
@@ -0,0 +1,46 @@
+#include "server.h"
+
+static void packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet)
+{
+ (void)userdata;
+ (void)stream;
+ aki_packet_free(packet);
+}
+
+static void connection_closed_callback(void *userdata, struct aki_packet_stream *stream)
+{
+ (void)userdata;
+ (void)stream;
+}
+
+static void packet_sent_callback(void *userdata, struct aki_packet *packet)
+{
+ (void)userdata;
+ (void)packet;
+}
+
+static void flushed_callback(void *userdata)
+{
+ (void)userdata;
+}
+
+static void connection_callback(void *userdata, struct aki_packet_stream *stream)
+{
+ stream->userdata = userdata;
+ stream->packet_callback = packet_callback;
+ stream->packet_sent_callback = packet_sent_callback;
+ stream->connection_closed_callback = connection_closed_callback;
+ stream->flushed_callback = flushed_callback;
+}
+
+bool bmu_server_init(struct bmu_server *server)
+{
+ return aki_packet_stream_init(&server->server, AKI_SOCKET_TCP, connection_callback,
+ NULL, NULL, NULL, server);
+}
+
+void bmu_server_listen(struct bmu_server *server, struct aki_event_loop *loop, str *addr, s32 port)
+{
+ server->loop = loop;
+ aki_packet_stream_listen(&server->server, server->loop, addr, port);
+}
diff --git a/src/bimu/server.h b/src/bimu/server.h
new file mode 100644
index 0000000..ba3093e
--- /dev/null
+++ b/src/bimu/server.h
@@ -0,0 +1,21 @@
+#pragma once
+
+#include <aki/packet_stream.h>
+
+#include "../cache/entry.h"
+
+struct bmu_node_connection {
+ struct cch_handle handle;
+};
+
+struct bmu_node {
+ struct cch_entry *entry;
+};
+
+struct bmu_server {
+ struct aki_packet_stream server;
+ struct aki_event_loop *loop;
+};
+
+bool bmu_server_init(struct bmu_server *server);
+void bmu_server_listen(struct bmu_server *server, struct aki_event_loop *loop, str *addr, s32 port);