scripts: livekit participant (Phase 1) + proper-auth token script
This commit is contained in:
1 parent
df4d00ee29
commit
7fabeeef91
2 files changed
+207
No files matched your search
@@ -0,0 +1,188 @@
|
||||
#!/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()
|
||||
Reference in new issue
Block a user