diff options
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 85 | ||||
| -rw-r--r-- | src/liana/client.h | 8 | ||||
| -rw-r--r-- | src/liana/handler.h | 8 | ||||
| -rw-r--r-- | src/liana/handlers.c | 2 | ||||
| -rw-r--r-- | src/liana/handlers/cdio.h | 2 | ||||
| -rw-r--r-- | src/liana/handlers/cdio_client.c | 12 | ||||
| -rw-r--r-- | src/liana/handlers/cdio_server.c | 42 | ||||
| -rw-r--r-- | src/liana/handlers/codec_client.c | 12 | ||||
| -rw-r--r-- | src/liana/handlers/codec_server.c | 44 | ||||
| -rw-r--r-- | src/liana/list.c | 14 | ||||
| -rw-r--r-- | src/liana/meson.build | 1 | ||||
| -rw-r--r-- | src/liana/server.c | 157 | ||||
| -rw-r--r-- | src/liana/server.h | 32 | ||||
| -rw-r--r-- | src/liana/vcr.c | 107 | ||||
| -rw-r--r-- | src/liana/vcr.h | 22 |
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); |