193 lines
9.0 KiB
C++
193 lines
9.0 KiB
C++
#include <chrono>
|
|
#include <cstring>
|
|
#include <string>
|
|
#include <vector>
|
|
#include <mxl/flow.h>
|
|
#include <mxl/time.h>
|
|
#include "NodeBase.hpp"
|
|
#include "FlowDef.hpp"
|
|
#include "NDIReceiver.hpp"
|
|
|
|
class NDIInNode : public dmf::NodeBase {
|
|
void run() override {
|
|
const auto source_num = static_cast<uint32_t>(config().value("source_num", 0));
|
|
|
|
dmf::NDIReceiver ndi;
|
|
dmf::NDIReceiver::SourceInfo src;
|
|
try {
|
|
auto sources = ndi.find_sources(5000);
|
|
log("Available NDI sources:");
|
|
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;
|
|
}
|
|
|
|
// --- video flow ---
|
|
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<std::string>();
|
|
const int width = src.width;
|
|
const int height = src.height;
|
|
const int fps_num = src.fps_num;
|
|
const int fps_den = src.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 st = 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 (st != MXL_STATUS_OK) { log("video mxlCreateFlowWriter failed (%s)", dmf::mxl_status_str(st)); 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<uint32_t>(height), video_cfg.discrete.grainCount);
|
|
|
|
// --- audio flow (optional — only created when graph wires audio_flow_id) ---
|
|
mxlFlowWriter audio_writer{};
|
|
mxlFlowConfigInfo audio_cfg{};
|
|
int sample_rate = 0;
|
|
int channels = 0;
|
|
int bit_depth = 32;
|
|
bool has_audio = config().contains("audio_flow_id");
|
|
|
|
if (has_audio) {
|
|
const auto audio_flow_info = config().at("audio_flow_id");
|
|
const auto audio_flow_id = audio_flow_info.at("id").get<std::string>();
|
|
bit_depth = audio_flow_info.value("bit_depth", 32);
|
|
|
|
log("probing audio format from NDI source...");
|
|
try {
|
|
const auto ap = ndi.probe_audio();
|
|
sample_rate = ap.sample_rate;
|
|
channels = ap.channels;
|
|
log("detected audio: %d Hz %dch", sample_rate, channels);
|
|
} catch (const std::runtime_error& e) {
|
|
log("audio probe failed (%s) — using config defaults", e.what());
|
|
sample_rate = audio_flow_info.value("sample_rate", 48000);
|
|
channels = audio_flow_info.value("channels", 2);
|
|
}
|
|
|
|
log("audio flow=%s %d Hz %dch %d-bit", audio_flow_id.c_str(), sample_rate, channels, bit_depth);
|
|
|
|
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));
|
|
has_audio = false;
|
|
} else {
|
|
log("audio channels=%u buffer=%u samples",
|
|
audio_cfg.continuous.channelCount, audio_cfg.continuous.bufferLength);
|
|
}
|
|
}
|
|
|
|
// --- main loop ---
|
|
const mxlRational video_rate = {fps_num, fps_den};
|
|
|
|
uint64_t video_index = mxlGetCurrentIndex(&video_rate);
|
|
uint64_t audio_index = 0;
|
|
if (has_audio) {
|
|
const mxlRational audio_rate = {sample_rate, 1};
|
|
audio_index = mxlGetCurrentIndex(&audio_rate);
|
|
}
|
|
log("start video_index=%llu", video_index);
|
|
|
|
std::vector<uint8_t> latest_video(video_stride * height);
|
|
bool have_video = false;
|
|
std::vector<float> audio_buf;
|
|
dmf::NDIReceiver::AudioInfo audio_info;
|
|
|
|
while (dmf::g_running.load(std::memory_order_relaxed)) {
|
|
// Drain available NDI frames before writing one MXL grain.
|
|
// With audio enabled, many audio packets may be queued ahead of a video frame.
|
|
// Budget: drain for at most half a frame period so the grain clock stays on time
|
|
// even when catching up from a startup audio backlog.
|
|
try {
|
|
const auto drain_deadline =
|
|
std::chrono::steady_clock::now() + std::chrono::milliseconds(20);
|
|
bool first = true;
|
|
while (dmf::g_running.load(std::memory_order_relaxed) &&
|
|
std::chrono::steady_clock::now() < drain_deadline) {
|
|
auto kind = ndi.capture(latest_video.data(), video_stride,
|
|
audio_buf, audio_info, has_audio,
|
|
first ? 5 : 0);
|
|
first = false;
|
|
if (kind == dmf::NDIReceiver::FrameKind::Video) {
|
|
have_video = true;
|
|
break; // got video — write grain now, remaining audio deferred to next window
|
|
} else if (kind == dmf::NDIReceiver::FrameKind::Audio && has_audio) {
|
|
mxlMutableWrappedMultiBufferSlice slice{};
|
|
if (mxlFlowWriterOpenSamples(audio_writer, audio_index,
|
|
static_cast<size_t>(audio_info.samples),
|
|
&slice) == MXL_STATUS_OK) {
|
|
const size_t frag0 = slice.base.fragments[0].size / sizeof(float);
|
|
const size_t frag1 = slice.base.fragments[1].size / sizeof(float);
|
|
for (int c = 0; c < audio_info.channels; ++c) {
|
|
const float* src = &audio_buf[c * audio_info.channel_stride];
|
|
auto* dst0 = reinterpret_cast<float*>(
|
|
static_cast<uint8_t*>(slice.base.fragments[0].pointer) + c * slice.stride);
|
|
std::memcpy(dst0, src, frag0 * sizeof(float));
|
|
if (frag1 > 0) {
|
|
auto* dst1 = reinterpret_cast<float*>(
|
|
static_cast<uint8_t*>(slice.base.fragments[1].pointer) + c * slice.stride);
|
|
std::memcpy(dst1, src + frag0, frag1 * sizeof(float));
|
|
}
|
|
}
|
|
mxlFlowWriterCommitSamples(audio_writer);
|
|
}
|
|
audio_index += audio_info.samples;
|
|
} else {
|
|
break; // None — queue drained
|
|
}
|
|
}
|
|
} catch (const std::runtime_error& e) {
|
|
log("NDI error: %s — stopping", e.what());
|
|
break;
|
|
}
|
|
|
|
// Write video grain whenever the MXL clock has reached video_index
|
|
const uint64_t current = mxlGetCurrentIndex(&video_rate);
|
|
if (current >= video_index) {
|
|
mxlGrainInfo grain{};
|
|
uint8_t* buf = nullptr;
|
|
st = mxlFlowWriterOpenGrain(video_writer, video_index, &grain, &buf);
|
|
if (st == MXL_STATUS_OK) {
|
|
if (have_video) {
|
|
std::memcpy(buf, latest_video.data(), latest_video.size());
|
|
grain.flags = 0;
|
|
} else {
|
|
grain.flags = MXL_GRAIN_FLAG_INVALID;
|
|
}
|
|
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 = mxlGetCurrentIndex(&video_rate);
|
|
} else {
|
|
log("video OpenGrain failed (%s) at index=%llu", dmf::mxl_status_str(st), video_index);
|
|
video_index = current + 1;
|
|
}
|
|
}
|
|
}
|
|
|
|
log("stopped at video_index=%llu", video_index);
|
|
mxlReleaseFlowWriter(instance(), video_writer);
|
|
if (has_audio) mxlReleaseFlowWriter(instance(), audio_writer);
|
|
}
|
|
};
|
|
|
|
int main() {
|
|
NDIInNode node;
|
|
return node.execute();
|
|
}
|