#include #include #include #include #include "NodeBase.hpp" #include "FlowDef.hpp" #include "V210.hpp" #include 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 reader{}; mxlStatus st = mxlCreateFlowReader(instance(), flow_id.c_str(), nullptr, &reader); if (st != MXL_STATUS_OK) { log("mxlCreateFlowReader failed (status=%d)", st); return; } mxlFlowConfigInfo cfg_info{}; mxlFlowReaderGetConfigInfo(reader, &cfg_info); const uint32_t mxl_stride = cfg_info.discrete.sliceSizes[0]; const mxlRational rate = {fps_num, fps_den}; uint64_t index = mxlGetCurrentIndex(&rate); uint64_t frame_count = 0; uint64_t invalid_count = 0; uint64_t late_count = 0; auto wall_start = std::chrono::steady_clock::now(); auto last_log_time = wall_start; // NDI part if (!NDIlib_initialize()) { // Cannot run NDI. Most likely because the CPU is not sufficient (see SDK documentation). log("Cannot run NDI"); if (!NDIlib_is_supported_CPU()) { log("CPU is not sufficient for NDI"); } return; } NDIlib_send_create_t NDI_send_create_desc; NDI_send_create_desc.p_ndi_name = flow_id.c_str(); NDIlib_send_instance_t pNDI_send = NDIlib_send_create(&NDI_send_create_desc); if (!pNDI_send) { log("Cannot create NDI send instance"); return; } NDIlib_video_frame_v2_t NDI_video_frame_10bit; NDI_video_frame_10bit.xres = width; NDI_video_frame_10bit.yres = height; NDI_video_frame_10bit.FourCC = (NDIlib_FourCC_video_type_e)NDI_LIB_FOURCC('V', '2', '1', '0'); NDI_video_frame_10bit.line_stride_in_bytes = mxl_stride; NDI_video_frame_10bit.p_data = (uint8_t*)malloc(NDI_video_frame_10bit.line_stride_in_bytes * NDI_video_frame_10bit.yres); NDIlib_video_frame_v2_t NDI_video_frame_16bit; NDI_video_frame_16bit.xres = NDI_video_frame_10bit.xres; NDI_video_frame_16bit.yres = NDI_video_frame_10bit.yres; NDI_video_frame_16bit.frame_rate_N = fps_num; NDI_video_frame_16bit.frame_rate_D = fps_den; NDI_video_frame_16bit.line_stride_in_bytes = NDI_video_frame_16bit.xres * sizeof(uint16_t); NDI_video_frame_16bit.p_data = (uint8_t*)malloc(NDI_video_frame_16bit.line_stride_in_bytes * 2 * NDI_video_frame_16bit.yres); // uint64_t ndi_frame_counter = 0; while (dmf::g_running.load(std::memory_order_relaxed)) { mxlGrainInfo grain{}; uint8_t* buf = nullptr; st = mxlFlowReaderGetGrainNonBlocking(reader, index, &grain, &buf); if (st == MXL_STATUS_OK) { frame_count++; if (grain.flags & MXL_GRAIN_FLAG_INVALID) invalid_count++; index++; if (NDIlib_send_get_no_connections(pNDI_send,0) == 0) { ndi_frame_counter = 0; continue; } memcpy(NDI_video_frame_10bit.p_data, buf, mxl_stride * height); NDIlib_util_V210_to_P216(&NDI_video_frame_10bit, &NDI_video_frame_16bit); NDIlib_send_send_video_v2(pNDI_send, &NDI_video_frame_16bit); ndi_frame_counter++; if (ndi_frame_counter == 1) log("NDI receiver connected"); } else if (st == MXL_ERR_OUT_OF_RANGE_TOO_EARLY) { mxlSleepForNs(1'000'000); // 1 ms poll } else if (st == MXL_ERR_OUT_OF_RANGE_TOO_LATE) { late_count++; // Jump to the most recent frame in the ring buffer mxlFlowRuntimeInfo ri{}; mxlFlowReaderGetRuntimeInfo(reader, &ri); index = ri.headIndex; } else { log("unexpected status=%d on index=%llu", st, 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(), reader); free(NDI_video_frame_10bit.p_data); free(NDI_video_frame_16bit.p_data); NDIlib_send_destroy(pNDI_send); NDIlib_destroy(); } }; int main() { NDIOutNode node; return node.execute(); }