Files

523 lines
23 KiB
C++

// native_bridge/juce_fx/JuceFxLoop.cpp
// juce_fx_bridge --juce-fx <job.json> --shm <name> [--parent <pid>]
// 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 <windows.h>
#include <mmsystem.h>
#include <avrt.h>
#endif
#include <algorithm>
#include <chrono>
#include <cstring>
#include <cstdio>
#include <fstream>
#include <iostream>
#include <iterator>
#include <mutex>
#include <string>
#include <thread>
#include <vector>
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<char>(f)),
std::istreambuf_iterator<char>());
}
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=<pid> t=<ms-tu-start>] - 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<std::mutex> lk(mu_);
if (!writeOne(ch)) return traits_type::eof();
return c;
}
std::streamsize xsputn(const char* s, std::streamsize n) override {
std::lock_guard<std::mutex> 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<std::mutex> 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::milliseconds>(
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<double>(32.0, std::min<double>(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<FxRealtimeIPC*>(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: <jobPath>.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<uint32_t>(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<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++;
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<double, std::milli>(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<double, std::milli>(perfIterT1 - perfT0).count();
perfT0 = perfIterT1; ++perfIter;
perfLoopSum += perfLoopMs; if (perfLoopMs > perfLoopMax) perfLoopMax = perfLoopMs;
if (perfIter % 100 == 0) {
const double runSec = std::chrono::duration<double>(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;
}