diff --git a/nodes/ndiin/main.cpp b/nodes/ndiin/main.cpp index fd106c8..0e1388f 100644 --- a/nodes/ndiin/main.cpp +++ b/nodes/ndiin/main.cpp @@ -1,37 +1,36 @@ +#include #include +#include #include #include #include "NodeBase.hpp" #include "FlowDef.hpp" -#include "V210.hpp" -#include "NDIHelper.hpp" -#include -#include +#include "NDIReceiver.hpp" class NDIInNode : public dmf::NodeBase { void run() override { - const uint32_t source_num = config().value("source_num", 0); + const auto source_num = static_cast(config().value("source_num", 0)); - dmf::NDIHelper ndi_helper; + dmf::NDIReceiver ndi; + dmf::NDIReceiver::SourceInfo src; try { - std::vector ndi_sources; - ndi_helper.find_sources(&ndi_sources, 5000); + auto sources = ndi.find_sources(5000); log("Available NDI sources:"); - for (const auto& source_name : ndi_sources) - log(" %s", source_name.c_str()); - ndi_helper.select_source(source_num); - ndi_helper.get_source_info(source_num); + for (const auto& name : sources) + log(" %s", name.c_str()); + ndi.connect(source_num); + src = ndi.probe(); } catch (const std::runtime_error& e) { log("Error: %s", e.what()); return; } - + const auto flow_info = config().at("flow_id"); const auto flow_id = flow_info.at("id").get(); - const int width = flow_info.value("width", ndi_helper.xres); - const int height = flow_info.value("height", ndi_helper.yres); - const int fps_num = flow_info.value("fps_num", ndi_helper.frame_N); - const int fps_den = flow_info.value("fps_den", ndi_helper.frame_D); + const int width = flow_info.value("width", src.width); + const int height = flow_info.value("height", src.height); + const int fps_num = flow_info.value("fps_num", src.fps_num); + const int fps_den = flow_info.value("fps_den", src.fps_den); log("flow=%s %dx%d @ %d/%d fps", flow_id.c_str(), width, height, fps_num, fps_den); @@ -56,23 +55,22 @@ class NDIInNode : public dmf::NodeBase { const mxlRational rate = {fps_num, fps_den}; uint64_t index = mxlGetCurrentIndex(&rate); log("start index=%llu", index); - - // because NDI can drop to 1 FPS for static frames, even if source is 29.97p - const size_t frame_bytes = stride * height; - uint8_t* latest_buffer = (uint8_t*)malloc(frame_bytes); - bool last_ndi_frame_valid = false; + + // NDI can drop to 1fps for static content — hold last valid frame + std::vector latest_frame(stride * height); + bool have_frame = false; while (dmf::g_running.load(std::memory_order_relaxed)) { try { - if (ndi_helper.getV210_video_frame(source_num, latest_buffer, stride)) - last_ndi_frame_valid = true; + if (ndi.capture_v210(latest_frame.data(), stride)) + have_frame = true; } catch (const std::runtime_error& e) { log("NDI error: %s — stopping", e.what()); break; } - uint8_t* buf = nullptr; mxlGrainInfo grain{}; + uint8_t* buf = nullptr; st = mxlFlowWriterOpenGrain(writer, index, &grain, &buf); if (st != MXL_STATUS_OK) { @@ -81,14 +79,14 @@ class NDIInNode : public dmf::NodeBase { continue; } - if (last_ndi_frame_valid) { - std::memcpy(buf, latest_buffer, frame_bytes); + if (have_frame) { + std::memcpy(buf, latest_frame.data(), latest_frame.size()); grain.flags = 0; } else { grain.flags = MXL_GRAIN_FLAG_INVALID; } - grain.validSlices = grain.totalSlices; // mark grain complete so readers can consume it + grain.validSlices = grain.totalSlices; mxlFlowWriterCommitGrain(writer, &grain); const uint64_t ns = mxlGetNsUntilIndex(index + 1, &rate); @@ -97,13 +95,11 @@ class NDIInNode : public dmf::NodeBase { } log("stopped at index=%llu", index); - free(latest_buffer); mxlReleaseFlowWriter(instance(), writer); } }; -int main() -{ +int main() { NDIInNode node; return node.execute(); -} \ No newline at end of file +} diff --git a/shared/NDIHelper.hpp b/shared/NDIHelper.hpp deleted file mode 100644 index 50ca5e3..0000000 --- a/shared/NDIHelper.hpp +++ /dev/null @@ -1,206 +0,0 @@ -#include -#include -#include -#include -#include "Signal.hpp" -#include "V210.hpp" - -namespace dmf { -class NDIHelper { - public: - - int xres = 0, yres = 0, frame_D = 0, frame_N = 0, stride = 0; - - NDIHelper() { - if (!NDIlib_is_supported_CPU()) { - throw std::runtime_error("CPU is not sufficient for NDI"); - } - if (!NDIlib_initialize()) { - throw std::runtime_error("NDI lib init failed"); - } - } - - ~NDIHelper() { - if (pNDI_recv) NDIlib_recv_destroy(pNDI_recv); - if (pNDI_find) NDIlib_find_destroy(pNDI_find); - NDIlib_destroy(); - } - - void find_sources(std::vector* sources, u_int32_t timeout_ms) { - pNDI_find = NDIlib_find_create_v2(); - if (!pNDI_find) { - throw std::runtime_error("Cannot create NDI finder"); - } - - const int max_attempts = 10; - uint32_t sources_amount = 0; - const NDIlib_source_t* p_sources = nullptr; - for (int attempt = 0; !sources_amount && attempt < max_attempts; ++attempt) { - if (!dmf::g_running.load(std::memory_order_relaxed)) { - NDIlib_find_destroy(pNDI_find); - pNDI_find = nullptr; - throw std::runtime_error("Interrupted while searching for NDI sources"); - } - NDIlib_find_wait_for_sources(pNDI_find, timeout_ms); - p_sources = NDIlib_find_get_current_sources(pNDI_find, &sources_amount); - } - if (!sources_amount) { - NDIlib_find_destroy(pNDI_find); - pNDI_find = nullptr; - throw std::runtime_error("No NDI sources found after timeout"); - } - - // Copy while finder is alive: p_sources points into finder-owned memory - for (uint32_t i = 0; i < sources_amount; ++i) { - sources->push_back(p_sources[i].p_ndi_name); - cached_names.emplace_back(p_sources[i].p_ndi_name); - cached_urls.emplace_back(p_sources[i].p_url_address); - } - cached_sources.reserve(cached_names.size()); - for (uint32_t i = 0; i < cached_names.size(); ++i) { - cached_sources.push_back({cached_names[i].c_str(), cached_urls[i].c_str()}); - } - - NDIlib_find_destroy(pNDI_find); - pNDI_find = nullptr; - } - - void select_source(uint32_t source_num) { - if (cached_sources.empty()) { - throw std::runtime_error("0 sources found"); - } else if (source_num >= cached_sources.size()) { - throw std::runtime_error("Source_num bigger that sources amount"); - } - - pNDI_recv = NDIlib_recv_create_v3(); - if (!pNDI_recv) { - throw std::runtime_error("Cannot create NDI receive instance"); - } - NDIlib_recv_connect(pNDI_recv, &cached_sources[source_num]); - } - - void get_source_info(uint32_t source_num) { - NDIlib_video_frame_v2_t video_frame; - NDIlib_frame_type_e frame_type; - bool is_got_info = false; - while(!is_got_info) - { - frame_type = NDIlib_recv_capture_v3(pNDI_recv, &video_frame, nullptr, nullptr, 1000); - switch(frame_type) - { - case NDIlib_frame_type_video: - is_got_info = true; - xres = video_frame.xres; - yres = video_frame.yres; - frame_D = video_frame.frame_rate_D; - frame_N = video_frame.frame_rate_N; - fourCC = video_frame.FourCC; - stride = video_frame.line_stride_in_bytes; - if (stride == 0) { - stride = xres * get_bytes_per_pixel(fourCC); - } - break; - case NDIlib_frame_type_error: - is_got_info = true; - throw std::runtime_error("Selected NDI source is lost"); - break; - } - } - NDIlib_recv_free_video_v2(pNDI_recv, &video_frame); - } - - int get_bytes_per_pixel(NDIlib_FourCC_video_type_e fourCC) { - switch (fourCC) { - case NDIlib_FourCC_video_type_UYVY: // Standard 8-bit YUV 4:2:2 - case NDIlib_FourCC_video_type_YV12: // 8-bit YUV 4:2:0 - case NDIlib_FourCC_video_type_I420: // 8-bit YUV 4:2:0 - case NDIlib_FourCC_video_type_NV12: // 8-bit YUV 4:2:0 - // These are 4:2:2 or 4:2:0 formats. - // On average, they use 2 bytes (16 bits) per pixel across the macroblock. - return 2; - - case NDIlib_FourCC_video_type_BGRA: // 8-bit RGB with Alpha - case NDIlib_FourCC_video_type_RGBA: // 8-bit RGB with Alpha - // 4 channels (Red, Green, Blue, Alpha) * 1 byte each - return 4; - - case NDIlib_FourCC_video_type_BGRX: // 8-bit RGB (Padding) - case NDIlib_FourCC_video_type_RGBX: // 8-bit RGB (Padding) - // 4 channels (Red, Green, Blue, Empty) * 1 byte each - return 4; - - case NDIlib_FourCC_video_type_UYVA: // 8-bit YUV 4:2:2 + Alpha channel - // 2 bytes for YUV + 1 byte for Alpha split - return 3; - - case NDIlib_FourCC_video_type_P216: // 16-bit YUV 4:2:2 (High bit depth) - // 2 channels packed at 2 bytes (16-bits) per sample = 4 bytes per pixel - return 4; - - case NDIlib_FourCC_video_type_PA16: // 16-bit YUV 4:2:2 + 16-bit Alpha - return 6; - - default: - return 2; // Safe NDI default fallback - } - } - - std::string fourCCtoStr() { - char fourcc_str[5]; - uint32_t fourcc = (uint32_t)fourCC; - fourcc_str[0] = (fourcc >> 0) & 0xFF; - fourcc_str[1] = (fourcc >> 8) & 0xFF; - fourcc_str[2] = (fourcc >> 16) & 0xFF; - fourcc_str[3] = (fourcc >> 24) & 0xFF; - fourcc_str[4] = '\0'; - return std::string(fourcc_str); - } - - bool getV210_video_frame(uint32_t source_num, uint8_t* frame_buffer, uint32_t frame_stride) { - NDIlib_video_frame_v2_t video_frame; - NDIlib_frame_type_e frame_type; - frame_type = NDIlib_recv_capture_v3(pNDI_recv, &video_frame, nullptr, nullptr, 5); - - if (frame_type == NDIlib_frame_type_error) - throw std::runtime_error("NDI source lost"); - - if (frame_type == NDIlib_frame_type_status_change) { - // Source changed resolution or framerate — refresh internal info, repeat last frame - get_source_info(source_num); - return false; - } - - if (frame_type != NDIlib_frame_type_video) - return false; - - switch (fourCC) { - case NDIlib_FourCC_type_UYVY: - v210::UYVYtoV210(video_frame.p_data, frame_buffer, xres, yres, stride, frame_stride); - break; - case NDIlib_FourCC_type_P216: { - // P216→V210: write directly into caller's buffer via a wrapper frame - NDIlib_video_frame_v2_t dst{}; - dst.p_data = frame_buffer; - dst.line_stride_in_bytes = frame_stride; - NDIlib_util_P216_to_V210(&video_frame, &dst); - break; - } - default: - NDIlib_recv_free_video_v2(pNDI_recv, &video_frame); - throw std::runtime_error("Unsupported NDI color format: " + fourCCtoStr()); - } - NDIlib_recv_free_video_v2(pNDI_recv, &video_frame); - return true; - } - - private: - // receive - NDIlib_find_instance_t pNDI_find = nullptr; - NDIlib_recv_instance_t pNDI_recv = nullptr; - NDIlib_FourCC_type_e fourCC; - // owned copies so finder can be destroyed early - std::vector cached_names; - std::vector cached_urls; - std::vector cached_sources; -}; -} \ No newline at end of file diff --git a/shared/NDIReceiver.hpp b/shared/NDIReceiver.hpp new file mode 100644 index 0000000..aea4654 --- /dev/null +++ b/shared/NDIReceiver.hpp @@ -0,0 +1,191 @@ +#pragma once + +#include +#include +#include +#include +#include "Signal.hpp" +#include "V210.hpp" + +namespace dmf { + +class NDIReceiver { +public: + struct SourceInfo { + int width = 0; + int height = 0; + int fps_num = 0; + int fps_den = 0; + int stride = 0; + NDIlib_FourCC_type_e fourcc{}; + }; + + NDIReceiver() { + if (!NDIlib_is_supported_CPU()) + throw std::runtime_error("CPU is not sufficient for NDI"); + if (!NDIlib_initialize()) + throw std::runtime_error("NDI lib init failed"); + } + + ~NDIReceiver() { + if (recv_) NDIlib_recv_destroy(recv_); + if (find_) NDIlib_find_destroy(find_); + NDIlib_destroy(); + } + + // Discovers NDI sources on the network. Polls up to max_attempts times, + // each waiting timeout_ms milliseconds. Respects g_running. + std::vector find_sources(uint32_t timeout_ms, int max_attempts = 10) { + find_ = NDIlib_find_create_v2(); + if (!find_) + throw std::runtime_error("Cannot create NDI finder"); + + uint32_t count = 0; + const NDIlib_source_t* p_sources = nullptr; + for (int i = 0; !count && i < max_attempts; ++i) { + if (!dmf::g_running.load(std::memory_order_relaxed)) { + NDIlib_find_destroy(find_); + find_ = nullptr; + throw std::runtime_error("Interrupted while searching for NDI sources"); + } + NDIlib_find_wait_for_sources(find_, timeout_ms); + p_sources = NDIlib_find_get_current_sources(find_, &count); + } + if (!count) { + NDIlib_find_destroy(find_); + find_ = nullptr; + throw std::runtime_error("No NDI sources found after timeout"); + } + + // Copy while finder is alive: p_sources points into finder-owned memory + std::vector names; + for (uint32_t i = 0; i < count; ++i) { + names.emplace_back(p_sources[i].p_ndi_name); + source_names_.emplace_back(p_sources[i].p_ndi_name); + source_urls_.emplace_back(p_sources[i].p_url_address); + } + sources_.reserve(source_names_.size()); + for (size_t i = 0; i < source_names_.size(); ++i) + sources_.push_back({source_names_[i].c_str(), source_urls_[i].c_str()}); + + NDIlib_find_destroy(find_); + find_ = nullptr; + return names; + } + + // Connects to a discovered source by index. + void connect(uint32_t source_num) { + if (sources_.empty()) + throw std::runtime_error("No sources available — call find_sources() first"); + if (source_num >= static_cast(sources_.size())) + throw std::runtime_error("source_num exceeds available source count"); + + recv_ = NDIlib_recv_create_v3(); + if (!recv_) + throw std::runtime_error("Cannot create NDI receive instance"); + NDIlib_recv_connect(recv_, &sources_[source_num]); + } + + // Captures the first video frame to determine resolution, frame rate, and format. + // Stores the result internally for use by capture_v210(). Respects g_running. + SourceInfo probe() { + while (dmf::g_running.load(std::memory_order_relaxed)) { + NDIlib_video_frame_v2_t frame; + auto type = NDIlib_recv_capture_v3(recv_, &frame, nullptr, nullptr, 1000); + if (type == NDIlib_frame_type_error) + throw std::runtime_error("NDI source lost during probe"); + if (type != NDIlib_frame_type_video) + continue; + + info_.width = frame.xres; + info_.height = frame.yres; + info_.fps_num = frame.frame_rate_N; + info_.fps_den = frame.frame_rate_D; + info_.fourcc = static_cast(frame.FourCC); + info_.stride = frame.line_stride_in_bytes; + if (info_.stride == 0) + info_.stride = info_.width * bytes_per_pixel(info_.fourcc); + + NDIlib_recv_free_video_v2(recv_, &frame); + return info_; + } + throw std::runtime_error("Interrupted during probe"); + } + + // Captures one video frame and converts it to V210 in frame_buffer. + // Returns false if no frame was available this tick (caller should repeat last frame). + // Throws on source lost or unsupported format. + bool capture_v210(uint8_t* frame_buffer, uint32_t frame_stride) { + NDIlib_video_frame_v2_t frame; + auto type = NDIlib_recv_capture_v3(recv_, &frame, nullptr, nullptr, 5); + + if (type == NDIlib_frame_type_error) + throw std::runtime_error("NDI source lost"); + + if (type == NDIlib_frame_type_status_change) { + info_ = probe(); // source changed resolution or framerate — refresh + return false; + } + + if (type != NDIlib_frame_type_video) + return false; + + switch (info_.fourcc) { + case NDIlib_FourCC_type_UYVY: + v210::UYVYtoV210(frame.p_data, frame_buffer, + info_.width, info_.height, info_.stride, frame_stride); + break; + case NDIlib_FourCC_type_P216: { + NDIlib_video_frame_v2_t dst{}; + dst.p_data = frame_buffer; + dst.line_stride_in_bytes = frame_stride; + NDIlib_util_P216_to_V210(&frame, &dst); + break; + } + default: + NDIlib_recv_free_video_v2(recv_, &frame); + throw std::runtime_error("Unsupported NDI color format: " + fourcc_str(info_.fourcc)); + } + NDIlib_recv_free_video_v2(recv_, &frame); + return true; + } + +private: + NDIlib_find_instance_t find_ = nullptr; + NDIlib_recv_instance_t recv_ = nullptr; + SourceInfo info_; + std::vector source_names_; + std::vector source_urls_; + std::vector sources_; + + static int bytes_per_pixel(NDIlib_FourCC_type_e fc) { + switch (fc) { + case NDIlib_FourCC_video_type_UYVY: + case NDIlib_FourCC_video_type_YV12: + case NDIlib_FourCC_video_type_I420: + case NDIlib_FourCC_video_type_NV12: return 2; + case NDIlib_FourCC_video_type_BGRA: + case NDIlib_FourCC_video_type_RGBA: + case NDIlib_FourCC_video_type_BGRX: + case NDIlib_FourCC_video_type_RGBX: return 4; + case NDIlib_FourCC_video_type_UYVA: return 3; + case NDIlib_FourCC_video_type_P216: return 4; + case NDIlib_FourCC_video_type_PA16: return 6; + default: return 2; + } + } + + static std::string fourcc_str(NDIlib_FourCC_type_e fc) { + uint32_t v = static_cast(fc); + char s[5] = { + static_cast((v >> 0) & 0xFF), + static_cast((v >> 8) & 0xFF), + static_cast((v >> 16) & 0xFF), + static_cast((v >> 24) & 0xFF), + '\0' + }; + return s; + } +}; + +} // namespace dmf