diff options
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 29 | ||||
| -rw-r--r-- | src/liana/handlers/codec_client.c | 44 | ||||
| -rw-r--r-- | src/liana/handlers/codec_server.c | 18 | ||||
| -rw-r--r-- | src/liana/list.h | 2 | ||||
| -rw-r--r-- | src/liana/server.c | 50 | ||||
| -rw-r--r-- | src/liana/server.h | 2 | ||||
| -rw-r--r-- | src/liana/vcr.c | 258 | ||||
| -rw-r--r-- | src/liana/vcr.h | 10 |
8 files changed, 286 insertions, 127 deletions
diff --git a/src/liana/client.c b/src/liana/client.c index bb39f3b..f1d2a61 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -1,4 +1,5 @@ #include <al/log.h> +#include <nnwt/multiplex.h> #include "../server/common.h" #ifdef CAMU_HAVE_FFMPEG @@ -12,7 +13,7 @@ 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; + al_assert(client->vcr.data == stream); lia_vcr_push_packet(&client->vcr, packet); } @@ -125,7 +126,7 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream struct lia_client *client = (struct lia_client *)userdata; client->connection_id = nn_packet_read_u32(packet); parse_info_packet(client, packet); - nn_packet_free(packet); + nn_packet_stream_return_packet(stream, packet); if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) { client->reconnect = false; nn_packet_stream_disconnect(&client->data); @@ -135,6 +136,7 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream struct nn_packet *rpacket = nn_packet_create(); nn_packet_write_s32(rpacket, client->mask); nn_packet_stream_send_packet(stream, rpacket); + lia_vcr_start(&client->vcr); } static void packet_sent_callback(void *userdata, struct nn_packet *packet) @@ -162,7 +164,6 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) } client->reconnect = false; } - lia_vcr_start(&client->vcr); stream->packet_sent_callback = packet_sent_callback; struct nn_packet *packet = nn_packet_create(); nn_packet_write_u32(packet, client->id); @@ -173,6 +174,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) stream->packet_callback = info_packet_callback; } else { stream->packet_callback = data_packet_callback; + lia_vcr_start(&client->vcr); } nn_packet_stream_send_packet(stream, packet); return true; @@ -191,7 +193,11 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * // any unexpected behavior. client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect); if (client->reconnect) { +#ifdef CAMU_DIRECT_MODE + nn_multiplex_direct_reconnect(stream); +#else nn_packet_stream_reconnect(stream, &client->addr, client->port); +#endif } else { client->callback(client->userdata, LIANA_CLIENT_CLOSED, NULL, NULL); } @@ -208,17 +214,14 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u lia_vcr_init(&client->vcr, client->loop, &client->data); al_str_clone(&client->addr, addr); client->port = port; - if (!nn_packet_stream_init(&client->data, type, CAMU_MULTIPLEX_LIANA, - connection_callback, connection_closed_callback, client)) { - connection_closed_callback(client, &client->data); - } - client->renderer = renderer; - 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) -{ + nn_packet_stream_init(&client->data, connection_callback, connection_closed_callback, client); client->renderer = renderer; +#ifdef CAMU_DIRECT_MODE + (void)type; + nn_multiplex_direct_connect(&client->data, CAMU_MULTIPLEX_LIANA); +#else + nn_packet_stream_connect(&client->data, client->loop, CAMU_MULTIPLEX_LIANA, type, &client->addr, client->port); +#endif } void lia_client_seek(struct lia_client *client, u64 pos, u64 at) diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index 3722197..2492533 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -1,8 +1,11 @@ #include "../../codec/codec.h" +#include "../../server/common.h" + #ifdef CAMU_HAVE_FFMPEG #include "../../codec/ffmpeg/decoder.h" #include "../../codec/ffmpeg/packet_ext.h" #endif + #include "../../codec/stb_image/decoder.h" #include "../../codec/spng/decoder.h" #include "../../codec/wuffs/decoder.h" @@ -37,9 +40,7 @@ static bool codec_client_init(struct lia_client_handler *handler, struct camu_re #ifdef CAMU_HAVE_FFMPEG static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt) { - struct camu_codec_packet packet; - packet.av.pkt = pkt; - s32 ret = codec->dec->push(codec->dec, &packet); + s32 ret = codec->dec->push_av_packet(codec->dec, pkt); return ret == CAMU_OK; } #endif @@ -60,7 +61,14 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc if (!packet) { if (codec->dec) { // Flush always returns success. - codec->dec->push(codec->dec, NULL); + switch (codec->dec->mode) { + case CAMU_NORMAL: + codec->dec->push(codec->dec, NULL); + break; + case CAMU_FFMPEG_COMPAT: + codec->dec->push_av_packet(codec->dec, NULL); + break; + } // process() could still error. s32 ret = codec->dec->process(codec->dec); codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); @@ -71,14 +79,19 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc } // Push packet. + u32 rindex = packet->rindex; bool success; - u8 type = nn_packet_read_u8(packet); - switch (type) { + u8 mode = nn_packet_read_u8(packet); + switch (mode) { case CAMU_NORMAL: { if (codec->dec) { +#ifdef CAMU_DIRECT_MODE + success = push_packet(codec, (struct nn_buffer *)packet->opaque); +#else struct nn_buffer buffer; nn_packet_read_buffer(packet, &buffer); success = push_packet(codec, &buffer); +#endif } else { success = true; } @@ -86,15 +99,27 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc } #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: { - AVPacket *pkt = nn_packet_read_av_packet(packet); + AVPacket *pkt; + if (packet->opaque) { + pkt = (AVPacket *)packet->opaque; + } else { + pkt = av_packet_alloc(); + nn_packet_read_av_packet(packet, pkt); + packet->opaque = pkt; + } + struct camu_codec_stream *stream = codec->handler.stream; if (codec->dec) { success = push_av_packet(codec, pkt); } else { - codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, codec->handler.stream, pkt); + codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, stream, pkt); success = true; } +#ifdef VCR_BUFFER_WHOLE_FILE + pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, stream->av.stream->time_base); +#else av_packet_unref(pkt); av_packet_free(&pkt); +#endif break; } default: @@ -102,6 +127,9 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc #endif } + // Restore packet rindex in case we resue it. + packet->rindex = rindex; + if (!success) { // Forcing in EOF on an error is not necessary but should exhibit less erratic behavior sink-side. codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c index 7167f97..24552bb 100644 --- a/src/liana/handlers/codec_server.c +++ b/src/liana/handlers/codec_server.c @@ -1,8 +1,11 @@ #include "../../codec/codec.h" +#include "../../server/common.h" + #ifdef CAMU_HAVE_FFMPEG #include "../../codec/ffmpeg/demuxer.h" #include "../../codec/ffmpeg/packet_ext.h" #endif + #include "../../codec/stb_image/demuxer.h" #include "../../codec/spng/demuxer.h" #include "../../codec/wuffs/demuxer.h" @@ -26,7 +29,9 @@ static bool codec_server_init(struct lia_server_handler *handler, struct cch_han break; #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: +#ifndef CAMU_DIRECT_MODE codec->packet.av.pkt = av_packet_alloc(); +#endif break; #endif } @@ -84,6 +89,9 @@ static bool codec_server_seek(struct lia_server_handler *handler, u64 pos) static void codec_server_step(struct lia_server_handler *handler) { struct lia_codec_server *codec = (struct lia_codec_server *)handler; +#ifdef CAMU_DIRECT_MODE + codec->packet.av.pkt = av_packet_alloc(); +#endif codec->handler.status = codec->demux->get_packet(codec->demux, &codec->packet); } @@ -96,7 +104,11 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct case CAMU_NORMAL: { nn_packet_write_s32(packet, 0); nn_packet_write_u8(packet, codec->packet.mode); +#ifdef CAMU_DIRECT_MODE + packet->opaque = codec->packet.buffer; +#else nn_packet_write_buffer(packet, codec->packet.buffer); +#endif break; } #ifdef CAMU_HAVE_FFMPEG @@ -104,8 +116,12 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct AVPacket *pkt = codec->packet.av.pkt; nn_packet_write_s32(packet, pkt->stream_index); nn_packet_write_u8(packet, codec->packet.mode); +#ifdef CAMU_DIRECT_MODE + packet->opaque = pkt; +#else nn_packet_write_av_packet(packet, pkt); av_packet_unref(pkt); +#endif break; } #endif @@ -125,7 +141,9 @@ static void codec_server_free(struct lia_server_handler **handler) break; #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: +#ifndef CAMU_DIRECT_MODE av_packet_free(&codec->packet.av.pkt); +#endif break; #endif } diff --git a/src/liana/list.h b/src/liana/list.h index bb4ac17..3750d09 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -14,7 +14,7 @@ #define LIANA_BUFFER_AHEAD 2 -#define LIANA_LIST_SCUFFED_LOOP +//#define LIANA_LIST_SCUFFED_LOOP enum { LIANA_SINK_SET = 0, diff --git a/src/liana/server.c b/src/liana/server.c index b7b8963..3b7bc2d 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -4,6 +4,7 @@ #include "handler.h" #include "handlers.h" #include "list.h" +#include "vcr.h" static inline u32 get_incremental_id(struct lia_server *server) { @@ -28,16 +29,29 @@ static void remove_zombie(struct lia_server *server, struct nn_packet_stream *st static void data_packet_sent_callback(void *userdata, struct nn_packet *packet) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + nn_packet_pool_lock(&conn->pool); nn_packet_pool_return(&conn->pool, packet); + nn_packet_pool_unlock(&conn->pool); } -static u8 packet_pool_callback(void *userdata, struct nn_packet *packet) +static void data_packets_sent_callback(void *userdata, struct nn_packet **packets, u32 count) +{ + struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + nn_packet_pool_lock(&conn->pool); + for (u32 i = 0; i < count; i++) { + if (packets[i]) { + nn_packet_pool_return(&conn->pool, packets[i]); + } + } + nn_packet_pool_unlock(&conn->pool); +} + +static void packet_pool_callback(void *userdata, struct nn_packet *packet) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; if (!nn_packet_stream_send_packet(conn->stream, packet)) { - return NNWT_PACKET_POOL_RETURN; + data_packet_sent_callback(conn, packet); } - return NNWT_PACKET_POOL_KEEP; } static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) @@ -58,8 +72,7 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { (void)userdata; - (void)stream; - nn_packet_free(packet); + nn_packet_stream_return_packet(stream, packet); // We should never be here. Although, we also shouldn't assert because // any erroneous connection can bring us here. al_assert(false); @@ -91,14 +104,14 @@ static void data_connection_closed_callback(void *userdata, struct nn_packet_str // This is joining handler_thread(), we will never be here if init_thread() blocks or fails. nn_thread_join(&conn->thread); + free_connection_stream(conn); + if (conn->disconnected) { free_connection(conn); } else { cch_handle_enable(&conn->handle); nn_packet_pool_enable(&conn->pool); } - - free_connection_stream(conn); } static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) @@ -108,9 +121,10 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s conn->handler->subscribe(conn->handler, mask); stream->packet_callback = discard_packet_callback; stream->packet_sent_callback = data_packet_sent_callback; + stream->packets_sent_callback = data_packets_sent_callback; stream->connection_closed_callback = data_connection_closed_callback; nn_thread_create(&conn->thread, handler_thread, conn); - nn_packet_free(packet); + nn_packet_stream_return_packet(stream, packet); } static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet) @@ -123,7 +137,6 @@ static void subscribe_connection_closed_callback(void *userdata, struct nn_packe { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; - free_connection(conn); free_connection_stream(conn); } @@ -152,8 +165,11 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet conn->handler->subscribe(conn->handler, mask); stream->packet_callback = discard_packet_callback; stream->packet_sent_callback = data_packet_sent_callback; + stream->packets_sent_callback = data_packets_sent_callback; stream->connection_closed_callback = data_connection_closed_callback; +#ifndef VCR_BUFFER_WHOLE_FILE nn_thread_create(&conn->thread, handler_thread, conn); +#endif } } @@ -183,13 +199,13 @@ static void signal_callback(void *userdata) // Connection was closed before init was done. return; } + struct nn_packet_stream *stream = conn->stream; if (!conn->errored) { conn->id = get_incremental_id(server); - nn_packet_pool_init(&conn->pool, 96, server->loop, packet_pool_callback, conn); + nn_packet_pool_init(&conn->pool, 512, server->loop, packet_pool_callback, conn); al_array_push(node->connections, conn); handle_connection(conn, packet); } else { - struct nn_packet_stream *stream = conn->stream; al_free(conn); conn = NULL; // This connection is now nothing but a packet stream. @@ -197,7 +213,7 @@ static void signal_callback(void *userdata) stream->connection_closed_callback = connection_closed_callback; nn_packet_stream_disconnect(stream); } - nn_packet_free(packet); + nn_packet_stream_return_packet(stream, packet); } static nn_thread_result NNWT_THREADCALL init_thread(void *userdata) @@ -231,10 +247,10 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id) 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; - nn_packet_free(conn->packet); + struct nn_packet *packet = conn->packet; // Checked in signal_callback and will signal to cleanup the connection. conn->packet = NULL; + nn_packet_stream_return_packet(stream, packet); } static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) @@ -273,7 +289,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str } else { nn_packet_stream_disconnect(stream); } - nn_packet_free(packet); + nn_packet_stream_return_packet(stream, packet); } } @@ -292,13 +308,11 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) return true; } -void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock) +void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *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; - nn_packet_stream_from_socket(stream, server->loop, sock); } struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry) diff --git a/src/liana/server.h b/src/liana/server.h index 7c155d0..849fd38 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -51,7 +51,7 @@ struct lia_server { }; 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); +void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *stream); 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 4131c80..c12c422 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -3,7 +3,7 @@ #include "vcr.h" #include "handler.h" -#define VCR_BUFFER_BUFFERED MB(8) +#define VCR_BUFFER_BUFFERED MB(6) enum { VCR_EXPAND_UNTOUCHED = 0, @@ -11,16 +11,22 @@ enum { VCR_EXPAND_COMPLETE }; -// Only track the buffered state of audio and video streams as we don't -// expect any other type of stream to ever call cork(). -#define TRACK_IGNORE_BUFFERED(track) \ - (!(track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO)) +enum { + VCR_TRACK_RUNNING = 0, + VCR_TRACK_STOPPED, + VCR_TRACK_CLOSED +}; + +#define VCR_TRACK_THREADED(track) \ + (track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO) +#ifndef CAMU_DIRECT_MODE static void signal_callback(void *userdata) { struct lia_vcr *vcr = (struct lia_vcr *)userdata; nn_packet_stream_cork(vcr->data, false); } +#endif static void reset_metrics(struct lia_vcr *vcr) { @@ -36,82 +42,146 @@ void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_pac vcr->mark.low = 0; vcr->expand = VCR_EXPAND_UNTOUCHED; vcr->data = data; +#ifndef CAMU_DIRECT_MODE nn_signal_init(&vcr->signal, loop, signal_callback, vcr); +#else + (void)loop; +#endif reset_metrics(vcr); } -void lia_vcr_start(struct lia_vcr *vcr) +static void return_entire_cache(struct lia_vcr_track *track) { - nn_signal_start(&vcr->signal); + struct lia_vcr *vcr = track->vcr; + u32 size = track->cache.cache.size; + nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), size); + track->cache.cache.size = 0; } static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) { + nn_thread_set_priority(NNWT_THREAD_SCHED_FIFO, 32); + struct lia_vcr_track *track = (struct lia_vcr_track *)userdata; struct lia_vcr *vcr = track->vcr; - u32 packets; + + s32 state; + bool corked; + u32 packets, index = 0; struct nn_packet *packet = NULL; while (nn_packet_cache_wait(&track->cache, &packets)) { - s32 state = 0; - for (u32 i = 0; i < packets; i++) { - packet = nn_packet_cache_pop(&track->cache); - state = 0; + al_assert(packets >= index); + corked = false; + for (; index < packets; index++) { + packet = nn_packet_cache_at(&track->cache, index); + +#ifdef VCR_BUFFER_WHOLE_FILE + if (!packet) { // Loop. + index = 0; + corked = false; + break; + } +#endif + if (packet) { - state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED); - if (state == LIANA_STREAM_CLOSED) { - // We were signaled to close, exit thread. - goto out; - } + // Check if we should uncork the packet stream. 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); +#ifndef CAMU_DIRECT_MODE + u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED); if (buffered && buffer <= vcr->mark.low) { nn_signal_send(&vcr->signal); } +#else + (void)buffer; +#endif } - if (!track->client->handle_packet(track->client, packet)) { - // Handler error, exit. - goto out; + + nn_mutex_lock(&track->mutex); + + // NULL packet means flush. + bool success = track->client->handle_packet(track->client, packet); + if (!success) { + al_log_error("liana", "Error handling packet, exiting track thread."); + nn_packet_cache_unlock(&track->cache); + nn_mutex_unlock(&track->mutex); + return 0; } - if (packet) { - nn_packet_free(packet); - packet = NULL; - if (state == LIANA_STREAM_STOPPED) { - // Unlock here to accumulate packets while waiting. - 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 { - // NULL packet, exit thread. - goto out; + + // Take state after handle_packet() because we may have been + // corked from within it. + state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED); + + // We wait if corked (TRACK_STOPPED) or EOF. + corked = (packet && state == VCR_TRACK_STOPPED) || !packet; + + // Don't wait, continue processing packets from the current set. + if (!corked) { + nn_mutex_unlock(&track->mutex); + continue; } + + nn_packet_cache_unlock(&track->cache); + + // Wait for uncork. + nn_cond_wait(&track->cond, &track->mutex); + + // Check for possibly updated state. + state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED); + + nn_mutex_unlock(&track->mutex); + + if (state == VCR_TRACK_CLOSED) { + // We already unlocked the cache. + return 0; + } + + al_assert(state == VCR_TRACK_RUNNING); + + // Wait on the packet cache again. + break; } - if (state != LIANA_STREAM_STOPPED) { - // If state = STOPPED, we already unlocked. + + if (corked) { + // The cache is already unlocked here. + index++; + } else { +#ifndef VCR_BUFFER_WHOLE_FILE + al_assert(index == packets); + index = 0; + nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), packets); + al_array_remove_range(track->cache.cache, 0, packets); +#endif 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) nn_packet_free(packet); - nn_packet_cache_unlock(&track->cache); + return 0; } +void lia_vcr_start(struct lia_vcr *vcr) +{ +#ifndef CAMU_DIRECT_MODE + nn_signal_start(&vcr->signal); +#endif + struct lia_vcr_track *track; + al_array_foreach(vcr->tracks, i, track) { + if (VCR_TRACK_THREADED(track) && !track->running) { + nn_thread_create(&track->thread, vcr_track_thread, track); + track->running = true; + } + } +} + void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track) { track->vcr = vcr; nn_cond_init(&track->cond); nn_mutex_init(&track->mutex); - al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED); + al_atomic_store(u8)(&track->buffered, !VCR_TRACK_THREADED(track), AL_ATOMIC_RELAXED); 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); + al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED); } bool lia_vcr_is_empty(struct lia_vcr *vcr) @@ -128,6 +198,7 @@ static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index return NULL; } +#ifndef CAMU_DIRECT_MODE static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) { u8 buffered = 1; @@ -151,6 +222,7 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) } } } +#endif static void update_metrics(struct lia_vcr *vcr, u32 size) { @@ -180,69 +252,85 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) struct lia_vcr_track *track; u8 op = nn_packet_read_u8(packet); switch (op) { - case LIANA_PACKET_DATA: + case LIANA_PACKET_DATA: { + u32 size = nn_packet_get_size(packet); + update_metrics(vcr, size); track = get_track_from_index(vcr, nn_packet_read_s32(packet)); if (!track) { al_log_warn("liana", "Received data from errored or unknown track."); - nn_packet_free(packet); - return; - } - if (!track->running) { - nn_thread_create(&track->thread, vcr_track_thread, track); - track->running = true; + break; } - u32 size = nn_packet_get_size(packet); - if (!nn_packet_cache_send_packet(&track->cache, packet)) { - nn_packet_free(packet); + if (VCR_TRACK_THREADED(track)) { + if (!nn_packet_cache_send_packet(&track->cache, packet)) { + break; + } + u64 buffer; + if ((buffer = al_atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED)) >= vcr->mark.buffered) { +#ifndef CAMU_DIRECT_MODE + cork_if_buffered(vcr, buffer); +#endif + } + // Keep packet. return; + } else { + if (!track->client->handle_packet(track->client, packet)) { + al_log_warn("liana", "Error handling non-buffered packet."); + } } - u64 buffer; - if ((buffer = al_atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED)) >= vcr->mark.buffered) { - cork_if_buffered(vcr, buffer); - } - update_metrics(vcr, size); break; + } case LIANA_PACKET_EOF: al_array_foreach(vcr->tracks, i, track) { nn_packet_cache_send_packet(&track->cache, NULL); } - nn_packet_free(packet); +#ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); +#endif break; case LIANA_PACKET_ERROR: al_log_warn("liana", "Unhandled error packet."); - nn_packet_free(packet); break; default: - al_assert(false); + al_log_warn("liana", "Erroneous packet."); + break; } + nn_packet_stream_return_packet(vcr->data, packet); +} + +void lia_vcr_set_buffered(struct lia_vcr_track *track) +{ + al_atomic_store(u8)(&track->buffered, 1, AL_ATOMIC_RELAXED); } void lia_vcr_cork(struct lia_vcr_track *track) { - al_atomic_store(s32)(&track->state, LIANA_STREAM_STOPPED, AL_ATOMIC_RELAXED); + al_atomic_store(s32)(&track->state, VCR_TRACK_STOPPED, AL_ATOMIC_RELAXED); al_atomic_store(u8)(&track->buffered, 1, AL_ATOMIC_RELAXED); } void lia_vcr_uncork(struct lia_vcr_track *track) { - if (al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED) != LIANA_STREAM_STOPPED) { - // This can be reached during normal operation. Whether or not that makes - // sense is up for consideration. + if (al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED) != VCR_TRACK_STOPPED) { + // We will get here during normal operation. Early returning is historically + // tricky in vcr_uncork(). If I'm understanding correctly, asserting that cond_is_waiting() + // just below means we are safe. + return; } - al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); + // Lock before setting track->state to avoid a race with cork(). nn_mutex_lock(&track->mutex); - if (nn_cond_is_waiting(&track->cond)) { - nn_cond_signal(&track->cond); - } + al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED); + al_assert(nn_cond_is_waiting(&track->cond)); + nn_cond_signal(&track->cond); 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); + // Calling packet_cache_disable() while holding the track mutex can + // very possibly deadlock. nn_packet_cache_disable(&track->cache); nn_mutex_lock(&track->mutex); + al_atomic_store(s32)(&track->state, VCR_TRACK_CLOSED, AL_ATOMIC_RELAXED); if (nn_cond_is_waiting(&track->cond)) { nn_cond_signal(&track->cond); } @@ -251,23 +339,28 @@ static void vcr_track_close_internal(struct lia_vcr_track *track) nn_thread_join(&track->thread); track->running = false; } - struct nn_packet *packet; - while ((packet = nn_packet_cache_pop(&track->cache))) { - nn_packet_free(packet); - } } void lia_vcr_flush(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - vcr_track_close_internal(track); - track->client->flush(track->client); - 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); + if (VCR_TRACK_THREADED(track)) { + vcr_track_close_internal(track); +#ifndef VCR_BUFFER_WHOLE_FILE + return_entire_cache(track); +#endif + al_atomic_store(u8)(&track->buffered, 0, AL_ATOMIC_RELAXED); + track->client->flush(track->client); + nn_packet_cache_enable(&track->cache); + al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED); + } else { + track->client->flush(track->client); + } } +#ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); +#endif al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); reset_metrics(vcr); if (vcr->expand == VCR_EXPAND_COMPLETE) { @@ -281,8 +374,11 @@ void lia_vcr_close_all(struct lia_vcr *vcr) struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { vcr_track_close_internal(track); + return_entire_cache(track); } +#ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); +#endif } void lia_vcr_free(struct lia_vcr *vcr) diff --git a/src/liana/vcr.h b/src/liana/vcr.h index ce47104..8aefaa2 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -6,12 +6,9 @@ #include <nnwt/signal.h> #include "../codec/codec.h" +#include "../server/common.h" -enum { - LIANA_STREAM_RUNNING = 0, - LIANA_STREAM_STOPPED, - LIANA_STREAM_CLOSED -}; +//#define VCR_BUFFER_WHOLE_FILE struct lia_vcr_track { s32 index; @@ -33,7 +30,9 @@ struct lia_vcr { struct { u64 buffered, low; } mark; u8 expand; struct nn_packet_stream *data; +#ifndef CAMU_DIRECT_MODE struct nn_signal signal; +#endif struct { u64 current_frame; u64 last_report_ts; @@ -45,6 +44,7 @@ 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 nn_packet *packet); +void lia_vcr_set_buffered(struct lia_vcr_track *track); 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); |