# 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 logging import os import subprocess import tempfile import threading import time import uuid from app.config import settings 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 # M4.5 watchdog: số lần AUTO restart tối đa sau crash (editor in-process chết # không được giết session mãi — Reaper chấp nhận gap audio khi plugin crash). FXRT_MAX_RESTARTS = 3 _log = logging.getLogger("fx_realtime") 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): import numpy as np 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} self.processed = None # G3: processed samples tu bridge (lat slot 0xFFFFFFFE) # 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, 'flush': 0} # M4.5 watchdog (chính sách FXRealtime đã duyệt): proc chết ngoài ý # muốn (crash view/plugin/loop) → AUTO restart; user kill / close() có # chủ đích → KHÔNG restart (cờ _closing / running=0 — user gọi lại). self._closing = False self._restarts = 0 # M8: bridge JUCE doc job 1 lan luc spawn (khong poll seq) -> set_chain # phai reload (terminate -> watchdog respawn). Co chong dem crash. self._juce_bridge = False self._reload_planned = False self._watch = threading.Thread(target=self._watch_loop, daemon=True) self._spawn() self._watch.start() # ── 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): from app.core.native_render import find_fx_bridge_exe, find_juce_fx_bridge_exe juce_exe = os.getenv("SF_JUCE_FX_BRIDGE_PATH", "").strip() # G3: mặc định dùng JUCE bridge (latency graph thật + processed samples); # fallback bridge cũ khi không có exe JUCE (regression an toàn). exe = juce_exe or find_juce_fx_bridge_exe() or find_fx_bridge_exe() if not exe: raise RuntimeError("Không tìm thấy bridge exe (juce_fx_bridge / fx_vst_bridge)") job = { "sample_rate": self.sample_rate, "block_size": self.block_size, "seq": self._chain_seq, "fx_chain": self._chain_json_for_bridge(), } # M4.5 respawn: job cũ không ai đọc (proc cũ đã chết) — dọn rác. _old_job = self._job_path if _old_job and os.path.exists(_old_job): try: os.remove(_old_job) except Exception: pass fd, self._job_path = tempfile.mkstemp(suffix=".json", prefix="fxrt_") with os.fdopen(fd, "w", encoding="utf-8") as f: json.dump(job, f) self._juce_bridge = os.path.basename(exe).lower().startswith("juce_fx_bridge") _log.info("[fxrt] spawn exe=%s juce=%s seq=%d shm=%s", os.path.basename(exe), self._juce_bridge, self._chain_seq, self.name) if self._juce_bridge: cmd = [exe, "--juce-fx", self._job_path, "--shm", self.name, "--parent", str(os.getpid())] else: cmd = [exe, "--realtime-fx", self._job_path, "--shm", self.name, "--parent", str(os.getpid())] # FXRT_BRIDGE_LOG=: bat stderr bridge (FxRTPerf) vao file — chan # doan outFull/drop; mac dinh DEVNULL nhu cu. # V34 diag: mac dinh ghi FxRTPerf vao fxrt_bridge.log de phan tich # stall (procMax=bridge/VST cham, outDepth=Python pump cham). _blog = os.environ.get("FXRT_BRIDGE_LOG") if not _blog: # M13: per-session log (fxrt_bridge_.log) — trước đây 1 # file chung cho MỌI session → log 2 process lẫn nhau. Env override # giữ nguyên để ép 1 file khi cần. _ad = os.environ.get("APPDATA") if _ad: _blog = os.path.join(_ad, "SonicForgeDAW", "logs", "fxrt_bridge_" + self.name + ".log") else: _blog = "" self.proc = subprocess.Popen( cmd, stdout=subprocess.DEVNULL, stderr=open(_blog, "ab", buffering=0) if _blog else subprocess.DEVNULL, creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0), ) _log.info("[fxrt] spawned pid=%d log=%s", self.proc.pid, _blog) 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 # M8: juce_fx_bridge doc job CHI 1 LAN luc spawn (JuceFxLoop.cpp khong # poll seq nhu RealtimeFxLoop cu) -> chi ghi job KHONG bao gio ap len # bridge: chain bridge giu nguyen trong khi engine chain da doi -> GUI # mo slot moi C++ bao "slot ngoai chain (M) - bo qua". Reload chu dong: # terminate -> watchdog respawn doc chain hien tai (khong tinh crash, # khong dung gioi han FXRT_MAX_RESTARTS). Legacy fx_vst_bridge poll seq # giu nguyen (async swap muot - khong reload). _log.info("[fxrt] set_chain seq=%d len=%d juce=%s proc_alive=%s reload_planned=%s", self._chain_seq, len(self.chain), self._juce_bridge, self.proc is not None and self.proc.poll() is None, self._reload_planned) if self._juce_bridge: self._request_chain_reload() _log.info("[fxrt] set_chain requested reload proc_pid=%s", self.proc.pid if self.proc else None) return True def _request_chain_reload(self): """M8: yeu cau watchdog respawn sach voi chain hien tai (chi juce bridge). Khong block: terminate proc -> vong watchdog (200ms) thay chet + co _reload_planned -> respawn doc job moi, khong tinh crash.""" with self._lock: if self._closing or self.proc is None or self.proc.poll() is not None: return self._reload_planned = True try: self.proc.terminate() except Exception: self._reload_planned = False def wait_chain_ready(self, timeout=20.0): """M8: cho chain moi thuc su ap tren bridge sau set_chain (juce reload async): het co reload + proc respawn READY. Khong co reload -> tra ngay True. False khi session dong/chet/qua han.""" deadline = time.time() + timeout while time.time() < deadline: with self._lock: if self._closing: return False p = self.proc reloading = self._reload_planned alive = p is not None and p.poll() is None ready = alive and self.h.state == FXRT_STATE_READY if not reloading and ready: return True if not alive and not reloading: return False time.sleep(0.1) return False 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): self._closing = True # M4.5: watchdog không auto-restart khi close chủ đích 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).""" import numpy as np 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). V24 NO-DROP: ring đầy → GIỮ pending (pump flush_input_stale 2ms retry) — không drop → hết content gap + hết race pointer (drop cũ cập in_read vượt in_write → avail > in_slots → ghi đè slot bridge đang đọc → torn frame). Cap 2048 block (~11s) chống tích lũy vô hạn khi client dồn liên tục.""" if not self._in_pending: return h = self.h if len(self._in_pending) > 2048: drop = len(self._in_pending) - 2048 del self._in_pending[:drop] # drop block CŨ NHẤT (front) self.stats['drop_blocks'] += drop frames = self._in_pending avail = h.in_slots - (h.in_write - h.in_read) self.stats['publish'] += len(frames) if avail <= 0: # ring đầy → giữ nguyên pending; pump retry khi bridge drain return n = min(len(frames), avail) publish = frames[:n] if n < len(frames): del frames[:n] # publish n block CŨ NHẤT, giữ phần thừa làm pending self._in_pending_t0 = time.time() if frames else None else: self._in_pending = [] self._in_pending_t0 = None for i in range(n): L, R = publish[i] 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 flush_all(self): """V24: reset toàn bộ pipeline khi client resume sau stall — xóa input/output ring + pending → bridge xử lí audio MỚI, không backlog cũ (lag cộng dồn ~1s sau stall). Gọi khi nhận {cmd:'flush'} qua WS.""" with self._lock: self._in_pending = [] self._in_pending_t0 = None self.h.in_read = self.h.in_write self.h.out_read = self.h.out_write self.stats['flush'] += 1 def read_output(self): import numpy as np """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 get_processed(self): """G3: processed samples (bridge da xu li) - playhead engine-derived.""" with self._lock: self._lat_drain() return self.processed 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) _slot = int(lats[ls].slot) _samples = int(lats[ls].samples) if _slot == 0xFFFFFFFE: # G3: processed - khong vao latencies dict (sum chain sai). self.processed = _samples else: self.latencies[_slot] = _samples h.lat_read += 1 except Exception: pass # ── editor control (Phase 2 — Reaper-style) ───────────────────────────── # M4: engine điều khiển editor của CHÍNH instance DSP trong juce_fx_bridge # qua SHM ctrl ring (key ed.*). Không spawn instance 2 / feeder. def send_editor_cmd(self, action, slot, w=0, h=0): """Gửi lệnh editor (ed.open/close/show/hide/resize/capture) cho slot chain. w/h encode w*4096+h (chính xác tuyệt đối w,h<=4095 — float 24-bit mantissa; xem M3); w,h<=0 → editor preferred size (bridge chỉ setSize khi w,h>0). Trả False nếu session không có SHM (fake object).""" key = "ed.%s" % (action or "").strip().lower() value = 0.0 try: if int(w) > 0 and int(h) > 0: value = float(int(w) * 4096 + int(h)) except Exception: value = 0.0 with self._lock: if getattr(self, "shm", None) is None: return False try: self._ctrl_enqueue_locked(int(slot), key, value) return True except Exception: return False def send_editor_rect(self, slot, x, y, w, h): """M7: anchor rect cho GUI native (window nằm ĐÚNG rect panel embed). Gửi 4 key SHM ctrl ed.rx/ry/rw/rh — MỖI key 1 số nguyên float (chính xác tuyệt đối |v| <= 2^24; x/y âm khi panel nằm màn hình trái OK vì C++ parse int trực tiếp — KHÔNG encode w*4096+h như ed.resize, giới hạn 4095). JuceFxGuiHost lưu anchor (thread-safe) + reconcile GUI thread 30ms SetWindowPos khi rect ĐỔI. Trả False nếu session không SHM.""" with self._lock: if getattr(self, "shm", None) is None: return False try: for f, v in (("x", int(x)), ("y", int(y)), ("w", int(w)), ("h", int(h))): self._ctrl_enqueue_locked(int(slot), "ed.r" + f, float(v)) return True except Exception: return False def chain_slot_for_path(self, plugin_path): """Slot index (theo chain bridge ĐANG CHẠY — đã lọc active=false, giữ thứ tự) của plugin path; None nếu plugin không trong chain. Index khớp editors[] của JuceFxEngine (C++) — ed.* gửi theo index này.""" norm = _norm_plugin_path(plugin_path or "") if not norm: return None for i, s in enumerate(self._chain_json_for_bridge()): p = (s.get("path") or "") if p and _norm_plugin_path(p) == norm: return i return None def capture_editor_preset(self, slot, timeout=3.0): """M5: capture state editor của CHÍNH instance DSP (ed.capture) → C++ ghi file JSON cạnh job: .preset.json {slot, preset_b64} (xem JuceFxLoop.cpp) → đọc + xoá file → trả preset_b64 ("" nếu plugin không có state / bridge chết / timeout). Không cần editor đang mở — state đọc từ processor sống trong graph; gọi TRƯỚC ed.close khi đóng GUI.""" cap_path = (self._job_path or "") + ".preset.json" if cap_path and os.path.exists(cap_path): try: os.remove(cap_path) except Exception: pass if not self.send_editor_cmd("capture", slot=slot): return "" deadline = time.time() + timeout while time.time() < deadline: if not self.alive(): return "" if cap_path and os.path.exists(cap_path): try: with open(cap_path, "r", encoding="utf-8") as f: data = json.load(f) os.remove(cap_path) return data.get("preset_b64", "") except Exception: try: os.remove(cap_path) except Exception: pass return "" time.sleep(0.05) return "" def _watch_loop(self): """M4.5 watchdog — thread riêng (ngoài audio). Proc chết ngoài ý muốn (crash view/plugin/loop) → AUTO restart tối đa FXRT_MAX_RESTARTS lần, log rc để đối chiếu CPU 100%. Không restart khi: close() chủ đích (_closing) / running=0 (user kill — user gọi lại). Respawn đúng cmd line (--juce-fx --shm --parent ) qua _spawn().""" while True: time.sleep(0.2) with self._lock: if self._closing: return proc = self.proc if proc is None: return rc = proc.poll() if rc is None: continue if self._closing: return if getattr(self.h, "running", 1) == 0: _log.warning("[fxrt] bridge thoát rc=%s sau running=0 (user kill — không restart)", rc) return if self._reload_planned: self._reload_planned = False _log.info("[fxrt] bridge rc=%s — reload chain chủ động (set_chain juce) → respawn đọc chain mới", rc) elif self._restarts >= FXRT_MAX_RESTARTS: _log.error("[fxrt] bridge crash rc=%s — restart tối đa %d lần đã đạt, bỏ", rc, FXRT_MAX_RESTARTS) return else: self._restarts += 1 _log.warning("[fxrt] bridge chết rc=%s — auto restart #%d", rc, self._restarts) try: self._spawn() except Exception as e: _log.error("[fxrt] respawn lỗi: %s", e) return # Chờ READY (thread riêng — không block API); crash lần nữa → vòng # ngoài restart tiếp (đếm _restarts đã tăng → dừng đúng giới hạn). deadline = time.time() + 8.0 while time.time() < deadline: time.sleep(0.1) with self._lock: if self._closing or self.proc is None: return if self.proc.poll() is not None: break if self.h.state == FXRT_STATE_READY: break # Registry toàn tiến trình (desktop app 1 user). _SESSIONS = {} _SESSIONS_LOCK = threading.Lock() def _norm_plugin_path(p): """Chuẩn hóa path plugin để so khớp (Windows case-insensitive).""" p = (p or "").strip() if not p: return "" return os.path.normcase(os.path.normpath(p)) # M4: mapping plugin path -> (session_id, slot) — GUI native đang mở (editor # instance DSP trong session fxrt). /fx-gui/{close,hide,show} chỉ nhận path → # tìm session qua registry này. Session chết → entry dọn khi dùng/stop. _FXRT_EDITORS = {} _FXRT_EDITORS_LOCK = threading.Lock() def editor_register(plugin_path, session_id, slot): key = _norm_plugin_path(plugin_path) if not key: return with _FXRT_EDITORS_LOCK: _FXRT_EDITORS[key] = (session_id, int(slot)) def editor_find(plugin_path): key = _norm_plugin_path(plugin_path) with _FXRT_EDITORS_LOCK: return _FXRT_EDITORS.get(key) def editor_unregister(plugin_path): key = _norm_plugin_path(plugin_path) with _FXRT_EDITORS_LOCK: _FXRT_EDITORS.pop(key, None) def editor_unregister_session(session_id): with _FXRT_EDITORS_LOCK: for k in [k for k, v in _FXRT_EDITORS.items() if v[0] == session_id]: _FXRT_EDITORS.pop(k, None) 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() editor_unregister_session(session_id) # M4: dọn mapping GUI native return sess is not None def stop_all(): with _SESSIONS_LOCK: ids = list(_SESSIONS.keys()) for sid in ids: stop_session(sid)