// native_bridge/juce_fx/JuceFxLoop.cpp // juce_fx_bridge --juce-fx --shm [--parent ] // G1 POC: vòng lặp SHM giống RealtimeFxLoop.cpp nhưng xử lý qua JuceFxEngine // (JUCE AudioProcessorGraph chain rỗng) thay vì RealtimeFxChain. Job format // giữ nguyên {sample_rate, block_size, fx_chain} — G1 bỏ qua fx_chain. // Validate header mỗi iteration: sampleRate/blockSize đổi giữa chừng → // teardown + prepareToPlay lại + report latency mới qua FxLatReport. #include "JuceFxEngine.h" #include "JuceFxGuiHost.h" #include "FxRealtimeIPC.h" #include "FxShm.h" #ifdef _WIN32 #ifndef NOMINMAX #define NOMINMAX #endif #include #include #endif #include #include #include #include #include #include #include #include #include namespace { 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()); } 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":[...] — như RealtimeFxLoop.cpp. 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_juce_fx_loop(const std::string& jobPath, const std::string& shmName, uint32_t parentPid) { std::cerr << "[JuceFxLoop] start job=" << jobPath << " shm=" << shmName << " parent=" << parentPid << std::endl; const std::string job = readFile(jobPath); if (job.empty()) { std::cerr << "[JuceFxLoop] 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)); // G2: đọc fx_chain -> build graph (VST3 theo chain). Rỗng = passthrough. const std::string chainJson = extractChainJson(job); if (chainJson.empty()) std::cerr << "[JuceFxLoop] no fx_chain in job (empty chain = passthrough)" << std::endl; fxshm::ShmView* v = fxshm::openShm(shmName, sizeof(FxRealtimeIPC)); if (!v) { std::cerr << "[JuceFxLoop] cannot open SHM: " << shmName << " (engine phải tạo trước)" << std::endl; return 2; } auto* ipc = static_cast(v->view); for (int i = 0; i < 200 && ipc->h.magic != FXRT_MAGIC; ++i) sleepMs(10); if (ipc->h.magic != FXRT_MAGIC) { std::cerr << "[JuceFxLoop] SHM magic mismatch (engine chưa init?)" << std::endl; fxshm::closeShm(v); return 2; } if (ipc->h.inSlots != FXRT_IN_SLOTS || ipc->h.outSlots != FXRT_OUT_SLOTS) { std::cerr << "[JuceFxLoop] slot count mismatch" << std::endl; fxshm::closeShm(v); return 2; } #ifdef _WIN32 timeBeginPeriod(1); SetThreadPriority(GetCurrentThread(), THREAD_PRIORITY_HIGHEST); #endif // Engine: prepare graph (chain theo fx_chain). sr/block từ job; header là // nguồn thật (engine ghi lúc tạo SHM) — nếu khác, lấy header. uint32_t curSr = sampleRate; uint32_t curBlock = block; JuceFxEngine engine; engine.setChain(chainJson); if (ipc->h.sampleRate) curSr = ipc->h.sampleRate; if (ipc->h.blockSize) curBlock = ipc->h.blockSize; // M2: GUI thread + host window (Reaper-style — editor của CHÍNH instance // DSP trong graph, không instance 2 / feeder). GUI thread cũng là thread // tạo plugin instance: JUCE plugin gắn MM vào thread tạo instance — editor // tạo trên thread khác sẽ treo (selftest: createEditorIfNeeded hang khi // instance tạo trên main). Vì vậy start() TRƯỚC prepare, prepare chạy qua // runOnGui; main thread sau này chỉ engine.process (JUCE cho phép // cross-thread process như host chuẩn). JuceFxGuiHost guiHost(engine); guiHost.start(); guiHost.runOnGui([&] { engine.prepare(curSr, curBlock); }); std::cerr << "[JuceFxLoop] prepared sr=" << curSr << " block=" << curBlock << " latency=" << engine.latencySamples() << std::endl; // Report latency ban đầu: từng slot (giữ giao thức FxLatReport cũ — engine // Python sum → total). Chain rỗng → slot 0 = 0 (khớp G1). auto reportLatencies = [&]() { const auto lats = engine.entryLatencies(); if (lats.empty()) { const uint32_t ls = ipc->h.latWrite & (FXRT_LAT_SLOTS - 1); ipc->lat[ls].slot = 0; ipc->lat[ls].samples = 0; MemoryBarrier(); ipc->h.latWrite++; } else { 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++; } } }; reportLatencies(); ipc->h.state = FXRT_STATE_READY; std::cerr << "[JuceFxLoop] ready" << std::endl; // (guiHost đã start TRƯỚC engine.prepare ở trên — instance phải tạo trên // GUI thread.) Editor chỉ mở khi có lệnh (M3 ctrl ring); idle ≈ 0 CPU. std::thread hb([&]() { while (ipc->h.running && ipc->h.state == FXRT_STATE_READY) { ipc->h.heartbeat++; sleepMs(100); } }); // M3: ctrl ring reader (engine → bridge, poll 2ms). Key "ed.*" = editor // command (Reaper-style: engine điều khiển editor của CHÍNH instance DSP // trong graph — không instance 2). Dispatch qua engine.requestEditor // (thread-safe queue); GUI tick 30ms pumpEditorQueue thực thi. Key khác // (SET_PARAM cũ: param name) — JuceFxEngine chưa có param API (params đi // trong chain JSON) → consume + bỏ qua an toàn (ring không đầy, không chặn // producer). w/h mã hoá trong value: w*4096 + h — chính xác tuyệt đối với // w,h <= 4095 (float 24-bit mantissa); value <= 0 → editor preferred size // (engine chỉ setSize khi w,h > 0). std::thread ctrl([&]() { while (ipc->h.running && ipc->h.state == FXRT_STATE_READY) { 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 uint32_t slot = cmd.slot; char key[17]; std::memcpy(key, cmd.key, 16); key[16] = '\0'; const float value = cmd.value; MemoryBarrier(); ipc->h.ctrlRead++; // consume trước dispatch — không block ring if (std::strncmp(key, "ed.", 3) != 0) continue; // SET_PARAM cũ — no-op (xem comment trên) JuceFxEngine::EditorAction act; bool known = true; if (std::strcmp(key, "ed.open") == 0) act = JuceFxEngine::EditorAction::Open; else if (std::strcmp(key, "ed.close") == 0) act = JuceFxEngine::EditorAction::Close; else if (std::strcmp(key, "ed.show") == 0) act = JuceFxEngine::EditorAction::Show; else if (std::strcmp(key, "ed.hide") == 0) act = JuceFxEngine::EditorAction::Hide; else if (std::strcmp(key, "ed.resize") == 0) act = JuceFxEngine::EditorAction::Resize; else if (std::strcmp(key, "ed.capture") == 0) act = JuceFxEngine::EditorAction::Capture; else known = false; if (!known) { std::cerr << "[JuceFxLoop] ctrl ed.* lạ: " << key << " slot=" << slot << std::endl; continue; } int w = 0, h = 0; if (value > 0.0f) { const int v = (int)(value + 0.5f); w = v / 4096; h = v % 4096; } engine.requestEditor(slot, act, w, h); std::cerr << "[JuceFxLoop] ctrl " << key << " slot=" << slot << " w=" << w << " h=" << h << std::endl; } sleepMs(2); } }); const uint32_t n = curBlock; const uint32_t inMask = ipc->h.inSlots - 1; const uint32_t outMask = ipc->h.outSlots - 1; float L[FXRT_BLOCK * FXRT_IN_SLOTS], R[FXRT_BLOCK * FXRT_IN_SLOTS]; uint64_t processed = 0; uint64_t perfIter = 0, perfProc = 0, perfIdle = 0, perfOutFull = 0; double perfProcSum = 0.0, perfProcMax = 0.0; double perfLoopSum = 0.0, perfLoopMax = 0.0; uint32_t perfTakeSum = 0; auto perfT0 = std::chrono::steady_clock::now(); auto perfRunStart = perfT0; while (ipc->h.running) { if (parentPid && !parentAlive(parentPid)) { std::cerr << "[JuceFxLoop] parent gone — exiting" << std::endl; break; } // Validate header mỗi iteration: sampleRate/blockSize đổi giữa chừng // (đổi thiết bị audio / session mới khác rate) → teardown graph + // prepareToPlay lại + report latency mới. WebAudio sampleRate bất // biến — check này phòng header bị ghi lại. if (ipc->h.sampleRate && (ipc->h.sampleRate != curSr || ipc->h.blockSize != curBlock)) { curSr = ipc->h.sampleRate; curBlock = ipc->h.blockSize; // M2: editor là JUCE Component — đóng trên GUI thread TRƯỚC khi // re-prepare; prepare() chạy trên GUI thread (thread tạo instance — // xem startup, M2b). guiHost.waitEditorsClosed(1500); guiHost.runOnGui([&] { engine.prepare(curSr, curBlock); }); reportLatencies(); std::cerr << "[JuceFxLoop] re-prepared sr=" << curSr << " block=" << curBlock << " latency=" << engine.latencySamples() << std::endl; } // OUT-ring backpressure (giữ nguyên cơ chế RealtimeFxLoop — pointer // collision → torn frame nếu ghi đè 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; } const uint32_t avail = ipc->h.inWrite - ipc->h.inRead; if (avail == 0) { ++perfIdle; std::this_thread::sleep_until(std::chrono::steady_clock::now() + std::chrono::milliseconds(1)); continue; } 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++; off += n; } const auto perfP0 = std::chrono::steady_clock::now(); engine.process(L, R, off); // 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; 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(); ipc->h.outWrite++; ++processed; // G3: báo processed samples (playhead engine-derived — app tính // position = (processed - anchor - latency)/sr) qua lat ring slot // đặc biệt 0xFFFFFFFE. Python _lat_drain tách riêng (không vào // latencies dict). Ring 8 slot, Python drain mỗi ~2ms — an toàn. const uint32_t pls = ipc->h.latWrite & (FXRT_LAT_SLOTS - 1); ipc->lat[pls].slot = 0xFFFFFFFEu; ipc->lat[pls].samples = (uint32_t)(processed * n); MemoryBarrier(); ipc->h.latWrite++; } 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 << "[JuceFxPerf] iter=" << perfIter << " procN=" << perfProc << " procAvg=" << (perfProc ? perfProcSum / perfProc : 0.0) << "ms" << " procMax=" << perfProcMax << "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; hb.join(); ctrl.join(); // dừng trước shutdown — không còn lệnh editor mới guiHost.runOnGui([&] { engine.shutdown(); }); // graph clear trên GUI thread guiHost.stop(); // M2: đóng editor (GUI thread) + dừng GUI thread + join fxshm::closeShm(v); #ifdef _WIN32 timeEndPeriod(1); #endif std::cerr << "[JuceFxLoop] exit processed=" << processed << std::endl; return 0; }