#!/usr/bin/env python3 """LiveKit participant for chaos-prime (Phase 1, Sept 23). Joins the MatrixRTC LiveKit room (token via scripts/livekit_token.sh) and publishes an audio track. Manual-trigger for now; m.call.member watch later. Modes: --wav FILE publish FILE (wav) as an audio track --say "text" synthesize via MeloTTS -> RVC (Yuliko voice), then publish --listen SECONDS join and print received audio track events (no publish) Usage: .venv-livekit/bin/python scripts/livekit_participant.py --say "hello pierre" """ import argparse import os import subprocess import sys import wave import tempfile import asyncio HERE = os.path.dirname(os.path.abspath(__file__)) TOKEN_SH = os.path.join(HERE, "livekit_token.sh") MELO_PY = os.path.expanduser("~/Projects/MeloTTS/.venv/bin/python") MELO_SERVER = os.path.join(HERE, "melo_server.py") RVC_ROOT = "/mnt/engram/RVC1006Nvidia" RVC_MODEL = "yuliko_proto.pth" RVC_INDEX = os.path.join(RVC_ROOT, "logs/yuliko_proto/added_IVF401_Flat_nprobe_1_yuliko_proto_v2.index") DEFAULT_ROOM = "!vp_58NnFUUE4dmniRdrFIuNbMaaqLrQLOzcC3UUpAi8:chaosmith.systems" def get_jwt(room: str) -> str: out = subprocess.run(["bash", TOKEN_SH, room], capture_output=True, text=True, check=True) jwt = out.stdout.strip() if not jwt or jwt == "null": raise RuntimeError(f"token script returned nothing: {out.stderr[:500]}") return jwt def synth_yuliko(text: str) -> str: """text -> MeloTTS PCM -> RVC (yuliko) -> wav path. Blocking, ~1-2 min.""" # 1) Melo synth via the long-running server protocol (v2, raw PCM) melo = subprocess.Popen( [MELO_PY, MELO_SERVER], stdin=subprocess.PIPE, stdout=subprocess.PIPE ) # Server handshake: READY on startup, then per-request WAV header + PCM. # Text must be sent BEFORE reading the header (server emits header only # after synthesizing). melo.stdin.write((text + "\n").encode()) melo.stdin.flush() first = melo.stdout.readline().decode().strip() if first == "READY": first = melo.stdout.readline().decode().strip() header = first.split() assert header and header[0] == "WAV", f"melo protocol error: {header}" n_bytes, rate = int(header[1]), int(header[2]) pcm = b"" while len(pcm) < n_bytes: chunk = melo.stdout.read(n_bytes - len(pcm)) if not chunk: break pcm += chunk melo.stdin.close() melo.wait(timeout=30) tmp_in = tempfile.NamedTemporaryFile(suffix=".wav", delete=False) with wave.open(tmp_in.name, "wb") as w: w.setnchannels(1) w.setsampwidth(2) w.setframerate(rate) w.writeframes(pcm) # 2) RVC convert (run from RVC root so relative asset paths resolve). # infer_batch_rvc works on DIRECTORIES: input_path = dir of wavs, # opt_path = output dir. in_dir = tempfile.mkdtemp(prefix="rvc_in_") out_dir = tempfile.mkdtemp(prefix="rvc_out_") os.rename(tmp_in.name, os.path.join(in_dir, "in.wav")) subprocess.run( [ MELO_PY, "tools/infer_batch_rvc.py", "--f0up_key", "0", "--input_path", in_dir, "--index_path", RVC_INDEX, "--f0method", "rmvpe", "--opt_path", out_dir, "--model_name", RVC_MODEL, "--index_rate", "0.66", "--device", "cuda:0", "--is_half", "True", "--filter_radius", "3", "--rms_mix_rate", "1", "--protect", "0.33", ], cwd=RVC_ROOT, check=True, capture_output=True, ) result = os.path.join(out_dir, "in.wav") return result def wav_source(path: str): """Blocking wav -> AudioSource generator (16kHz mono s16).""" from livekit import rtc with wave.open(path, "rb") as w: rate, n = w.getframerate(), w.getnframes() frames = w.readframes(n) source = rtc.AudioSource(rate, 1) frame_bytes = rate * 2 // 50 # 20ms chunks async def _run(): for i in range(0, len(frames), frame_bytes): chunk = frames[i:i + frame_bytes] f = rtc.AudioFrame( data=chunk, sample_rate=rate, num_channels=1, samples_per_channel=len(chunk) // 2, ) await source.capture_frame(f) await asyncio.sleep(0.02) return source, _run async def amain(args): from livekit import rtc url = "wss://livekit.chaosmith.systems" jwt = get_jwt(args.room) room = rtc.Room() @room.on("track_subscribed") def on_track(track, pub, participant): print(f"[participant] subscribed: {track.kind} from {participant.identity}") if track.kind == rtc.TrackKind.KIND_AUDIO and args.listen: stream = rtc.AudioStream(track) async def _drain(): async for ev in stream: pass # Phase 2: feed whisper STT here asyncio.ensure_future(_drain()) await room.connect(url, jwt) print(f"[participant] connected to {args.room} as {room.local_participant.identity}") published = None if args.wav: src_path = args.wav elif args.say: print("[participant] synthesizing Yuliko voice...") src_path = synth_yuliko(args.say) print(f"[participant] synthesized: {src_path}") else: src_path = None if src_path: source, runner = wav_source(src_path) track = rtc.LocalAudioTrack.create_audio_track("chaos-prime-voice", source) opts = rtc.TrackPublishOptions() opts.source = rtc.TrackSource.SOURCE_MICROPHONE published = await room.local_participant.publish_track(track, opts) print("[participant] audio track published") await runner() if args.listen: print(f"[participant] listening for {args.listen}s...") await asyncio.sleep(args.listen) elif published: await asyncio.sleep(2) # let the tail flush await room.disconnect() print("[participant] disconnected") def main(): p = argparse.ArgumentParser() p.add_argument("--room", default=DEFAULT_ROOM) p.add_argument("--wav") p.add_argument("--say") p.add_argument("--listen", type=int, default=0) args = p.parse_args() if not (args.wav or args.say or args.listen): p.error("one of --wav / --say / --listen required") asyncio.run(amain(args)) if __name__ == "__main__": main()