// 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 #include #endif #include #include #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); } // M5: base64 encode (RFC 4648) cho editor state capture -> file preset. static std::string base64Encode(const std::string& in) { static const char* tbl = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/"; std::string out; out.reserve(((in.size() + 2) / 3) * 4); for (size_t i = 0; i < in.size(); i += 3) { unsigned a = (unsigned char)in[i]; unsigned b = i + 1 < in.size() ? (unsigned char)in[i + 1] : 0; unsigned c = i + 2 < in.size() ? (unsigned char)in[i + 2] : 0; out += tbl[a >> 2]; out += tbl[((a & 3) << 4) | (b >> 4)]; out += i + 1 < in.size() ? tbl[((b & 15) << 2) | (c >> 6)] : '='; out += i + 2 < in.size() ? tbl[c & 63] : '='; } return out; } // M13: prefix moi dong std::cerr bang [p= t=] - log cua // nhieu bridge process tron 1 file (FXRT_BRIDGE_LOG default) nen khong tach // duoc event theo process/thoi gian. Wrapper stream buffer: phat prefix o dau // dong, khong buffer ca dong (cerr unitbuf flush tung token - emit o sync() // se tao prefix gia); mutex bao an toan thread (main/hb/ctrl/gui thread cung // cerr). Cai dau run_juce_fx_loop truoc khi thread khac chay. class PfxBuf : public std::streambuf { public: explicit PfxBuf(std::streambuf* sink) : sink_(sink), lineStart_(true), t0_(std::chrono::steady_clock::now()) { pid_ = (unsigned long)::GetCurrentProcessId(); } ~PfxBuf() override { if (sink_) std::cerr.rdbuf(sink_); } protected: int_type overflow(int_type c) override { if (traits_type::eq_int_type(c, traits_type::eof())) return traits_type::not_eof(c); char ch = traits_type::to_char_type(c); std::lock_guard lk(mu_); if (!writeOne(ch)) return traits_type::eof(); return c; } std::streamsize xsputn(const char* s, std::streamsize n) override { std::lock_guard lk(mu_); std::streamsize done = 0; while (done < n) { const char c = s[done]; if (c == '\n') { if (sink_->sputc('\n') == traits_type::eof()) break; ++done; lineStart_ = true; continue; } if (!emitPrefix()) break; std::streamsize j = done; while (j < n && s[j] != '\n') ++j; const std::streamsize chunk = j - done; if (chunk > 0) { const std::streamsize w = sink_->sputn(s + done, chunk); if (w != chunk) break; done += chunk; } } return done; } int sync() override { std::lock_guard lk(mu_); return sink_->pubsync() == 0 ? 0 : -1; } private: bool emitPrefix() { if (!lineStart_) return true; lineStart_ = false; const long long ms = std::chrono::duration_cast( std::chrono::steady_clock::now() - t0_).count(); char buf[64]; const int len = std::snprintf(buf, sizeof(buf), "[p=%lu t=%lld] ", pid_, ms); if (len <= 0) return true; return sink_->sputn(buf, len) == (std::streamsize)len; } bool writeOne(char c) { if (c == '\n') { if (sink_->sputc('\n') == traits_type::eof()) return false; lineStart_ = true; return true; } if (!emitPrefix()) return false; return sink_->sputc(c) != traits_type::eof(); } std::streambuf* sink_; std::mutex mu_; bool lineStart_; unsigned long pid_; std::chrono::steady_clock::time_point t0_; }; } // namespace int run_juce_fx_loop(const std::string& jobPath, const std::string& shmName, uint32_t parentPid) { PfxBuf pfx(std::cerr.rdbuf()); std::cerr.rdbuf(&pfx); 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); // MMCSS "Pro Audio": Windows đảm bảo quantum + ưu tiên cho thread này khi // CPU cạnh tranh — chống spike deschedule (loopMax 430ms / 5756ms ở log). // AvSetMmThreadCharacteristicsW (Avrt.lib). Process HIGH_PRIORITY_CLASS — // không cần admin (REALTIME_PRIORITY_CLASS cần SeIncreaseBasePriorityPrivilege). DWORD mmTaskIndex = 0; HANDLE hMm = AvSetMmThreadCharacteristicsW(L"Pro Audio", &mmTaskIndex); if (hMm) { SetPriorityClass(GetCurrentProcess(), HIGH_PRIORITY_CLASS); std::cerr << "[JuceFxLoop] MMCSS Pro Audio ok" << std::endl; } else { std::cerr << "[JuceFxLoop] MMCSS unavailable err=" << GetLastError() << std::endl; } #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) // M7: ed.rx/ry/rw/rh = anchor rect (frontend đo panel embed, // gửi 4 key mỗi 250ms) → guiHost.setAnchorField (thread-safe, // không qua engine queue); reconcile GUI thread SetWindowPos // khi rect đổi → window native nằm đúng rect panel. if (key[3] == 'r' && (key[4] == 'x' || key[4] == 'y' || key[4] == 'w' || key[4] == 'h') && key[5] == '\0') { guiHost.setAnchorField(slot, key[4], (int)value); continue; } 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; } // M16: KHÔNG clearAnchor khi Open — backend retry ed.open sau // ~700ms (plugins.py) clear rect fresh JS vừa gửi → window mới // tạo không anchor được (lần đầu mở panel). Rect cũ vẫn đúng khi // panel yên (M7 dedupe); panel dời → JS gửi rect mới ≤250ms sửa. // Chỉ clear khi Close (window xoá — rect cũ vô dụng). engine.requestEditor(slot, act, w, h); std::cerr << "[JuceFxLoop] ctrl " << key << " slot=" << slot << " w=" << w << " h=" << h << std::endl; // M5: capture state editor (GUI đóng → preset về project như // WM_FXGUI_CAPTURE cũ). PumpEditorQueue (GUI tick 30ms) thực thi // Capture → captured[slot] (raw). Ctrl thread đợi tối đa 3s rồi // base64 + ghi file JSON cạnh job: .preset.json // {slot, preset_b64}. Python /fx-gui/close poll file này (xoá file // cũ trước — chỉ đọc file MỚI cho capture lần này). if (act == JuceFxEngine::EditorAction::Capture) { const std::string capFile = jobPath + ".preset.json"; std::remove(capFile.c_str()); std::string raw; for (int i = 0; i < 300 && raw.empty(); ++i) { sleepMs(10); if (!ipc->h.running || ipc->h.state != FXRT_STATE_READY) break; raw = engine.takeCapturedState(slot); } if (!raw.empty()) { const std::string b64 = base64Encode(raw); std::ofstream f(capFile, std::ios::binary); f << "{\"slot\":" << slot << ",\"preset_b64\":\"" << b64 << "\"}" << std::endl; std::cerr << "[JuceFxLoop] capture slot=" << slot << " -> " << capFile << " (" << raw.size() << " bytes, " << b64.size() << " b64)" << std::endl; } else { std::cerr << "[JuceFxLoop] capture slot=" << slot << " rỗng — plugin không có state" << 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 if (hMm) AvRevertMmThreadCharacteristics(hMm); timeEndPeriod(1); #endif std::cerr << "[JuceFxLoop] exit processed=" << processed << std::endl; return 0; }