#include #include #include #include #include #include #include #include "NodeBase.hpp" #include // RAII wrapper: init NDI, create sender, destroy both on scope exit. struct NDIContext { NDIlib_send_instance_t sender = nullptr; explicit NDIContext(const char* ndi_name) { if (!NDIlib_is_supported_CPU()) throw std::runtime_error("CPU not sufficient for NDI"); if (!NDIlib_initialize()) throw std::runtime_error("NDI lib init failed"); NDIlib_send_create_t desc{}; desc.p_ndi_name = ndi_name; sender = NDIlib_send_create(&desc); if (!sender) { NDIlib_destroy(); throw std::runtime_error("Cannot create NDI send instance"); } } ~NDIContext() { NDIlib_send_destroy(sender); NDIlib_destroy(); } }; class NDIOutNode : public dmf::NodeBase { void run() override { // --- video flow (optional) --- bool has_video = config().contains("flow_id"); int width = 1920; int height = 1080; int fps_num = 25; int fps_den = 1; std::string flow_id; mxlFlowReader video_reader{}; uint32_t video_stride = 0; if (has_video) { const auto flow_info = config().at("flow_id"); flow_id = flow_info.at("id").get(); width = flow_info.value("width", 1920); height = flow_info.value("height", 1080); fps_num = flow_info.value("fps_num", 25); fps_den = flow_info.value("fps_den", 1); log("video flow=%s %dx%d @ %d/%d fps", flow_id.c_str(), width, height, fps_num, fps_den); log("waiting for flow to become active..."); bool active = false; while (!active && dmf::g_running.load(std::memory_order_relaxed)) { mxlIsFlowActive(instance(), flow_id.c_str(), &active); if (!active) mxlSleepForNs(100'000'000); } if (!dmf::g_running) return; log("flow active — starting read"); mxlFlowConfigInfo video_cfg{}; mxlStatus vst = mxlCreateFlowReader(instance(), flow_id.c_str(), "", &video_reader); if (vst != MXL_STATUS_OK) { log("video mxlCreateFlowReader failed (%s)", dmf::mxl_status_str(vst)); return; } mxlFlowReaderGetConfigInfo(video_reader, &video_cfg); video_stride = video_cfg.discrete.sliceSizes[0]; } // --- audio flow (optional) --- mxlFlowReader audio_reader{}; mxlFlowConfigInfo audio_cfg{}; int sample_rate = 0; int channels = 0; int samples_per_frame = 0; bool has_audio = config().contains("audio_flow_id"); std::string audio_flow_id; if (has_audio) { const auto audio_flow_info = config().at("audio_flow_id"); audio_flow_id = audio_flow_info.at("id").get(); sample_rate = audio_flow_info.value("sample_rate", 48000); channels = audio_flow_info.value("channels", 2); samples_per_frame = sample_rate * fps_den / fps_num; log("audio flow=%s %d Hz %dch %d samples/frame", audio_flow_id.c_str(), sample_rate, channels, samples_per_frame); mxlStatus ast = mxlCreateFlowReader(instance(), audio_flow_id.c_str(), "", &audio_reader); if (ast != MXL_STATUS_OK) { log("audio mxlCreateFlowReader failed (%s) — continuing without audio", dmf::mxl_status_str(ast)); has_audio = false; } else { mxlFlowReaderGetConfigInfo(audio_reader, &audio_cfg); log("audio channels=%u buffer=%u samples", audio_cfg.continuous.channelCount, audio_cfg.continuous.bufferLength); } } if (!has_video && !has_audio) { log("no flows configured — exiting"); return; } NDIContext ndi(node_id().c_str()); // Video buffers and NDI frames (only when has_video) std::vector v210_buf, p216_buf; NDIlib_video_frame_v2_t v210_frame{}, p216_frame{}; if (has_video) { v210_buf.resize(video_stride * height); p216_buf.resize(width * sizeof(uint16_t) * 2 * height); v210_frame.xres = width; v210_frame.yres = height; v210_frame.frame_rate_N = fps_num; v210_frame.frame_rate_D = fps_den; v210_frame.FourCC = static_cast(NDI_LIB_FOURCC('V','2','1','0')); v210_frame.line_stride_in_bytes = video_stride; v210_frame.p_data = v210_buf.data(); p216_frame.xres = width; p216_frame.yres = height; p216_frame.FourCC = NDIlib_FourCC_video_type_P216; p216_frame.frame_rate_N = fps_num; p216_frame.frame_rate_D = fps_den; p216_frame.line_stride_in_bytes = width * static_cast(sizeof(uint16_t)); p216_frame.p_data = p216_buf.data(); } // Audio buffer and NDI frame (only when has_audio) std::vector audio_planar(static_cast(channels) * samples_per_frame); NDIlib_audio_frame_v3_t ndi_audio{}; if (has_audio) { ndi_audio.sample_rate = sample_rate; ndi_audio.no_channels = channels; ndi_audio.no_samples = samples_per_frame; ndi_audio.FourCC = NDIlib_FourCC_audio_type_FLTP; ndi_audio.channel_stride_in_bytes = samples_per_frame * sizeof(float); ndi_audio.p_data = reinterpret_cast(audio_planar.data()); } // --- main loop --- uint64_t video_index = 0; uint64_t audio_index = 0; uint64_t frame_count = 0; uint64_t invalid_count = 0; uint64_t late_count = 0; uint64_t ndi_frame_count = 0; auto wall_start = std::chrono::steady_clock::now(); auto last_log_time = wall_start; if (has_video) { const mxlRational video_rate = {fps_num, fps_den}; video_index = mxlGetCurrentIndex(&video_rate); } if (has_audio) { const mxlRational audio_rate = {sample_rate, 1}; audio_index = mxlGetCurrentIndex(&audio_rate); } while (dmf::g_running.load(std::memory_order_relaxed)) { // --- audio: non-blocking, one chunk per video frame (or free-running) --- bool audio_advanced = false; if (has_audio) { mxlWrappedMultiBufferSlice audio_slices{}; mxlStatus ast = mxlFlowReaderGetSamplesNonBlocking( audio_reader, audio_index, samples_per_frame, &audio_slices); if (ast == MXL_STATUS_OK) { const size_t frag0 = audio_slices.base.fragments[0].size / sizeof(float); const size_t frag1 = audio_slices.base.fragments[1].size / sizeof(float); for (int c = 0; c < channels; ++c) { float* dst = audio_planar.data() + c * samples_per_frame; const auto* src0 = reinterpret_cast( static_cast(audio_slices.base.fragments[0].pointer) + c * audio_slices.stride); std::memcpy(dst, src0, frag0 * sizeof(float)); if (frag1 > 0) { const auto* src1 = reinterpret_cast( static_cast(audio_slices.base.fragments[1].pointer) + c * audio_slices.stride); std::memcpy(dst + frag0, src1, frag1 * sizeof(float)); } } NDIlib_send_send_audio_v3(ndi.sender, &ndi_audio); audio_index += samples_per_frame; audio_advanced = true; } else if (ast == MXL_ERR_OUT_OF_RANGE_TOO_LATE) { mxlFlowRuntimeInfo ari{}; mxlFlowReaderGetRuntimeInfo(audio_reader, &ari); audio_index = ari.headIndex; } else if (ast == MXL_ERR_FLOW_INVALID) { log("audio flow invalidated — reconnecting..."); mxlReleaseFlowReader(instance(), audio_reader); audio_reader = nullptr; mxlSleepForNs(100'000'000); if (mxlCreateFlowReader(instance(), audio_flow_id.c_str(), "", &audio_reader) == MXL_STATUS_OK) { log("audio flow reconnected"); const mxlRational r = {sample_rate, 1}; audio_index = mxlGetCurrentIndex(&r); } } } // --- video --- if (has_video) { mxlGrainInfo video_grain{}; uint8_t* video_buf = nullptr; mxlStatus vst = mxlFlowReaderGetGrainNonBlocking( video_reader, video_index, &video_grain, &video_buf); if (vst == MXL_STATUS_OK) { frame_count++; if (video_grain.flags & MXL_GRAIN_FLAG_INVALID) invalid_count++; if (NDIlib_send_get_no_connections(ndi.sender, 0) > 0) { std::memcpy(v210_frame.p_data, video_buf, video_stride * height); NDIlib_util_V210_to_P216(&v210_frame, &p216_frame); NDIlib_send_send_video_v2(ndi.sender, &p216_frame); if (++ndi_frame_count == 1) log("first NDI receiver connected"); } else { ndi_frame_count = 0; } video_index++; } else if (vst == MXL_ERR_OUT_OF_RANGE_TOO_EARLY) { mxlSleepForNs(1'000'000); } else if (vst == MXL_ERR_OUT_OF_RANGE_TOO_LATE) { late_count++; mxlFlowRuntimeInfo ri{}; mxlFlowReaderGetRuntimeInfo(video_reader, &ri); video_index = ri.headIndex; } else if (vst == MXL_ERR_FLOW_INVALID) { log("video flow invalidated — reconnecting..."); mxlReleaseFlowReader(instance(), video_reader); video_reader = nullptr; mxlSleepForNs(100'000'000); if (mxlCreateFlowReader(instance(), flow_id.c_str(), "", &video_reader) == MXL_STATUS_OK) { log("video flow reconnected"); const mxlRational r = {fps_num, fps_den}; video_index = mxlGetCurrentIndex(&r); } } else { log("unexpected video status (%s) on index=%llu", dmf::mxl_status_str(vst), video_index); break; } auto now = std::chrono::steady_clock::now(); if (std::chrono::duration(now - last_log_time).count() >= 1.0) { const double elapsed = std::chrono::duration(now - wall_start).count(); log("frames=%llu invalid=%llu late=%llu avg=%.2f fps", frame_count, invalid_count, late_count, static_cast(frame_count) / elapsed); last_log_time = now; } } else if (!audio_advanced) { // audio-only and nothing was ready — avoid busy spin mxlSleepForNs(1'000'000); } } if (has_video) log("stopped — total frames=%llu invalid=%llu late=%llu", frame_count, invalid_count, late_count); else log("stopped"); if (has_video && video_reader) mxlReleaseFlowReader(instance(), video_reader); if (has_audio && audio_reader) mxlReleaseFlowReader(instance(), audio_reader); } }; int main() { NDIOutNode node; return node.execute(); }