#define AL_LOG_SECTION "vcr" //#define AL_LOG_ENABLE_TRACE #include #include #include "handlers/handler.h" #include "vcr.h" #include "client.h" #include "common.h" // packet_cache_v2 // - Each track has it's own cond. // - Tracks get signaled in order of stream index. // - Cache has to handle case where the next packet in order is // from a track that has yet to call packet_cache_wait() (Immediate return). // - Once track has read as much as it can: // - If it read all avaliable packets, call packet_cache_wait() again and that entire // section will be consumed. // - If it was partial, call packet_cache_yield() for that many packets // to be consumed. // Considerations: // - Works after client disconnect. // - Optimizes seek and reseek. // - Supports multiple AVPackets per nn_packet. #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 { VCR_EXPAND_UNTOUCHED = 0, VCR_EXPAND_GROWN, VCR_EXPAND_COMPLETE }; enum { VCR_TRACK_RUNNING = 0, VCR_TRACK_STOPPED, VCR_TRACK_ERRORED, VCR_TRACK_CLOSED }; enum { VCR_NOT_EOF = 0, VCR_EOF, VCR_EOF_ERRORED }; #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; if (vcr->corked) { nn_packet_stream_cork(vcr->data, false); vcr->corked = false; log_trace("Uncorked."); u64 now = nn_get_timestamp(); // This keeps the difference between `mark` and `now` equal to the amount // of time we were actually receiving packets for. al_assert(vcr->metrics.last_report_mark != LIANA_TIMESTAMP_INVALID); vcr->metrics.last_report_mark += (now - vcr->metrics.last_cork_ts); al_assert(vcr->metrics.last_report_mark < now); vcr->metrics.last_cork_ts = LIANA_TIMESTAMP_INVALID; } } #endif static void reset_metrics(struct lia_vcr *vcr) { vcr->metrics.current_frame = 0; vcr->metrics.last_cork_ts = LIANA_TIMESTAMP_INVALID; vcr->metrics.last_report_ts = LIANA_TIMESTAMP_INVALID; vcr->metrics.last_report_mark = LIANA_TIMESTAMP_INVALID; vcr->metrics.average_kbps = 0.f; } static void update_metrics(struct lia_vcr *vcr, u64 size) { vcr->metrics.current_frame += size; log_trace("current_frame: %lu, corked: %s.", vcr->metrics.current_frame, BOOLSTR(vcr->corked)); u64 now = nn_get_timestamp(); if (vcr->metrics.last_report_mark == LIANA_TIMESTAMP_INVALID) { vcr->metrics.last_report_ts = now; vcr->metrics.last_report_mark = now; return; } u64 diff = now - vcr->metrics.last_report_ts; u64 mark = now - vcr->metrics.last_report_mark; u64 frame = vcr->metrics.current_frame; if (mark > 750000 || (diff > 5000000 && vcr->metrics.current_frame >= KB(500)) || (!size && frame > 0)) { al_assert(frame > 0); f32 kbps = (frame / 125.f) / (mark / 1000000.f); f32 average_kbps = vcr->metrics.average_kbps; average_kbps = average_kbps == 0.f ? kbps : (average_kbps + kbps) / 2.f; f32 buffered = atomic_load(u64)(&vcr->size, AL_ATOMIC_RELAXED) / (f32)MB(1); f32 capacity = vcr->mark.buffered / (f32)MB(1); log_info("Receiving packets at %.2fkbps (%.2f/%.2fMB).", average_kbps, buffered, capacity); vcr->metrics.average_kbps = average_kbps; vcr->metrics.last_report_ts = now; vcr->metrics.last_report_mark = now; vcr->metrics.current_frame = 0; } } void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data, u16 node_id) { vcr->data = data; vcr->node_id = node_id; al_array_init(vcr->tracks); atomic_store(u64)(&vcr->size, 0, AL_ATOMIC_RELAXED); vcr->mark.buffered = VCR_BUFFER_BUFFERED; atomic_store(u64)(&vcr->mark.low, 0, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_UNTOUCHED; vcr->started = false; vcr->corked = false; #ifndef CAMU_DIRECT_MODE nn_signal_init(&vcr->signal, loop, signal_callback, vcr); #else (void)loop; #endif reset_metrics(vcr); } #define count_minus_eof(packets, count) (packets[count - 1] ? count : count - 1) static void return_entire_cache(struct lia_vcr_track *track) { u32 count = track->cache.cache.count; if (count > 0) { struct lia_vcr *vcr = track->vcr; struct nn_packet **packets = al_array_offset(track->cache.cache, 0); count = count_minus_eof(packets, count); #ifdef VCR_BUFFER_WHOLE_FILE for (u32 i = 0; i < count; i++) { disown_packet(packets[i]); } #endif nn_packet_stream_return_packets(vcr->data, packets, count); track->cache.cache.count = 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; char thread_name[16] = "\0"; // 16 = limit. al_snprintf(thread_name, sizeof(thread_name), "vcr:%hu_%d", vcr->node_id, track->stream->index); nn_thread_set_name(thread_name); bool corked; u32 count, index = 0; struct nn_packet *packet = NULL; while (nn_packet_cache_wait(&track->cache, &count)) { al_assert(count >= index); corked = false; for (; index < count; index++) { packet = nn_packet_cache_at(&track->cache, index); #ifdef VCR_BUFFER_WHOLE_FILE if (!packet && count > 2) { // Loop. index = 0; corked = false; break; } #endif if (packet) { // Check if we should uncork the packet stream. u32 size = nn_packet_get_size(packet); u64 buffer = atomic_sub(u64)(&vcr->size, size, AL_ATOMIC_RELAXED); #ifndef CAMU_DIRECT_MODE bool buffered = atomic_load(bool)(&track->buffered, AL_ATOMIC_RELAXED); u64 low = atomic_load(u64)(&vcr->mark.low, AL_ATOMIC_RELAXED); if (buffered && (low && buffer <= low)) { nn_signal_send(&vcr->signal); } #else (void)buffer; #endif } nn_mutex_lock(&track->lock); struct lia_client_handler *handler = track->handler; bool success = handler->handle_packet(handler, packet); // NULL packet = flush. if (!success || (!packet && track->eof == VCR_EOF_ERRORED)) { log_error("Error handling packet, exiting track thread."); handler->callback(handler->userdata, LIANA_CLIENT_ERRORED, handler->stream, NULL); return_entire_cache(track); track->cache.disabled = true; nn_packet_cache_unlock(&track->cache); // Let flush()/close() know we forcefully exited the thread. track->state = VCR_TRACK_ERRORED; nn_mutex_unlock(&track->lock); return 0; } if (!packet) { // Wait on EOF. track->state = VCR_TRACK_STOPPED; } // 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->lock); continue; } index++; // Count of packets consumed, also increment for WHOLE_FILE mode. #ifndef VCR_BUFFER_WHOLE_FILE struct nn_packet **packets = al_array_offset(track->cache.cache, 0); nn_packet_stream_return_packets(vcr->data, packets, count_minus_eof(packets, index)); al_array_remove_range(track->cache.cache, 0, index); #endif nn_packet_cache_unlock(&track->cache); // Wait for uncork. nn_cond_wait(&track->cond, &track->lock); // Check for possibly updated state. u32 state = track->state; nn_mutex_unlock(&track->lock); 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 (corked) { // The cache is already unlocked here. #ifndef VCR_BUFFER_WHOLE_FILE index = 0; #endif } else { #ifndef VCR_BUFFER_WHOLE_FILE al_assert(index == count); index = 0; if (count > 0) { struct nn_packet **packets = al_array_offset(track->cache.cache, 0); nn_packet_stream_return_packets(vcr->data, packets, count_minus_eof(packets, count)); al_array_remove_range(track->cache.cache, 0, count); } #endif 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) { al_assert(!track->started); if (VCR_TRACK_THREADED(track)) { nn_thread_create(&track->thread, vcr_track_thread, track); track->started = true; } } vcr->started = 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->lock); atomic_store(bool)(&track->buffered, !VCR_TRACK_THREADED(track), AL_ATOMIC_RELAXED); track->started = false; if (VCR_TRACK_THREADED(track)) { nn_packet_cache_init(&track->cache, 256); track->state = VCR_TRACK_RUNNING; } track->eof = VCR_NOT_EOF; al_array_push(vcr->tracks, track); } bool lia_vcr_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_stream *stream) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (track->stream == stream) { al_array_remove_at(vcr->tracks, i); if (VCR_TRACK_THREADED(track)) { nn_packet_cache_free(&track->cache); } track->handler->free(&track->handler); al_free(track); return true; } } return false; } bool lia_vcr_is_empty(struct lia_vcr *vcr) { return !vcr->tracks.count; } static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (track->stream->index == index) { return track; } } return NULL; } static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) { bool buffered = true; struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { // The audio/video buffers that a track supplies call vcr_set_buffered(track) // when they decide that they're buffered. buffered &= atomic_load(bool)(&track->buffered, AL_ATOMIC_RELAXED); } if (buffered) { #ifndef CAMU_DIRECT_MODE 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)); return; } if (!vcr->corked) { log_trace("Corked."); nn_packet_stream_cork(vcr->data, true); vcr->corked = true; vcr->metrics.last_cork_ts = nn_get_timestamp(); } // Don't set low until after we corked so that vcr_track_thread() will never try uncorking // until we know what the low mark is. if (vcr->expand == VCR_EXPAND_GROWN) { atomic_store(u64)(&vcr->mark.low, vcr->mark.buffered - VCR_BUFFER_LOW_OFFSET, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_COMPLETE; } al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { nn_packet_cache_flush(&track->cache); } } #else (void)buffer; #endif } } bool 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: { u32 size = nn_packet_get_size(packet); update_metrics(vcr, size); s32 index = nn_packet_read_s32(packet); if (!(track = get_track_from_index(vcr, index))) { log_error("Received data from an errored or unknown track (index: %d).", index); break; } if (VCR_TRACK_THREADED(track)) { u64 buffer = 0; bool send_to_cache = !nn_packet_cache_disabled(&track->cache); if (send_to_cache) { buffer = atomic_add(u64)(&vcr->size, size, AL_ATOMIC_RELAXED); nn_packet_cache_send_packet(&track->cache, packet); } nn_packet_cache_unlock(&track->cache); if (send_to_cache) { if (buffer >= vcr->mark.buffered) { cork_if_buffered(vcr, buffer); } return true; // Keep packet. } } else if (!track->handler->handle_packet(track->handler, packet)) { log_warn("Error handling non-buffered packet."); } break; } case LIANA_PACKET_EOF: case LIANA_PACKET_ERROR: { bool error = op == LIANA_PACKET_ERROR; update_metrics(vcr, 0); // Flush. al_array_foreach(vcr->tracks, i, track) { al_assert(track->eof == VCR_NOT_EOF); track->eof = error ? VCR_EOF_ERRORED : VCR_EOF; if (VCR_TRACK_THREADED(track)) { if (!nn_packet_cache_disabled(&track->cache)) { nn_packet_cache_send_packet(&track->cache, NULL); } nn_packet_cache_unlock(&track->cache); } } #ifndef CAMU_DIRECT_MODE // A corked stream won't close after an unexpected disconnect. This results in better // behavior for the sink (e.g., an image buffer won't immediately be removed). nn_packet_stream_cork(vcr->data, true); nn_signal_stop(&vcr->signal); #endif if (error) { log_warn("Forcing EOF due to an error packet."); } else { log_info("Received EOF."); } break; } default: log_warn("Erroneous packet."); break; } return false; } 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) { al_assert(track->state != VCR_TRACK_ERRORED); track->state = VCR_TRACK_STOPPED; } void lia_vcr_uncork(struct lia_vcr_track *track) { // 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. // Possible race with vcr_track_thread(): // -> Track corks while holding lock. // -> Buffer goes under MARK_LOW, tries to uncork(), evaluates track->state = STOPPED // then waits on the lock. // -> handle_packet() errors and sets state to ERRORED. // -> uncork() aquires the lock and causes an invalid state. nn_mutex_lock(&track->lock); if (track->state != VCR_TRACK_STOPPED) { nn_mutex_unlock(&track->lock); return; } track->state = VCR_TRACK_RUNNING; al_assert(nn_cond_is_waiting(&track->cond)); nn_cond_signal(&track->cond); nn_mutex_unlock(&track->lock); } static void vcr_threaded_track_close(struct lia_vcr_track *track) { struct lia_vcr *vcr = track->vcr; // Calling packet_cache_disable() while holding track->lock can very possibly deadlock. nn_packet_cache_disable(&track->cache); 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)); } track->state = VCR_TRACK_CLOSED; nn_mutex_unlock(&track->lock); if (vcr->started) { al_assert(track->started); nn_thread_join(&track->thread); track->started = false; } } static void vcr_close_all(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { vcr_threaded_track_close(track); return_entire_cache(track); } } vcr->started = false; vcr->corked = false; #ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); #endif reset_metrics(vcr); } void lia_vcr_flush(struct lia_vcr *vcr) { vcr_close_all(vcr); struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { atomic_store(bool)(&track->buffered, false, AL_ATOMIC_RELAXED); track->handler->flush(track->handler); nn_packet_cache_enable(&track->cache); track->state = VCR_TRACK_RUNNING; } else { track->handler->flush(track->handler); } track->eof = VCR_NOT_EOF; } atomic_store(u64)(&vcr->size, 0, AL_ATOMIC_RELAXED); if (vcr->expand == VCR_EXPAND_COMPLETE) { atomic_store(u64)(&vcr->mark.low, 0, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_GROWN; } } void lia_vcr_close_all(struct lia_vcr *vcr) { vcr_close_all(vcr); } void lia_vcr_free(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { nn_packet_cache_free(&track->cache); } track->handler->free(&track->handler); al_free(track); } al_array_free(vcr->tracks); }