280 lines
10 KiB
Python
280 lines
10 KiB
Python
# 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 12 u32 (48B): magic state sampleRate blockSize running heartbeat
|
|
inWrite inRead inSlots outWrite outRead outSlots
|
|
inL[4][256] inR[4][256] outL[8][256] outR[8][256]
|
|
|
|
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_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),
|
|
]
|
|
|
|
|
|
HEADER_SIZE = ctypes.sizeof(FxRTHeader) # 48
|
|
IN_L_OFF = HEADER_SIZE # 48
|
|
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
|
|
TOTAL_SIZE = OUT_R_OFF + FXRT_OUT_SLOTS * FXRT_BLOCK * 4
|
|
|
|
|
|
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._np = np.frombuffer(self.shm.buf, dtype=np.float32)
|
|
self._lock = threading.Lock()
|
|
self.proc = None
|
|
self._job_path = ""
|
|
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}
|
|
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("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,
|
|
"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 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):
|
|
"""Ghi 1 block stereo float32 (2048 bytes) vào input ring. Ring đầy →
|
|
drop block 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:
|
|
h = self.h
|
|
while True:
|
|
if h.in_write - h.in_read < h.in_slots:
|
|
slot = h.in_write & (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] = arr[:, 0]
|
|
self._np[br:br + self.block_size] = arr[:, 1]
|
|
h.in_write += 1
|
|
return
|
|
h.in_read += 1 # full → drop oldest input
|
|
|
|
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
|
|
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
|
|
|
|
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:
|
|
return dict(self.latencies)
|
|
|
|
|
|
# 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)
|