diff --git a/PLAN.md b/PLAN.md new file mode 100644 index 0000000..b8bebac --- /dev/null +++ b/PLAN.md @@ -0,0 +1,122 @@ +# DMF Studio — Roadmap + +## Next steps (in order) + +### 1. WebSocket API in studio-manager +Allow the graph to be changed at runtime without restarting. +- Add a WebSocket server to `studio-manager` +- API: load/reload graph, start/stop individual nodes, query status +- `studio-manager` already has `load_graph()` — the WS layer calls it on demand + and diffs against the running set (stop removed nodes, fork new ones) +- **Required before**: frontend, live source switching, PiP (otherwise every + graph change is a full restart) + +### 2. Processing nodes — PiP / mixer +First node that takes multiple input flows and produces an output flow. +- Uses `mxlFlowSynchronizationGroup` to align grains from two inputs +- Reference implementation: `nodes/testpattern` (writer) + `nodes/fakesink` (reader) +- Only becomes useful with the WebSocket API (so you can switch sources live) + +### 3. Vue.js frontend +Visual graph editor that drives the WebSocket API. + +--- + +## Redundancy + +Key constraint: **one writer per MXL flow** — can't run two identical nodes writing +the same flow simultaneously. Redundancy lives at the pipeline level, not the node level. + +### Dual pipeline on separate machines + +``` +Machine 1 (k8s node A) Machine 2 (k8s node B) + decklinkin → [MXL] → ndiout decklinkin → [MXL] → ndiout + ↓ ↓ + (primary path) (backup path) + \ / + └──────→ [selector node] ←──────────┘ + ↓ + ndiout (final) +``` + +**Selector node** — reads two input flows, monitors grain validity flags, switches +to backup when primary fails. Fits the existing node model; uses +`mxlFlowSynchronizationGroup` to watch both flows. Key processing node to build +once redundancy becomes a requirement. + +MXL shared memory requires all pods in a pipeline to be co-located on the same +physical machine. Redundant pipelines naturally go on *different* machines — which +is exactly right for hardware failure redundancy. + +### What k8s gives for free + +- Stateless processing nodes (PiP, denoise, format convert): k8s restarts on crash, + ~1-2 s gap — acceptable for non-critical path +- `PodDisruptionBudget`: ensures critical nodes survive cluster maintenance +- Leader election (k8s lease objects): two studio-managers, one active, one standby; + automatic failover with no node code changes + +--- + +## Kubernetes integration (mxl-k8s) + +Source: `~/codeproj/mxl-k8s` — not official, treat as reference, not truth. + +mxl-k8s is a full k8s control plane for MXL flows. Four runtime pieces: + +- **Operator** (Deployment): watches `MxlReceiver` CRDs, creates `MxlFlowMirror` per target node +- **Agent** (DaemonSet): watches each node's MXL domain via `fanotify`, publishes `MxlFlow` CRDs with where flows live +- **Gateway** (DaemonSet, `hostNetwork`): drives libmxl-fabrics RDMA/TCP between nodes — zero-copy grain transfer via registered mmap regions +- **Shim** (`libmxl-intent.so`, LD_PRELOAD): intercepts `openat`/`stat`/`access` on `.mxl-flow/` paths in consumer pods; when a flow isn't local, asks the agent's UDS socket (`/run/mxl/agent.sock`) to materialize it via mirror, then retries — transparent to node code + +### What changes for our nodes in k8s + +**Producer pods** (decklinkin, ndiin, testpattern): **zero code change**. +- Add `hostPath: /run/mxl/domain` volume + `IPC_LOCK`, `SYS_RESOURCE` capabilities +- `NODE_CONFIG` → Pod env var (from ConfigMap) +- `MXL_DOMAIN` → `/run/mxl/domain` (standardized in k8s context) + +**Consumer pods, same node**: same as above, no code change. + +**Consumer pods, different node**: still no code change. +- Add `initContainer` copying `libmxl-intent.so` from shim image +- Set `LD_PRELOAD=/opt/mxl-intent/libmxl-intent.so` +- Mount `/run/mxl` (whole dir, not just `/domain`) so agent socket is accessible +- Create an `MxlReceiver` CRD pointing at the flow — operator handles the mirror plumbing + +### What studio-manager becomes in k8s + +Currently: fork/exec child processes. In k8s: apply/delete Pods (or Deployments) with `NODE_CONFIG` env vars. For cross-node flows: create `MxlReceiver` CRDs instead of wiring flows manually. + +Same-node pipeline: all pods get `nodeAffinity: requiredDuringScheduling → same host`. +Cross-node: add shim + `MxlReceiver`; mxl-k8s handles the rest. + +--- + +## Architecture decisions + +### Separate audio and video threads in nodes + +**Decision**: sink nodes (`decklinkout`, `ndiout`) and likely source nodes should +process audio and video on separate threads. + +**Why**: audio and video have different timing granularities. +- Video: one grain every ~40 ms (at 25 fps) — coarse, can block +- Audio: must flow continuously at sample granularity — any stall causes dropout + +In the current single-thread model, video stalls (e.g. `TOO_EARLY` retries) +pause audio too. In `decklinkout` this is especially bad — DeckLink's timestamped +audio buffer underruns if it isn't fed consistently. + +**What it looks like**: +- **Audio thread**: tight loop, continuously drains MXL audio ring buffer and + pushes to output (DeckLink `ScheduleAudioSamples` / NDI send). No video logic. +- **Video thread**: current main loop, handles grain read → process → output + at frame rate. +- **Shared state**: only `g_running` and the output handle. No frame data crosses + the boundary — each thread reads its own MXL flow independently. + MXL clock keeps them in sync without explicit A/V coordination. + +**When**: after the WebSocket API, when running real content and audio quality matters. +Current single-thread model is acceptable for development. diff --git a/ndi-ndi.json b/ndi-ndi.json index 037af62..ba182ec 100644 --- a/ndi-ndi.json +++ b/ndi-ndi.json @@ -1,7 +1,7 @@ { "nodes": [ { "id": "ndiin", "type": "ndiin", "params": {} }, - { "id": "ndiout", "type": "ndiout", "params": { "device_index": 0 } } + { "id": "ndiout", "type": "ndiout", "params": {} } ], "edges": [ { diff --git a/nodes/decklinkout/main.cpp b/nodes/decklinkout/main.cpp index 7c17eef..7e02987 100644 --- a/nodes/decklinkout/main.cpp +++ b/nodes/decklinkout/main.cpp @@ -1,38 +1,34 @@ -#include -#include "DeckLinkSender.hpp" - +#include +#include +#include +#include #include #include -#include -#include -#include -#include "Signal.hpp" #include "NodeBase.hpp" +#include "DeckLinkSender.hpp" -namespace dmf { class DeckLinkOutNode : public dmf::NodeBase { void run() override { - // --- common params (optional) --- - const uint32_t device_index = config().value("device_index", 0u); + const uint32_t device_index = config().value("device_index", 0u); + // --- video flow (optional) --- bool has_video = config().contains("video_flow_id"); - int width = 1920; - int height = 1080; - int fps_num = 25; - int fps_den = 1; + int width = 1920; + int height = 1080; + int fps_num = 25; + int fps_den = 1; std::string flow_id; - mxlFlowReader video_reader{}; - uint32_t video_stride = 0; + mxlFlowReader video_reader{}; + uint32_t video_stride = 0; if (has_video) { const auto flow_info = config().at("video_flow_id"); - flow_id = flow_info.at("id").get(); - width = flow_info.value("width", 1920); - height = flow_info.value("height", 1080); - fps_num = flow_info.value("fps_num", 25); - fps_den = flow_info.value("fps_den", 1); - + flow_id = flow_info.at("id").get(); + width = flow_info.value("width", 1920); + height = flow_info.value("height", 1080); + fps_num = flow_info.value("fps_num", 25); + fps_den = flow_info.value("fps_den", 1); log("video flow=%s %dx%d @ %d/%d fps", flow_id.c_str(), width, height, fps_num, fps_den); log("waiting for flow to become active..."); @@ -46,15 +42,17 @@ class DeckLinkOutNode : public dmf::NodeBase { mxlFlowConfigInfo video_cfg{}; mxlStatus vst = mxlCreateFlowReader(instance(), flow_id.c_str(), "", &video_reader); - if (vst != MXL_STATUS_OK) { log("video mxlCreateFlowReader failed (%s)", dmf::mxl_status_str(vst)); return; } + if (vst != MXL_STATUS_OK) { + log("video mxlCreateFlowReader failed (%s)", dmf::mxl_status_str(vst)); + return; + } mxlFlowReaderGetConfigInfo(video_reader, &video_cfg); video_stride = video_cfg.discrete.sliceSizes[0]; } // --- audio flow (optional) --- - mxlFlowReader audio_reader{}; - mxlFlowConfigInfo audio_cfg{}; - int sample_rate = 0; + mxlFlowReader audio_reader{}; + int sample_rate = 48000; int channels = 0; int samples_per_frame = 0; bool has_audio = config().contains("audio_flow_id"); @@ -62,19 +60,20 @@ class DeckLinkOutNode : public dmf::NodeBase { if (has_audio) { const auto audio_flow_info = config().at("audio_flow_id"); - audio_flow_id = audio_flow_info.at("id").get(); + audio_flow_id = audio_flow_info.at("id").get(); sample_rate = audio_flow_info.value("sample_rate", 48000); channels = audio_flow_info.value("channels", 2); samples_per_frame = sample_rate * fps_den / fps_num; - log("audio flow=%s %d Hz %dch %d samples/frame", audio_flow_id.c_str(), sample_rate, channels, samples_per_frame); mxlStatus ast = mxlCreateFlowReader(instance(), audio_flow_id.c_str(), "", &audio_reader); if (ast != MXL_STATUS_OK) { - log("audio mxlCreateFlowReader failed (%s) — continuing without audio", dmf::mxl_status_str(ast)); + log("audio mxlCreateFlowReader failed (%s) — continuing without audio", + dmf::mxl_status_str(ast)); has_audio = false; } else { + mxlFlowConfigInfo audio_cfg{}; mxlFlowReaderGetConfigInfo(audio_reader, &audio_cfg); log("audio channels=%u buffer=%u samples", audio_cfg.continuous.channelCount, audio_cfg.continuous.bufferLength); @@ -85,25 +84,24 @@ class DeckLinkOutNode : public dmf::NodeBase { dmf::DeckLinkSender sender; try { - log("Available DeckLink output devices:\n"); - for (const auto& d : sender.devices) { + log("Available DeckLink output devices:"); + for (const auto& d : sender.devices) log(" %u) %s", d.index, d.name.c_str()); - } sender.start_output(device_index, width, height, fps_num, fps_den, channels); } catch (const std::runtime_error& e) { log("DeckLink init error: %s", e.what()); + if (has_video && video_reader) mxlReleaseFlowReader(instance(), video_reader); + if (has_audio && audio_reader) mxlReleaseFlowReader(instance(), audio_reader); return; } - // --- main loop --- - uint64_t video_index = 0; - uint64_t audio_index = 0; - 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; + // Pre-allocate audio staging buffer (planar float32) + std::vector audio_planar( + static_cast(channels) * static_cast(samples_per_frame)); + // --- Clock init --- + uint64_t video_index = 0; + uint64_t audio_index = 0; if (has_video) { const mxlRational video_rate = {fps_num, fps_den}; video_index = mxlGetCurrentIndex(&video_rate); @@ -113,19 +111,21 @@ class DeckLinkOutNode : public dmf::NodeBase { audio_index = mxlGetCurrentIndex(&audio_rate); } - while (dmf::g_running.load(std::memory_order_relaxed)) { - // --- audio: non-blocking, one chunk per video frame (or free-running) --- - bool audio_advanced = false; + 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; + // --- Main loop --- + while (dmf::g_running.load(std::memory_order_relaxed)) { + // Audio: non-blocking, one chunk per video frame + bool audio_advanced = false; if (has_audio) { mxlWrappedMultiBufferSlice audio_slices{}; mxlStatus ast = mxlFlowReaderGetSamplesNonBlocking( audio_reader, audio_index, samples_per_frame, &audio_slices); if (ast == MXL_STATUS_OK) { - // Extract planar float32 samples from MXL ring buffer - // (per-channel: ring buffer can wrap, so 2 fragments) - std::vector audio_planar( - static_cast(channels) * samples_per_frame); const size_t frag0 = audio_slices.base.fragments[0].size / sizeof(float); const size_t frag1 = audio_slices.base.fragments[1].size / sizeof(float); for (int c = 0; c < channels; ++c) { @@ -142,7 +142,7 @@ class DeckLinkOutNode : public dmf::NodeBase { } } sender.submit_audio(audio_planar.data(), samples_per_frame); - audio_index += samples_per_frame; + audio_index += samples_per_frame; audio_advanced = true; } else if (ast == MXL_ERR_OUT_OF_RANGE_TOO_LATE) { mxlFlowRuntimeInfo ari{}; @@ -161,7 +161,7 @@ class DeckLinkOutNode : public dmf::NodeBase { } } - // --- video --- + // Video if (has_video) { mxlGrainInfo video_grain{}; uint8_t* video_buf = nullptr; @@ -191,7 +191,8 @@ class DeckLinkOutNode : public dmf::NodeBase { video_index = mxlGetCurrentIndex(&r); } } else { - log("unexpected video status (%s) on index=%llu", dmf::mxl_status_str(vst), video_index); + log("unexpected video status (%s) on index=%llu", + dmf::mxl_status_str(vst), static_cast(video_index)); break; } @@ -199,24 +200,31 @@ class DeckLinkOutNode : public dmf::NodeBase { 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), + static_cast(invalid_count), + static_cast(late_count), static_cast(frame_count) / elapsed); last_log_time = now; } } else if (!audio_advanced) { - // audio-only and nothing was ready — avoid busy spin mxlSleepForNs(1'000'000); } } + + if (has_video) + log("stopped — total frames=%llu invalid=%llu late=%llu", + static_cast(frame_count), + static_cast(invalid_count), + static_cast(late_count)); + else + log("stopped"); + if (has_video && video_reader) mxlReleaseFlowReader(instance(), video_reader); if (has_audio && audio_reader) mxlReleaseFlowReader(instance(), audio_reader); } }; -} -int main() -{ - dmf::DeckLinkOutNode node; - node.execute(); - return 0; -} \ No newline at end of file +int main() { + DeckLinkOutNode node; + return node.execute(); +} diff --git a/nodes/ndiin/main.cpp b/nodes/ndiin/main.cpp index 68a94b8..ab14cb4 100644 --- a/nodes/ndiin/main.cpp +++ b/nodes/ndiin/main.cpp @@ -147,11 +147,11 @@ class NDIInNode : public dmf::NodeBase { 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++; + video_index = mxlGetCurrentIndex(&video_rate); } else { log("video OpenGrain failed (%s) at index=%llu", dmf::mxl_status_str(st), video_index); + video_index = current + 1; } - video_index = current + 1; } } diff --git a/shared/DeckLinkSender.hpp b/shared/DeckLinkSender.hpp index 036e6cf..ad2d37a 100644 --- a/shared/DeckLinkSender.hpp +++ b/shared/DeckLinkSender.hpp @@ -1,11 +1,17 @@ #pragma once +#include #include +#include +#include +#include +#include #include #include #include #include +#include "Signal.hpp" #include "V210.hpp" namespace dmf { @@ -33,38 +39,30 @@ public: VideoInfo video_info{}; AudioInfo audio_info{}; bool has_audio = false; - int64_t frame_count = 0; - std::atomic audio_stream_time{0}; DeckLinkSender() { enumerate_devices(); } + ~DeckLinkSender() { if (decklink_output) { - decklink_output->StopScheduledPlayback(frame_count * video_info.fps_den, - nullptr, - video_info.fps_num); + decklink_output->StopScheduledPlayback(0, nullptr, 1); decklink_output->DisableVideoOutput(); - if (has_audio) decklink_output->DisableAudioOutput();; + if (has_audio) decklink_output->DisableAudioOutput(); decklink_output->SetScheduledFrameCompletionCallback(nullptr); decklink_output->Release(); - for (auto* vf : frame_pool) vf->Release(); - frame_pool.clear(); } delete output_callback; + for (auto* vf : all_frames) vf->Release(); for (auto* d : raw_devices) d->Release(); if (selected_device) selected_device->Release(); } // channels > 0 enables audio output at 48 kHz / 32-bit int. - void start_output(uint32_t device_index, int width, int height, int fps_num, int fps_den, int channels = 0) { + void start_output(uint32_t device_index, int width, int height, + int fps_num, int fps_den, int channels = 0) { if (device_index >= raw_devices.size()) throw std::runtime_error("Device index out of range"); - video_info.width = width; - video_info.height = height; - video_info.fps_num = fps_num; - video_info.fps_den = fps_den; - - BMDDisplayMode bm_display_mode = bmdModeHD1080p25; //need to create a func for detection + video_info = {width, height, fps_num, fps_den}; selected_device = raw_devices[device_index]; selected_device->AddRef(); @@ -76,7 +74,23 @@ public: r = decklink_output->SetScheduledFrameCompletionCallback(output_callback); if (r != S_OK) throw std::runtime_error("Could not set output callback"); - r = decklink_output->EnableVideoOutput(bmdModeHD1080p25, bmdVideoOutputFlagDefault); + const BMDDisplayMode mode = pick_display_mode(width, height, fps_num, fps_den); + + // Verify the card supports this mode in 10-bit YUV before enabling + bool is_supported = false; + BMDDisplayMode actual_mode = mode; + r = decklink_output->DoesSupportVideoMode( + bmdVideoConnectionUnspecified, + mode, + bmdFormat10BitYUV, + bmdNoVideoOutputConversion, + bmdSupportedVideoModeDefault, + &actual_mode, + &is_supported); + if (r != S_OK || !is_supported) + throw std::runtime_error("Display mode not supported in 10-bit YUV"); + + r = decklink_output->EnableVideoOutput(actual_mode, bmdVideoOutputFlagDefault); if (r != S_OK) throw std::runtime_error("Could not enable video output"); if (channels > 0) { @@ -85,36 +99,24 @@ public: static_cast(channels), bmdAudioOutputStreamTimestamped); if (r != S_OK) throw std::runtime_error("Could not enable audio output"); - // r = decklink_output->SetAudioCallback - has_audio = true; audio_info.channels = channels; + has_audio = true; } - - bool is_supported = false; - BMDDisplayMode actual_mode; - r = decklink_output->DoesSupportVideoMode( - bmdVideoConnectionUnspecified, // TODO: create selection between sdi, hdmi, etc. - bmdModeHD1080p25, - bmdFormat10BitYUV, - bmdNoVideoOutputConversion, - bmdSupportedVideoModeDefault, - &actual_mode, - &is_supported - ); - if (r != S_OK) throw std::runtime_error("Selected mode is not supported for bmdFormat10BitYUV"); - int32_t row_bytes; + int32_t row_bytes = 0; r = decklink_output->RowBytesForPixelFormat(bmdFormat10BitYUV, width, &row_bytes); - if (r != S_OK) throw std::runtime_error("Could not get row bytes for display mode"); + if (r != S_OK) throw std::runtime_error("Could not get row bytes for pixel format"); - // prefill frame pool - const int64_t duration = fps_den; - const int64_t timescale = fps_num; - const size_t preroll_pool_size = 3; // min=3, cause 1 displaying, 1 queued, 1 writable - for (size_t i = 0; i < preroll_pool_size; ++i) { + // Pre-fill 3 black frames and schedule them to prime the pipeline. + // The card fires ScheduledFrameCompleted once each is done, at which point + // submit_frame() can reclaim and reschedule them with live video. + constexpr size_t kPreroll = 3; + for (size_t i = 0; i < kPreroll; ++i) { IDeckLinkMutableVideoFrame* vf = nullptr; - r = decklink_output->CreateVideoFrame(width, height, row_bytes, bmdFormat10BitYUV, bmdFrameFlagDefault, &vf); - if (r != S_OK) throw std::runtime_error("Could not create a video frame"); + r = decklink_output->CreateVideoFrame( + width, height, row_bytes, bmdFormat10BitYUV, bmdFrameFlagDefault, &vf); + if (r != S_OK) throw std::runtime_error("Could not create video frame"); + IDeckLinkVideoBuffer* buf = nullptr; vf->QueryInterface(IID_IDeckLinkVideoBuffer, (void**)&buf); buf->StartAccess(bmdBufferAccessWrite); @@ -124,20 +126,30 @@ public: buf->EndAccess(bmdBufferAccessWrite); buf->Release(); - decklink_output->ScheduleVideoFrame(vf, i * duration, duration, timescale); - frame_pool.push_back(vf); + decklink_output->ScheduleVideoFrame( + vf, static_cast(i * fps_den), fps_den, fps_num); + all_frames.push_back(vf); } - output_callback->next_time.store(preroll_pool_size * duration); + scheduled_time = static_cast(kPreroll) * fps_den; - // start playback - r = decklink_output->StartScheduledPlayback(0, timescale, 1.0); - if (r != S_OK) throw std::runtime_error("Could not start streams"); + r = decklink_output->StartScheduledPlayback(0, fps_num, 1.0); + if (r != S_OK) throw std::runtime_error("Could not start scheduled playback"); } + // Blocks until a frame slot is free (returned by the card via callback), + // then copies src into it and schedules it for display. void submit_frame(const uint8_t* src, uint32_t stride) { - // Grab a free frame from the pool (round-robin) - IDeckLinkMutableVideoFrame* vf = frame_pool[frame_count % frame_pool.size()]; - frame_count++; + IDeckLinkMutableVideoFrame* vf = nullptr; + { + std::unique_lock lk(pool_mutex); + pool_cv.wait(lk, [this] { + return !free_frames.empty() || + !dmf::g_running.load(std::memory_order_relaxed); + }); + if (free_frames.empty()) return; + vf = free_frames.back(); + free_frames.pop_back(); + } IDeckLinkVideoBuffer* buf = nullptr; if (vf->QueryInterface(IID_IDeckLinkVideoBuffer, (void**)&buf) != S_OK) return; @@ -145,83 +157,108 @@ public: void* ptr = nullptr; buf->GetBytes(&ptr); if (ptr) { - // Row-by-row copy with stride adaptation - const uint32_t dst_stride = vf->GetRowBytes(); - const uint32_t copy_row = std::min(stride, dst_stride); - for (int y = 0; y < video_info.height; ++y) { + const uint32_t dst_stride = static_cast(vf->GetRowBytes()); + const uint32_t copy_stride = std::min(stride, dst_stride); + for (int y = 0; y < video_info.height; ++y) std::memcpy(static_cast(ptr) + y * dst_stride, - src + y * stride, copy_row); - } + src + y * stride, copy_stride); } buf->EndAccess(bmdBufferAccessWrite); buf->Release(); - // Schedule for display - BMDTimeValue display_time = output_callback->next_time.fetch_add(video_info.fps_den); - decklink_output->ScheduleVideoFrame(vf, display_time, - video_info.fps_den, video_info.fps_num); - + const int64_t t = scheduled_time; + scheduled_time += video_info.fps_den; + decklink_output->ScheduleVideoFrame(vf, t, video_info.fps_den, video_info.fps_num); } + // Converts float32 planar → interleaved int32 and pushes to DeckLink audio buffer. void submit_audio(const float* planar, int samples) { if (!has_audio || samples <= 0) return; - // float32 planar → interleaved int32 (4-byte PCM) scaled to int32 range - const int channels = audio_info.channels; - std::vector interleaved( - static_cast(channels) * static_cast(samples)); - for (int s = 0; s < samples; ++s) { - for (int c = 0; c < channels; ++c) { - float v = planar[c * samples + s]; // planar: channel-major - // clamp to safe range and scale - if (v > 1.0f) v = 1.0f; - else if (v < -1.0f) v = -1.0f; - interleaved[static_cast(s) * channels + c] = + const int ch = audio_info.channels; + const size_t total = static_cast(ch) * static_cast(samples); + if (audio_convert_buf.size() < total) audio_convert_buf.resize(total); + auto& interleaved = audio_convert_buf; + for (int s = 0; s < samples; ++s) + for (int c = 0; c < ch; ++c) { + float v = std::max(-1.0f, std::min(1.0f, planar[c * samples + s])); + interleaved[static_cast(s) * ch + c] = static_cast(v * 2147483647.0f); } - } - const int64_t stream_time = audio_stream_time.fetch_add(samples); + const int64_t t = audio_stream_time; + audio_stream_time += samples; uint32_t written = 0; decklink_output->ScheduleAudioSamples( interleaved.data(), static_cast(samples), - stream_time, audio_info.sample_rate, &written); - (void)written; + t, audio_info.sample_rate, &written); } - private: - class OutputCallback: public IDeckLinkVideoOutputCallback { - public: + class OutputCallback : public IDeckLinkVideoOutputCallback { + public: explicit OutputCallback(DeckLinkSender& owner) : owner(owner) {} - HRESULT ScheduledFrameCompleted (IDeckLinkVideoFrame* completedFrame, BMDOutputFrameCompletionResult result) override { - // Frame is done displaying — return it to the pool. - // submit_frame will overwrite and reschedule it with fresh MXL data. + HRESULT ScheduledFrameCompleted(IDeckLinkVideoFrame* completed, + BMDOutputFrameCompletionResult) override { + { + std::lock_guard lk(owner.pool_mutex); + owner.free_frames.push_back( + static_cast(completed)); + } + owner.pool_cv.notify_one(); return S_OK; } - HRESULT ScheduledPlaybackHasStopped (void) override { - return S_OK; - } + HRESULT ScheduledPlaybackHasStopped() override { return S_OK; } HRESULT STDMETHODCALLTYPE QueryInterface(REFIID, LPVOID*) override { return E_NOINTERFACE; } ULONG STDMETHODCALLTYPE AddRef() override { return ++ref_count; } ULONG STDMETHODCALLTYPE Release() override { return --ref_count; } - std::atomic next_time{0}; - - private: + private: DeckLinkSender& owner; - std::atomic next_frame_idx{0}; std::atomic ref_count{1}; }; // DeckLink SDK objects - std::vector raw_devices; - IDeckLink* selected_device = nullptr; - IDeckLinkOutput* decklink_output = nullptr; - OutputCallback* output_callback = nullptr; + std::vector raw_devices; + IDeckLink* selected_device = nullptr; + IDeckLinkOutput* decklink_output = nullptr; + OutputCallback* output_callback = nullptr; + std::vector all_frames; // for destructor cleanup - std::vector frame_pool; + // Frame pool — frames returned by ScheduledFrameCompleted land here + std::mutex pool_mutex; + std::condition_variable pool_cv; + std::vector free_frames; + + // Scheduling state — only written from submit_frame/submit_audio (single-threaded caller) + int64_t scheduled_time = 0; // next video frame position (fps_num units) + int64_t audio_stream_time = 0; // next audio batch position (sample units) + std::vector audio_convert_buf; // reused across submit_audio calls + + // Maps width/height/fps to a BMDDisplayMode. Uses float comparison to handle + // any representation of drop-frame rates (e.g. 30000/1001 or 2997/100). + static BMDDisplayMode pick_display_mode(int width, int height, int fps_num, int fps_den) { + const double fps = static_cast(fps_num) / static_cast(fps_den); + if (width == 1920 && height == 1080) { + if (std::abs(fps - 23.976) < 0.01) return bmdModeHD1080p2398; + else if (std::abs(fps - 24.0) < 0.01) return bmdModeHD1080p24; + else if (std::abs(fps - 25.0) < 0.01) return bmdModeHD1080p25; + else if (std::abs(fps - 29.97) < 0.01) return bmdModeHD1080p2997; + else if (std::abs(fps - 30.0) < 0.01) return bmdModeHD1080p30; + else if (std::abs(fps - 50.0) < 0.01) return bmdModeHD1080p50; + else if (std::abs(fps - 59.94) < 0.01) return bmdModeHD1080p5994; + else if (std::abs(fps - 60.0) < 0.01) return bmdModeHD1080p6000; + } else if (width == 1280 && height == 720) { + if (std::abs(fps - 50.0) < 0.01) return bmdModeHD720p50; + else if (std::abs(fps - 59.94) < 0.01) return bmdModeHD720p5994; + else if (std::abs(fps - 60.0) < 0.01) return bmdModeHD720p60; + } + throw std::runtime_error( + "No DeckLink display mode for " + std::to_string(width) + "x" + + std::to_string(height) + " @ " + std::to_string(fps_num) + + "/" + std::to_string(fps_den) + " fps"); + } void enumerate_devices() { IDeckLinkIterator* it = CreateDeckLinkIteratorInstance(); @@ -230,9 +267,9 @@ private: IDeckLink* device = nullptr; uint32_t index = 0; while (it->Next(&device) == S_OK) { - IDeckLinkInput* inp = nullptr; - if (device->QueryInterface(IID_IDeckLinkOutput, (void**)&inp) == S_OK) { - inp->Release(); + IDeckLinkOutput* out = nullptr; + if (device->QueryInterface(IID_IDeckLinkOutput, (void**)&out) == S_OK) { + out->Release(); const char* name = nullptr; device->GetDisplayName(&name); devices.push_back({index, name ? name : "?"}); @@ -247,4 +284,5 @@ private: throw std::runtime_error("No DeckLink output devices found"); } }; -} \ No newline at end of file + +} // namespace dmf diff --git a/shared/NodeBase.hpp b/shared/NodeBase.hpp index 62e588a..2040eae 100644 --- a/shared/NodeBase.hpp +++ b/shared/NodeBase.hpp @@ -37,8 +37,6 @@ inline const char* mxl_status_str(mxlStatus s) noexcept { } } - - // Base class for all DMF node binaries. // // Handles the boilerplate every node needs: diff --git a/shared/V210.hpp b/shared/V210.hpp index 461a105..391e1d0 100644 --- a/shared/V210.hpp +++ b/shared/V210.hpp @@ -60,29 +60,6 @@ inline void pack_block( w[3] = (y4 & 0x3FFu) | ((p45.cr & 0x3FFu) << 10) | ((y5 & 0x3FFu) << 20); } -// Write one horizontal line of bars. -// `stride` is the line size in bytes as returned by MXL (configInfo.discrete.sliceSizes[0]). -// Bytes beyond the active pixels are already zeroed by the mmap, so no explicit padding needed. -inline void write_bar_line(uint8_t* line, int width, uint32_t /*stride*/) -{ - const int n = static_cast(SMPTE_BARS.size()); - const int blocks = width / 6; // one V210 block = 6 pixels = 16 bytes - - for (int b = 0; b < blocks; b++) { - int x = b * 6; - auto color = [&](int px) -> const Color& { - return SMPTE_BARS[static_cast(px * n / width)]; - }; - const Color& c01 = color(x); - const Color& c23 = color(x + 2); - const Color& c45 = color(x + 4); - pack_block(line + b * 16, - c01, c01.y, color(x+1).y, - c23, c23.y, color(x+3).y, - c45, c45.y, color(x+5).y); - } -} - // Write one horizontal line of an arbitrary bar palette. template inline void write_palette_line(uint8_t* line, int width, const std::array& palette) @@ -140,12 +117,6 @@ inline void fill_white(uint8_t* buf, int width, int height, uint32_t stride) fill_solid(buf, width, height, stride, {940, 512, 512}); } -// Kept for backward compatibility. -inline void fill_frame(uint8_t* buf, int width, int height, uint32_t stride) -{ - fill_colorbars(buf, width, height, stride); -} - inline void UYVYtoV210(uint8_t* src_buf, uint8_t* dst_buf, int width, int height, uint32_t src_stride, uint32_t dst_stride) { const uint8_t* src = src_buf; @@ -160,49 +131,46 @@ inline void UYVYtoV210(uint8_t* src_buf, uint8_t* dst_buf, int width, int height // mp[8]=U2, mp[9]=Y4, mp[10]=V2, mp[11]=Y5 dmf::v210::pack_block(dst + b * 16, - {0, (uint16_t)(mp[0]<<2), (uint16_t)(mp[2]<<2)}, (uint16_t)(mp[1]<<2), (uint16_t)(mp[3]<<2), - {0, (uint16_t)(mp[4]<<2), (uint16_t)(mp[6]<<2)}, (uint16_t)(mp[5]<<2), (uint16_t)(mp[7]<<2), - {0, (uint16_t)(mp[8]<<2), (uint16_t)(mp[10]<<2)}, (uint16_t)(mp[9]<<2), (uint16_t)(mp[11]<<2) - ); + {0, static_cast(mp[0]<<2), static_cast(mp[2]<<2)}, + static_cast(mp[1]<<2), static_cast(mp[3]<<2), + {0, static_cast(mp[4]<<2), static_cast(mp[6]<<2)}, + static_cast(mp[5]<<2), static_cast(mp[7]<<2), + {0, static_cast(mp[8]<<2), static_cast(mp[10]<<2)}, + static_cast(mp[9]<<2), static_cast(mp[11]<<2)); } src += src_stride; dst += dst_stride; } } -inline void YUV422P10toV210(const uint16_t* y, const uint16_t* u, const uint16_t* v, +inline void YUV422P10toV210( + const uint16_t* y, const uint16_t* u, const uint16_t* v, uint8_t* dst, int width, int height, - int y_stride, int u_stride, int v_stride, // bytes between rows + int y_stride, int u_stride, int v_stride, uint32_t dst_stride) - { - for (int row = 0; row < height; row++) { - const uint16_t* y_row = reinterpret_cast( - reinterpret_cast(y) + row * y_stride - ); - const uint16_t* u_row = reinterpret_cast( - reinterpret_cast(u) + row * u_stride - ); - const uint16_t* v_row = reinterpret_cast( - reinterpret_cast(v) + row * v_stride - ); - - uint8_t* dst_row = dst + static_cast(row) * dst_stride; - const int blocks = width / 6; - - for (int b = 0; b < blocks; b++) { - int x = b * 6; - const uint16_t cb0 = u_row[x/2], cb1 = u_row[x/2+1], cb2 = u_row[x/2+2]; - const uint16_t cr0 = v_row[x/2], cr1 = v_row[x/2+1], cr2 = v_row[x/2+2]; - const uint16_t y0 = y_row[x], y1 = y_row[x+1], y2 = y_row[x+2]; - const uint16_t y3 = y_row[x+3], y4 = y_row[x+4], y5 = y_row[x+5]; - - auto* w = reinterpret_cast(dst_row + b * 16); - w[0] = (cb0 & 0x3FFu) | ((y0 & 0x3FFu) << 10) | ((cr0 & 0x3FFu) << 20); - w[1] = (y1 & 0x3FFu) | ((cb1 & 0x3FFu) << 10) | ((y2 & 0x3FFu) << 20); - w[2] = (cr1 & 0x3FFu) | ((y3 & 0x3FFu) << 10) | ((cb2 & 0x3FFu) << 20); - w[3] = (y4 & 0x3FFu) | ((cr2 & 0x3FFu) << 10) | ((y5 & 0x3FFu) << 20); - } +{ + for (int row = 0; row < height; row++) { + const uint16_t* y_row = reinterpret_cast( + reinterpret_cast(y) + row * y_stride); + const uint16_t* u_row = reinterpret_cast( + reinterpret_cast(u) + row * u_stride); + const uint16_t* v_row = reinterpret_cast( + reinterpret_cast(v) + row * v_stride); + uint8_t* dst_row = dst + static_cast(row) * dst_stride; + const int blocks = width / 6; + for (int b = 0; b < blocks; b++) { + const int x = b * 6; + const uint16_t cb0 = u_row[x/2], cb1 = u_row[x/2+1], cb2 = u_row[x/2+2]; + const uint16_t cr0 = v_row[x/2], cr1 = v_row[x/2+1], cr2 = v_row[x/2+2]; + const uint16_t y0 = y_row[x], y1 = y_row[x+1], y2 = y_row[x+2]; + const uint16_t y3 = y_row[x+3], y4 = y_row[x+4], y5 = y_row[x+5]; + auto* w = reinterpret_cast(dst_row + b * 16); + w[0] = (cb0 & 0x3FFu) | ((y0 & 0x3FFu) << 10) | ((cr0 & 0x3FFu) << 20); + w[1] = (y1 & 0x3FFu) | ((cb1 & 0x3FFu) << 10) | ((y2 & 0x3FFu) << 20); + w[2] = (cr1 & 0x3FFu) | ((y3 & 0x3FFu) << 10) | ((cb2 & 0x3FFu) << 20); + w[3] = (y4 & 0x3FFu) | ((cr2 & 0x3FFu) << 10) | ((y5 & 0x3FFu) << 20); } } +} } // namespace dmf::v210