summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c85
-rw-r--r--src/liana/client.h8
-rw-r--r--src/liana/handler.h8
-rw-r--r--src/liana/handlers.c2
-rw-r--r--src/liana/handlers/cdio.h2
-rw-r--r--src/liana/handlers/cdio_client.c12
-rw-r--r--src/liana/handlers/cdio_server.c42
-rw-r--r--src/liana/handlers/codec_client.c12
-rw-r--r--src/liana/handlers/codec_server.c44
-rw-r--r--src/liana/list.c14
-rw-r--r--src/liana/meson.build1
-rw-r--r--src/liana/server.c157
-rw-r--r--src/liana/server.h32
-rw-r--r--src/liana/vcr.c107
-rw-r--r--src/liana/vcr.h22
15 files changed, 273 insertions, 275 deletions
diff --git a/src/liana/client.c b/src/liana/client.c
index ad365b3..1ba09cc 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -9,24 +9,24 @@
#include "handlers.h"
#include "list.h"
-static void data_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet)
+static void data_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_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)
+static void parse_info_packet(struct lia_client *client, struct nn_packet *packet)
{
str liana;
- aki_packet_read_str(packet, &liana);
- client->duration = aki_packet_read_u64(packet);
- u32 count = aki_packet_read_u32(packet);
+ nn_packet_read_str(packet, &liana);
+ client->duration = nn_packet_read_u64(packet);
+ u32 count = nn_packet_read_u32(packet);
for (u32 i = 0; i < count; i++) {
- u8 mode = aki_packet_read_u8(packet);
- u8 type = aki_packet_read_u8(packet);
- u64 duration = aki_packet_read_u64(packet);
- s32 index = aki_packet_read_s32(packet);
+ u8 mode = nn_packet_read_u8(packet);
+ u8 type = nn_packet_read_u8(packet);
+ u64 duration = nn_packet_read_u64(packet);
+ s32 index = nn_packet_read_s32(packet);
al_assert(index < 32);
struct lia_vcr_track *track = NULL;
switch (mode) {
@@ -35,17 +35,17 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack
track = al_alloc_object(struct lia_vcr_track);
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);
+ fmt->format = nn_packet_read_s32(packet);
+ fmt->sample_rate = nn_packet_read_s32(packet);
+ fmt->channel_count = nn_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) {
struct camu_video_format *fmt = &track->stream.video.fmt;
- fmt->width = aki_packet_read_s32(packet);
- fmt->height = aki_packet_read_s32(packet);
- fmt->format = aki_packet_read_s32(packet);
+ fmt->width = nn_packet_read_s32(packet);
+ fmt->height = nn_packet_read_s32(packet);
+ fmt->format = nn_packet_read_s32(packet);
}
break;
}
@@ -67,10 +67,10 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack
default:
continue;
}
- enum AVCodecID codec_id = aki_packet_read_av_codec_id(packet);
+ enum AVCodecID codec_id = nn_packet_read_av_codec_id(packet);
const AVCodec *codec = avcodec_find_decoder(codec_id);
AVFormatContext *format_context = avformat_alloc_context();
- AVStream *stream = aki_packet_read_av_stream(format_context, codec, packet);
+ AVStream *stream = nn_packet_read_av_stream(format_context, codec, packet);
if (type == CAMU_STREAM_SUBTITLE && codec_id != AV_CODEC_ID_ASS) {
client->mask &= ~(1 << index);
continue;
@@ -120,30 +120,30 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack
}
}
-static void info_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet)
+static void info_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
struct lia_client *client = (struct lia_client *)userdata;
- client->connection_id = aki_packet_read_u32(packet);
+ client->connection_id = nn_packet_read_u32(packet);
parse_info_packet(client, packet);
- aki_packet_free(packet);
+ nn_packet_free(packet);
if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) {
client->reconnect = false;
- aki_packet_stream_disconnect(&client->data);
+ nn_packet_stream_disconnect(&client->data);
return;
}
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);
+ struct nn_packet *rpacket = nn_packet_create();
+ nn_packet_write_s32(rpacket, client->mask);
+ nn_packet_stream_send_packet(stream, rpacket);
}
-static void packet_sent_callback(void *userdata, struct aki_packet *packet)
+static void packet_sent_callback(void *userdata, struct nn_packet *packet)
{
(void)userdata;
- aki_packet_free(packet);
+ nn_packet_free(packet);
}
-static void connection_callback(void *userdata, struct aki_packet_stream *stream)
+static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_client *client = (struct lia_client *)userdata;
if (client->reconnect) {
@@ -164,20 +164,21 @@ static void connection_callback(void *userdata, struct aki_packet_stream *stream
}
lia_vcr_start(&client->vcr);
stream->packet_sent_callback = packet_sent_callback;
- struct aki_packet *packet = aki_packet_create();
- aki_packet_write_u32(packet, client->id);
- aki_packet_write_u32(packet, 0);
- aki_packet_write_s32(packet, client->mask);
- aki_packet_write_u64(packet, client->pos);
+ struct nn_packet *packet = nn_packet_create();
+ nn_packet_write_u32(packet, client->id);
+ nn_packet_write_u32(packet, 0);
+ nn_packet_write_s32(packet, client->mask);
+ nn_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);
+ nn_packet_stream_send_packet(stream, packet);
+ return true;
}
-static void connection_closed_callback(void *userdata, struct aki_packet_stream *stream)
+static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_client *client = (struct lia_client *)userdata;
if (client->reconnect) {
@@ -190,13 +191,13 @@ static void connection_closed_callback(void *userdata, struct aki_packet_stream
// any unexpected behavior.
client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect);
if (client->reconnect) {
- aki_packet_stream_reconnect(stream, &client->addr, client->port);
+ nn_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,
+void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type,
str *addr, u16 port, u32 id, u64 pos, struct camu_renderer *renderer)
{
client->loop = loop;
@@ -207,12 +208,12 @@ void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop,
lia_vcr_init(&client->vcr, client->loop, &client->data);
al_str_clone(&client->addr, addr);
client->port = port;
- if (!aki_packet_stream_init(&client->data, type, connection_callback, connection_closed_callback, client)) {
+ if (!nn_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);
+ nn_packet_stream_set_multiplex(&client->data, CAMU_MULTIPLEX_LIANA);
client->renderer = renderer;
- aki_packet_stream_connect(&client->data, client->loop, &client->addr, client->port);
+ nn_packet_stream_connect(&client->data, client->loop, &client->addr, client->port);
}
void lia_client_set_renderer(struct lia_client *client, struct camu_renderer *renderer)
@@ -226,7 +227,7 @@ void lia_client_seek(struct lia_client *client, u64 pos, u64 at)
client->at = at;
if (!client->reconnect) {
client->reconnect = true;
- aki_packet_stream_disconnect(&client->data);
+ nn_packet_stream_disconnect(&client->data);
}
}
@@ -238,12 +239,12 @@ void lia_client_reseek(struct lia_client *client)
void lia_client_disconnect(struct lia_client *client)
{
client->reconnect = false;
- aki_packet_stream_disconnect(&client->data);
+ nn_packet_stream_disconnect(&client->data);
}
void lia_client_free(struct lia_client *client)
{
lia_vcr_free(&client->vcr);
- aki_packet_stream_free(&client->data);
+ nn_packet_stream_free(&client->data);
al_str_free(&client->addr);
}
diff --git a/src/liana/client.h b/src/liana/client.h
index cd1297e..55c5b3a 100644
--- a/src/liana/client.h
+++ b/src/liana/client.h
@@ -1,13 +1,13 @@
#pragma once
-#include <aki/packet_stream.h>
+#include <nnwt/packet_stream.h>
#include "../codec/codec.h"
#include "vcr.h"
struct lia_client {
- struct aki_event_loop *loop;
+ struct nn_event_loop *loop;
u32 id;
s32 mask;
u64 pos;
@@ -16,7 +16,7 @@ struct lia_client {
str addr;
u16 port;
u32 connection_id;
- struct aki_packet_stream data;
+ struct nn_packet_stream data;
u64 duration;
struct lia_vcr vcr;
struct camu_renderer *renderer;
@@ -24,7 +24,7 @@ struct lia_client {
void *userdata;
};
-void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop, u8 type,
+void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type,
str *addr, u16 port, u32 id, u64 pos, struct camu_renderer *renderer);
void lia_client_seek(struct lia_client *client, u64 pos, u64 at);
void lia_client_reseek(struct lia_client *client);
diff --git a/src/liana/handler.h b/src/liana/handler.h
index 21437c4..9f93e5c 100644
--- a/src/liana/handler.h
+++ b/src/liana/handler.h
@@ -1,18 +1,18 @@
#pragma once
-#include <aki/packet_pool.h>
+#include <nnwt/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 (*write_info)(struct lia_server_handler *, struct nn_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 (*write_packet)(struct lia_server_handler *, struct nn_packet *);
void (*free)(struct lia_server_handler **);
s32 status;
};
@@ -49,7 +49,7 @@ struct lia_seek_req {
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 *);
+ bool (*handle_packet)(struct lia_client_handler *, struct nn_packet *);
void (*flush)(struct lia_client_handler *);
void (*free)(struct lia_client_handler **);
struct camu_codec_stream *stream;
diff --git a/src/liana/handlers.c b/src/liana/handlers.c
index 936f801..1a28024 100644
--- a/src/liana/handlers.c
+++ b/src/liana/handlers.c
@@ -13,7 +13,7 @@ struct lia_handler_entry liana_handlers[] = {
.create_client_handler = lia_codec_client_create
#endif
},
-#ifdef LIANA_HAVE_CDIO
+#ifdef CACHE_HAVE_CDIO
{
.name = al_str_c("cdio"),
#ifdef LIANA_SERVER
diff --git a/src/liana/handlers/cdio.h b/src/liana/handlers/cdio.h
index 2727616..4c40986 100644
--- a/src/liana/handlers/cdio.h
+++ b/src/liana/handlers/cdio.h
@@ -6,8 +6,8 @@ struct lia_cdio_server {
struct lia_server_handler handler;
struct cch_handle *handle;
struct camu_audio_format fmt;
+ struct nn_buffer buffer;
f64 pts;
- struct aki_buffer buffer;
};
struct lia_cdio_client {
diff --git a/src/liana/handlers/cdio_client.c b/src/liana/handlers/cdio_client.c
index 93b1a27..665858f 100644
--- a/src/liana/handlers/cdio_client.c
+++ b/src/liana/handlers/cdio_client.c
@@ -9,7 +9,7 @@ static bool cdio_client_init(struct lia_client_handler *handler, struct camu_ren
return true;
}
-static bool cdio_client_handle_packet(struct lia_client_handler *handler, struct aki_packet *packet)
+static bool cdio_client_handle_packet(struct lia_client_handler *handler, struct nn_packet *packet)
{
struct lia_cdio_client *cdio = (struct lia_cdio_client *)handler;
if (!packet) {
@@ -18,12 +18,12 @@ static bool cdio_client_handle_packet(struct lia_client_handler *handler, struct
}
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->pts = nn_packet_read_f64(packet);
+ frame->audio.sample_count = nn_packet_read_s32(packet);
+ struct nn_buffer buffer;
+ nn_packet_read_buffer(packet, &buffer);
frame->data = al_malloc(buffer.size);
- al_memcpy(frame->data, aki_buffer_get_ptr(&buffer, 0), buffer.size);
+ al_memcpy(frame->data, nn_buffer_get_ptr(&buffer, 0), buffer.size);
cdio->handler.callback(cdio->handler.userdata, LIANA_CLIENT_DATA, cdio->handler.stream, frame);
return true;
}
diff --git a/src/liana/handlers/cdio_server.c b/src/liana/handlers/cdio_server.c
index 641ed92..01a9265 100644
--- a/src/liana/handlers/cdio_server.c
+++ b/src/liana/handlers/cdio_server.c
@@ -15,8 +15,8 @@ static bool cdio_server_init(struct lia_server_handler *handler, struct cch_hand
cdio->fmt.format = CAMU_SAMPLE_FORMAT_S16;
cdio->fmt.sample_rate = 44100;
cdio->fmt.channel_count = 2;
+ nn_buffer_init(&cdio->buffer);
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);
@@ -25,18 +25,18 @@ static bool cdio_server_init(struct lia_server_handler *handler, struct cch_hand
return true;
}
-static void cdio_server_write_info(struct lia_server_handler *handler, struct aki_packet *packet)
+static void cdio_server_write_info(struct lia_server_handler *handler, struct nn_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_u64(packet, cdio->handler.get_duration(&cdio->handler));
- aki_packet_write_s32(packet, 0);
- 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);
+ nn_packet_write_u64(packet, 0);
+ nn_packet_write_u32(packet, 1);
+ nn_packet_write_u8(packet, CAMU_NORMAL);
+ nn_packet_write_u8(packet, CAMU_STREAM_AUDIO);
+ nn_packet_write_u64(packet, cdio->handler.get_duration(&cdio->handler));
+ nn_packet_write_s32(packet, 0);
+ nn_packet_write_s32(packet, cdio->fmt.format);
+ nn_packet_write_s32(packet, cdio->fmt.sample_rate);
+ nn_packet_write_s32(packet, cdio->fmt.channel_count);
}
static void cdio_server_subscribe(struct lia_server_handler *handler, s32 mask)
@@ -69,26 +69,26 @@ static void cdio_server_step(struct lia_server_handler *handler)
(void)cdio;
}
-static void cdio_server_write_packet(struct lia_server_handler *handler, struct aki_packet *packet)
+static void cdio_server_write_packet(struct lia_server_handler *handler, struct nn_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);
+ nn_buffer_ensure_space(&cdio->buffer, size);
+ s32 ret = cch_handle_read(cdio->handle, nn_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);
+ nn_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);
+ nn_packet_write_u8(packet, LIANA_PACKET_DATA);
+ nn_packet_write_s32(packet, 0);
+ nn_packet_write_f64(packet, cdio->pts);
s32 sample_count = ret / (camu_audio_format_bytes_per_sample(&cdio->fmt) * cdio->fmt.channel_count);
+ nn_packet_write_s32(packet, sample_count);
+ nn_buffer_set_size(&cdio->buffer, ret);
+ nn_packet_write_buffer(packet, &cdio->buffer);
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)
diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c
index dc53129..3722197 100644
--- a/src/liana/handlers/codec_client.c
+++ b/src/liana/handlers/codec_client.c
@@ -44,7 +44,7 @@ static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt)
}
#endif
-static bool push_packet(struct lia_codec_client *codec, struct aki_buffer *buffer)
+static bool push_packet(struct lia_codec_client *codec, struct nn_buffer *buffer)
{
struct camu_codec_packet packet;
packet.buffer = buffer;
@@ -52,7 +52,7 @@ static bool push_packet(struct lia_codec_client *codec, struct aki_buffer *buffe
return ret == CAMU_OK;
}
-static bool codec_client_handle_packet(struct lia_client_handler *handler, struct aki_packet *packet)
+static bool codec_client_handle_packet(struct lia_client_handler *handler, struct nn_packet *packet)
{
struct lia_codec_client *codec = (struct lia_codec_client *)handler;
@@ -72,12 +72,12 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
// Push packet.
bool success;
- u8 type = aki_packet_read_u8(packet);
+ u8 type = nn_packet_read_u8(packet);
switch (type) {
case CAMU_NORMAL: {
if (codec->dec) {
- struct aki_buffer buffer;
- aki_packet_read_buffer(packet, &buffer);
+ struct nn_buffer buffer;
+ nn_packet_read_buffer(packet, &buffer);
success = push_packet(codec, &buffer);
} else {
success = true;
@@ -86,7 +86,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
}
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
- AVPacket *pkt = aki_packet_read_av_packet(packet);
+ AVPacket *pkt = nn_packet_read_av_packet(packet);
if (codec->dec) {
success = push_av_packet(codec, pkt);
} else {
diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c
index 94b18df..7167f97 100644
--- a/src/liana/handlers/codec_server.c
+++ b/src/liana/handlers/codec_server.c
@@ -33,30 +33,30 @@ static bool codec_server_init(struct lia_server_handler *handler, struct cch_han
return true;
}
-static void codec_server_write_info(struct lia_server_handler *handler, struct aki_packet *packet)
+static void codec_server_write_info(struct lia_server_handler *handler, struct nn_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);
+ nn_packet_write_u64(packet, duration);
+ nn_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);
- aki_packet_write_u8(packet, stream->type);
- aki_packet_write_u64(packet, stream->duration);
- aki_packet_write_s32(packet, i);
+ nn_packet_write_u8(packet, stream->mode);
+ nn_packet_write_u8(packet, stream->type);
+ nn_packet_write_u64(packet, stream->duration);
+ nn_packet_write_s32(packet, i);
switch (stream->mode) {
case CAMU_NORMAL: {
struct camu_video_format *fmt = &stream->video.fmt;
- aki_packet_write_s32(packet, fmt->width);
- aki_packet_write_s32(packet, fmt->height);
- aki_packet_write_s32(packet, fmt->format);
+ nn_packet_write_s32(packet, fmt->width);
+ nn_packet_write_s32(packet, fmt->height);
+ nn_packet_write_s32(packet, fmt->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);
+ nn_packet_write_av_codec_id(packet, stream->av.stream->codecpar->codec_id);
+ nn_packet_write_av_stream(packet, stream->av.stream);
break;
#endif
}
@@ -87,33 +87,33 @@ static void codec_server_step(struct lia_server_handler *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)
+static void codec_server_write_packet(struct lia_server_handler *handler, struct nn_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);
+ nn_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);
+ nn_packet_write_s32(packet, 0);
+ nn_packet_write_u8(packet, codec->packet.mode);
+ nn_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);
+ nn_packet_write_s32(packet, pkt->stream_index);
+ nn_packet_write_u8(packet, codec->packet.mode);
+ nn_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);
+ nn_packet_write_u8(packet, LIANA_PACKET_EOF);
} else {
- aki_packet_write_u8(packet, LIANA_PACKET_ERROR);
+ nn_packet_write_u8(packet, LIANA_PACKET_ERROR);
}
}
diff --git a/src/liana/list.c b/src/liana/list.c
index 073cdf0..7c4979b 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -1,4 +1,4 @@
-#include <aki/thread.h>
+#include <nnwt/thread.h>
#include <al/random.h>
#include <al/lib.h>
#include <al/log.h>
@@ -112,7 +112,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
if (error) pump_queue(list);
return false;
}
- u64 now = aki_get_timestamp();
+ u64 now = nn_get_timestamp();
u8 pause;
u64 at = LIANA_TIMESTAMP_INVALID;
u64 seek_pos = current->offset;
@@ -168,7 +168,7 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
}
list->current++;
list->idle = false;
- entry->start = aki_get_timestamp() + LIANA_BASE_DELAY;
+ entry->start = nn_get_timestamp() + LIANA_BASE_DELAY;
struct lia_timing time = {
.at = entry->start,
.seek_pos = entry->offset,
@@ -299,7 +299,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
return false;
}
- u64 now = aki_get_timestamp();
+ u64 now = nn_get_timestamp();
u64 at = now + LIANA_BASE_DELAY;
u8 pause;
@@ -405,7 +405,7 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
struct lia_list_entry *entry = get_entry_from_sequence(list, sequence);
al_assert(entry && !entry->held);
- u64 now = aki_get_timestamp();
+ u64 now = nn_get_timestamp();
u8 pause = entry->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_PAUSE : LIANA_PAUSE_RESUME;
u64 at;
bool ended = assume_ended(entry, now);
@@ -461,7 +461,7 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent
}
entry->ended = false;
- u64 now = aki_get_timestamp();
+ u64 now = nn_get_timestamp();
u64 pos = (u64)(entry->duration * percent);
u64 at = now + LIANA_BASE_DELAY;
u8 pause = entry->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE;
@@ -527,10 +527,12 @@ static bool handle_end(struct lia_list *list, s32 id)
return false;
} else {
list->idle = true;
+ /*
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
sink->set = -1;
}
+ */
}
}
diff --git a/src/liana/meson.build b/src/liana/meson.build
index 32e6a00..b68ba8b 100644
--- a/src/liana/meson.build
+++ b/src/liana/meson.build
@@ -16,7 +16,6 @@ liana_args = []
if cache_have_cdio
liana_server_src += ['handlers/cdio_server.c']
liana_client_src += ['handlers/cdio_client.c']
- liana_args += ['-DLIANA_HAVE_CDIO']
endif
liana_server = declare_dependency(sources: liana_server_src,
diff --git a/src/liana/server.c b/src/liana/server.c
index 44b3f12..27c253b 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -5,7 +5,7 @@
#include "handlers.h"
#include "list.h"
-bool lia_server_init(struct lia_server *server, struct aki_event_loop *loop)
+bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop)
{
server->loop = loop;
server->increment = 1;
@@ -14,52 +14,46 @@ bool lia_server_init(struct lia_server *server, struct aki_event_loop *loop)
return true;
}
-static void remove_zombie(struct lia_server *server, struct aki_packet_stream *stream)
+static void remove_zombie(struct lia_server *server, struct nn_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;
- }
- }
+ al_array_remove(server->zombies, stream);
}
-static void data_packet_sent_callback(void *userdata, struct aki_packet *packet)
+static void data_packet_sent_callback(void *userdata, struct nn_packet *packet)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
- aki_packet_pool_return(&conn->pool, packet);
+ nn_packet_pool_return(&conn->pool, packet);
}
-static u8 packet_pool_callback(void *userdata, struct aki_packet *packet)
+static u8 packet_pool_callback(void *userdata, struct nn_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;
+ if (!nn_packet_stream_send_packet(conn->stream, packet)) {
+ return NNWT_PACKET_POOL_RETURN;
}
- return AKI_PACKET_POOL_KEEP;
+ return NNWT_PACKET_POOL_KEEP;
}
-static aki_thread_result AKI_THREADCALL handler_thread(void *userdata)
+static nn_thread_result NNWT_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);
+ struct nn_packet *packet = nn_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);
+ nn_packet_pool_submit(&conn->pool, packet);
if (conn->handler->status != CAMU_OK) break;
} while (1);
- aki_packet_pool_flush(&conn->pool);
+ nn_packet_pool_flush(&conn->pool);
return 0;
}
-static void discard_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet)
+static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
(void)userdata;
(void)stream;
- aki_packet_free(packet);
+ nn_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);
@@ -70,83 +64,83 @@ 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);
+ nn_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)
+static void data_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
(void)stream;
- aki_packet_pool_disable(&conn->pool);
+ nn_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);
+ nn_thread_join(&conn->thread);
close_connection_internal(conn);
- aki_packet_pool_free(&conn->pool);
+ nn_packet_pool_free(&conn->pool);
}
-static void subscribe_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet)
+static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
- s32 mask = aki_packet_read_s32(packet);
+ s32 mask = nn_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);
+ nn_thread_create(&conn->thread, handler_thread, conn);
+ nn_packet_free(packet);
}
-static void subscribe_packet_sent_callback(void *userdata, struct aki_packet *packet)
+static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet)
{
(void)userdata;
- aki_packet_free(packet);
+ nn_packet_free(packet);
}
-static void subscribe_connection_closed_callback(void *userdata, struct aki_packet_stream *stream)
+static void subscribe_connection_closed_callback(void *userdata, struct nn_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);
+ nn_packet_pool_free(&conn->pool);
}
-static void handle_connection(struct lia_node_connection *conn, struct aki_packet *packet)
+static void handle_connection(struct lia_node_connection *conn, struct nn_packet *packet)
{
- struct aki_packet_stream *stream = conn->stream;
+ struct nn_packet_stream *stream = conn->stream;
- s32 mask = aki_packet_read_s32(packet);
- u64 seek_pos = aki_packet_read_u64(packet);
+ s32 mask = nn_packet_read_s32(packet);
+ u64 seek_pos = nn_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_u32(rpacket, conn->id);
- aki_packet_write_str(rpacket, cch_entry_get_liana(conn->node->entry));
+ struct nn_packet *rpacket = nn_packet_create();
+ nn_packet_write_u32(rpacket, conn->id);
+ nn_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);
+ nn_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);
+ nn_thread_create(&conn->thread, handler_thread, conn);
}
}
-static void connection_closed_callback(void *userdata, struct aki_packet_stream *stream)
+static void connection_closed_callback(void *userdata, struct nn_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);
+ nn_packet_stream_free(stream);
al_free(stream);
}
@@ -155,9 +149,9 @@ 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;
+ nn_signal_stop(&conn->signal);
+ nn_thread_join(&conn->thread);
+ struct nn_packet *packet = conn->packet;
if (!packet) {
// Connection was closed before init was done.
close_connection_internal(conn);
@@ -166,30 +160,30 @@ static void signal_callback(void *userdata)
if (!conn->errored) {
conn->id = server->increment;
server->increment = al_u32_inc_wrap(server->increment);
- aki_packet_pool_init(&conn->pool, 96, server->loop, packet_pool_callback, conn);
+ nn_packet_pool_init(&conn->pool, 96, 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;
+ struct nn_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);
+ nn_packet_stream_disconnect(stream);
}
- aki_packet_free(packet);
+ nn_packet_free(packet);
}
-static aki_thread_result AKI_THREADCALL init_thread(void *userdata)
+static nn_thread_result NNWT_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);
+ nn_signal_send(&conn->signal);
return 0;
}
@@ -211,24 +205,24 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id)
return NULL;
}
-static void pre_init_connection_closed_callback(void *userdata, struct aki_packet_stream *stream)
+static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
(void)stream;
- aki_packet_free(conn->packet);
+ nn_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)
+static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_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);
- u32 node_id = aki_packet_read_u32(packet);
- u32 connection_id = aki_packet_read_u32(packet);
+ u32 node_id = nn_packet_read_u32(packet);
+ u32 connection_id = nn_packet_read_u32(packet);
struct lia_node *node = get_node_from_id(server, node_id);
struct lia_node_connection *conn = NULL;
@@ -238,46 +232,47 @@ static void packet_callback(void *userdata, struct aki_packet_stream *stream, st
conn->node = node;
conn->stream = stream;
conn->packet = packet;
- aki_signal_init(&conn->signal, server->loop, signal_callback, conn);
- aki_signal_start(&conn->signal);
+ nn_signal_init(&conn->signal, server->loop, signal_callback, conn);
+ nn_signal_start(&conn->signal);
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);
+ nn_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);
+ nn_packet_stream_disconnect(stream);
}
- aki_packet_free(packet);
+ nn_packet_free(packet);
}
}
-static void packet_sent_callback(void *userdata, struct aki_packet *packet)
+static void packet_sent_callback(void *userdata, struct nn_packet *packet)
{
(void)userdata;
- aki_packet_free(packet);
+ nn_packet_free(packet);
}
-static void connection_callback(void *userdata, struct aki_packet_stream *stream)
+static bool connection_callback(void *userdata, struct nn_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);
+ return true;
}
-void lia_server_add_socket(struct lia_server *server, struct aki_socket *sock)
+void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock)
{
- struct aki_packet_stream *stream = al_alloc_object(struct aki_packet_stream);
+ struct nn_packet_stream *stream = al_alloc_object(struct nn_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);
+ nn_packet_stream_from_socket(stream, server->loop, sock);
}
struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry)
@@ -292,7 +287,7 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en
return node;
}
-static aki_thread_result AKI_THREADCALL init_duration_thread(void *userdata)
+static nn_thread_result NNWT_THREADCALL init_duration_thread(void *userdata)
{
struct lia_node *node = (struct lia_node *)userdata;
if (!node->handler->init(node->handler, &node->handle)) {
@@ -301,15 +296,15 @@ static aki_thread_result AKI_THREADCALL init_duration_thread(void *userdata)
} else {
node->duration = node->handler->get_duration(node->handler);
}
- aki_signal_send(&node->signal);
+ nn_signal_send(&node->signal);
return 0;
}
static void duration_signal_callback(void *userdata)
{
struct lia_node *node = (struct lia_node *)userdata;
- aki_signal_stop(&node->signal);
- aki_thread_join(&node->thread);
+ nn_signal_stop(&node->signal);
+ nn_thread_join(&node->thread);
node->handler->free(&node->handler);
cch_entry_return_handle(node->entry, &node->handle);
node->callback(node->userdata, LIANA_NODE_DURATION, node->duration);
@@ -317,11 +312,11 @@ static void duration_signal_callback(void *userdata)
void lia_node_get_duration(struct lia_node *node)
{
- aki_signal_init(&node->signal, node->server->loop, duration_signal_callback, node);
- aki_signal_start(&node->signal);
+ nn_signal_init(&node->signal, node->server->loop, duration_signal_callback, node);
+ nn_signal_start(&node->signal);
cch_entry_get_handle(node->entry, &node->handle);
node->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler();
- aki_thread_create(&node->thread, init_duration_thread, node);
+ nn_thread_create(&node->thread, init_duration_thread, node);
}
void lia_server_close(struct lia_server *server)
@@ -331,13 +326,13 @@ void lia_server_close(struct lia_server *server)
struct lia_node_connection *connection;
al_array_foreach(node->connections, j, connection) {
if (connection->stream) {
- aki_packet_stream_disconnect(connection->stream);
+ nn_packet_stream_disconnect(connection->stream);
}
}
}
- struct aki_packet_stream *zombie;
+ struct nn_packet_stream *zombie;
al_array_foreach_rev(server->zombies, i, zombie) {
- aki_packet_stream_disconnect(zombie);
+ nn_packet_stream_disconnect(zombie);
}
}
@@ -353,7 +348,7 @@ void lia_server_free(struct lia_server *server)
al_free(node);
}
al_array_free(server->nodes);
- struct aki_packet_stream *zombie;
+ struct nn_packet_stream *zombie;
al_array_foreach(server->zombies, i, zombie) {
al_free(zombie);
}
diff --git a/src/liana/server.h b/src/liana/server.h
index c62eb29..b4ada58 100644
--- a/src/liana/server.h
+++ b/src/liana/server.h
@@ -1,23 +1,23 @@
#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 <nnwt/event_loop.h>
+#include <nnwt/packet_stream.h>
+#include <nnwt/packet_pool.h>
+#include <nnwt/socket.h>
+#include <nnwt/signal.h>
#include "../cache/entry.h"
struct lia_node_connection {
u32 id;
- struct aki_packet *packet;
- struct aki_packet_stream *stream;
+ struct nn_packet *packet;
+ struct nn_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 nn_thread thread;
+ struct nn_signal signal;
+ struct nn_packet_pool pool;
struct lia_node *node;
};
@@ -36,21 +36,21 @@ struct lia_node {
struct lia_server_handler *handler;
bool errored;
struct cch_handle handle;
- struct aki_thread thread;
- struct aki_signal signal;
+ struct nn_thread thread;
+ struct nn_signal signal;
void (*callback)(void *, u8, u64);
void *userdata;
};
struct lia_server {
- struct aki_event_loop *loop;
+ struct nn_event_loop *loop;
u32 increment;
array(struct lia_node *) nodes;
- array(struct aki_packet_stream *) zombies;
+ array(struct nn_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);
+bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop);
+void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock);
struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry);
void lia_node_get_duration(struct lia_node *node);
void lia_server_close(struct lia_server *server);
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index c4bf583..4131c80 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -19,7 +19,7 @@ enum {
static void signal_callback(void *userdata)
{
struct lia_vcr *vcr = (struct lia_vcr *)userdata;
- aki_packet_stream_cork(vcr->data, false);
+ nn_packet_stream_cork(vcr->data, false);
}
static void reset_metrics(struct lia_vcr *vcr)
@@ -28,7 +28,7 @@ static void reset_metrics(struct lia_vcr *vcr)
vcr->metric.last_report_ts = 0Lu;
}
-void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_packet_stream *data)
+void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data)
{
al_array_init(vcr->tracks);
al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED);
@@ -36,25 +36,25 @@ void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_p
vcr->mark.low = 0;
vcr->expand = VCR_EXPAND_UNTOUCHED;
vcr->data = data;
- aki_signal_init(&vcr->signal, loop, signal_callback, vcr);
+ nn_signal_init(&vcr->signal, loop, signal_callback, vcr);
reset_metrics(vcr);
}
void lia_vcr_start(struct lia_vcr *vcr)
{
- aki_signal_start(&vcr->signal);
+ nn_signal_start(&vcr->signal);
}
-static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata)
+static nn_thread_result NNWT_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)) {
+ struct nn_packet *packet = NULL;
+ while (nn_packet_cache_wait(&track->cache, &packets)) {
s32 state = 0;
for (u32 i = 0; i < packets; i++) {
- packet = aki_packet_cache_pop(&track->cache);
+ packet = nn_packet_cache_pop(&track->cache);
state = 0;
if (packet) {
state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
@@ -62,11 +62,11 @@ static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata)
// We were signaled to close, exit thread.
goto out;
}
- u32 size = aki_packet_get_size(packet);
+ u32 size = nn_packet_get_size(packet);
u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED);
u64 buffer = al_atomic_sub(u64)(&vcr->count, size, AL_ATOMIC_RELAXED);
if (buffered && buffer <= vcr->mark.low) {
- aki_signal_send(&vcr->signal);
+ nn_signal_send(&vcr->signal);
}
}
if (!track->client->handle_packet(track->client, packet)) {
@@ -74,14 +74,14 @@ static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata)
goto out;
}
if (packet) {
- aki_packet_free(packet);
+ nn_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);
+ nn_packet_cache_unlock(&track->cache);
+ nn_mutex_lock(&track->mutex);
+ nn_cond_wait(&track->cond, &track->mutex);
+ nn_mutex_unlock(&track->mutex);
break;
}
} else {
@@ -91,25 +91,25 @@ static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata)
}
if (state != LIANA_STREAM_STOPPED) {
// If state = STOPPED, we already unlocked.
- aki_packet_cache_unlock(&track->cache);
+ nn_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);
+ if (packet) nn_packet_free(packet);
+ nn_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);
+ nn_cond_init(&track->cond);
+ nn_mutex_init(&track->mutex);
al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED);
- aki_packet_cache_init(&track->cache, 256);
+ nn_packet_cache_init(&track->cache, 256);
al_array_push(vcr->tracks, track);
al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
}
@@ -145,9 +145,9 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer)
vcr->mark.low = vcr->mark.buffered - MB(2);
vcr->expand = VCR_EXPAND_COMPLETE;
}
- aki_packet_stream_cork(vcr->data, true);
+ nn_packet_stream_cork(vcr->data, true);
al_array_foreach(vcr->tracks, i, track) {
- aki_packet_cache_flush(&track->cache);
+ nn_packet_cache_flush(&track->cache);
}
}
}
@@ -155,7 +155,7 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer)
static void update_metrics(struct lia_vcr *vcr, u32 size)
{
vcr->metric.current_frame += size;
- u64 now = aki_get_timestamp();
+ u64 now = nn_get_timestamp();
if (!vcr->metric.last_report_ts) {
vcr->metric.last_report_ts = now;
return;
@@ -166,6 +166,7 @@ static void update_metrics(struct lia_vcr *vcr, u32 size)
u64 frame = vcr->metric.current_frame;
vcr->metric.current_frame = 0Lu;
if (diff > 2500000Lu) {
+ // We are buffering fast enough for it to not matter.
al_log_debug("vcr", "Ignoring %lu bytes in metrics.", frame);
return;
}
@@ -174,25 +175,25 @@ static void update_metrics(struct lia_vcr *vcr, u32 size)
}
}
-void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet)
+void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet)
{
struct lia_vcr_track *track;
- u8 op = aki_packet_read_u8(packet);
+ u8 op = nn_packet_read_u8(packet);
switch (op) {
case LIANA_PACKET_DATA:
- track = get_track_from_index(vcr, aki_packet_read_s32(packet));
+ track = get_track_from_index(vcr, nn_packet_read_s32(packet));
if (!track) {
al_log_warn("liana", "Received data from errored or unknown track.");
- aki_packet_free(packet);
+ nn_packet_free(packet);
return;
}
if (!track->running) {
- aki_thread_create(&track->thread, vcr_track_thread, track);
+ nn_thread_create(&track->thread, vcr_track_thread, track);
track->running = true;
}
- u32 size = aki_packet_get_size(packet);
- if (!aki_packet_cache_send_packet(&track->cache, packet)) {
- aki_packet_free(packet);
+ u32 size = nn_packet_get_size(packet);
+ if (!nn_packet_cache_send_packet(&track->cache, packet)) {
+ nn_packet_free(packet);
return;
}
u64 buffer;
@@ -203,14 +204,14 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet)
break;
case LIANA_PACKET_EOF:
al_array_foreach(vcr->tracks, i, track) {
- aki_packet_cache_send_packet(&track->cache, NULL);
+ nn_packet_cache_send_packet(&track->cache, NULL);
}
- aki_packet_free(packet);
- aki_signal_stop(&vcr->signal);
+ nn_packet_free(packet);
+ nn_signal_stop(&vcr->signal);
break;
case LIANA_PACKET_ERROR:
al_log_warn("liana", "Unhandled error packet.");
- aki_packet_free(packet);
+ nn_packet_free(packet);
break;
default:
al_assert(false);
@@ -230,29 +231,29 @@ void lia_vcr_uncork(struct lia_vcr_track *track)
// sense is up for consideration.
}
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);
+ nn_mutex_lock(&track->mutex);
+ if (nn_cond_is_waiting(&track->cond)) {
+ nn_cond_signal(&track->cond);
}
- aki_mutex_unlock(&track->mutex);
+ nn_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);
+ nn_packet_cache_disable(&track->cache);
+ nn_mutex_lock(&track->mutex);
+ if (nn_cond_is_waiting(&track->cond)) {
+ nn_cond_signal(&track->cond);
}
- aki_mutex_unlock(&track->mutex);
+ nn_mutex_unlock(&track->mutex);
if (track->running) {
- aki_thread_join(&track->thread);
+ nn_thread_join(&track->thread);
track->running = false;
}
- struct aki_packet *packet;
- while ((packet = aki_packet_cache_pop(&track->cache))) {
- aki_packet_free(packet);
+ struct nn_packet *packet;
+ while ((packet = nn_packet_cache_pop(&track->cache))) {
+ nn_packet_free(packet);
}
}
@@ -262,11 +263,11 @@ void lia_vcr_flush(struct lia_vcr *vcr)
al_array_foreach(vcr->tracks, i, track) {
vcr_track_close_internal(track);
track->client->flush(track->client);
- aki_packet_cache_enable(&track->cache);
+ nn_packet_cache_enable(&track->cache);
al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED);
al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
}
- aki_signal_stop(&vcr->signal);
+ nn_signal_stop(&vcr->signal);
al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED);
reset_metrics(vcr);
if (vcr->expand == VCR_EXPAND_COMPLETE) {
@@ -281,14 +282,14 @@ void lia_vcr_close_all(struct lia_vcr *vcr)
al_array_foreach(vcr->tracks, i, track) {
vcr_track_close_internal(track);
}
- aki_signal_stop(&vcr->signal);
+ nn_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);
+ nn_packet_cache_free(&track->cache);
track->client->free(&track->client);
#ifdef CAMU_HAVE_FFMPEG
if (track->stream.mode == CAMU_FFMPEG_COMPAT) {
diff --git a/src/liana/vcr.h b/src/liana/vcr.h
index 8041d96..ce47104 100644
--- a/src/liana/vcr.h
+++ b/src/liana/vcr.h
@@ -1,9 +1,9 @@
#pragma once
#include <al/atomic.h>
-#include <aki/packet_cache.h>
-#include <aki/packet_stream.h>
-#include <aki/signal.h>
+#include <nnwt/packet_cache.h>
+#include <nnwt/packet_stream.h>
+#include <nnwt/signal.h>
#include "../codec/codec.h"
@@ -19,11 +19,11 @@ struct lia_vcr_track {
struct lia_client_handler *client;
atomic(s32) state;
atomic(u8) buffered;
- struct aki_packet_cache cache;
- struct aki_cond cond;
- struct aki_mutex mutex;
+ struct nn_packet_cache cache;
+ struct nn_cond cond;
+ struct nn_mutex mutex;
bool running;
- struct aki_thread thread;
+ struct nn_thread thread;
struct lia_vcr *vcr;
};
@@ -32,19 +32,19 @@ struct lia_vcr {
atomic(u64) count;
struct { u64 buffered, low; } mark;
u8 expand;
- struct aki_packet_stream *data;
- struct aki_signal signal;
+ struct nn_packet_stream *data;
+ struct nn_signal signal;
struct {
u64 current_frame;
u64 last_report_ts;
} metric;
};
-void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_packet_stream *data);
+void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data);
void lia_vcr_start(struct lia_vcr *vcr);
void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track);
bool lia_vcr_is_empty(struct lia_vcr *vcr);
-void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet);
+void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_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);