433 lines
17 KiB
Python
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()
|