0acc3ee7ee
- worklet gom 1024 mẫu (4 block) vào 1 frame interleaved, main thread tách 4 WS frame 512 floats - bridge take=min(avail,4) xử lí 1 call process cho burst → overhead VST3 amortize 4 lần - engine _pump_output pace theo deadline thay vì sleep cố định (outRead 131→187.5/s) - FX_RT_ROUNDTRIP_SAMPLES 512→1792 (RTT batching ~7 block)
240 lines
9.5 KiB
C++
240 lines
9.5 KiB
C++
// native_bridge/src/RealtimeFxLoop.cpp
|
|
// fx_vst_bridge --realtime-fx <job.json> --shm <name> [--parent <pid>]
|
|
// Vòng lặp SHM realtime: engine (Python) tạo mapping FxRealtimeIPC, bridge mở,
|
|
// đọc block input ring (engine->bridge), chạy qua RealtimeFxChain (VST3 FX,
|
|
// SEH-guarded, in-place), ghi block output ring (bridge->engine). Heartbeat
|
|
// 100ms riêng (watchdog). Thoát khi running==0 hoặc parent chết.
|
|
//
|
|
// job.json: { "sample_rate": 48000, "block_size": 256, "fx_chain": [ ... ] }
|
|
// fx_chain: mảng {type,path,name,preset_b64,bypass} — RealtimeFxChain::setChain
|
|
// nhận JSON mảng {path,bypass,preset_b64} (type/name bỏ qua).
|
|
#include "RenderFxJob.h"
|
|
#include "FxRealtimeIPC.h"
|
|
|
|
#ifdef _WIN32
|
|
#ifndef NOMINMAX
|
|
#define NOMINMAX
|
|
#endif
|
|
#include <windows.h>
|
|
#endif
|
|
|
|
#include <algorithm>
|
|
#include <cctype>
|
|
#include <chrono>
|
|
#include <cstring>
|
|
#include <fstream>
|
|
#include <iostream>
|
|
#include <iterator>
|
|
#include <string>
|
|
#include <thread>
|
|
|
|
namespace {
|
|
|
|
struct ShmView {
|
|
void* map = nullptr;
|
|
void* view = nullptr;
|
|
};
|
|
|
|
ShmView* openShm(const std::string& name, size_t size) {
|
|
int wlen = MultiByteToWideChar(CP_UTF8, 0, name.c_str(), -1, nullptr, 0);
|
|
std::wstring wname(wlen, L'\0');
|
|
MultiByteToWideChar(CP_UTF8, 0, name.c_str(), -1, &wname[0], wlen);
|
|
HANDLE map = OpenFileMappingW(FILE_MAP_ALL_ACCESS, FALSE, wname.c_str());
|
|
if (!map) return nullptr;
|
|
void* view = MapViewOfFile(map, FILE_MAP_ALL_ACCESS, 0, 0, size);
|
|
if (!view) { CloseHandle(map); return nullptr; }
|
|
ShmView* v = new ShmView();
|
|
v->map = map;
|
|
v->view = view;
|
|
return v;
|
|
}
|
|
|
|
void closeShm(ShmView* v) {
|
|
if (!v) return;
|
|
if (v->view) UnmapViewOfFile(v->view);
|
|
if (v->map) CloseHandle((HANDLE)v->map);
|
|
delete v;
|
|
}
|
|
|
|
static bool parentAlive(uint32_t pid) {
|
|
if (pid == 0) return true;
|
|
HANDLE h = OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, FALSE, pid);
|
|
if (!h) return false;
|
|
CloseHandle(h);
|
|
return true;
|
|
}
|
|
|
|
static void sleepMs(uint32_t ms) { Sleep(ms); }
|
|
|
|
std::string readFile(const std::string& path) {
|
|
std::ifstream f(path, std::ios::binary);
|
|
return std::string((std::istreambuf_iterator<char>(f)),
|
|
std::istreambuf_iterator<char>());
|
|
}
|
|
|
|
// Tách giá trị số khỏi JSON flat: "key":<num> hoặc "key": <num>
|
|
double jsonNumber(const std::string& s, const std::string& key, double def) {
|
|
std::string k = "\"" + key + "\"";
|
|
size_t p = s.find(k);
|
|
if (p == std::string::npos) return def;
|
|
p += k.size();
|
|
while (p < s.size() && (s[p] == ':' || s[p] == ' ' || s[p] == '\t' ||
|
|
s[p] == '\r' || s[p] == '\n'))
|
|
++p;
|
|
size_t e = p;
|
|
while (e < s.size() && (std::isdigit((unsigned char)s[e]) || s[e] == '.' ||
|
|
s[e] == '-' || s[e] == '+' || s[e] == 'e' ||
|
|
s[e] == 'E'))
|
|
++e;
|
|
if (e == p) return def;
|
|
try { return std::stod(s.substr(p, e - p)); } catch (...) { return def; }
|
|
}
|
|
|
|
// Cắt substring mảng JSON "fx_chain":[...] (base64 không chứa [ ] nên depth
|
|
// scan an toàn). Trả rỗng nếu không có.
|
|
std::string extractChainJson(const std::string& s) {
|
|
std::string k = "\"fx_chain\"";
|
|
size_t p = s.find(k);
|
|
if (p == std::string::npos) return "";
|
|
p = s.find('[', p);
|
|
if (p == std::string::npos) return "";
|
|
size_t depth = 0, i = p;
|
|
for (; i < s.size(); ++i) {
|
|
if (s[i] == '[') ++depth;
|
|
else if (s[i] == ']') { --depth; if (depth == 0) break; }
|
|
}
|
|
if (i >= s.size()) return "";
|
|
return s.substr(p, i - p + 1);
|
|
}
|
|
|
|
} // namespace
|
|
|
|
int run_realtime_fx_loop(const std::string& jobPath, const std::string& shmName,
|
|
uint32_t parentPid) {
|
|
std::cerr << "[RealtimeFxLoop] start job=" << jobPath << " shm=" << shmName
|
|
<< " parent=" << parentPid << std::endl;
|
|
const std::string job = readFile(jobPath);
|
|
if (job.empty()) {
|
|
std::cerr << "[RealtimeFxLoop] cannot read job file" << std::endl;
|
|
return 1;
|
|
}
|
|
const double srD = jsonNumber(job, "sample_rate", 44100.0);
|
|
const uint32_t sampleRate = (uint32_t)srD;
|
|
const double blkD = jsonNumber(job, "block_size", (double)FXRT_BLOCK);
|
|
const uint32_t block = (uint32_t)std::max<double>(32.0, std::min<double>(blkD, (double)FXRT_BLOCK));
|
|
const std::string chainJson = extractChainJson(job);
|
|
if (chainJson.empty()) {
|
|
std::cerr << "[RealtimeFxLoop] no fx_chain in job (empty chain = passthrough)" << std::endl;
|
|
}
|
|
|
|
// Batching (fix lag/giật/crackle): gom tối đa FXRT_IN_SLOTS block input
|
|
// (256 mẫu) → 1 call chain.process(L,R,batchN). Overhead VST3 (memcpy,
|
|
// lock param, setup processData, memset output) amortize 4 lần → engine
|
|
// theo kịp 187.5 block/s. setChain với batchN để plugin buffer đủ cho
|
|
// process n lớn (không overflow khi xử lí burst).
|
|
const uint32_t batchN = block * FXRT_IN_SLOTS;
|
|
RealtimeFxChain chain;
|
|
chain.setChain(chainJson, sampleRate, (int32_t)batchN);
|
|
|
|
ShmView* v = openShm(shmName, sizeof(FxRealtimeIPC));
|
|
if (!v) {
|
|
std::cerr << "[RealtimeFxLoop] cannot open SHM: " << shmName
|
|
<< " (engine phải tạo trước)" << std::endl;
|
|
return 2;
|
|
}
|
|
auto* ipc = static_cast<FxRealtimeIPC*>(v->view);
|
|
|
|
// Chờ engine khởi tạo header (magic) — tối đa 2s.
|
|
for (int i = 0; i < 200 && ipc->h.magic != FXRT_MAGIC; ++i) sleepMs(10);
|
|
if (ipc->h.magic != FXRT_MAGIC) {
|
|
std::cerr << "[RealtimeFxLoop] SHM magic mismatch (engine chưa init?)" << std::endl;
|
|
closeShm(v);
|
|
return 2;
|
|
}
|
|
if (ipc->h.inSlots != FXRT_IN_SLOTS || ipc->h.outSlots != FXRT_OUT_SLOTS) {
|
|
std::cerr << "[RealtimeFxLoop] slot count mismatch" << std::endl;
|
|
closeShm(v);
|
|
return 2;
|
|
}
|
|
ipc->h.state = FXRT_STATE_READY;
|
|
std::cerr << "[RealtimeFxLoop] ready sr=" << sampleRate << " block=" << block
|
|
<< " chain=" << (chainJson.empty() ? "(empty)" : "loaded") << std::endl;
|
|
|
|
// Heartbeat thread riêng — watchdog của engine (loop chậm không = chết).
|
|
std::thread hb([&]() {
|
|
while (ipc->h.running && ipc->h.state == FXRT_STATE_READY) {
|
|
ipc->h.heartbeat++;
|
|
sleepMs(100);
|
|
}
|
|
});
|
|
|
|
const uint32_t n = block;
|
|
float L[FXRT_BLOCK * FXRT_IN_SLOTS], R[FXRT_BLOCK * FXRT_IN_SLOTS];
|
|
uint32_t inMask = ipc->h.inSlots - 1;
|
|
uint32_t outMask = ipc->h.outSlots - 1;
|
|
uint64_t processed = 0;
|
|
uint64_t lastGen = 0; // chain generation đã báo latency
|
|
while (ipc->h.running) {
|
|
if (parentPid && !parentAlive(parentPid)) {
|
|
std::cerr << "[RealtimeFxLoop] parent gone — exiting" << std::endl;
|
|
break;
|
|
}
|
|
// SET_PARAM ring (engine -> bridge, Phase 2.8): drain mỗi iteration —
|
|
// áp dụng vào chain TRƯỚC block kế tiếp. Ring đầy → drop lệnh cũ.
|
|
while (ipc->h.ctrlRead < ipc->h.ctrlWrite) {
|
|
const uint32_t cs = ipc->h.ctrlRead & (FXRT_CTRL_SLOTS - 1);
|
|
const FxCtrlCmd& cmd = ipc->ctrl[cs];
|
|
const char* kend = std::find(cmd.key, cmd.key + sizeof(cmd.key), '\0');
|
|
chain.setParam((int)cmd.slot,
|
|
std::string(cmd.key, (size_t)(kend - cmd.key)),
|
|
(double)cmd.value);
|
|
MemoryBarrier();
|
|
ipc->h.ctrlRead++;
|
|
}
|
|
// REPORT_LATENCY (bridge -> engine): chain swap (gen đổi) → báo lại
|
|
// latency từng slot (builtin=0; VST3=getLatencySamples). Ring 8 —
|
|
// engine drain định kỳ, đầy thì overwrite (drop cũ).
|
|
const uint64_t gen = chain.chainGen();
|
|
if (gen != lastGen) {
|
|
const std::vector<uint32_t> lats = chain.entryLatencies();
|
|
for (size_t i = 0; i < lats.size(); ++i) {
|
|
const uint32_t ls = ipc->h.latWrite & (FXRT_LAT_SLOTS - 1);
|
|
ipc->lat[ls].slot = (uint32_t)i;
|
|
ipc->lat[ls].samples = lats[i];
|
|
MemoryBarrier();
|
|
ipc->h.latWrite++;
|
|
}
|
|
lastGen = gen;
|
|
std::cerr << "[RealtimeFxLoop] latency reported: " << lats.size()
|
|
<< " slot(s)" << std::endl;
|
|
}
|
|
uint32_t avail = ipc->h.inWrite - ipc->h.inRead;
|
|
if (avail == 0) { sleepMs(1); continue; }
|
|
// Batch: tiêu thụ tối đa FXRT_IN_SLOTS block, xử lí 1 call. Input về
|
|
// đều → take=1 (latency thấp nhất); dồn burst → take>1 gom lại.
|
|
const uint32_t take = std::min<uint32_t>(avail, FXRT_IN_SLOTS);
|
|
uint32_t off = 0;
|
|
for (uint32_t i = 0; i < take; ++i) {
|
|
const uint32_t slot = ipc->h.inRead & inMask;
|
|
std::memcpy(L + off, ipc->inL[slot], n * sizeof(float));
|
|
std::memcpy(R + off, ipc->inR[slot], n * sizeof(float));
|
|
ipc->h.inRead++; // consume
|
|
off += n;
|
|
}
|
|
chain.process(L, R, off); // SEH-guarded; chain rỗng = passthrough
|
|
for (uint32_t i = 0, o = 0; i < take; ++i, o += n) {
|
|
const uint32_t oslot = ipc->h.outWrite & outMask;
|
|
std::memcpy(ipc->outL[oslot], L + o, n * sizeof(float));
|
|
std::memcpy(ipc->outR[oslot], R + o, n * sizeof(float));
|
|
MemoryBarrier(); // dữ liệu trước index (đúng thứ tự trên ARM64)
|
|
ipc->h.outWrite++; // publish
|
|
++processed;
|
|
}
|
|
}
|
|
ipc->h.state = FXRT_STATE_STARTING; // đã dừng
|
|
hb.join();
|
|
closeShm(v);
|
|
std::cerr << "[RealtimeFxLoop] exit processed=" << processed << std::endl;
|
|
return 0;
|
|
}
|