#include #include #include #include #include "NodeBase.hpp" #include "FlowDef.hpp" #include "DeckLinkReceiver.hpp" class DeckLinkInNode : public dmf::NodeBase { void run() override { const uint32_t device_index = config().value("device_index", 0u); const bool want_audio = config().contains("audio_flow_id"); const int channels = want_audio ? config().at("audio_flow_id").value("channels", 2) : 0; dmf::DeckLinkReceiver receiver; try { log("Available DeckLink devices:"); for (const auto& d : receiver.devices) log(" %u) %s", d.index, d.name.c_str()); receiver.start_capture(device_index, channels); log("Capturing from: %s", receiver.devices[device_index].name.c_str()); } catch (const std::runtime_error& e) { log("DeckLink init error: %s", e.what()); return; } if (!receiver.wait_for_format(5000)) { log("Timeout waiting for format detection"); return; } const auto& vi = receiver.video_info; if (vi.width == 0 || vi.fps_num == 0) { log("Invalid format detected"); return; } log("Detected: %dx%d @ %d/%d fps", vi.width, vi.height, vi.fps_num, vi.fps_den); // --- Video writer --- if (!config().contains("video_flow_id")) { log("no video output connected — exiting"); 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 = vi.width; const int height = vi.height; const int fps_num = vi.fps_num; const int fps_den = vi.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 = nullptr; 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("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); // --- Audio writer (optional) --- mxlFlowWriter audio_writer = nullptr; int max_audio_samples = 0; const bool has_audio = want_audio && receiver.has_audio; if (has_audio) { const auto audio_flow_info = config().at("audio_flow_id"); const auto audio_flow_id = audio_flow_info.at("id").get(); const int sample_rate = receiver.audio_info.sample_rate; const int bit_depth = 32; log("audio flow=%s %d Hz %dch %d-bit", audio_flow_id.c_str(), sample_rate, channels, bit_depth); mxlFlowConfigInfo audio_cfg{}; mxlStatus ast = mxlCreateFlowWriter( instance(), dmf::make_audio_flow_def(audio_flow_id, node_id(), sample_rate, channels, bit_depth, fps_num, fps_den).c_str(), "", &audio_writer, &audio_cfg, &created); if (ast != MXL_STATUS_OK) { log("audio mxlCreateFlowWriter failed (%s) — continuing without audio", dmf::mxl_status_str(ast)); } else { log("audio channels=%u buffer=%u samples", audio_cfg.continuous.channelCount, audio_cfg.continuous.bufferLength); size_t max_write = 0; mxlFlowWriterGetMaxWriteLengthSamples(audio_writer, &max_write); max_audio_samples = static_cast(max_write); } } // --- Buffers --- std::vector frame_buf(static_cast(video_stride) * static_cast(height)); // Planar float32: channel c at audio_buf[c * max_audio_samples] std::vector audio_buf(static_cast(max_audio_samples) * static_cast(channels)); // --- Clock --- const mxlRational video_rate = {fps_num, fps_den}; const mxlRational audio_rate = {receiver.audio_info.sample_rate, 1}; uint64_t video_index = mxlGetCurrentIndex(&video_rate); uint64_t audio_index = has_audio ? mxlGetCurrentIndex(&audio_rate) : 0; log("start video_index=%llu", static_cast(video_index)); // --- Capture loop --- while (dmf::g_running.load(std::memory_order_relaxed)) { int samples_written = 0; // DeckLink delivers one video frame + accompanying audio per callback. if (!receiver.wait_for_frame( frame_buf.data(), video_stride, width, height, (has_audio && audio_writer) ? audio_buf.data() : nullptr, max_audio_samples, (has_audio && audio_writer) ? &samples_written : nullptr)) break; // Video grain mxlGrainInfo grain{}; uint8_t* video_buf_ptr = nullptr; vst = mxlFlowWriterOpenGrain(video_writer, video_index, &grain, &video_buf_ptr); if (vst == MXL_STATUS_OK) { std::memcpy(video_buf_ptr, frame_buf.data(), frame_buf.size()); grain.flags = 0; grain.validSlices = grain.totalSlices; mxlFlowWriterCommitGrain(video_writer, &grain); } // Audio samples (same fragment-wrap pattern as videoin) if (has_audio && audio_writer && samples_written > 0) { mxlMutableWrappedMultiBufferSlice slice{}; mxlStatus ast = mxlFlowWriterOpenSamples( audio_writer, audio_index, static_cast(samples_written), &slice); if (ast == MXL_STATUS_OK) { for (int ch = 0; ch < channels; ++ch) { const uint8_t* src = reinterpret_cast( audio_buf.data() + ch * max_audio_samples); uint8_t* dst0 = static_cast( slice.base.fragments[0].pointer) + ch * slice.stride; const size_t frag0_bytes = slice.base.fragments[0].size; const size_t total_bytes = static_cast(samples_written) * sizeof(float); if (total_bytes <= frag0_bytes) { std::memcpy(dst0, src, total_bytes); } else { std::memcpy(dst0, src, frag0_bytes); uint8_t* dst1 = static_cast( slice.base.fragments[1].pointer) + ch * slice.stride; std::memcpy(dst1, src + frag0_bytes, total_bytes - frag0_bytes); } } mxlFlowWriterCommitSamples(audio_writer); } else { log("audio OpenSamples failed (%s) index=%llu — skipping", dmf::mxl_status_str(ast), static_cast(audio_index)); } audio_index += static_cast(samples_written); } // Pace video to the MXL clock const uint64_t ns = mxlGetNsUntilIndex(video_index + 1, &video_rate); if (ns > 0 && ns < 2'000'000'000ULL) mxlSleepForNs(ns); video_index = mxlGetCurrentIndex(&video_rate); } log("stopped at video_index=%llu", static_cast(video_index)); mxlReleaseFlowWriter(instance(), video_writer); if (audio_writer) mxlReleaseFlowWriter(instance(), audio_writer); } }; int main() { DeckLinkInNode node; return node.execute(); }