Compare commits

..
5 Commits
Author SHA1 Message Date
pierre ac4571e751 harness: full-date timestamps on forwarded matrix messages (Pierre De Lancre) 2026-10-02 02:53:22 +03:00
pierre 9f7b573d01 matrix: authed v1 media download — homeserver has unauthenticated media disabled (MSC3916) (Pierre De Lancre) 2026-10-02 02:41:12 +03:00
pierre 7fabeeef91 scripts: livekit participant (Phase 1) + proper-auth token script 2026-09-26 03:53:30 +03:00
pierre df4d00ee29 harness: throttle bubble edits — cap 12/bubble, min 8s interval
Element X chokes loading rooms with long m.replace edit chains (the tool
bubble got an edit per tool call — dozens per turn). Throttle pushes; the
final state still flushes via the reply/closing path.
2026-09-24 18:09:49 +03:00
pierreandLetta Code 890970de81 harness: fix double free of streamed chunk text (chunk payloads are borrowed per letta.zig contract)
onProgressEvent freed .chunk text that handleLine also frees via defer —
double free corrupted the heap mid-turn; DebugAllocator caught it and
the reply tail was silently lost (delivered messages missing their final
characters while the CLI transcript held the full text).

👾 Generated with [Letta Code](https://letta.com)

Co-Authored-By: Letta Code <noreply@letta.com>
2026-09-21 16:10:20 +03:00
4 changed files with 233 additions and 4 deletions

No files matched your search

+188
View File
@@ -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()
+19
View File
@@ -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
View File
@@ -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
View File
@@ -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.