// native_bridge/src/RealtimeFxLoop.cpp // fx_vst_bridge --realtime-fx --shm [--parent ] // 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 #include #endif #include #include #include #include #include #include #include #include #include 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(f)), std::istreambuf_iterator()); } // Tách giá trị số khỏi JSON flat: "key": hoặc "key": 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(32.0, std::min(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(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; } // Windows timer resolution 1ms — bỏ granularity Sleep 15.6ms cho toàn // process (idle sleep, chain wait, heartbeat). timeEndPeriod đối ứng ở // cuối. (winmm đã link sẵn — CMakeLists fx_vst_bridge.) #ifdef _WIN32 timeBeginPeriod(1); // Thread audio realtime ABOVE_NORMAL: khong bi tien trinh NORMAL khac // (fx-gui spawn + VST load instance thu 2, native VSTi GUI open...) preempt // -> underrun -> am bi "bop nghet"/ngat quan. ChannelWorker (main.cpp) // da dung THREAD_PRIORITY_ABOVE_NORMAL. SetThreadPriority(GetCurrentThread(), THREAD_PRIORITY_ABOVE_NORMAL); #endif // Chờ chain ĐẦU TIÊN được worker build xong (gen>=1) trước khi READY. // setChain() là ASYNC — worker thread load VST3 + preset mất 1-4s; nếu set // READY ngay, engine mở WS và audio chảy qua chain RỖNG (passthrough), rồi // worker swap chain GIỮA STREAM khi load xong → jump dry→VST + latency đổi // đột ngột → crackle. Chờ swap xong: engine không gửi audio tới khi chain // thật sự sẵn sàng. (Engine wait_ready timeout = 30s > deadline 28s.) const auto chainDeadline = std::chrono::steady_clock::now() + std::chrono::seconds(28); while (chain.chainGen() == 0 && std::chrono::steady_clock::now() < chainDeadline) sleepMs(10); if (chain.chainGen() == 0) std::cerr << "[RealtimeFxLoop] chain load TIMEOUT 28s — READY anyway" << std::endl; else std::cerr << "[RealtimeFxLoop] chain ready (gen=" << chain.chainGen() << ")" << std::endl; 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); } }); // Chain update (set_chain): engine ghi lại job file seq++ — poll mỗi // 100ms; seq đổi → setChain ASYNC (worker load chain mới, swap atomic // dưới mutex — audio thread copy shared_ptr, không dừng) → load/đổi VST // FX + đổi preset GUI không ngắt âm/crackle. uint64_t lastSeq = (uint64_t)jsonNumber(job, "seq", 0); auto lastPoll = std::chrono::steady_clock::now(); 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 uint32_t fadeLeft = 0; // output fade-in còn lại sau chain swap (che gap latency jump) const uint32_t fadeTotal = 512; // ~11.6ms @44.1k // ---- FXRT_PERF instrumentation (tạm, đo crackle) ---- uint64_t perfIter = 0, perfProc = 0, perfIdle = 0, perfOutFull = 0; double perfProcSum = 0.0, perfProcMax = 0.0, perfProcMin = 1e9; double perfLoopSum = 0.0, perfLoopMax = 0.0; uint32_t perfTakeSum = 0; auto perfT0 = std::chrono::steady_clock::now(); auto perfRunStart = perfT0; uint64_t perfLastOut = 0; while (ipc->h.running) { if (parentPid && !parentAlive(parentPid)) { std::cerr << "[RealtimeFxLoop] parent gone — exiting" << std::endl; break; } { const auto nowT = std::chrono::steady_clock::now(); if (std::chrono::duration_cast(nowT - lastPoll).count() >= 100) { lastPoll = nowT; const std::string j2 = readFile(jobPath); if (!j2.empty()) { const uint64_t s2 = (uint64_t)jsonNumber(j2, "seq", 0); if (s2 != lastSeq) { const std::string c2 = extractChainJson(j2); if (!c2.empty()) { lastSeq = s2; std::cerr << "[RealtimeFxLoop] chain update seq=" << s2 << std::endl; chain.setChain(c2, sampleRate, (int32_t)batchN); } } } } } // 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 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; fadeLeft = fadeTotal; std::cerr << "[RealtimeFxLoop] latency reported: " << lats.size() << " slot(s)" << std::endl; } // OUT-ring backpressure (fix crackle "wobble" — pointer collision): // kiểm tra chỗ trống TRƯỚC khi consume input. Nếu out ring đầy (engine // pump chậm/stall: WS send bị block khi main-thread browser busy) mà // bridge vẫn ghi → ghi đè slot engine đang đọc = read/write pointer // giẫm nhau → mẫu torn → méo "wobble" kéo dài tới khi reset. Khi đầy: // KHÔNG xử lí batch này (nhường engine drain), input backlog dồn về // engine → engine drop batch cũ nhất (drop sạch, worklet hold+fade che // gap) — không bao giờ ghi đè out slot chưa đọc. const uint32_t outFree = ipc->h.outSlots - (ipc->h.outWrite - ipc->h.outRead); if (outFree < std::min(ipc->h.inWrite - ipc->h.inRead, FXRT_IN_SLOTS)) { ++perfOutFull; std::this_thread::sleep_until(std::chrono::steady_clock::now() + std::chrono::milliseconds(1)); continue; } uint32_t avail = ipc->h.inWrite - ipc->h.inRead; if (avail == 0) { ++perfIdle; // sleep_until: pacing chính xác theo steady_clock (Sleep(1) bị // granularity 15.6ms khi không có timeBeginPeriod → jitter loop // tới 32ms, phá headroom FILL của worklet). timeBeginPeriod(1) // ở đầu loop + sleep_until giữ loopMax ~2ms. std::this_thread::sleep_until(std::chrono::steady_clock::now() + std::chrono::milliseconds(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(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; } const auto perfP0 = std::chrono::steady_clock::now(); chain.process(L, R, off); // SEH-guarded; chain rỗng = passthrough const auto perfP1 = std::chrono::steady_clock::now(); const double perfProcMs = std::chrono::duration(perfP1 - perfP0).count(); ++perfProc; perfProcSum += perfProcMs; perfTakeSum += take; if (perfProcMs > perfProcMax) perfProcMax = perfProcMs; if (perfProcMs < perfProcMin) perfProcMin = perfProcMs; for (uint32_t i = 0, o = 0; i < take; ++i, o += n) { if (fadeLeft > 0) { const uint32_t f = std::min(fadeLeft, n); const float start = (float)(fadeTotal - fadeLeft); for (uint32_t k = 0; k < f; ++k) { const float g = (start + (float)k) / (float)fadeTotal; L[o + k] *= g; R[o + k] *= g; } fadeLeft -= f; } 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; } const auto perfIterT1 = std::chrono::steady_clock::now(); const double perfLoopMs = std::chrono::duration(perfIterT1 - perfT0).count(); perfT0 = perfIterT1; ++perfIter; perfLoopSum += perfLoopMs; if (perfLoopMs > perfLoopMax) perfLoopMax = perfLoopMs; if (perfIter % 100 == 0) { const double runSec = std::chrono::duration(std::chrono::steady_clock::now() - perfRunStart).count(); std::cerr << "[FxRTPerf] iter=" << perfIter << " procN=" << perfProc << " procAvg=" << (perfProc ? perfProcSum / perfProc : 0.0) << "ms" << " procMax=" << perfProcMax << "ms" << " procMin=" << (perfProc ? perfProcMin : 0.0) << "ms" << " loopAvg=" << (perfIter ? perfLoopSum / perfIter : 0.0) << "ms" << " loopMax=" << perfLoopMax << "ms" << " idle=" << perfIdle << " outFull=" << perfOutFull << " takeAvg=" << (perfProc ? (double)perfTakeSum / perfProc : 0.0) << " procBlocks=" << processed << " rate=" << (runSec > 0 ? processed / runSec : 0.0) << "blk/s" << " outDepth=" << (ipc->h.outWrite - ipc->h.outRead) << std::endl; perfProc = 0; perfProcSum = 0.0; perfTakeSum = 0; perfIdle = 0; perfOutFull = 0; } } ipc->h.state = FXRT_STATE_STARTING; // đã dừng hb.join(); closeShm(v); #ifdef _WIN32 timeEndPeriod(1); #endif std::cerr << "[RealtimeFxLoop] exit processed=" << processed << std::endl; return 0; }