#include "NodeBase.hpp" #include "DeckLinkReceiver.hpp" #include "Signal.hpp" #include "FlowDef.hpp" #include #include #include #include #include #include #include #include class NodeDeckLinkIn: public dmf::NodeBase { void run() override { dmf::DeckLinkReceiver decklink_receiver; log("Available DeckLink input devices:"); for (auto device : decklink_receiver.devices_list) { log("%i) %s", device.index, device.display_name.c_str()); } int selected_device = 0; decklink_receiver.start_capture(selected_device); log( "DeckLink feed for device '%s':", decklink_receiver.devices_list.at(selected_device).display_name.c_str() ); dmf::SourceInfo video_source_info{}; if (!decklink_receiver.wait_for_format(5000)) { log("Timeout waiting for format detection"); return; } video_source_info = decklink_receiver.get_input_callback()->video_info; if (video_source_info.width == 0 || video_source_info.fps_num == 0) { log("Invalid format detected"); return; } const auto video_flow_info = config().at("video_flow_id"); const auto video_flow_id = video_flow_info.at("id").get(); const int width = video_flow_info.value("width", video_source_info.width); const int height = video_flow_info.value("height", video_source_info.height); const int fps_num = video_flow_info.value("fps_num", video_source_info.fps_num); const int fps_den = video_flow_info.value("fps_den", video_source_info.fps_den); log("video flow=%s %dx%d @ %d/%d fps", video_flow_id.c_str(), width, height, fps_num, fps_den); mxlFlowWriter video_writer{}; mxlFlowConfigInfo video_cfg{}; bool created = false; mxlStatus vst = mxlCreateFlowWriter( instance(), dmf::make_video_flow_def(video_flow_id, node_id(), width, height, fps_num, fps_den).c_str(), "", &video_writer, &video_cfg, &created); if (vst != MXL_STATUS_OK) { log("video mxlCreateFlowWriter failed (%s)", dmf::mxl_status_str(vst)); return; } const uint32_t video_stride = video_cfg.discrete.sliceSizes[0]; log("video stride=%u B/line grain=%u B ring=%u grains", video_stride, video_stride * static_cast(height), video_cfg.discrete.grainCount); const mxlRational video_rate = {fps_num, fps_den}; uint64_t video_index = mxlGetCurrentIndex(&video_rate); log("start video_index=%lu", static_cast(video_index)); auto* cb = decklink_receiver.get_input_callback(); const uint32_t dst_row = video_stride; std::vector local_frame(dst_row * height, 0); while(dmf::g_running.load(std::memory_order_relaxed)) { // Wait for a fresh frame from the callback uint32_t src_row; int src_w, src_h; { std::unique_lock lk(cb->frame_mutex); cb->frame_cv.wait(lk, [&] { return cb->frame_ready.load(std::memory_order_acquire) || !dmf::g_running.load(std::memory_order_relaxed); }); if (!dmf::g_running.load(std::memory_order_relaxed)) break; src_row = cb->frame_row_bytes ? cb->frame_row_bytes : dst_row; src_w = cb->frame_width; src_h = cb->frame_height; // Row-by-row copy with zero padding for stride difference const uint8_t* src = cb->frame_buffer.data(); uint8_t* dst = local_frame.data(); const uint32_t copy_row = std::min(src_row, dst_row); const int rows = std::min(src_h, height); std::memset(local_frame.data(), 0, local_frame.size()); for (int y = 0; y < rows; ++y) { std::memcpy(dst, src, copy_row); src += src_row; dst += dst_row; } cb->frame_ready.store(false, std::memory_order_release); } if (src_w != width || src_h != height) { log("frame dim mismatch: src=%dx%d mxl=%dx%d (skipping)", src_w, src_h, width, height); continue; } // MXL write — no lock held, callback can fill next frame in parallel mxlGrainInfo grain{}; uint8_t* video_buf = nullptr; vst = mxlFlowWriterOpenGrain(video_writer, video_index, &grain, &video_buf); if (vst == MXL_STATUS_OK) { std::memcpy(video_buf, local_frame.data(), local_frame.size()); grain.flags = 0; grain.validSlices = grain.totalSlices; mxlFlowWriterCommitGrain(video_writer, &grain); const uint64_t ns = mxlGetNsUntilIndex(video_index + 1, &video_rate); if (ns > 0 && ns < 2'000'000'000ULL) mxlSleepForNs(ns); video_index++; } } } }; int main () { NodeDeckLinkIn node; node.execute(); return 0; }