#define AL_LOG_SECTION "vcr" //#define AL_LOG_ENABLE_TRACE #include #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) 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_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; 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; } } 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->count, 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; } } #endif 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->count, 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; #ifndef CAMU_DIRECT_MODE vcr->corked = false; nn_signal_init(&vcr->signal, loop, signal_callback, vcr); reset_metrics(vcr); #endif } static void return_entire_cache(struct lia_vcr_track *track) { struct lia_vcr *vcr = track->vcr; u32 count = track->cache.cache.count; nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), count); track->cache.cache.count = 0; } // 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. 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; const char thread_name[16] = "\0"; // 16 = limit. 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; while (nn_packet_cache_wait(&track->cache, &packets)) { 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 && packets > 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->count, 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->mutex); // NULL packet means flush. bool success = track->client->handle_packet(track->client, packet); if (!success) { // In the case of codec_client, an EOF will be sent on an error. log_error("Error handling packet, exiting track thread."); return_entire_cache(track); track->cache.disabled = true; nn_packet_cache_unlock(&track->cache); nn_mutex_unlock(&track->mutex); 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); // 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; } index++; // Count of packets consumed, also increment for WHOLE_FILE mode. #ifndef VCR_BUFFER_WHOLE_FILE nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), 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->mutex); // Check for possibly updated state. state = 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 (corked) { // The cache is already unlocked here. #ifndef VCR_BUFFER_WHOLE_FILE index = 0; #endif } 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); } } 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->running); if (VCR_TRACK_THREADED(track)) { nn_thread_create(&track->thread, vcr_track_thread, track); track->running = 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->mutex); 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); } 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); nn_packet_cache_free(&track->cache); track->client->free(&track->client); 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; } #ifndef CAMU_DIRECT_MODE 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) { buffered &= atomic_load(bool)(&track->buffered, AL_ATOMIC_RELAXED); } if (buffered) { 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) { nn_packet_cache_flush(&track->cache); } } } #endif 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: { u32 size = nn_packet_get_size(packet); #ifndef CAMU_DIRECT_MODE update_metrics(vcr, size); #endif 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 can_send = nn_packet_cache_available(&track->cache); if (can_send) { buffer = atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED); nn_packet_cache_send_packet(&track->cache, packet); } nn_packet_cache_unlock(&track->cache); if (can_send) { if (buffer >= vcr->mark.buffered) { #ifndef CAMU_DIRECT_MODE cork_if_buffered(vcr, buffer); #endif } return; // Keep packet. } } else { if (!track->client->handle_packet(track->client, packet)) { log_warn("Error handling non-buffered packet."); } } break; } case LIANA_PACKET_EOF: case LIANA_PACKET_ERROR: { // @TODO: Should ERROR be passed down to LIANA_CLIENT_ERRORED? #ifndef CAMU_DIRECT_MODE update_metrics(vcr, 0); // Flush. #endif al_array_foreach(vcr->tracks, i, track) { if (nn_packet_cache_available(&track->cache)) { nn_packet_cache_send_packet(&track->cache, NULL); } nn_packet_cache_unlock(&track->cache); } #ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); #endif if (op == LIANA_PACKET_EOF) { log_info("Received EOF."); } else if (op == LIANA_PACKET_ERROR) { log_warn("Forcing EOF due to an error packet."); } break; } default: log_warn("Erroneous packet."); break; } nn_packet_stream_return_packet(vcr->data, packet); } void lia_vcr_set_buffered(struct lia_vcr_track *track) { atomic_store(bool)(&track->buffered, true, AL_ATOMIC_RELAXED); } void lia_vcr_cork(struct lia_vcr_track *track) { atomic_store(s32)(&track->state, VCR_TRACK_STOPPED, AL_ATOMIC_RELAXED); } 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. 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); 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) { struct lia_vcr *vcr = track->vcr; // Calling packet_cache_disable() while holding the track mutex 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_cond_signal(&track->cond); } nn_mutex_unlock(&track->mutex); if (vcr->started) { al_assert(track->running); nn_thread_join(&track->thread); track->running = false; } } static void vcr_close_all_internal(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { vcr_track_close_internal(track); return_entire_cache(track); } } vcr->started = false; #ifndef CAMU_DIRECT_MODE vcr->corked = false; nn_signal_stop(&vcr->signal); reset_metrics(vcr); #endif } void lia_vcr_flush(struct lia_vcr *vcr) { vcr_close_all_internal(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->client->flush(track->client); nn_packet_cache_enable(&track->cache); atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED); } else { track->client->flush(track->client); } } atomic_store(u64)(&vcr->count, 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_internal(vcr); } void lia_vcr_free(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { nn_packet_cache_free(&track->cache); track->client->free(&track->client); al_free(track); } al_array_free(vcr->tracks); }