G1: JUCE engine SHM chain rong POC - probe pass, khong lag them
This commit is contained in:
@@ -0,0 +1,241 @@
|
||||
// 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 "FxRealtimeIPC.h"
|
||||
#include "FxShm.h"
|
||||
|
||||
#ifdef _WIN32
|
||||
#ifndef NOMINMAX
|
||||
#define NOMINMAX
|
||||
#endif
|
||||
#include <windows.h>
|
||||
#include <mmsystem.h>
|
||||
#endif
|
||||
|
||||
#include <algorithm>
|
||||
#include <chrono>
|
||||
#include <cstring>
|
||||
#include <fstream>
|
||||
#include <iostream>
|
||||
#include <iterator>
|
||||
#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; }
|
||||
}
|
||||
|
||||
} // 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<double>(32.0, std::min<double>(blkD, (double)FXRT_BLOCK));
|
||||
// G1: fx_chain bỏ qua (chain rỗng). G2: đọc fx_chain -> build graph.
|
||||
|
||||
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_ABOVE_NORMAL);
|
||||
#endif
|
||||
|
||||
// Engine: prepare graph (chain rỗng). 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;
|
||||
if (ipc->h.sampleRate) curSr = ipc->h.sampleRate;
|
||||
if (ipc->h.blockSize) curBlock = ipc->h.blockSize;
|
||||
engine.prepare(curSr, curBlock);
|
||||
std::cerr << "[JuceFxLoop] prepared sr=" << curSr << " block=" << curBlock
|
||||
<< " latency=" << engine.latencySamples() << std::endl;
|
||||
|
||||
// Report latency ban đầu (slot 0 = tổng chain).
|
||||
{
|
||||
const uint32_t ls = ipc->h.latWrite & (FXRT_LAT_SLOTS - 1);
|
||||
ipc->lat[ls].slot = 0;
|
||||
ipc->lat[ls].samples = engine.latencySamples();
|
||||
MemoryBarrier();
|
||||
ipc->h.latWrite++;
|
||||
}
|
||||
ipc->h.state = FXRT_STATE_READY;
|
||||
std::cerr << "[JuceFxLoop] ready" << std::endl;
|
||||
|
||||
std::thread hb([&]() {
|
||||
while (ipc->h.running && ipc->h.state == FXRT_STATE_READY) {
|
||||
ipc->h.heartbeat++;
|
||||
sleepMs(100);
|
||||
}
|
||||
});
|
||||
|
||||
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;
|
||||
engine.prepare(curSr, curBlock);
|
||||
const uint32_t ls = ipc->h.latWrite & (FXRT_LAT_SLOTS - 1);
|
||||
ipc->lat[ls].slot = 0;
|
||||
ipc->lat[ls].samples = engine.latencySamples();
|
||||
MemoryBarrier();
|
||||
ipc->h.latWrite++;
|
||||
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;
|
||||
}
|
||||
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();
|
||||
engine.shutdown();
|
||||
fxshm::closeShm(v);
|
||||
#ifdef _WIN32
|
||||
timeEndPeriod(1);
|
||||
#endif
|
||||
std::cerr << "[JuceFxLoop] exit processed=" << processed << std::endl;
|
||||
return 0;
|
||||
}
|
||||
Reference in New Issue
Block a user