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
This commit is contained in:
2026-08-24 10:40:24 +07:00
parent 21130858d6
commit 3c9f1c5809
2 changed files with 544 additions and 0 deletions
+166
View File
@@ -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_<tag>.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()
+378
View File
@@ -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()