From 3c9f1c58096b459025e2729e3846c5b8d6fbfac3 Mon Sep 17 00:00:00 2001 From: locphamtran Date: Mon, 24 Aug 2026 10:40:24 +0700 Subject: [PATCH] chore(probe): them fxrt probe diagnostics - fxrt_probe.py: in-process bridge loopMax/jitter do - fxrt_ws_probe.py: pipeline that (uvicorn :8011 + WS + worklet sim), dem lost/silence/peak@silence, stalls heavy|light|none --- fxrt_probe.py | 166 +++++++++++++++++++++ fxrt_ws_probe.py | 378 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 544 insertions(+) create mode 100644 fxrt_probe.py create mode 100644 fxrt_ws_probe.py diff --git a/fxrt_probe.py b/fxrt_probe.py new file mode 100644 index 0000000..9b27691 --- /dev/null +++ b/fxrt_probe.py @@ -0,0 +1,166 @@ +# fxrt_probe.py — đo bridge realtime với VST thật (tạm, chẩn đoán crackle). +# Chạy: cd repo && python fxrt_probe.py [--secs 20] +# - tạo session realtime (bắt stderr bridge ra fxrt_bridge_.log) +# - producer: publish batch 4 block mỗi 21.333ms (timeBeginPeriod(1)) +# - consumer: drain output, ghi thời điểm từng block tới +# - báo cáo: engine drops, bridge perf, output gap +import ctypes +import os +import sys +import threading +import time + +import numpy as np + +sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) + +from app.core import fx_realtime # noqa: E402 + +SR = 48000 +BLOCK = 256 +BATCH = 4 +BATCH_PERIOD = BATCH * BLOCK / SR # 21.333ms +BLOCK_PERIOD = BLOCK / SR # 5.333ms + +IZ = r"C:\Program Files\Common Files\VST3\iZotope" +CHAIN_OZONE_MASTER = [ + {"type": "vst3", "path": IZ + r"\Ozone 11 Equalizer.vst3", "active": True}, + {"type": "vst3", "path": IZ + r"\Ozone 11 Dynamics.vst3", "active": True}, + {"type": "vst3", "path": IZ + r"\Ozone 11 Maximizer.vst3", "active": True}, +] +CHAIN_OZONE_1 = [ + {"type": "vst3", "path": IZ + r"\Ozone 11 Maximizer.vst3", "active": True}, +] +CHAIN_RX_DENOISE = [ + {"type": "vst3", "path": IZ + r"\RX 11 Spectral De-noise.vst3", "active": True}, +] +CHAIN_EMPTY = [] + + +def set_timer_res(ms=1): + try: + ctypes.windll.winmm.timeBeginPeriod(ms) + except Exception: + pass + + +def run_test(tag, chain, secs): + log_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), + "fxrt_bridge_%s.log" % tag) + real_popen = fx_realtime.subprocess.Popen + + def capped_popen(cmd, **kw): + kw.pop("stderr", None) + f = open(log_path, "ab", buffering=0) + return real_popen(cmd, stderr=f, **kw) + + fx_realtime.subprocess.Popen = capped_popen + try: + sess = fx_realtime.start_session(chain, SR) + finally: + fx_realtime.subprocess.Popen = real_popen + + print("== %s: session=%s chain=%d slot(s)" % (tag, sess.name, len(chain))) + + stop = threading.Event() + arr_times = [] # (block_idx, perf_counter) khi consumer đọc ra + arr_lock = threading.Lock() + + def producer(): + # 1 block = 256 sample stereo = 512 float; gửi 4 block liên tiếp + # (mô phỏng worklet burst), engine gom 4 → publish 1 batch. + t = np.arange(BLOCK) / SR + sig = 0.2 * np.sin(2 * np.pi * 220 * t) + inter = np.empty(BLOCK * 2, dtype=np.float32) + inter[0::2] = sig + inter[1::2] = 0.9 * sig + data = inter.tobytes() + start = time.perf_counter() + k = 0 + while not stop.is_set(): + for _ in range(BATCH): + sess.write_input(data) + k += 1 + nxt = start + k * BATCH_PERIOD + d = nxt - time.perf_counter() + if d > 0: + time.sleep(d) + + def consumer(): + idx = 0 + while not stop.is_set(): + out = sess.read_output() + if not out: + time.sleep(0.001) + continue + now = time.perf_counter() + with arr_lock: + for _ in out: + arr_times.append((idx, now)) + idx += 1 + + pt = threading.Thread(target=producer) + ct = threading.Thread(target=consumer) + t0 = time.perf_counter() + pt.start(); ct.start() + time.sleep(secs) + stop.set() + pt.join(timeout=5); ct.join(timeout=5) + + stats = dict(sess.stats) + sess.close() + time.sleep(0.3) + + with arr_lock: + times = list(arr_times) + if len(times) > 1: + gaps = [b - a for (i, a), (j, b) in zip(times[:-1], times[1:]) if b > a] + dur = time.perf_counter() - t0 + rate = len(times) / dur + exp = 1.0 / (BATCH / BLOCK_PERIOD) # batch period + print(" engine stats: %s" % stats) + print(" output: %d block / %.1fs = %.1f blk/s (expect %.1f)" + % (len(times), dur, rate, BLOCK / BLOCK_PERIOD)) + print(" gap: n=%d avg=%.3fms p95=%.3fms max=%.3fms" + % (len(gaps), np.mean(gaps) * 1e3, np.percentile(gaps, 95) * 1e3, + max(gaps) * 1e3)) + over = [g * 1e3 for g in gaps if g > 1.35 * BLOCK_PERIOD] + print(" gaps>%.1fms (underrun candidate): %d (%.2f%%)" + % (1.35 * BLOCK_PERIOD * 1e3, len(over), + 100.0 * len(over) / len(gaps) if gaps else 0.0)) + else: + print(" NO OUTPUT (engine stats: %s)" % stats) + + print(" bridge log tail (%s):" % log_path) + try: + with open(log_path, "r", encoding="utf-8", errors="replace") as f: + lines = f.readlines() + for ln in lines[-12:]: + print(" " + ln.rstrip()) + except Exception as e: + print(" (no log: %s)" % e) + + +if __name__ == "__main__": + set_timer_res(1) + secs = 20 + if "--secs" in sys.argv: + secs = int(sys.argv[sys.argv.index("--secs") + 1]) + tests = [] + if "--test" in sys.argv: + t = sys.argv[sys.argv.index("--test") + 1] + tests = [(t, {"ozone1": CHAIN_OZONE_1, + "ozone3": CHAIN_OZONE_MASTER, + "rx": CHAIN_RX_DENOISE, + "empty": CHAIN_EMPTY}[t])] + else: + tests = [("empty", CHAIN_EMPTY), + ("ozone1", CHAIN_OZONE_1), + ("ozone3", CHAIN_OZONE_MASTER), + ("rx", CHAIN_RX_DENOISE)] + for tag, chain in tests: + try: + run_test(tag, chain, secs) + except Exception as e: + print("== %s FAILED: %r" % (tag, e)) + import traceback; traceback.print_exc() diff --git a/fxrt_ws_probe.py b/fxrt_ws_probe.py new file mode 100644 index 0000000..5ba4eed --- /dev/null +++ b/fxrt_ws_probe.py @@ -0,0 +1,378 @@ +# fxrt_ws_probe.py — đo toàn pipeline THẬT (engine uvicorn + WS + bridge) với +# stall model giống main-thread browser. Tạm, chẩn đoán crackle master FX. +# python fxrt_ws_probe.py [--secs 20] [--chain empty|adelay|ozone1|ozone3] +# [--stalls heavy|light|none] +# Kết quả: số block output tới client, gap phân bố, và MÔ PHỎNG CHÍNH XÁC +# worklet outQueue (FILL/HOLD/CAP) → số underrun/silence (crackle thật) + +# peak block cuối trước mỗi silence (fade-out engine → ~0, hết click). +import asyncio +import ctypes +import json +import os +import random +import subprocess +import sys +import tempfile +import time + +import numpy as np + +REPO = os.path.dirname(os.path.abspath(__file__)) +sys.path.insert(0, REPO) + +SR = 48000 +BLOCK = 256 +BATCH = 4 +BATCH_PERIOD = BATCH * BLOCK / SR # 21.333ms +BLOCK_PERIOD = BLOCK / SR # 5.333ms +FILL = 24 +CAP = 48 +HOLD = 2 # block lặp khi underrun (~10.7ms) + +IZ = r"C:\Program Files\Common Files\VST3\iZotope" +ADELAY = os.path.join(REPO, "native_bridge", "build", "VST3", "Release", "adelay.vst3") +CHAINS = { + "empty": [], + "adelay": [{"type": "vst3", "path": ADELAY, "active": True}], + "ozone1": [{"type": "vst3", "path": IZ + r"\Ozone 11 Maximizer.vst3", "active": True}], + "ozone3": [ + {"type": "vst3", "path": IZ + r"\Ozone 11 Equalizer.vst3", "active": True}, + {"type": "vst3", "path": IZ + r"\Ozone 11 Dynamics.vst3", "active": True}, + {"type": "vst3", "path": IZ + r"\Ozone 11 Maximizer.vst3", "active": True}, + ], +} + + +def stall_model(name): + # (xác suất stall mỗi giây, min_ms, max_ms) — mô phỏng main-thread busy + # (GUI decode, React render, mở mastering panel, rAF). + return { + "none": (0.0, 0, 0), + "light": (0.5, 60, 120), + "heavy": (1.5, 120, 280), + }[name] + + +class WorkletSim: + """Mô phỏng chính xác SfFxRealtimeProcessor outQueue từ arrival times.""" + + def __init__(self): + self.arrivals = [] # (block_seq, time, peak) + self.underruns = 0 + self.silences = 0 + self.silent_ms = 0.0 + self.hold_events = 0 + self.silence_peaks = [] # peak block cuối trước mỗi silence (fade -> ~0) + + def run(self, duration): + queue = [] + play = None # thời điểm consume block kế + started = False + idx = 0 + last_peak = 1.0 + for item in self.arrivals: + bi, t = item[0], item[1] + peak = item[2] if len(item) > 2 else 1.0 + queue.append((bi, peak)) + if len(queue) > CAP: + queue.pop(0) + if play is None: + if len(queue) >= FILL: + started = True + play = t + BLOCK_PERIOD + last_peak = queue.pop(0)[1] + continue + # consume các block đến hạn trước arrival này + while play is not None and play <= t: + if queue: + last_peak = queue.pop(0)[1] + play += BLOCK_PERIOD + else: + # underrun tại play — hold HOLD block rồi silence + self.underruns += 1 + hold_end = play + HOLD * BLOCK_PERIOD + if t >= hold_end: + # silence từ hold_end tới khi có block mới (t) + self.silences += 1 + self.silent_ms += (t - hold_end) * 1e3 + self.silence_peaks.append(last_peak) + else: + self.hold_events += 1 + play = None + idx += 1 + if play is None and queue: + play = t + BLOCK_PERIOD + last_peak = queue.pop(0)[1] + return self + + def report(self): + r = ("worklet sim: underruns=%d hold_only=%d silences=%d silent_ms=%.1f" + % (self.underruns, self.hold_events, self.silences, self.silent_ms)) + if self.silence_peaks: + arr = np.array(self.silence_peaks) + hard = int((arr > 0.05).sum()) + r += (" | peak@silence min=%.3f med=%.3f max=%.3f hard(>0.05)=%d" + % (arr.min(), np.median(arr), arr.max(), hard)) + return r + + +def make_block(): + t = np.arange(BLOCK) / SR + sig = 0.2 * np.sin(2 * np.pi * 220 * t) + 0.1 * np.sin(2 * np.pi * 331 * t) + inter = np.empty(BLOCK * 2, dtype=np.float32) + inter[0::2] = sig + inter[1::2] = 0.85 * sig + return inter.tobytes() + + +async def run_ws_probe(port, chain, secs, stall): + import websockets + url = "http://127.0.0.1:%d" % port + import urllib.request + req = urllib.request.Request( + url + "/api/v1/plugins/fx-realtime/start", + data=json.dumps({"fx_chain": chain, "sample_rate": SR}).encode(), + headers={"Content-Type": "application/json"}, method="POST") + try: + with urllib.request.urlopen(req, timeout=60) as r: + st = json.loads(r.read().decode()) + except Exception as e: + print(" start session FAILED:", e) + return None + sid = st["session_id"] + print(" session=%s chain=%d slot(s)" % (sid, len(chain))) + + ws_url = "ws://127.0.0.1:%d/api/v1/plugins/ws/fx-realtime/%s" % (port, sid) + s_prob, s_min, s_max = stall_model(stall) + sent_bursts = 0 + recv_times = [] # (block_seq, time, peak) + recv_lock = asyncio.Lock() + stop = asyncio.Event() + produced = asyncio.Queue() # burst buffers + backlog = [] + sent_times = [] + send_time_stats = [] + produce_count = 0 + pump_iter_stats = [] + + block_bytes = make_block() + + # sanity: đo resolution asyncio.sleep ngay trong loop này + _t0 = time.perf_counter() + for _ in range(50): + await asyncio.sleep(0.001) + print(" [sanity] 50x sleep(0.001)=%.1fms -> %.2fms/iter" + % ((time.perf_counter() - _t0) * 1e3, (time.perf_counter() - _t0) / 50 * 1e3)) + + async def producer(): + # worklet: sản xuất burst mỗi 21.333ms KHÔNG ngừng dù main thread busy + next_t = time.perf_counter() + 0.2 + n = 0 + while not stop.is_set(): + # to_thread: asyncio.sleep bi cap 15.6ms (GQCS Win11) -> pacing 21.3ms + # thành ~31ms -> probe không đại diện. time.sleep giữ precision. + await asyncio.to_thread(time.sleep, max(0, next_t - time.perf_counter())) + await produced.put(block_bytes) + n += 1 + next_t += BATCH_PERIOD + nonlocal produce_count + produce_count = n + + async def stall_scheduler(): + # main-thread busy: báo stall_until cho pump qua queue + while not stop.is_set(): + if s_prob: + await asyncio.sleep(random.expovariate(s_prob)) + dur = random.uniform(s_min, s_max) / 1000.0 + await stall_q.put(time.perf_counter() + dur) + else: + await asyncio.sleep(60) + + async def pump(): + # fxRtInPump JS: tối đa 1 burst/21ms; stall → không gửi + nonlocal backlog + last_sent = 0.0 + stall_until = 0.0 + send_times = [] + iter_times = [] + while not stop.is_set(): + now = time.perf_counter() + iter_times.append(now) + if now < stall_until: + await asyncio.to_thread(time.sleep, 0.002) + continue + while not produced.empty(): + backlog.append(await produced.get()) + if backlog and now - last_sent >= 0.021: + burst = backlog.pop(0) + sent_times.append(now) + t0 = time.perf_counter() + for i in range(BATCH): + await ws.send(burst) + send_times.append(time.perf_counter() - t0) + last_sent = now + await asyncio.to_thread(time.sleep, 0.001) + # cập nhật stall_until từ scheduler + if not stall_q.empty(): + stall_until = stall_q.get_nowait() + nonlocal send_time_stats, pump_iter_stats + send_time_stats = send_times + pump_iter_stats = iter_times + + stall_q = asyncio.Queue() + + async def consumer(): + if os.environ.get("FXRT_PROBE_PLAIN_RECV"): + while not stop.is_set(): + try: + msg = await ws.recv() + except Exception: + break + if isinstance(msg, bytes) and len(msg) == BLOCK * 2 * 4: + now = time.perf_counter() + pk = float(np.frombuffer(msg, dtype=np.float32)[::2].max()) + async with recv_lock: + recv_times.append((len(recv_times), now, pk)) + return + while not stop.is_set(): + try: + msg = await asyncio.wait_for(ws.recv(), timeout=1.0) + except asyncio.TimeoutError: + continue + except Exception: + break + if isinstance(msg, bytes) and len(msg) == BLOCK * 2 * 4: + now = time.perf_counter() + pk = float(np.frombuffer(msg, dtype=np.float32)[::2].max()) + async with recv_lock: + recv_times.append((len(recv_times), now, pk)) + # text frame (latency) bỏ qua + + async with websockets.connect(ws_url) as ws: + t0 = time.perf_counter() + tasks = [asyncio.create_task(producer()), asyncio.create_task(pump()), + asyncio.create_task(stall_scheduler()), asyncio.create_task(consumer())] + if os.environ.get("FXRT_PROBE_NO_CONSUMER"): + for t in tasks: + t.cancel() + tasks = [asyncio.create_task(producer()), asyncio.create_task(pump()), + asyncio.create_task(stall_scheduler())] + await asyncio.sleep(secs) + stop.set() + await asyncio.sleep(0.5) + for t in tasks: + t.cancel() + dur = time.perf_counter() - t0 + + async with recv_lock: + times = list(recv_times) + # stop session + try: + req2 = urllib.request.Request( + url + "/api/v1/plugins/fx-realtime/stop", + data=json.dumps({"session_id": sid}).encode(), + headers={"Content-Type": "application/json"}, method="POST") + urllib.request.urlopen(req2, timeout=10) + except Exception: + pass + + n_sent = len(sent_times) * BATCH + n_recv = len(times) + print(" producer bursts=%d pump sends=%d (backlog at end=%d)" + % (produce_count, len(sent_times), len(backlog))) + if pump_iter_stats: + it = np.diff(np.array(pump_iter_stats)) * 1e3 + it = it[it > 0] + print(" pump iter: avg=%.2fms p50=%.2f p95=%.2f max=%.2fms" + % (it.mean(), np.percentile(it, 50), np.percentile(it, 95), it.max())) + if len(sent_times) > 3: + si = np.diff(np.array(sent_times)) * 1e3 + si = si[si > 0] + print(" send interval: avg=%.2fms p50=%.2f p95=%.2f max=%.2fms" + % (si.mean(), np.percentile(si, 50), np.percentile(si, 95), si.max())) + if send_time_stats: + st = np.array(send_time_stats) * 1e3 + print(" ws send burst: avg=%.2fms p50=%.2f p95=%.2f max=%.2fms" + % (st.mean(), np.percentile(st, 50), np.percentile(st, 95), st.max())) + rate = n_recv / dur + print(" sent=%d bursts (%d blocks) recv=%d blocks over %.1fs = %.1f blk/s (expect 187.5)" + % (len(sent_times), n_sent, n_recv, dur, rate)) + print(" lost blocks (sent-recv): %d" % max(0, n_sent - n_recv)) + if len(times) > 1: + gaps = [b - a for (i, a, _), (j, b, _) in zip(times[:-1], times[1:])] + gaps = [g for g in gaps if g > 0] + g = np.array(gaps) * 1e3 + print(" gap: avg=%.2fms p50=%.2f p95=%.2f p99=%.2f max=%.2fms" + % (g.mean(), np.percentile(g, 50), np.percentile(g, 95), + np.percentile(g, 99), g.max())) + over = g[g > 1.35 * BLOCK_PERIOD * 1e3] + print(" gaps>%.1fms: %d (%.2f%%)" % (1.35 * BLOCK_PERIOD * 1e3, len(over), + 100.0 * len(over) / len(g) if len(g) else 0)) + sim = WorkletSim() + sim.arrivals = times + sim.run(dur) + print(" " + sim.report()) + + +def set_timer_res(ms=1): + try: + ctypes.windll.winmm.timeBeginPeriod(ms) + except Exception: + pass + + +def main(): + set_timer_res(1) + secs = 20 + chain = "ozone1" + stall = "heavy" + port = 8011 + if "--secs" in sys.argv: + secs = int(sys.argv[sys.argv.index("--secs") + 1]) + if "--chain" in sys.argv: + chain = sys.argv[sys.argv.index("--chain") + 1] + if "--stalls" in sys.argv: + stall = sys.argv[sys.argv.index("--stalls") + 1] + if "--port" in sys.argv: + port = int(sys.argv[sys.argv.index("--port") + 1]) + + # engine process + import app.main # noqa: F401 — kiểm tra import được + storage = tempfile.mkdtemp(prefix="sf_wsprobe_") + env = dict(os.environ) + env["SONICFORGE_STORAGE_DIR"] = storage + logf = open(os.path.join(REPO, "fxrt_engine_%d.log" % port), "ab", buffering=0) + boot = ("import ctypes; ctypes.windll.winmm.timeBeginPeriod(1); " + "import uvicorn, app.main; " + "uvicorn.run(app.main.app, host='127.0.0.1', port=%d, log_level='warning')" % port) + engine = subprocess.Popen( + [sys.executable, "-c", boot], + cwd=REPO, env=env, stdout=logf, stderr=subprocess.STDOUT, + creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0)) + try: + import urllib.request + deadline = time.time() + 60 + while time.time() < deadline: + try: + with urllib.request.urlopen("http://127.0.0.1:%d/health" % port, timeout=2) as r: + if r.status == 200: + break + except Exception: + time.sleep(0.5) + else: + print("engine not up; log tail:") + print(open(logf.name, "rb").read().decode(errors="replace")[-2000:]) + return + print("engine up on :%d" % port) + asyncio.run(run_ws_probe(port, CHAINS[chain], secs, stall)) + finally: + engine.terminate() + try: + engine.wait(timeout=5) + except Exception: + engine.kill() + logf.close() + + +if __name__ == "__main__": + main()