# 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 = 48 CAP = 64 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. mode='cut' — hiện tại: underrun → hold HOLD block rồi CẮT cứng silence (click nếu nội dung khác 0). mode='holdfade' — đề xuất (fix worklet): underrun → giữ block cuối + fade gain 1→0 qua FADE_HOLD block (~21ms) rồi silence ~0 (hết click); giữa hold, block mới tới → resume ngay. """ FADE_HOLD = 4 # block fade khi underrun (proposal) def __init__(self, mode="cut"): self.mode = mode self.arrivals = [] # (block_seq, time, peak) self.underruns = 0 self.silences = 0 self.silent_ms = 0.0 self.hold_events = 0 self.hold_fades = 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 hold_fade = 0 # còn N block fade đang phát (mode holdfade) 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: if hold_fade > 0 and self.mode == "holdfade": # có block mới giữa hold-fade → resume ngay (không đợi # hết fade, không cắt): fade chỉ còn là chỗi ~0. hold_fade = 0 last_peak = queue.pop(0)[1] play += BLOCK_PERIOD else: # underrun tại play self.underruns += 1 if self.mode == "holdfade": # giữ block cuối + fade 1→0 qua FADE_HOLD block hold_fade = self.FADE_HOLD self.hold_fades += 1 while play <= t and hold_fade > 0: last_peak *= (hold_fade - 1) / max(self.FADE_HOLD - 1, 1) hold_fade -= 1 play += BLOCK_PERIOD if hold_fade == 0 and play <= t: # đã fade về ~0 → silence êm (không click) self.silences += 1 self.silent_ms += (t - play) * 1e3 self.silence_peaks.append(last_peak) play = None idx += 1 else: 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[%s]: underruns=%d hold_only=%d hold_fades=%d silences=%d silent_ms=%.1f" % (self.mode, self.underruns, self.hold_events, self.hold_fades, 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 amp = float(os.environ.get("FXRT_PROBE_AMP", "1.0")) sig = amp * (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) out_peaks = [] # output peak per block (cause-5 clipping track) out_flat = [] # flat-top samples per block (|x|>0.9999) 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: set stall_until — pump VÀ consumer đều đọc (cùng # main thread bị block → cả 2 hướng đều đứng; app thật đúng vậy). while not stop.is_set(): if s_prob: await asyncio.sleep(random.expovariate(s_prob)) dur = random.uniform(s_min, s_max) / 1000.0 stall_state["until"] = 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 send_times = [] iter_times = [] while not stop.is_set(): now = time.perf_counter() iter_times.append(now) if now < stall_state["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) nonlocal send_time_stats, pump_iter_stats send_time_stats = send_times pump_iter_stats = iter_times stall_state = {"until": 0.0} async def param_burst(): # Mô phỏng fxRtPushAllParams (app gửi loạt set_param ngay sau WS open): # numeric ParamID 0..3 → VST3 setParamNormalized THẬT trên bridge. # Fix: setParam defer worker → audio loop không block → send interval ổn. if not os.environ.get("FXRT_PROBE_PARAMS"): return await asyncio.sleep(0.8) nslots = len(chain) ids = [s.strip() for s in os.environ.get("FXRT_PROBE_PARAM_IDS", "0,1,2,3").split(",")] for rep in range(4): for slot in range(nslots): for i, pid in enumerate(ids): await ws.send(json.dumps({"cmd": "set_param", "slot": slot, "key": pid, "value": 0.3 + 0.1 * i})) await asyncio.sleep(0.4) async def consumer(): while not stop.is_set(): if os.environ.get("FXRT_PROBE_CONSUMER_STALL") and \ time.perf_counter() < stall_state["until"]: # main-thread busy: WS message nằm trong browser queue, main # không decode/postMessage → worklet không nhận gì. await asyncio.to_thread(time.sleep, 0.002) continue try: msg = await ws.recv() except Exception: break if isinstance(msg, bytes) and len(msg) == BLOCK * 2 * 4: now = time.perf_counter() f = np.frombuffer(msg, dtype=np.float32) pk = float(f[::2].max()) flat = int(np.count_nonzero(np.abs(f) > 0.9999)) async with recv_lock: recv_times.append((len(recv_times), now, pk)) out_peaks.append(pk) out_flat.append(flat) # 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(param_burst()), 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)) if out_peaks: op = np.array(out_peaks) of = np.array(out_flat) print(" OUT peak: min=%.4f p50=%.4f p95=%.4f max=%.4f" % (op.min(), np.percentile(op, 50), np.percentile(op, 95), op.max())) print(" OUT flat-top(|x|>0.9999): blocks=%d samples=%d (%.4f%%)" % (int(np.count_nonzero(of > 0)), int(of.sum()), 100.0 * of.sum() / max(1, len(of) * BLOCK * 2))) 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(mode=os.environ.get("FXRT_PROBE_WORKLET", "cut")) 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()