Feature/ndi out node #2
+120
-90
@@ -33,31 +33,42 @@ struct NDIContext {
|
|||||||
|
|
||||||
class NDIOutNode : public dmf::NodeBase {
|
class NDIOutNode : public dmf::NodeBase {
|
||||||
void run() override {
|
void run() override {
|
||||||
// --- video flow ---
|
// --- video flow (optional) ---
|
||||||
const auto flow_info = config().at("flow_id");
|
bool has_video = config().contains("flow_id");
|
||||||
const auto flow_id = flow_info.at("id").get<std::string>();
|
|
||||||
const int width = flow_info.value("width", 1920);
|
|
||||||
const int height = flow_info.value("height", 1080);
|
|
||||||
const int fps_num = flow_info.value("fps_num", 25);
|
|
||||||
const int fps_den = flow_info.value("fps_den", 1);
|
|
||||||
|
|
||||||
log("flow=%s %dx%d @ %d/%d fps", flow_id.c_str(), width, height, fps_num, fps_den);
|
|
||||||
|
|
||||||
log("waiting for flow to become active...");
|
|
||||||
bool active = false;
|
|
||||||
while (!active && dmf::g_running.load(std::memory_order_relaxed)) {
|
|
||||||
mxlIsFlowActive(instance(), flow_id.c_str(), &active);
|
|
||||||
if (!active) mxlSleepForNs(100'000'000);
|
|
||||||
}
|
|
||||||
if (!dmf::g_running) return;
|
|
||||||
log("flow active — starting read");
|
|
||||||
|
|
||||||
|
int width = 1920;
|
||||||
|
int height = 1080;
|
||||||
|
int fps_num = 25;
|
||||||
|
int fps_den = 1;
|
||||||
|
std::string flow_id;
|
||||||
mxlFlowReader video_reader{};
|
mxlFlowReader video_reader{};
|
||||||
mxlFlowConfigInfo video_cfg{};
|
uint32_t video_stride = 0;
|
||||||
mxlStatus vst = mxlCreateFlowReader(instance(), flow_id.c_str(), nullptr, &video_reader);
|
|
||||||
if (vst != MXL_STATUS_OK) { log("mxlCreateFlowReader failed (status=%d)", vst); return; }
|
if (has_video) {
|
||||||
mxlFlowReaderGetConfigInfo(video_reader, &video_cfg);
|
const auto flow_info = config().at("flow_id");
|
||||||
const uint32_t video_stride = video_cfg.discrete.sliceSizes[0];
|
flow_id = flow_info.at("id").get<std::string>();
|
||||||
|
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...");
|
||||||
|
bool active = false;
|
||||||
|
while (!active && dmf::g_running.load(std::memory_order_relaxed)) {
|
||||||
|
mxlIsFlowActive(instance(), flow_id.c_str(), &active);
|
||||||
|
if (!active) mxlSleepForNs(100'000'000);
|
||||||
|
}
|
||||||
|
if (!dmf::g_running) return;
|
||||||
|
log("flow active — starting read");
|
||||||
|
|
||||||
|
mxlFlowConfigInfo video_cfg{};
|
||||||
|
mxlStatus vst = mxlCreateFlowReader(instance(), flow_id.c_str(), nullptr, &video_reader);
|
||||||
|
if (vst != MXL_STATUS_OK) { log("video mxlCreateFlowReader failed (status=%d)", vst); return; }
|
||||||
|
mxlFlowReaderGetConfigInfo(video_reader, &video_cfg);
|
||||||
|
video_stride = video_cfg.discrete.sliceSizes[0];
|
||||||
|
}
|
||||||
|
|
||||||
// --- audio flow (optional) ---
|
// --- audio flow (optional) ---
|
||||||
mxlFlowReader audio_reader{};
|
mxlFlowReader audio_reader{};
|
||||||
@@ -88,31 +99,35 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
NDIContext ndi(flow_id.c_str());
|
if (!has_video && !has_audio) { log("no flows configured — exiting"); return; }
|
||||||
|
|
||||||
// V210 intermediate and P216 send buffers for NDI video
|
NDIContext ndi(node_id().c_str());
|
||||||
std::vector<uint8_t> v210_buf(video_stride * height);
|
|
||||||
std::vector<uint8_t> p216_buf(width * sizeof(uint16_t) * 2 * height);
|
|
||||||
|
|
||||||
NDIlib_video_frame_v2_t v210_frame{};
|
// Video buffers and NDI frames (only when has_video)
|
||||||
v210_frame.xres = width;
|
std::vector<uint8_t> v210_buf, p216_buf;
|
||||||
v210_frame.yres = height;
|
NDIlib_video_frame_v2_t v210_frame{}, p216_frame{};
|
||||||
v210_frame.frame_rate_N = fps_num;
|
if (has_video) {
|
||||||
v210_frame.frame_rate_D = fps_den;
|
v210_buf.resize(video_stride * height);
|
||||||
v210_frame.FourCC = static_cast<NDIlib_FourCC_video_type_e>(NDI_LIB_FOURCC('V','2','1','0'));
|
p216_buf.resize(width * sizeof(uint16_t) * 2 * height);
|
||||||
v210_frame.line_stride_in_bytes = video_stride;
|
|
||||||
v210_frame.p_data = v210_buf.data();
|
|
||||||
|
|
||||||
NDIlib_video_frame_v2_t p216_frame{};
|
v210_frame.xres = width;
|
||||||
p216_frame.xres = width;
|
v210_frame.yres = height;
|
||||||
p216_frame.yres = height;
|
v210_frame.frame_rate_N = fps_num;
|
||||||
p216_frame.FourCC = NDIlib_FourCC_video_type_P216;
|
v210_frame.frame_rate_D = fps_den;
|
||||||
p216_frame.frame_rate_N = fps_num;
|
v210_frame.FourCC = static_cast<NDIlib_FourCC_video_type_e>(NDI_LIB_FOURCC('V','2','1','0'));
|
||||||
p216_frame.frame_rate_D = fps_den;
|
v210_frame.line_stride_in_bytes = video_stride;
|
||||||
p216_frame.line_stride_in_bytes = width * static_cast<int>(sizeof(uint16_t));
|
v210_frame.p_data = v210_buf.data();
|
||||||
p216_frame.p_data = p216_buf.data();
|
|
||||||
|
|
||||||
// Float planar buffer for NDI audio (ch0 samples, ch1 samples, ...)
|
p216_frame.xres = width;
|
||||||
|
p216_frame.yres = height;
|
||||||
|
p216_frame.FourCC = NDIlib_FourCC_video_type_P216;
|
||||||
|
p216_frame.frame_rate_N = fps_num;
|
||||||
|
p216_frame.frame_rate_D = fps_den;
|
||||||
|
p216_frame.line_stride_in_bytes = width * static_cast<int>(sizeof(uint16_t));
|
||||||
|
p216_frame.p_data = p216_buf.data();
|
||||||
|
}
|
||||||
|
|
||||||
|
// Audio buffer and NDI frame (only when has_audio)
|
||||||
std::vector<float> audio_planar(static_cast<size_t>(channels) * samples_per_frame);
|
std::vector<float> audio_planar(static_cast<size_t>(channels) * samples_per_frame);
|
||||||
NDIlib_audio_frame_v3_t ndi_audio{};
|
NDIlib_audio_frame_v3_t ndi_audio{};
|
||||||
if (has_audio) {
|
if (has_audio) {
|
||||||
@@ -125,14 +140,8 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// --- main loop ---
|
// --- main loop ---
|
||||||
const mxlRational video_rate = {fps_num, fps_den};
|
uint64_t video_index = 0;
|
||||||
|
|
||||||
uint64_t video_index = mxlGetCurrentIndex(&video_rate);
|
|
||||||
uint64_t audio_index = 0;
|
uint64_t audio_index = 0;
|
||||||
if (has_audio) {
|
|
||||||
const mxlRational audio_rate = {sample_rate, 1};
|
|
||||||
audio_index = mxlGetCurrentIndex(&audio_rate);
|
|
||||||
}
|
|
||||||
uint64_t frame_count = 0;
|
uint64_t frame_count = 0;
|
||||||
uint64_t invalid_count = 0;
|
uint64_t invalid_count = 0;
|
||||||
uint64_t late_count = 0;
|
uint64_t late_count = 0;
|
||||||
@@ -140,8 +149,18 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
auto wall_start = std::chrono::steady_clock::now();
|
auto wall_start = std::chrono::steady_clock::now();
|
||||||
auto last_log_time = wall_start;
|
auto last_log_time = wall_start;
|
||||||
|
|
||||||
|
if (has_video) {
|
||||||
|
const mxlRational video_rate = {fps_num, fps_den};
|
||||||
|
video_index = mxlGetCurrentIndex(&video_rate);
|
||||||
|
}
|
||||||
|
if (has_audio) {
|
||||||
|
const mxlRational audio_rate = {sample_rate, 1};
|
||||||
|
audio_index = mxlGetCurrentIndex(&audio_rate);
|
||||||
|
}
|
||||||
|
|
||||||
while (dmf::g_running.load(std::memory_order_relaxed)) {
|
while (dmf::g_running.load(std::memory_order_relaxed)) {
|
||||||
// --- audio: non-blocking, one chunk per video frame ---
|
// --- audio: non-blocking, one chunk per video frame (or free-running) ---
|
||||||
|
bool audio_advanced = false;
|
||||||
if (has_audio) {
|
if (has_audio) {
|
||||||
mxlWrappedMultiBufferSlice audio_slices{};
|
mxlWrappedMultiBufferSlice audio_slices{};
|
||||||
mxlStatus ast = mxlFlowReaderGetSamplesNonBlocking(
|
mxlStatus ast = mxlFlowReaderGetSamplesNonBlocking(
|
||||||
@@ -164,6 +183,7 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
}
|
}
|
||||||
NDIlib_send_send_audio_v3(ndi.sender, &ndi_audio);
|
NDIlib_send_send_audio_v3(ndi.sender, &ndi_audio);
|
||||||
audio_index += samples_per_frame;
|
audio_index += samples_per_frame;
|
||||||
|
audio_advanced = true;
|
||||||
} else if (ast == MXL_ERR_OUT_OF_RANGE_TOO_LATE) {
|
} else if (ast == MXL_ERR_OUT_OF_RANGE_TOO_LATE) {
|
||||||
mxlFlowRuntimeInfo ari{};
|
mxlFlowRuntimeInfo ari{};
|
||||||
mxlFlowReaderGetRuntimeInfo(audio_reader, &ari);
|
mxlFlowReaderGetRuntimeInfo(audio_reader, &ari);
|
||||||
@@ -172,53 +192,63 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// --- video ---
|
// --- video ---
|
||||||
mxlGrainInfo video_grain{};
|
if (has_video) {
|
||||||
uint8_t* video_buf = nullptr;
|
mxlGrainInfo video_grain{};
|
||||||
vst = mxlFlowReaderGetGrainNonBlocking(video_reader, video_index, &video_grain, &video_buf);
|
uint8_t* video_buf = nullptr;
|
||||||
|
mxlStatus vst = mxlFlowReaderGetGrainNonBlocking(
|
||||||
|
video_reader, video_index, &video_grain, &video_buf);
|
||||||
|
|
||||||
if (vst == MXL_STATUS_OK) {
|
if (vst == MXL_STATUS_OK) {
|
||||||
frame_count++;
|
frame_count++;
|
||||||
if (video_grain.flags & MXL_GRAIN_FLAG_INVALID) invalid_count++;
|
if (video_grain.flags & MXL_GRAIN_FLAG_INVALID) invalid_count++;
|
||||||
|
|
||||||
|
if (NDIlib_send_get_no_connections(ndi.sender, 0) > 0) {
|
||||||
|
std::memcpy(v210_frame.p_data, video_buf, video_stride * height);
|
||||||
|
NDIlib_util_V210_to_P216(&v210_frame, &p216_frame);
|
||||||
|
NDIlib_send_send_video_v2(ndi.sender, &p216_frame);
|
||||||
|
if (++ndi_frame_count == 1)
|
||||||
|
log("first NDI receiver connected");
|
||||||
|
} else {
|
||||||
|
ndi_frame_count = 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
video_index++;
|
||||||
|
|
||||||
|
} else if (vst == MXL_ERR_OUT_OF_RANGE_TOO_EARLY) {
|
||||||
|
mxlSleepForNs(1'000'000);
|
||||||
|
|
||||||
|
} else if (vst == MXL_ERR_OUT_OF_RANGE_TOO_LATE) {
|
||||||
|
late_count++;
|
||||||
|
mxlFlowRuntimeInfo ri{};
|
||||||
|
mxlFlowReaderGetRuntimeInfo(video_reader, &ri);
|
||||||
|
video_index = ri.headIndex;
|
||||||
|
|
||||||
if (NDIlib_send_get_no_connections(ndi.sender, 0) > 0) {
|
|
||||||
std::memcpy(v210_frame.p_data, video_buf, video_stride * height);
|
|
||||||
NDIlib_util_V210_to_P216(&v210_frame, &p216_frame);
|
|
||||||
NDIlib_send_send_video_v2(ndi.sender, &p216_frame);
|
|
||||||
if (++ndi_frame_count == 1)
|
|
||||||
log("first NDI receiver connected");
|
|
||||||
} else {
|
} else {
|
||||||
ndi_frame_count = 0;
|
log("unexpected video status=%d on index=%llu", vst, video_index);
|
||||||
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
video_index++;
|
auto now = std::chrono::steady_clock::now();
|
||||||
|
if (std::chrono::duration<double>(now - last_log_time).count() >= 1.0) {
|
||||||
} else if (vst == MXL_ERR_OUT_OF_RANGE_TOO_EARLY) {
|
const double elapsed = std::chrono::duration<double>(now - wall_start).count();
|
||||||
|
log("frames=%llu invalid=%llu late=%llu avg=%.2f fps",
|
||||||
|
frame_count, invalid_count, late_count,
|
||||||
|
static_cast<double>(frame_count) / elapsed);
|
||||||
|
last_log_time = now;
|
||||||
|
}
|
||||||
|
} else if (!audio_advanced) {
|
||||||
|
// audio-only and nothing was ready — avoid busy spin
|
||||||
mxlSleepForNs(1'000'000);
|
mxlSleepForNs(1'000'000);
|
||||||
|
|
||||||
} else if (vst == MXL_ERR_OUT_OF_RANGE_TOO_LATE) {
|
|
||||||
late_count++;
|
|
||||||
mxlFlowRuntimeInfo ri{};
|
|
||||||
mxlFlowReaderGetRuntimeInfo(video_reader, &ri);
|
|
||||||
video_index = ri.headIndex;
|
|
||||||
|
|
||||||
} else {
|
|
||||||
log("unexpected status=%d on index=%llu", vst, video_index);
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
auto now = std::chrono::steady_clock::now();
|
|
||||||
if (std::chrono::duration<double>(now - last_log_time).count() >= 1.0) {
|
|
||||||
const double elapsed = std::chrono::duration<double>(now - wall_start).count();
|
|
||||||
log("frames=%llu invalid=%llu late=%llu avg=%.2f fps",
|
|
||||||
frame_count, invalid_count, late_count,
|
|
||||||
static_cast<double>(frame_count) / elapsed);
|
|
||||||
last_log_time = now;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
log("stopped — total frames=%llu invalid=%llu late=%llu",
|
if (has_video)
|
||||||
frame_count, invalid_count, late_count);
|
log("stopped — total frames=%llu invalid=%llu late=%llu",
|
||||||
mxlReleaseFlowReader(instance(), video_reader);
|
frame_count, invalid_count, late_count);
|
||||||
|
else
|
||||||
|
log("stopped");
|
||||||
|
|
||||||
|
if (has_video) mxlReleaseFlowReader(instance(), video_reader);
|
||||||
if (has_audio) mxlReleaseFlowReader(instance(), audio_reader);
|
if (has_audio) mxlReleaseFlowReader(instance(), audio_reader);
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|||||||
Reference in New Issue
Block a user