diff options
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 64 | ||||
| -rw-r--r-- | src/liana/client.h | 5 | ||||
| -rw-r--r-- | src/liana/handlers/cdio_server.c | 2 | ||||
| -rw-r--r-- | src/liana/handlers/codec_server.c | 2 | ||||
| -rw-r--r-- | src/liana/handlers/dvd_server.c | 2 | ||||
| -rw-r--r-- | src/liana/handlers/handler.h | 2 | ||||
| -rw-r--r-- | src/liana/server.c | 24 | ||||
| -rw-r--r-- | src/liana/server.h | 1 | ||||
| -rw-r--r-- | src/liana/vcr.c | 80 | ||||
| -rw-r--r-- | src/liana/vcr.h | 4 |
10 files changed, 123 insertions, 63 deletions
diff --git a/src/liana/client.c b/src/liana/client.c index 4d5689f..2b62612 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -13,8 +13,10 @@ enum { RECONNECT_NONE = 0, - RECONNECT_ON_CONNECTION_CLOSED, + RECONNECT_RECOVER, + RECONNECT_SEEK, RECONNECT_SIGNAL_CLIENT, + RECONNECT_DISREGUARD, RECONNECT_DISCONNECTED }; @@ -48,7 +50,7 @@ static void collect_streams(struct lia_client *client, struct nn_packet *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); + al_assert(index < 64); switch (mode) { case CAMU_NORMAL: { if (type == CAMU_STREAM_AUDIO) { @@ -147,10 +149,6 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe struct camu_codec_stream *stream; al_array_foreach_ptr(client->streams, i, stream) { u8 type = stream->type; - if (type == CAMU_STREAM_SUBTITLE && client->mask == 0) { - log_warn("Ignoring subtitle-only resource."); - goto out; - } if ((selected & (1 << type)) || !(prefs->enabled & (1 << type))) { continue; } @@ -198,6 +196,10 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe lia_vcr_add_track(&client->vcr, track); } } + if (client->mask == (1 << CAMU_STREAM_SUBTITLE)) { + log_warn("Ignoring subtitle-only resource."); + client->mask = 0; + } out: client->callback(client->userdata, LIANA_CLIENT_CONFIGURE_COMPLETE, NULL, NULL); } @@ -211,13 +213,15 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) { log_warn("Discarding resource with no applicable streams."); al_assert(client->reconnect == RECONNECT_NONE); + client->reconnect = RECONNECT_DISREGUARD; nn_packet_stream_disconnect(&client->data); return; } stream->packet_callback = data_packet_callback; struct nn_packet *rpacket = nn_packet_create(); - nn_packet_write_u32(rpacket, client->mask); + nn_packet_write_u64(rpacket, client->mask); nn_packet_stream_send_packet(stream, rpacket); + client->reconnect = RECONNECT_RECOVER; lia_vcr_start(&client->vcr); } @@ -239,6 +243,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) .pos = client->pos, .pause = LIANA_PAUSE_NONE }; + client->at = LIANA_TIMESTAMP_INVALID; client->callback(client->userdata, LIANA_CLIENT_RESUME_AT, NULL, &time); // The value of client->mask will not have changed since connection_closed_callback(). if (client->rec.unconfigured) { @@ -247,18 +252,19 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) } client->callback(client->userdata, LIANA_CLIENT_RECONNECTED, NULL, &client->rec); } else { - al_assert(client->connection_id == 0); + al_assert(client->connection_id == 0 && client->reconnect == RECONNECT_NONE); } stream->packet_sent_callback = packet_sent_callback; struct nn_packet *packet = nn_packet_create(); nn_packet_write_u32(packet, client->node_id); nn_packet_write_u32(packet, client->connection_id); - nn_packet_write_u32(packet, client->mask); + nn_packet_write_u64(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; + client->reconnect = RECONNECT_RECOVER; lia_vcr_start(&client->vcr); } nn_packet_stream_send_packet(stream, packet); @@ -268,15 +274,26 @@ static bool connection_callback(void *userdata, struct nn_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 == RECONNECT_ON_CONNECTION_CLOSED) { + + bool reconnect = client->reconnect == RECONNECT_RECOVER || client->reconnect == RECONNECT_SEEK; + if (client->reconnect == RECONNECT_RECOVER) { + al_assert(client->at == LIANA_TIMESTAMP_INVALID); + struct lia_timing time; + client->callback(client->userdata, LIANA_CLIENT_RECOVER_TO, NULL, &time); + client->pos = time.pos; + client->at = time.at; + } + + if (reconnect) { lia_vcr_flush(&client->vcr); } else { lia_vcr_close_all(&client->vcr); } + // If reconnect = SIGNAL_CLIENT, we either never connected or recursed at the reconnect step. if (client->reconnect != RECONNECT_SIGNAL_CLIENT) { client->rec = (struct lia_reconnect_info){ - .reconnect = client->reconnect == RECONNECT_ON_CONNECTION_CLOSED, + .reconnect = reconnect, .unconfigured = client->mask == 0, .mask = client->mask }; @@ -295,10 +312,13 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * } al_array_free(client->rec.detached); } - if (client->reconnect == RECONNECT_ON_CONNECTION_CLOSED) { - // If stream_reconnect() errors, the client will be closed on recursion. + + if (reconnect) { + // If stream_reconnect() errors or is aborted, the client will be closed on recursion. client->reconnect = RECONNECT_SIGNAL_CLIENT; - if (!client->mask) { + // A client being seeked before an info packet is another reason mask may + // be unset here. In that case we don't want to forcefully close. + if (!client->rec.unconfigured && !client->mask) { connection_closed_callback(userdata, stream); } else { #ifdef CAMU_DIRECT_MODE @@ -318,6 +338,7 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, client->loop = loop; client->node_id = node_id; client->pos = pos; + client->at = LIANA_TIMESTAMP_INVALID; client->mask = 0; al_array_init(client->streams); client->reconnect = RECONNECT_NONE; @@ -336,13 +357,16 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, void lia_client_seek(struct lia_client *client, u64 pos, u64 at) { - if (client->reconnect == RECONNECT_DISCONNECTED) return; - // If reconnect = ON_CONNECTION_CLOSED or SIGNAL_CLIENT, we are safe to edit pos + u8 reconnect = client->reconnect; + if (reconnect == RECONNECT_DISCONNECTED || reconnect == RECONNECT_DISREGUARD) { + return; + } + // If reconnect = SEEK, RECOVER or SIGNAL_CLIENT, we are safe to edit pos // and at in-place because they aren't evaluated until connection_callback(). client->pos = pos; client->at = at; - if (client->reconnect == RECONNECT_NONE) { - client->reconnect = RECONNECT_ON_CONNECTION_CLOSED; + if (reconnect == RECONNECT_NONE || reconnect == RECONNECT_RECOVER) { + client->reconnect = RECONNECT_SEEK; nn_packet_stream_disconnect(&client->data); } } @@ -352,7 +376,9 @@ void lia_client_disconnect(struct lia_client *client) u8 reconnect = client->reconnect; al_assert(reconnect != RECONNECT_DISCONNECTED); client->reconnect = RECONNECT_DISCONNECTED; - if (reconnect != RECONNECT_ON_CONNECTION_CLOSED) { + // If reconnect == SIGNAL_CLIENT, stream_disconnect() needs to ensure + // connection_callback() is never called. + if (reconnect != RECONNECT_SEEK && reconnect != RECONNECT_DISREGUARD) { nn_packet_stream_disconnect(&client->data); } } diff --git a/src/liana/client.h b/src/liana/client.h index 681b134..452150d 100644 --- a/src/liana/client.h +++ b/src/liana/client.h @@ -12,6 +12,7 @@ enum { LIANA_CLIENT_DATA, LIANA_CLIENT_SUBTITLE, LIANA_CLIENT_REMOVE_BUFFERS, + LIANA_CLIENT_RECOVER_TO, LIANA_CLIENT_RESUME_AT, LIANA_CLIENT_RECONNECTED, LIANA_CLIENT_EOF, @@ -22,7 +23,7 @@ enum { struct lia_reconnect_info { bool reconnect; bool unconfigured; - u32 mask; + u64 mask; array(struct camu_codec_stream *) detached; }; @@ -39,7 +40,7 @@ struct lia_client { u32 node_id; struct lia_prefs prefs; array(struct camu_codec_stream) streams; - u32 mask; + u64 mask; u64 pos; u64 at; u8 reconnect; diff --git a/src/liana/handlers/cdio_server.c b/src/liana/handlers/cdio_server.c index 54c63ab..c67330d 100644 --- a/src/liana/handlers/cdio_server.c +++ b/src/liana/handlers/cdio_server.c @@ -42,7 +42,7 @@ static void cdio_server_write_info(struct lia_server_handler *handler, struct nn nn_packet_write_s32(packet, cdio->fmt.channel_count); } -static void cdio_server_subscribe(struct lia_server_handler *handler, u32 mask) +static void cdio_server_subscribe(struct lia_server_handler *handler, u64 mask) { (void)handler; (void)mask; diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c index fbb1e72..a7fc587 100644 --- a/src/liana/handlers/codec_server.c +++ b/src/liana/handlers/codec_server.c @@ -79,7 +79,7 @@ static void codec_server_write_info(struct lia_server_handler *handler, struct n } } -static void codec_server_subscribe(struct lia_server_handler *handler, u32 mask) +static void codec_server_subscribe(struct lia_server_handler *handler, u64 mask) { struct lia_codec_server *codec = (struct lia_codec_server *)handler; codec->demux->subscribed = mask; diff --git a/src/liana/handlers/dvd_server.c b/src/liana/handlers/dvd_server.c index 5a9c793..53ef323 100644 --- a/src/liana/handlers/dvd_server.c +++ b/src/liana/handlers/dvd_server.c @@ -53,7 +53,7 @@ static void dvd_server_write_info(struct lia_server_handler *handler, struct nn_ (void)packet; } -static void dvd_server_subscribe(struct lia_server_handler *handler, u32 mask) +static void dvd_server_subscribe(struct lia_server_handler *handler, u64 mask) { (void)handler; (void)mask; diff --git a/src/liana/handlers/handler.h b/src/liana/handlers/handler.h index 48fcef3..ec2f3fe 100644 --- a/src/liana/handlers/handler.h +++ b/src/liana/handlers/handler.h @@ -14,7 +14,7 @@ enum { struct lia_server_handler { bool (*init)(struct lia_server_handler *, struct cch_handle *); void (*write_info)(struct lia_server_handler *, struct nn_packet *); - void (*subscribe)(struct lia_server_handler *, u32); + void (*subscribe)(struct lia_server_handler *, u64); u64 (*get_duration)(struct lia_server_handler *); bool (*seek)(struct lia_server_handler *, u64); void (*step)(struct lia_server_handler *); diff --git a/src/liana/server.c b/src/liana/server.c index 04ef54b..5065425 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -1,3 +1,6 @@ +#define AL_LOG_SECTION "liana" +#include <al/log.h> + #include "server.h" #include "handlers.h" #include "list.h" @@ -166,7 +169,7 @@ static void data_connection_closed_callback(void *userdata, struct nn_packet_str } } -static void start_connection_handler(struct lia_node_connection *conn, u32 mask) +static void start_connection_handler(struct lia_node_connection *conn, u64 mask) { struct nn_packet_stream *stream = conn->stream; conn->handler->subscribe(conn->handler, mask); @@ -180,7 +183,7 @@ static void start_connection_handler(struct lia_node_connection *conn, u32 mask) 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; - u32 mask = nn_packet_read_u32(packet); + u64 mask = nn_packet_read_u64(packet); nn_packet_stream_return_packet(stream, packet); start_connection_handler(conn, mask); } @@ -203,7 +206,7 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet { struct nn_packet_stream *stream = conn->stream; - u32 mask = nn_packet_read_u32(packet); + u64 mask = nn_packet_read_u64(packet); u64 seek_pos = nn_packet_read_u64(packet); al_assert(!conn->ref); @@ -487,7 +490,7 @@ void lia_node_close(struct lia_node *node) al_array_foreach_rev(node->connections, i, conn) { if (conn->stream) { conn->disconnected = true; - nn_packet_stream_disconnect(conn->stream); + //nn_packet_stream_disconnect(conn->stream); } else { free_connection(conn); } @@ -503,6 +506,19 @@ void lia_server_close(struct lia_server *server) } } +void lia_server_force_disconnect_nodes(struct lia_server *server) +{ + struct lia_node *node; + al_array_foreach(server->nodes, i, node) { + struct lia_node_connection *conn; + al_array_foreach_rev(node->connections, j, conn) { + if (conn->stream) { + nn_packet_stream_disconnect(conn->stream); + } + } + } +} + void lia_server_free(struct lia_server *server) { // Assuming we joined on the event loop, server->nodes should be empty. diff --git a/src/liana/server.h b/src/liana/server.h index 8035ee1..bafe95b 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -61,4 +61,5 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en void lia_node_get_duration(struct lia_node *node); void lia_node_close(struct lia_node *node); void lia_server_close(struct lia_server *server); +void lia_server_force_disconnect_nodes(struct lia_server *server); void lia_server_free(struct lia_server *server); diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 4fa548c..f9dad7b 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -1,15 +1,20 @@ #define AL_LOG_SECTION "vcr" //#define AL_LOG_ENABLE_TRACE #include <al/log.h> +#include <nnwt/time.h> #include "handlers/handler.h" #include "vcr.h" #include "common.h" -#define VCR_BUFFER_BUFFERED MB(6LL) -#define VCR_BUFFER_GROW_FACTOR 8LL -#define VCR_BUFFER_LOW_OFFSET KB(500LL) +#define VCR_BUFFER_BUFFERED MB((u64)6) +#ifdef CAMU_HUGE_VIDEO_BUFFER +#define VCR_BUFFER_GROW_FACTOR ((u64)24) +#else +#define VCR_BUFFER_GROW_FACTOR ((u64)8) +#endif +#define VCR_BUFFER_LOW_OFFSET KB((u64)500) AL_STATIC_ASSERT(buf_gt_low_offset, VCR_BUFFER_BUFFERED * VCR_BUFFER_GROW_FACTOR, >, VCR_BUFFER_LOW_OFFSET); enum { @@ -97,6 +102,8 @@ void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_pac vcr->corked = false; nn_signal_init(&vcr->signal, loop, signal_callback, vcr); reset_metrics(vcr); +#else + (void)loop; #endif } @@ -128,7 +135,6 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) al_snprintf((char *)thread_name, sizeof(thread_name), "vcr:%hu_%d", vcr->node_id, track->stream->index); nn_thread_set_name(thread_name); - s32 state; bool corked; u32 packets, index = 0; struct nn_packet *packet = NULL; @@ -161,7 +167,7 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) #endif } - nn_mutex_lock(&track->mutex); + nn_mutex_lock(&track->lock); // NULL packet means flush. bool success = track->client->handle_packet(track->client, packet); @@ -171,19 +177,21 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) return_entire_cache(track); track->cache.disabled = true; nn_packet_cache_unlock(&track->cache); - nn_mutex_unlock(&track->mutex); + nn_mutex_unlock(&track->lock); return 0; } - // Take state again because handle_packet() could have caused the track to be corked. - state = atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED); + if (!packet) { + // Wait on EOF. + track->state = VCR_TRACK_STOPPED; + } - // We wait if corked (TRACK_STOPPED) or EOF. - corked = (packet && state == VCR_TRACK_STOPPED) || !packet; + // Track could have been corked from within handle_packet(). + corked = track->state == VCR_TRACK_STOPPED; // Don't wait, continue processing packets from the current set. if (!corked) { - nn_mutex_unlock(&track->mutex); + nn_mutex_unlock(&track->lock); continue; } @@ -196,12 +204,12 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) nn_packet_cache_unlock(&track->cache); // Wait for uncork. - nn_cond_wait(&track->cond, &track->mutex); + nn_cond_wait(&track->cond, &track->lock); // Check for possibly updated state. - state = atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED); + u32 state = track->state; - nn_mutex_unlock(&track->mutex); + nn_mutex_unlock(&track->lock); if (state == VCR_TRACK_CLOSED) { // We already unlocked the cache. @@ -253,12 +261,12 @@ 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); + nn_mutex_init(&track->lock); atomic_store(bool)(&track->buffered, !VCR_TRACK_THREADED(track), AL_ATOMIC_RELAXED); track->running = false; nn_packet_cache_init(&track->cache, 256); al_array_push(vcr->tracks, track); - atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED); + track->state = VCR_TRACK_RUNNING; } bool lia_vcr_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_stream *stream) @@ -304,7 +312,7 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) if (vcr->expand == VCR_EXPAND_UNTOUCHED) { vcr->mark.buffered = buffer * VCR_BUFFER_GROW_FACTOR; vcr->expand = VCR_EXPAND_GROWN; - log_debug("Expanded buffer to size %.2fMB.", vcr->mark.buffered / (f32)MB(1)); + log_info("Expanded buffer to size %.2fMB.", vcr->mark.buffered / (f32)MB(1)); return; } if (!vcr->corked) { @@ -398,38 +406,46 @@ void lia_vcr_set_buffered(struct lia_vcr_track *track) atomic_store(bool)(&track->buffered, true, AL_ATOMIC_RELAXED); } +// cork() is only ever called from within handle_packet() in vcr_track_thread(). +// Meaning track->lock will be held. void lia_vcr_cork(struct lia_vcr_track *track) { - atomic_store(s32)(&track->state, VCR_TRACK_STOPPED, AL_ATOMIC_RELAXED); + track->state = VCR_TRACK_STOPPED; } void lia_vcr_uncork(struct lia_vcr_track *track) { - if (atomic_load(s32)(&track->state, AL_ATOMIC_ACQUIRE) != 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. + // We need to avoid a race with cork() _and_ flush(). + // Possible race with flush() if locking after checking if state != STOPPED: + // - cork() -> flush() -> uncork(). + // - In uncork() we evaluate track->state to be STOPPED then wait on the lock + // being held by flush(). flush() sets the state to CLOSED and signals the + // cond. Now nn_cond_is_waiting() is false at the point uncork() acquires the lock. + nn_mutex_lock(&track->lock); + if (track->state != VCR_TRACK_STOPPED) { + nn_mutex_unlock(&track->lock); return; } - // Lock before setting track->state to avoid a race with cork(). - nn_mutex_lock(&track->mutex); - atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELEASE); + track->state = VCR_TRACK_RUNNING; al_assert(nn_cond_is_waiting(&track->cond)); nn_cond_signal(&track->cond); - nn_mutex_unlock(&track->mutex); + nn_mutex_unlock(&track->lock); } static void vcr_track_close_internal(struct lia_vcr_track *track) { struct lia_vcr *vcr = track->vcr; - // Calling packet_cache_disable() while holding the track mutex can very possibly deadlock. + // Calling packet_cache_disable() while holding track->lock can very possibly deadlock. nn_packet_cache_disable(&track->cache); - nn_mutex_lock(&track->mutex); - atomic_store(s32)(&track->state, VCR_TRACK_CLOSED, AL_ATOMIC_RELAXED); - if (nn_cond_is_waiting(&track->cond)) { + nn_mutex_lock(&track->lock); + if (track->state == VCR_TRACK_STOPPED) { + al_assert(nn_cond_is_waiting(&track->cond)); nn_cond_signal(&track->cond); + } else { + al_assert(!nn_cond_is_waiting(&track->cond)); } - nn_mutex_unlock(&track->mutex); + track->state = VCR_TRACK_CLOSED; + nn_mutex_unlock(&track->lock); if (vcr->started) { al_assert(track->running); nn_thread_join(&track->thread); @@ -463,7 +479,7 @@ void lia_vcr_flush(struct lia_vcr *vcr) atomic_store(bool)(&track->buffered, false, AL_ATOMIC_RELAXED); track->client->flush(track->client); nn_packet_cache_enable(&track->cache); - atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED); + track->state = VCR_TRACK_RUNNING; } else { track->client->flush(track->client); } diff --git a/src/liana/vcr.h b/src/liana/vcr.h index 1ab58da..5b8b9df 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -14,11 +14,11 @@ struct lia_vcr_track { struct camu_codec_stream *stream; struct lia_client_handler *client; struct nn_packet_cache cache; - atomic(s32) state; + u32 state; atomic(bool) buffered; bool running; struct nn_cond cond; - struct nn_mutex mutex; + struct nn_mutex lock; struct nn_thread thread; struct lia_vcr *vcr; }; |