From 4973d0f9fcddcf6d6e488c0623807fd284666763 Mon Sep 17 00:00:00 2001 From: JohannesItten Date: Thu, 9 Jul 2026 20:48:07 +0300 Subject: [PATCH] gaindb node --- nodes/gaindb/main.cpp | 181 ++++++++++++++++++++++++++++++++---------- 1 file changed, 138 insertions(+), 43 deletions(-) diff --git a/nodes/gaindb/main.cpp b/nodes/gaindb/main.cpp index 6dc8db1..06e2fcb 100644 --- a/nodes/gaindb/main.cpp +++ b/nodes/gaindb/main.cpp @@ -1,58 +1,153 @@ +#include +#include #include -#include "NodeBase.hpp" -#include "FlowDef.hpp" +#include #include #include +#include "NodeBase.hpp" +#include "FlowDef.hpp" -namespace dmf { -class GainDbNode: public dmf::NodeBase { +class GainDbNode : public dmf::NodeBase { void run() override { - if (!config().contains("audio_in_flow_id") || !config().contains("audio_out_flow_id")) { - log("no audio connected to in or out"); + if (!config().contains("audio_in_flow_id")) { log("no audio input connected"); return; } + if (!config().contains("audio_out_flow_id")) { log("no audio output connected"); return; } + + const auto in_id = config().at("audio_in_flow_id").at("id").get(); + const auto out_id = config().at("audio_out_flow_id").at("id").get(); + + const float gain_db = config().value("gain_db", 0.0f); + const float gain_linear = std::pow(10.0f, gain_db / 20.0f); + log("gain=%.2f dB (x%.4f linear)", gain_db, gain_linear); + + // --- wait for input flow --- + log("waiting for flow %s...", in_id.c_str()); + bool active = false; + while (!active && dmf::g_running.load(std::memory_order_relaxed)) { + mxlIsFlowActive(instance(), in_id.c_str(), &active); + if (!active) mxlSleepForNs(100'000'000); + } + if (!dmf::g_running) return; + + // --- create reader --- + mxlFlowReader in_reader{}; + if (mxlCreateFlowReader(instance(), in_id.c_str(), "", &in_reader) != MXL_STATUS_OK) { + log("mxlCreateFlowReader failed"); return; + } + + // --- read format from upstream flow_def --- + const auto fi = dmf::read_audio_flow_info(domain(), in_id); + const int sample_rate = fi.sample_rate; + const int channels = fi.channels; + const int samples_per_grain = fi.samples_per_grain; + log("audio: %d Hz %dch %d samples/grain", sample_rate, channels, samples_per_grain); + + // --- create output writer (same format as input) --- + mxlFlowWriter out_writer{}; + mxlFlowConfigInfo out_cfg{}; + bool created = false; + // gr_num/gr_den = sample_rate/samples_per_grain (e.g. 48000/1920 = 25/1) + const mxlStatus wst = mxlCreateFlowWriter( + instance(), + dmf::make_audio_flow_def(out_id, node_id(), + sample_rate, channels, /*bit_depth=*/32, + /*gr_num=*/sample_rate, /*gr_den=*/samples_per_grain).c_str(), + "", &out_writer, &out_cfg, &created); + if (wst != MXL_STATUS_OK) { + log("mxlCreateFlowWriter failed (%s)", dmf::mxl_status_str(wst)); + mxlReleaseFlowReader(instance(), in_reader); return; } - const std::string in_flow_id = config().at("audio_in_flow_id").at("id").get(); - dmf::AudioFlowInfo in_flow_info = dmf::read_audio_flow_info(domain(), in_flow_id); - int gain = 0; - if (config().contains("gain")) { - gain = config().at("gain").get(); - } - // create writer - int sample_rate = in_flow_info.sample_rate; - const std::string out_flow_id = config().at("audio_out_flow_id").at("id").get(); - std::string flow_def = dmf::make_audio_flow_def( - out_flow_id, - node_id(), - sample_rate, - in_flow_info.channels, - 32, - sample_rate, - in_flow_info.samples_per_grain - ); - mxlFlowWriter writer{}; - mxlFlowConfigInfo cfg{}; - bool created = false; - mxlStatus st = mxlCreateFlowWriter(instance(), flow_def.c_str(), nullptr, &writer, &cfg, &created); - if (st != MXL_STATUS_OK) { - log("audio mxlCreateFlowWriter failed (%s) — continuing without audio", dmf::mxl_status_str(st)); - } else { - log("audio channels=%u buffer=%u samples", - cfg.continuous.channelCount, cfg.continuous.bufferLength); - } + log("output ready buffer=%u samples", out_cfg.continuous.bufferLength); + + // temp flat buffer for one channel — handles ring wrap on both in and out slices + std::vector temp(static_cast(samples_per_grain)); + + const mxlRational audio_rate = {sample_rate, 1}; + uint64_t audio_index = mxlGetCurrentIndex(&audio_rate); + log("start index=%llu", audio_index); + + uint64_t grain_count = 0, late_count = 0; - const mxlRational rate = {sample_rate, 1}; - uint64_t index = mxlGetCurrentIndex(&rate); while (dmf::g_running.load(std::memory_order_relaxed)) { - mxlMutableBufferSlice slice{}; - uint8_t* buf = nullptr; + mxlWrappedMultiBufferSlice in_slice{}; + const mxlStatus rst = mxlFlowReaderGetSamplesNonBlocking( + in_reader, audio_index, + static_cast(samples_per_grain), &in_slice); + if (rst == MXL_STATUS_OK) { + mxlMutableWrappedMultiBufferSlice out_slice{}; + const mxlStatus ost = mxlFlowWriterOpenSamples( + out_writer, audio_index, + static_cast(samples_per_grain), &out_slice); + + if (ost == MXL_STATUS_OK) { + const size_t in_f0 = in_slice.base.fragments[0].size / sizeof(float); + const size_t in_f1 = in_slice.base.fragments[1].size / sizeof(float); + const size_t out_f0 = out_slice.base.fragments[0].size / sizeof(float); + const size_t out_f1 = out_slice.base.fragments[1].size / sizeof(float); + + for (int c = 0; c < channels; ++c) { + // flatten input channel c into temp + const auto* s0 = reinterpret_cast( + static_cast(in_slice.base.fragments[0].pointer) + + c * in_slice.stride); + std::memcpy(temp.data(), s0, in_f0 * sizeof(float)); + if (in_f1 > 0) { + const auto* s1 = reinterpret_cast( + static_cast(in_slice.base.fragments[1].pointer) + + c * in_slice.stride); + std::memcpy(temp.data() + in_f0, s1, in_f1 * sizeof(float)); + } + + // apply gain + for (size_t i = 0; i < static_cast(samples_per_grain); ++i) + temp[i] *= gain_linear; + + // scatter to output channel c (may also wrap) + auto* d0 = reinterpret_cast( + static_cast(out_slice.base.fragments[0].pointer) + + c * out_slice.stride); + std::memcpy(d0, temp.data(), out_f0 * sizeof(float)); + if (out_f1 > 0) { + auto* d1 = reinterpret_cast( + static_cast(out_slice.base.fragments[1].pointer) + + c * out_slice.stride); + std::memcpy(d1, temp.data() + out_f0, out_f1 * sizeof(float)); + } + } + mxlFlowWriterCommitSamples(out_writer); + grain_count++; + } else { + log("OpenSamples failed (%s) at index=%llu", + dmf::mxl_status_str(ost), audio_index); + } + + audio_index += static_cast(samples_per_grain); + const uint64_t ns = mxlGetNsUntilIndex(audio_index, &audio_rate); + if (ns > 0 && ns < 200'000'000ULL) mxlSleepForNs(ns); + + } else if (rst == MXL_ERR_OUT_OF_RANGE_TOO_EARLY) { + mxlSleepForNs(1'000'000); + + } else if (rst == MXL_ERR_OUT_OF_RANGE_TOO_LATE) { + late_count++; + mxlFlowRuntimeInfo ri{}; + mxlFlowReaderGetRuntimeInfo(in_reader, &ri); + audio_index = ri.headIndex; + + } else { + log("read error (%s) at index=%llu", dmf::mxl_status_str(rst), audio_index); + break; + } } - mxlReleaseFlowWriter(instance(), writer); - } + + log("stopped grains=%llu late=%llu", grain_count, late_count); + mxlReleaseFlowReader(instance(), in_reader); + mxlReleaseFlowWriter(instance(), out_writer); + } }; -} int main() { - dmf::GainDbNode node; + GainDbNode node; return node.execute(); -} \ No newline at end of file +}