gaindb node

This commit is contained in:
JohannesItten
2026-07-09 20:48:07 +03:00
parent 8f28b2f768
commit 4973d0f9fc
+138 -43
View File
@@ -1,58 +1,153 @@
#include <cmath>
#include <cstring>
#include <string>
#include "NodeBase.hpp"
#include "FlowDef.hpp"
#include <vector>
#include <mxl/flow.h>
#include <mxl/time.h>
#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<std::string>();
const auto out_id = config().at("audio_out_flow_id").at("id").get<std::string>();
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<std::string>();
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<int>();
}
// 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>();
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<float> temp(static_cast<size_t>(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<size_t>(samples_per_grain), &in_slice);
if (rst == MXL_STATUS_OK) {
mxlMutableWrappedMultiBufferSlice out_slice{};
const mxlStatus ost = mxlFlowWriterOpenSamples(
out_writer, audio_index,
static_cast<size_t>(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<const float*>(
static_cast<const uint8_t*>(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<const float*>(
static_cast<const uint8_t*>(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<size_t>(samples_per_grain); ++i)
temp[i] *= gain_linear;
// scatter to output channel c (may also wrap)
auto* d0 = reinterpret_cast<float*>(
static_cast<uint8_t*>(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<float*>(
static_cast<uint8_t*>(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<uint64_t>(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();
}
}