Compare commits
5
Commits
27327e94b0
...
master
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ac4571e751 | ||
|
|
9f7b573d01 | ||
|
|
7fabeeef91 | ||
|
|
df4d00ee29 | ||
|
|
890970de81 |
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()
|
||||
Executable
+19
@@ -0,0 +1,19 @@
|
||||
#!/usr/bin/env bash
|
||||
# Proper-auth LiveKit token for the chaos-prime agent (MatrixRTC flow).
|
||||
# Usage: livekit_token.sh [room_id] (defaults to the chaos-prime DM)
|
||||
# Requires MATRIX_PASSWORD in env (see .matrix-env). Prints the JWT to stdout.
|
||||
set -a
|
||||
source /home/pierre/symbol-inferral/.matrix-env
|
||||
set +a
|
||||
HS=https://matrix.chaosmith.systems
|
||||
ROOM="${1:-!vp_58NnFUUE4dmniRdrFIuNbMaaqLrQLOzcC3UUpAi8:chaosmith.systems}"
|
||||
|
||||
TOKEN=$(curl -s -m 10 -X POST "$HS/_matrix/client/v3/login" \
|
||||
-d "{\"type\":\"m.login.password\",\"user\":\"chaos-prime\",\"password\":\"$MATRIX_PASSWORD\"}" | jq -r .access_token)
|
||||
USER_ID=$(curl -s -m 10 "$HS/_matrix/client/v3/account/whoami?access_token=$TOKEN" | jq -r .user_id)
|
||||
OID=$(curl -s -m 10 -X POST "$HS/_matrix/client/v3/user/$USER_ID/openid/request_token?access_token=$TOKEN" -d "{}")
|
||||
BODY=$(jq -nc --argjson oid "$OID" --arg room "$ROOM" --arg uid "$USER_ID" \
|
||||
'{room_id:$room, slot_id:"0", openid_token:$oid,
|
||||
member:{id:"chaos-prime-agent", claimed_user_id:$uid, claimed_device_id:"CHAOSPRIMEAGENT"}}')
|
||||
curl -s -m 20 -X POST "https://livekit.chaosmith.systems/get_token" \
|
||||
-H "Content-Type: application/json" --data-binary "$BODY" | jq -r .jwt
|
||||
+3
-1
@@ -120,7 +120,9 @@ pub fn downloadFile(alloc: std.mem.Allocator, io: Io, token: []const u8, mxc: []
|
||||
const slash = std.mem.indexOfScalar(u8, rest, '/') orelse return error.BadMxc;
|
||||
const server = rest[0..slash];
|
||||
const media_id = rest[slash + 1 ..];
|
||||
const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/media/v3/download/{s}/{s}?allow_redirect=true", .{ HOMESERVER, server, media_id });
|
||||
// MSC3916: homeserver has unauthenticated media disabled — use the
|
||||
// authed client v1 endpoint (v3 download returns M_FORBIDDEN).
|
||||
const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v1/media/download/{s}/{s}", .{ HOMESERVER, server, media_id });
|
||||
defer alloc.free(url);
|
||||
const resp = try httpRaw(alloc, io, "GET", url, token, "", null);
|
||||
// Error responses are small JSON with errcode — detect and report.
|
||||
|
||||
+23
-3
@@ -201,6 +201,8 @@ const PROGRESS_ARG_CHARS = 120;
|
||||
const PROGRESS_RET_CHARS = 400;
|
||||
/// Char cap for the bubble body (edits re-send the full body every time).
|
||||
const PROGRESS_MAX_CHARS = 3500;
|
||||
/// Edit throttling: Element X chokes on long m.replace edit chains.
|
||||
const PROGRESS_MAX_EDITS = 12;
|
||||
|
||||
const Progress = struct {
|
||||
alloc: std.mem.Allocator,
|
||||
@@ -212,6 +214,8 @@ const Progress = struct {
|
||||
reply_event_id: ?[]u8 = null, // the streamed-reply message (edited live)
|
||||
reply: std.ArrayListUnmanaged(u8) = .empty, // accumulated reply text
|
||||
replied: bool = false, // true once the full reply was edited in
|
||||
last_edit_ms: i64 = 0,
|
||||
edits_done: usize = 0,
|
||||
|
||||
fn deinit(self: *Progress) void {
|
||||
if (self.event_id) |e| self.alloc.free(e);
|
||||
@@ -309,6 +313,14 @@ const Progress = struct {
|
||||
}
|
||||
|
||||
fn push(self: *Progress) void {
|
||||
// Throttle: Element X chokes on long m.replace edit chains, so cap
|
||||
// total edits per bubble and rate-limit to ~1 per 8 seconds. The
|
||||
// final state is always flushed by the caller (flush()).
|
||||
const now_ms: i64 = @intCast(@divTrunc(std.Io.Clock.now(.real, self.io).nanoseconds, std.time.ns_per_ms));
|
||||
if (self.edits_done >= PROGRESS_MAX_EDITS) return;
|
||||
if (now_ms - self.last_edit_ms < 8000 and self.edits_done > 0) return;
|
||||
self.last_edit_ms = now_ms;
|
||||
self.edits_done += 1;
|
||||
const id = self.event_id orelse {
|
||||
std.debug.print("DBG push: no event_id, blocks={d}\n", .{self.blocks.items.len});
|
||||
return;
|
||||
@@ -360,7 +372,9 @@ fn onProgressEvent(ctx: ?*anyopaque, ev: letta.Event) void {
|
||||
p.ret(out);
|
||||
},
|
||||
.chunk => |c| {
|
||||
defer p.alloc.free(c);
|
||||
// NOTE: chunk/step payloads are BORROWED — handleLine frees them
|
||||
// via defer (same contract as the overlay UI). Freeing here too
|
||||
// is a double free (heap corruption, silently lost reply tails).
|
||||
_ = p.replyChunk(c, false);
|
||||
},
|
||||
.step => {}, // legacy pre-formatted lines: superseded by structured
|
||||
@@ -547,13 +561,19 @@ pub fn main(init: std.process.Init) !void {
|
||||
if (!std.mem.eql(u8, ev.sender, peer)) continue;
|
||||
if (ev.body.len == 0) continue;
|
||||
|
||||
var ts_buf: [16]u8 = undefined;
|
||||
var ts_buf: [24]u8 = undefined;
|
||||
const ts_str: []const u8 = blk: {
|
||||
if (ev.ts <= 0) break :blk "";
|
||||
const secs: u64 = @intCast(@divTrunc(ev.ts, 1000));
|
||||
const epoch = std.time.epoch.EpochSeconds{ .secs = secs };
|
||||
const day = epoch.getEpochDay();
|
||||
const year_day = day.calculateYearDay();
|
||||
const month_day = year_day.calculateMonthDay();
|
||||
const sod = epoch.getDaySeconds();
|
||||
break :blk std.fmt.bufPrint(&ts_buf, " {d:0>2}:{d:0>2}Z", .{ sod.getHoursIntoDay(), sod.getMinutesIntoHour() }) catch "";
|
||||
break :blk std.fmt.bufPrint(&ts_buf, " {d:0>4}-{d:0>2}-{d:0>2} {d:0>2}:{d:0>2}Z", .{
|
||||
year_day.year, month_day.month.numeric(), month_day.day_index + 1,
|
||||
sod.getHoursIntoDay(), sod.getMinutesIntoHour(),
|
||||
}) catch "";
|
||||
};
|
||||
|
||||
// Attachment: download locally and forward the path to the agent.
|
||||
|
||||
Reference in new issue
Block a user