#include #include #include #include #include #include #include #include "NodeBase.hpp" #include "FlowDef.hpp" #include "V210.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 { const auto flow_info = config().at("flow_id"); const auto flow_id = flow_info.at("id").get(); const int width = flow_info.value("width", 1920); const int height = flow_info.value("height", 1080); const int fps_num = flow_info.value("fps_num", 25); const int fps_den = flow_info.value("fps_den", 1); log("flow=%s", flow_id.c_str()); 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"); mxlFlowReader video_reader{}; mxlStatus vst = mxlCreateFlowReader(instance(), flow_id.c_str(), nullptr, &video_reader); if (vst != MXL_STATUS_OK) { log("mxlCreateFlowReader failed (status=%d)", vst); return; } mxlFlowConfigInfo video_cfg{}; mxlFlowReaderGetConfigInfo(video_reader, &video_cfg); const uint32_t mxl_stride = video_cfg.discrete.sliceSizes[0]; NDIContext ndi(flow_id.c_str()); // V210 (10-bit) intermediate and P216 (16-bit) send buffers std::vector buf_10bit(mxl_stride * height); std::vector buf_16bit(width * sizeof(uint16_t) * 2 * height); NDIlib_video_frame_v2_t ndi_frame_10bit{}; ndi_frame_10bit.xres = width; ndi_frame_10bit.yres = height; ndi_frame_10bit.frame_rate_N = fps_num; ndi_frame_10bit.frame_rate_D = fps_den; ndi_frame_10bit.FourCC = static_cast(NDI_LIB_FOURCC('V','2','1','0')); ndi_frame_10bit.line_stride_in_bytes = mxl_stride; ndi_frame_10bit.p_data = buf_10bit.data(); NDIlib_video_frame_v2_t ndi_frame_16bit{}; ndi_frame_16bit.xres = width; ndi_frame_16bit.yres = height; ndi_frame_16bit.frame_rate_N = fps_num; ndi_frame_16bit.frame_rate_D = fps_den; ndi_frame_16bit.line_stride_in_bytes = width * static_cast(sizeof(uint16_t)); ndi_frame_16bit.p_data = buf_16bit.data(); // audio mxlFlowReader audio_reader{}; mxlFlowConfigInfo audio_cfg{}; mxlStatus ast; int sample_rate = 0; int channels = 0; int bit_depth = 32; int no_samples = 0; bool has_audio = config().contains("audio_flow_id"); if (has_audio) { const auto audio_flow_info = config().at("audio_flow_id"); const auto audio_flow_id = audio_flow_info.at("id").get(); ast = mxlCreateFlowReader(instance(), audio_flow_id.c_str(), nullptr, &audio_reader); if (ast != MXL_STATUS_OK) { log("audio mxlCreateFlowReader failed (status=%d) — continuing without audio", ast); has_audio = false; } else { mxlFlowReaderGetConfigInfo(audio_reader, &audio_cfg); sample_rate = audio_flow_info.at("sample_rate").get(); channels = audio_cfg.continuous.channelCount; no_samples = sample_rate / fps_num; log("audio channels: %i, samples: %i", channels, audio_cfg.continuous.bufferLength); } } NDIlib_audio_frame_v3_t ndi_audio_frame; ndi_audio_frame.sample_rate = sample_rate; ndi_audio_frame.no_channels = channels; ndi_audio_frame.no_samples = no_samples; ndi_audio_frame.FourCC = NDIlib_FourCC_audio_type_FLTP; ndi_audio_frame.channel_stride_in_bytes = ndi_audio_frame.no_samples * sizeof(float); // core loop const mxlRational video_rate = {fps_num, fps_den}; const mxlRational audio_rate = {sample_rate, 1}; uint64_t video_index = mxlGetCurrentIndex(&video_rate); uint64_t audio_index = has_audio ? mxlGetCurrentIndex(&audio_rate) : 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; std::vector audio_planar(static_cast(channels) * no_samples); ndi_audio_frame.p_data = reinterpret_cast(audio_planar.data()); while (dmf::g_running.load(std::memory_order_relaxed)) { mxlGrainInfo video_grain{}; mxlGrainInfo audio_grain{}; uint8_t* video_buf = nullptr; mxlWrappedMultiBufferSlice audio_slices; vst = mxlFlowReaderGetGrainNonBlocking(video_reader, video_index, &video_grain, &video_buf); bool ndi_has_connections = NDIlib_send_get_no_connections(ndi.sender, 0) > 0; if (has_audio) { ast = mxlFlowReaderGetSamplesNonBlocking(audio_reader, audio_index, no_samples, &audio_slices); if (ast == MXL_STATUS_OK) { for (size_t c = 0; c < audio_slices.count; c++) { float* dst = audio_planar.data() + c * 1920; size_t frag0_samples = audio_slices.base.fragments[0].size / sizeof(float); const uint8_t* src0 = static_cast(audio_slices.base.fragments[0].pointer) + c * audio_slices.stride; std::memcpy(dst, src0, frag0_samples * sizeof(float)); if (audio_slices.base.fragments[1].size > 0) { size_t frag1_samples = audio_slices.base.fragments[1].size / sizeof(float); const uint8_t* src1 = static_cast(audio_slices.base.fragments[1].pointer) + c * audio_slices.stride; std::memcpy(dst + frag0_samples, src1, frag1_samples * sizeof(float)); } } NDIlib_send_send_audio_v3(ndi.sender, &ndi_audio_frame); audio_index += no_samples; } else if (ast == MXL_ERR_OUT_OF_RANGE_TOO_LATE) { mxlFlowRuntimeInfo ari{}; mxlFlowReaderGetRuntimeInfo(audio_reader, &ari); audio_index = ari.headIndex; } } if (vst == MXL_STATUS_OK) { frame_count++; if (video_grain.flags & MXL_GRAIN_FLAG_INVALID) invalid_count++; if (ndi_has_connections) { std::memcpy(ndi_frame_10bit.p_data, video_buf, mxl_stride * height); NDIlib_util_V210_to_P216(&ndi_frame_10bit, &ndi_frame_16bit); NDIlib_send_send_video_v2(ndi.sender, &ndi_frame_16bit); if (++ndi_frame_count == 1) log("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 { log("unexpected status=%d on index=%llu", 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; } } log("stopped — total frames=%llu invalid=%llu late=%llu", frame_count, invalid_count, late_count); mxlReleaseFlowReader(instance(), video_reader); if (has_audio) mxlReleaseFlowReader(instance(), audio_reader); } }; int main() { NDIOutNode node; return node.execute(); }