fix: MXL_ERR_FLOW_INVALID reconnect, nullptr→"" options, null guards
Readers (ndiout, fakesink) now handle MXL_ERR_FLOW_INVALID by releasing and recreating the flow reader, then realigning the index to current time. This lets consumer nodes survive a producer restart without exiting. All mxlCreateInstance/FlowWriter/FlowReader options args changed from nullptr to "" to match MXL reference implementation style. mxlReleaseFlowReader calls guarded with null checks so cleanup is safe when a mid-run reconnect attempt fails. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
+12
-2
@@ -23,7 +23,7 @@ class FakeSinkNode : public dmf::NodeBase {
|
|||||||
log("flow active — starting read");
|
log("flow active — starting read");
|
||||||
|
|
||||||
mxlFlowReader reader{};
|
mxlFlowReader reader{};
|
||||||
mxlStatus st = mxlCreateFlowReader(instance(), flow_id.c_str(), nullptr, &reader);
|
mxlStatus st = mxlCreateFlowReader(instance(), flow_id.c_str(), "", &reader);
|
||||||
if (st != MXL_STATUS_OK) {
|
if (st != MXL_STATUS_OK) {
|
||||||
log("mxlCreateFlowReader failed (%s)", dmf::mxl_status_str(st));
|
log("mxlCreateFlowReader failed (%s)", dmf::mxl_status_str(st));
|
||||||
return;
|
return;
|
||||||
@@ -59,6 +59,16 @@ class FakeSinkNode : public dmf::NodeBase {
|
|||||||
mxlFlowReaderGetRuntimeInfo(reader, &ri);
|
mxlFlowReaderGetRuntimeInfo(reader, &ri);
|
||||||
index = ri.headIndex;
|
index = ri.headIndex;
|
||||||
|
|
||||||
|
} else if (st == MXL_ERR_FLOW_INVALID) {
|
||||||
|
log("flow invalidated — reconnecting...");
|
||||||
|
if (reader) mxlReleaseFlowReader(instance(), reader);
|
||||||
|
reader = nullptr;
|
||||||
|
mxlSleepForNs(100'000'000);
|
||||||
|
if (mxlCreateFlowReader(instance(), flow_id.c_str(), "", &reader) == MXL_STATUS_OK) {
|
||||||
|
log("flow reconnected");
|
||||||
|
index = mxlGetCurrentIndex(&rate);
|
||||||
|
}
|
||||||
|
|
||||||
} else {
|
} else {
|
||||||
log("unexpected status (%s) on index=%llu", dmf::mxl_status_str(st), index);
|
log("unexpected status (%s) on index=%llu", dmf::mxl_status_str(st), index);
|
||||||
break;
|
break;
|
||||||
@@ -77,7 +87,7 @@ class FakeSinkNode : public dmf::NodeBase {
|
|||||||
|
|
||||||
log("stopped — total frames=%llu invalid=%llu late=%llu",
|
log("stopped — total frames=%llu invalid=%llu late=%llu",
|
||||||
frame_count, invalid_count, late_count);
|
frame_count, invalid_count, late_count);
|
||||||
mxlReleaseFlowReader(instance(), reader);
|
if (reader) mxlReleaseFlowReader(instance(), reader);
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -41,7 +41,7 @@ class NDIInNode : public dmf::NodeBase {
|
|||||||
mxlStatus st = mxlCreateFlowWriter(
|
mxlStatus st = mxlCreateFlowWriter(
|
||||||
instance(),
|
instance(),
|
||||||
dmf::make_video_flow_def(video_flow_id, node_id(), width, height, fps_num, fps_den).c_str(),
|
dmf::make_video_flow_def(video_flow_id, node_id(), width, height, fps_num, fps_den).c_str(),
|
||||||
nullptr, &video_writer, &video_cfg, &created);
|
"", &video_writer, &video_cfg, &created);
|
||||||
if (st != MXL_STATUS_OK) { log("video mxlCreateFlowWriter failed (%s)", dmf::mxl_status_str(st)); return; }
|
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];
|
const uint32_t video_stride = video_cfg.discrete.sliceSizes[0];
|
||||||
@@ -69,7 +69,7 @@ class NDIInNode : public dmf::NodeBase {
|
|||||||
instance(),
|
instance(),
|
||||||
dmf::make_audio_flow_def(audio_flow_id, node_id(), sample_rate, channels, bit_depth,
|
dmf::make_audio_flow_def(audio_flow_id, node_id(), sample_rate, channels, bit_depth,
|
||||||
fps_num, fps_den).c_str(),
|
fps_num, fps_den).c_str(),
|
||||||
nullptr, &audio_writer, &audio_cfg, &created);
|
"", &audio_writer, &audio_cfg, &created);
|
||||||
if (ast != MXL_STATUS_OK) {
|
if (ast != MXL_STATUS_OK) {
|
||||||
log("audio mxlCreateFlowWriter failed (%s) — continuing without audio", dmf::mxl_status_str(ast));
|
log("audio mxlCreateFlowWriter failed (%s) — continuing without audio", dmf::mxl_status_str(ast));
|
||||||
has_audio = false;
|
has_audio = false;
|
||||||
|
|||||||
+27
-5
@@ -64,7 +64,7 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
log("flow active — starting read");
|
log("flow active — starting read");
|
||||||
|
|
||||||
mxlFlowConfigInfo video_cfg{};
|
mxlFlowConfigInfo video_cfg{};
|
||||||
mxlStatus vst = mxlCreateFlowReader(instance(), flow_id.c_str(), nullptr, &video_reader);
|
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);
|
mxlFlowReaderGetConfigInfo(video_reader, &video_cfg);
|
||||||
video_stride = video_cfg.discrete.sliceSizes[0];
|
video_stride = video_cfg.discrete.sliceSizes[0];
|
||||||
@@ -77,10 +77,11 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
int channels = 0;
|
int channels = 0;
|
||||||
int samples_per_frame = 0;
|
int samples_per_frame = 0;
|
||||||
bool has_audio = config().contains("audio_flow_id");
|
bool has_audio = config().contains("audio_flow_id");
|
||||||
|
std::string audio_flow_id;
|
||||||
|
|
||||||
if (has_audio) {
|
if (has_audio) {
|
||||||
const auto audio_flow_info = config().at("audio_flow_id");
|
const auto audio_flow_info = config().at("audio_flow_id");
|
||||||
const auto audio_flow_id = audio_flow_info.at("id").get<std::string>();
|
audio_flow_id = audio_flow_info.at("id").get<std::string>();
|
||||||
sample_rate = audio_flow_info.value("sample_rate", 48000);
|
sample_rate = audio_flow_info.value("sample_rate", 48000);
|
||||||
channels = audio_flow_info.value("channels", 2);
|
channels = audio_flow_info.value("channels", 2);
|
||||||
samples_per_frame = sample_rate * fps_den / fps_num;
|
samples_per_frame = sample_rate * fps_den / fps_num;
|
||||||
@@ -88,7 +89,7 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
log("audio flow=%s %d Hz %dch %d samples/frame",
|
log("audio flow=%s %d Hz %dch %d samples/frame",
|
||||||
audio_flow_id.c_str(), sample_rate, channels, samples_per_frame);
|
audio_flow_id.c_str(), sample_rate, channels, samples_per_frame);
|
||||||
|
|
||||||
mxlStatus ast = mxlCreateFlowReader(instance(), audio_flow_id.c_str(), nullptr, &audio_reader);
|
mxlStatus ast = mxlCreateFlowReader(instance(), audio_flow_id.c_str(), "", &audio_reader);
|
||||||
if (ast != MXL_STATUS_OK) {
|
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;
|
has_audio = false;
|
||||||
@@ -188,6 +189,16 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
mxlFlowRuntimeInfo ari{};
|
mxlFlowRuntimeInfo ari{};
|
||||||
mxlFlowReaderGetRuntimeInfo(audio_reader, &ari);
|
mxlFlowReaderGetRuntimeInfo(audio_reader, &ari);
|
||||||
audio_index = ari.headIndex;
|
audio_index = ari.headIndex;
|
||||||
|
} else if (ast == MXL_ERR_FLOW_INVALID) {
|
||||||
|
log("audio flow invalidated — reconnecting...");
|
||||||
|
mxlReleaseFlowReader(instance(), audio_reader);
|
||||||
|
audio_reader = nullptr;
|
||||||
|
mxlSleepForNs(100'000'000);
|
||||||
|
if (mxlCreateFlowReader(instance(), audio_flow_id.c_str(), "", &audio_reader) == MXL_STATUS_OK) {
|
||||||
|
log("audio flow reconnected");
|
||||||
|
const mxlRational r = {sample_rate, 1};
|
||||||
|
audio_index = mxlGetCurrentIndex(&r);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -223,6 +234,17 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
mxlFlowReaderGetRuntimeInfo(video_reader, &ri);
|
mxlFlowReaderGetRuntimeInfo(video_reader, &ri);
|
||||||
video_index = ri.headIndex;
|
video_index = ri.headIndex;
|
||||||
|
|
||||||
|
} else if (vst == MXL_ERR_FLOW_INVALID) {
|
||||||
|
log("video flow invalidated — reconnecting...");
|
||||||
|
mxlReleaseFlowReader(instance(), video_reader);
|
||||||
|
video_reader = nullptr;
|
||||||
|
mxlSleepForNs(100'000'000);
|
||||||
|
if (mxlCreateFlowReader(instance(), flow_id.c_str(), "", &video_reader) == MXL_STATUS_OK) {
|
||||||
|
log("video flow reconnected");
|
||||||
|
const mxlRational r = {fps_num, fps_den};
|
||||||
|
video_index = mxlGetCurrentIndex(&r);
|
||||||
|
}
|
||||||
|
|
||||||
} else {
|
} 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), video_index);
|
||||||
break;
|
break;
|
||||||
@@ -248,8 +270,8 @@ class NDIOutNode : public dmf::NodeBase {
|
|||||||
else
|
else
|
||||||
log("stopped");
|
log("stopped");
|
||||||
|
|
||||||
if (has_video) mxlReleaseFlowReader(instance(), video_reader);
|
if (has_video && video_reader) mxlReleaseFlowReader(instance(), video_reader);
|
||||||
if (has_audio) mxlReleaseFlowReader(instance(), audio_reader);
|
if (has_audio && audio_reader) mxlReleaseFlowReader(instance(), audio_reader);
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -27,7 +27,7 @@ class TestPatternNode : public dmf::NodeBase {
|
|||||||
mxlStatus vst = mxlCreateFlowWriter(
|
mxlStatus vst = mxlCreateFlowWriter(
|
||||||
instance(),
|
instance(),
|
||||||
dmf::make_video_flow_def(flow_id, node_id(), width, height, fps_num, fps_den).c_str(),
|
dmf::make_video_flow_def(flow_id, node_id(), width, height, fps_num, fps_den).c_str(),
|
||||||
nullptr, &video_writer, &video_cfg, &created);
|
"", &video_writer, &video_cfg, &created);
|
||||||
if (vst != MXL_STATUS_OK) { log("video mxlCreateFlowWriter failed (%s)", dmf::mxl_status_str(vst)); return; }
|
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];
|
const uint32_t video_stride = video_cfg.discrete.sliceSizes[0];
|
||||||
@@ -53,7 +53,7 @@ class TestPatternNode : public dmf::NodeBase {
|
|||||||
instance(),
|
instance(),
|
||||||
dmf::make_audio_flow_def(audio_flow_id, node_id(), sample_rate, channels, 32,
|
dmf::make_audio_flow_def(audio_flow_id, node_id(), sample_rate, channels, 32,
|
||||||
fps_num, fps_den).c_str(),
|
fps_num, fps_den).c_str(),
|
||||||
nullptr, &audio_writer, &audio_cfg, &created);
|
"", &audio_writer, &audio_cfg, &created);
|
||||||
if (ast != MXL_STATUS_OK) {
|
if (ast != MXL_STATUS_OK) {
|
||||||
log("audio mxlCreateFlowWriter failed (%s) — continuing without audio", dmf::mxl_status_str(ast));
|
log("audio mxlCreateFlowWriter failed (%s) — continuing without audio", dmf::mxl_status_str(ast));
|
||||||
has_audio = false;
|
has_audio = false;
|
||||||
|
|||||||
+1
-1
@@ -80,7 +80,7 @@ public:
|
|||||||
|
|
||||||
log("domain=%s", domain_.c_str());
|
log("domain=%s", domain_.c_str());
|
||||||
|
|
||||||
inst_ = mxlCreateInstance(domain_.c_str(), nullptr);
|
inst_ = mxlCreateInstance(domain_.c_str(), "");
|
||||||
if (!inst_) {
|
if (!inst_) {
|
||||||
log("mxlCreateInstance failed at %s", domain_.c_str());
|
log("mxlCreateInstance failed at %s", domain_.c_str());
|
||||||
return 1;
|
return 1;
|
||||||
|
|||||||
@@ -155,7 +155,7 @@ int main(int argc, char* argv[]) {
|
|||||||
|
|
||||||
// Clean up stale flow directories left by any previous crashed run
|
// Clean up stale flow directories left by any previous crashed run
|
||||||
{
|
{
|
||||||
mxlInstance gc = mxlCreateInstance(domain.c_str(), nullptr);
|
mxlInstance gc = mxlCreateInstance(domain.c_str(), "");
|
||||||
if (gc) {
|
if (gc) {
|
||||||
mxlGarbageCollectFlows(gc);
|
mxlGarbageCollectFlows(gc);
|
||||||
mxlDestroyInstance(gc);
|
mxlDestroyInstance(gc);
|
||||||
|
|||||||
Reference in New Issue
Block a user