# app/core/fx_realtime.py """VST FX SHM realtime (cloud-FX loop): engine tạo mapping FxRealtimeIPC, spawn fx_vst_bridge --realtime-fx, đẩy/kéo block audio qua SHM ring. Layout mirror chính xác native_bridge/include/FxRealtimeIPC.h (mọi field u32/float 4-byte, không padding): header 18 u32 (72B): magic state sampleRate blockSize running heartbeat inWrite inRead inSlots outWrite outRead outSlots ctrlWrite ctrlRead ctrlSlots latWrite latRead latSlots inL[4][256] inR[4][256] outL[8][256] outR[8][256] ctrl[8] (FxCtrlCmd 24B) — SET_PARAM ring (engine -> bridge) lat[8] (FxLatReport 8B) — REPORT_LATENCY ring (bridge -> engine) Luồng dữ liệu: browser → WS → write_input() → [bridge xử lý] → read_output() → WS → browser. Block = FXRT_BLOCK=256 stereo float32 = 2048 bytes. """ import ctypes import json import os import subprocess import tempfile import threading import time import uuid import numpy as np from app.config import settings from app.core.native_render import find_fx_bridge_exe FXRT_MAGIC = 0x46585254 FXRT_BLOCK = 256 FXRT_IN_SLOTS = 4 FXRT_OUT_SLOTS = 8 FXRT_CTRL_SLOTS = 8 FXRT_LAT_SLOTS = 8 FXRT_STATE_STARTING = 0 FXRT_STATE_READY = 1 FXRT_STATE_ERROR = 2 class FxRTHeader(ctypes.Structure): _fields_ = [ ("magic", ctypes.c_uint32), ("state", ctypes.c_uint32), ("sample_rate", ctypes.c_uint32), ("block_size", ctypes.c_uint32), ("running", ctypes.c_uint32), ("heartbeat", ctypes.c_uint32), ("in_write", ctypes.c_uint32), ("in_read", ctypes.c_uint32), ("in_slots", ctypes.c_uint32), ("out_write", ctypes.c_uint32), ("out_read", ctypes.c_uint32), ("out_slots", ctypes.c_uint32), ("ctrl_write", ctypes.c_uint32), ("ctrl_read", ctypes.c_uint32), ("ctrl_slots", ctypes.c_uint32), ("lat_write", ctypes.c_uint32), ("lat_read", ctypes.c_uint32), ("lat_slots", ctypes.c_uint32), ] class FxCtrlCmd(ctypes.Structure): _fields_ = [ ("slot", ctypes.c_uint32), ("key", ctypes.c_char * 16), ("value", ctypes.c_float), ] class FxLatReport(ctypes.Structure): _fields_ = [ ("slot", ctypes.c_uint32), ("samples", ctypes.c_uint32), ] HEADER_SIZE = ctypes.sizeof(FxRTHeader) # 72 IN_L_OFF = HEADER_SIZE # 72 IN_R_OFF = IN_L_OFF + FXRT_IN_SLOTS * FXRT_BLOCK * 4 OUT_L_OFF = IN_R_OFF + FXRT_IN_SLOTS * FXRT_BLOCK * 4 OUT_R_OFF = OUT_L_OFF + FXRT_OUT_SLOTS * FXRT_BLOCK * 4 CTRL_OFF = OUT_R_OFF + FXRT_OUT_SLOTS * FXRT_BLOCK * 4 LAT_OFF = CTRL_OFF + FXRT_CTRL_SLOTS * ctypes.sizeof(FxCtrlCmd) TOTAL_SIZE = LAT_OFF + FXRT_LAT_SLOTS * ctypes.sizeof(FxLatReport) class FxRealtimeSession: """Một session realtime: SHM + subprocess fx_vst_bridge + I/O block.""" def __init__(self, fx_chain, sample_rate, block_size=FXRT_BLOCK): from multiprocessing import shared_memory self.fx_chain = fx_chain or [] self.sample_rate = int(sample_rate) self.block_size = int(block_size or FXRT_BLOCK) if self.block_size != FXRT_BLOCK: raise ValueError("block_size phải = 256 (khớp C++ FxRealtimeIPC)") self.name = "SonicForge_FXRT_" + uuid.uuid4().hex[:12] self.shm = shared_memory.SharedMemory(name=self.name, create=True, size=TOTAL_SIZE) self.shm.buf[:TOTAL_SIZE] = b"\x00" * TOTAL_SIZE self.h = FxRTHeader.from_buffer(self.shm.buf) self.h.magic = FXRT_MAGIC self.h.state = FXRT_STATE_STARTING self.h.sample_rate = self.sample_rate self.h.block_size = self.block_size self.h.running = 1 self.h.in_slots = FXRT_IN_SLOTS self.h.out_slots = FXRT_OUT_SLOTS self.h.ctrl_slots = FXRT_CTRL_SLOTS self.h.lat_slots = FXRT_LAT_SLOTS self._np = np.frombuffer(self.shm.buf, dtype=np.float32) self._lock = threading.Lock() # Input batching (fix crackle): gom FXRT_IN_SLOTS frame rồi publish 1 lần # (in_write += n atomic) → bridge luôn thấy bội số của FXRT_IN_SLOTS → # take=4 → batch xử lí. Không phụ thuộc client gửi burst hay lẻ. self._in_pending = [] self._in_pending_t0 = None self.proc = None self._job_path = "" self._chain_seq = 0 self._started_at = time.time() # Chain hợp nhất (spec PLAN_DAW_A.md Phase 1): giữ đúng thứ tự slot. self.chain = list(self.fx_chain or []) # Control protocol (SET_PARAM / REPORT_LATENCY) — Phase 1: giữ ở engine; # Phase 2.8: C++ consume qua SHM control ring. self.pending_params = {} # {slot_idx: {key: value}} self.latencies = {} # {slot_idx: samples} # Diagnostics (FXRT crackle hunt): drop/publish counters. self.stats = {'publish': 0, 'drop_batch': 0, 'drop_blocks': 0, 'stale_flush': 0, 'read_blocks': 0, 'read_empty': 0} self._spawn() # ── process ───────────────────────────────────────────────────────────── def _chain_json_for_bridge(self): """JSON mảng slot theo spec Phase 1 — giữ NGUYÊN thứ tự + type/params; RealtimeFxChain::buildChain (C++) đọc {path,bypass,preset_b64} cho vst3 và {type,params} cho builtin (Phase 2). Slot active=false bị loại trước khi gửi (UI tắt hẳn — không gửi engine).""" out = [] for s in self.chain: if not s or not isinstance(s, dict): continue if s.get("active") is False: continue slot = { "type": s.get("type", ""), "bypass": bool(s.get("bypass")), } if s.get("id"): slot["id"] = s["id"] if s.get("path"): slot["path"] = s["path"] if s.get("preset_b64"): slot["preset_b64"] = s["preset_b64"] if isinstance(s.get("params"), dict) and s["params"]: slot["params"] = s["params"] out.append(slot) return out def _spawn(self): exe = find_fx_bridge_exe() if not exe: raise RuntimeError("Không tìm thấy fx_vst_bridge.exe") job = { "sample_rate": self.sample_rate, "block_size": self.block_size, "seq": self._chain_seq, "fx_chain": self._chain_json_for_bridge(), } fd, self._job_path = tempfile.mkstemp(suffix=".json", prefix="fxrt_") with os.fdopen(fd, "w", encoding="utf-8") as f: json.dump(job, f) cmd = [exe, "--realtime-fx", self._job_path, "--shm", self.name, "--parent", str(os.getpid())] self.proc = subprocess.Popen( cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0), ) def set_chain(self, fx_chain): """Cập nhật chain trên session ĐANG CHẠY (không restart session): ghi lại job file seq++ — bridge poll seq mỗi 500ms → setChain async swap (worker load chain mới, audio giữ chain cũ tới khi swap) → load/ đổi VST FX + đổi preset GUI không ngắt âm/crackle.""" self.fx_chain = list(fx_chain or []) self.chain = self.fx_chain self._chain_seq += 1 if not self._job_path or not os.path.exists(self._job_path): return False try: with open(self._job_path, "w", encoding="utf-8") as f: json.dump({ "sample_rate": self.sample_rate, "block_size": self.block_size, "seq": self._chain_seq, "fx_chain": self._chain_json_for_bridge(), }, f) except Exception: return False return True def wait_ready(self, timeout=30.0): deadline = time.time() + timeout while time.time() < deadline: if self.proc.poll() is not None: raise RuntimeError( f"fx_vst_bridge thoát sớm (rc={self.proc.returncode})") if self.h.state == FXRT_STATE_READY: return True time.sleep(0.02) raise RuntimeError("fx_vst_bridge không sẵn sàng trong %ss" % timeout) def alive(self): return self.proc is not None and self.proc.poll() is None def close(self): try: self.h.running = 0 except Exception: pass if self.proc is not None: try: self.proc.wait(timeout=3) except Exception: try: self.proc.kill() except Exception: pass self.proc = None try: self.shm.close() self.shm.unlink() except Exception: pass if self._job_path and os.path.exists(self._job_path): try: os.remove(self._job_path) except Exception: pass # ── audio I/O ─────────────────────────────────────────────────────────── def write_input(self, interleaved): """Gom FXRT_IN_SLOTS frame vào pending, đủ 4 → publish atomic vào ring (ghi data toàn bộ slot trước, sau đó in_write += n 1 lần). Ring đầy khi publish → drop batch cũ nhất (giữ latency thấp).""" arr = np.frombuffer(interleaved, dtype=np.float32) if arr.size != self.block_size * 2: return arr = arr.reshape(self.block_size, 2) with self._lock: if not self._in_pending: self._in_pending_t0 = time.time() self._in_pending.append((arr[:, 0].copy(), arr[:, 1].copy())) if len(self._in_pending) >= FXRT_IN_SLOTS: self._flush_input_locked() def _flush_input_locked(self): """Publish pending vào input ring. Ghi data trước, in_write += n sau → bridge đọc in_write 1 lần → thấy nguyên batch (take=4).""" if not self._in_pending: return h = self.h frames = self._in_pending self._in_pending = [] self._in_pending_t0 = None avail = h.in_slots - (h.in_write - h.in_read) self.stats['publish'] += len(frames) if avail <= 0: h.in_read += len(frames) # ring full → drop cả batch self.stats['drop_batch'] += len(frames) return n = min(len(frames), avail) drop = len(frames) - n if drop > 0: h.in_read += drop # drop frame cũ nhất self.stats['drop_blocks'] += drop for i in range(n): L, R = frames[i + drop] slot = (h.in_write + i) & (h.in_slots - 1) bl = IN_L_OFF // 4 + slot * self.block_size br = IN_R_OFF // 4 + slot * self.block_size self._np[bl:bl + self.block_size] = L self._np[br:br + self.block_size] = R h.in_write += n # publish atomic — bridge thấy đủ n block cùng lúc def flush_input_stale(self, timeout=0.004): """Flush pending nếu frame đầu chờ > timeout (client gửi lẻ/ngừng giữa burst) → không kẹt latency. Gọi từ _pump_output mỗi vòng.""" with self._lock: if self._in_pending and self._in_pending_t0 is not None and time.time() - self._in_pending_t0 > timeout: self.stats['stale_flush'] += 1 self._flush_input_locked() def read_output(self): """Drain output ring → list bytes (mỗi block 2048 bytes interleaved).""" out = [] with self._lock: h = self.h while h.out_read < h.out_write: slot = h.out_read & (h.out_slots - 1) bl = OUT_L_OFF // 4 + slot * self.block_size br = OUT_R_OFF // 4 + slot * self.block_size inter = np.empty(self.block_size * 2, dtype=np.float32) inter[0::2] = self._np[bl:bl + self.block_size] inter[1::2] = self._np[br:br + self.block_size] out.append(inter.tobytes()) h.out_read += 1 self.stats['read_blocks'] += len(out) if not out: self.stats['read_empty'] += 1 return out def heartbeat_age(self): return time.time() - self._started_at # placeholder: engine đọc h.heartbeat # ── control protocol (SET_PARAM / REPORT_LATENCY) ──────────────────────── # Phase 1: giữ ở engine, client gọi qua WS text frame. Phase 2.8: bridge # C++ đọc pending_params từ SHM control ring + ghi latency vào ring. def set_param(self, slot, key, value): with self._lock: self.pending_params.setdefault(slot, {})[str(key)] = value # Phase 2.8: đẩy vào SHM control ring (bridge drain mỗi # iteration). pending_params giữ làm fallback — fake test session # (object.__new__) không có shm; drain_params vẫn hoạt động. if getattr(self, "shm", None) is not None: try: self._ctrl_enqueue_locked(slot, str(key), value) except Exception: pass def _ctrl_enqueue_locked(self, slot, key, value): h = self.h kb = key.encode("utf-8")[:15] cmds = (FxCtrlCmd * FXRT_CTRL_SLOTS).from_buffer(self.shm.buf, CTRL_OFF) while True: if h.ctrl_write - h.ctrl_read < h.ctrl_slots: cs = h.ctrl_write & (h.ctrl_slots - 1) cmds[cs].slot = int(slot) & 0xFFFFFFFF cmds[cs].key = kb + b"\00" * (16 - len(kb)) cmds[cs].value = float(value) h.ctrl_write += 1 return h.ctrl_read += 1 # ring đầy → drop lệnh cũ nhất def report_latency(self, slot, samples): with self._lock: self.latencies[int(slot)] = int(samples) def drain_params(self): """Trả + xóa toàn bộ pending param (bridge consume 1 lần).""" with self._lock: out = self.pending_params self.pending_params = {} return out def get_latencies(self): with self._lock: self._lat_drain() return dict(self.latencies) def _lat_drain(self): shm = getattr(self, "shm", None) if shm is None: return try: h = self.h lats = (FxLatReport * FXRT_LAT_SLOTS).from_buffer(shm.buf, LAT_OFF) while h.lat_read < h.lat_write: ls = h.lat_read & (h.lat_slots - 1) self.latencies[int(lats[ls].slot)] = int(lats[ls].samples) h.lat_read += 1 except Exception: pass # Registry toàn tiến trình (desktop app 1 user). _SESSIONS = {} _SESSIONS_LOCK = threading.Lock() def start_session(fx_chain, sample_rate, block_size=FXRT_BLOCK): sess = FxRealtimeSession(fx_chain, sample_rate, block_size) try: sess.wait_ready(timeout=30) except Exception: sess.close() raise with _SESSIONS_LOCK: _SESSIONS[sess.name] = sess return sess def get_session(session_id): with _SESSIONS_LOCK: return _SESSIONS.get(session_id) def stop_session(session_id): with _SESSIONS_LOCK: sess = _SESSIONS.pop(session_id, None) if sess: sess.close() return sess is not None def stop_all(): with _SESSIONS_LOCK: ids = list(_SESSIONS.keys()) for sid in ids: stop_session(sid)