Files

433 lines
17 KiB
Python

# 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()