Compare commits
21
Commits
187bac9aa4
...
master
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ac4571e751 | ||
|
|
9f7b573d01 | ||
|
|
7fabeeef91 | ||
|
|
df4d00ee29 | ||
|
|
890970de81 | ||
|
|
27327e94b0 | ||
|
|
81fcbabefc | ||
|
|
8d89e3baeb | ||
|
|
58a66700c6 | ||
|
|
65e2ee083e | ||
|
|
b958eaab4c | ||
|
|
bf40b50f52 | ||
|
|
2f0890a61c | ||
|
|
861ac201cb | ||
|
|
cbe99e89d0 | ||
|
|
e212e3e120 | ||
|
|
dcf6d8df75 | ||
|
|
9f1a839a22 | ||
|
|
77f3150af4 | ||
|
|
a6d4e29589 | ||
|
|
da50dea0d1 |
No files matched your search
+5
-4
@@ -1,4 +1,5 @@
|
||||
zig-pkg
|
||||
zig-out
|
||||
.letta
|
||||
.zig-cache
|
||||
.matrix-env
|
||||
.letta/
|
||||
.zig-cache/
|
||||
zig-out/
|
||||
zig-pkg/
|
||||
@@ -58,8 +58,28 @@ pub fn build(b: *std.Build) void {
|
||||
|
||||
configureQtExeRootModule(b, exe, .{}) catch @panic("Qt configuration failed");
|
||||
|
||||
exe.use_llvm = true; // native lld chokes on GCC16 .sframe relocations
|
||||
b.installArtifact(exe);
|
||||
|
||||
// Matrix bridge test harness (independent of the tray daemon).
|
||||
const harness = b.addExecutable(.{
|
||||
.name = "matrix_harness",
|
||||
.root_module = b.createModule(.{
|
||||
.root_source_file = b.path("src/matrix_harness.zig"),
|
||||
.target = target,
|
||||
.optimize = optimize,
|
||||
.link_libc = true,
|
||||
.imports = &.{
|
||||
.{ .name = "libqt6zig", .module = qt6zig.module("libqt6zig") },
|
||||
},
|
||||
}),
|
||||
});
|
||||
for (qt_libraries) |lib|
|
||||
harness.root_module.linkLibrary(qt6zig.artifact(lib));
|
||||
configureQtExeRootModule(b, harness, .{}) catch @panic("Qt configuration failed");
|
||||
harness.use_llvm = true;
|
||||
b.installArtifact(harness);
|
||||
|
||||
const run_step = b.step("run", "Run inferon");
|
||||
const run_cmd = b.addRunArtifact(exe);
|
||||
run_step.dependOn(&run_cmd.step);
|
||||
|
||||
@@ -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
|
||||
+71
-15
@@ -3,6 +3,11 @@ const Io = std.Io;
|
||||
|
||||
const LETTA_PATH = "/home/pierre/.bun/bin/letta";
|
||||
|
||||
/// Conversation the matrix bridge routes messages into. Pinned to a specific
|
||||
/// conversation ID (not "default") so a tainted/default conversation can be
|
||||
/// swapped without touching the bridge.
|
||||
pub const CONVERSATION_ID = "local-conv-487";
|
||||
|
||||
/// Runs `letta -n <agent> -p <prompt>`, captures everything it writes to
|
||||
/// stdout, and returns it as an owned slice (caller frees it).
|
||||
pub fn infer(
|
||||
@@ -11,8 +16,10 @@ pub fn infer(
|
||||
agent: []const u8,
|
||||
prompt: []const u8,
|
||||
) ![]u8 {
|
||||
_ = agent; // unused: pinned conversation implies the agent on the CLI side
|
||||
var child = try std.process.spawn(io, .{
|
||||
.argv = &.{ LETTA_PATH, "--agent", agent, "-p", prompt, "--conversation", "default", "--toolset", "default" },
|
||||
// Pinned conversation: --agent must be omitted for non-default conversations.
|
||||
.argv = &.{ LETTA_PATH, "-p", prompt, "--conversation", CONVERSATION_ID, "--toolset", "default" },
|
||||
.stdout = .pipe,
|
||||
.stderr = .inherit, // letta's error output goes straight to your terminal
|
||||
});
|
||||
@@ -47,8 +54,12 @@ pub fn infer(
|
||||
// ------------------------------------------------------------- streaming ---
|
||||
|
||||
pub const Event = union(enum) {
|
||||
step: []const u8, // tool call / return, pre-formatted line
|
||||
step: []const u8, // tool call / return, pre-formatted line (legacy)
|
||||
chunk: []const u8, // streamed reply text
|
||||
/// Structured tool event (caller frees both fields).
|
||||
tool_call: struct { name: []u8, args: []u8 },
|
||||
/// Structured tool output (caller frees).
|
||||
tool_return: []u8,
|
||||
};
|
||||
|
||||
/// Like infer(), but spawns letta with --output-format stream-json and calls
|
||||
@@ -109,8 +120,12 @@ pub fn streamInferConv(
|
||||
var line_len: usize = 0;
|
||||
var buf: [16 * 1024]u8 = undefined;
|
||||
|
||||
var stream_ok = false; // did we see a proper result/completion line?
|
||||
while (true) {
|
||||
const n = stdout_file.readStreaming(io, &.{&buf}) catch break;
|
||||
const n = stdout_file.readStreaming(io, &.{&buf}) catch |e| {
|
||||
std.debug.print("inferon/letta: stdout read error {s} — reply may be TRUNCATED\n", .{@errorName(e)});
|
||||
break;
|
||||
};
|
||||
if (n == 0) break;
|
||||
raw_tail.appendSlice(allocator, buf[0..n]) catch {};
|
||||
if (raw_tail.items.len > 4096) {
|
||||
@@ -122,6 +137,7 @@ pub fn streamInferConv(
|
||||
if (ch == '\n') {
|
||||
if (line_len > 0) {
|
||||
handleLine(allocator, line_buf[0..line_len], ctx, onEvent, &reply);
|
||||
if (std.mem.indexOf(u8, line_buf[0..line_len], "\"type\":\"result\"") != null) stream_ok = true;
|
||||
line_len = 0;
|
||||
}
|
||||
} else if (line_len < line_buf.len) {
|
||||
@@ -130,8 +146,21 @@ pub fn streamInferConv(
|
||||
}
|
||||
}
|
||||
}
|
||||
// CRITICAL: the last line may arrive WITHOUT a trailing newline (EOF
|
||||
// right after the final JSON). It is still a complete line — dropping
|
||||
// it silently truncated replies (the text of the last assistant message
|
||||
// or the result line was lost).
|
||||
if (line_len > 0) {
|
||||
handleLine(allocator, line_buf[0..line_len], ctx, onEvent, &reply);
|
||||
if (std.mem.indexOf(u8, line_buf[0..line_len], "\"type\":\"result\"") != null) stream_ok = true;
|
||||
}
|
||||
_ = child.wait(io) catch {};
|
||||
|
||||
if (!stream_ok) {
|
||||
std.debug.print("inferon/letta: stream ended WITHOUT result line — reply likely TRUNCATED ({d} bytes collected). tail:\n{s}\n", .{ reply.items.len, raw_tail.items });
|
||||
}
|
||||
std.debug.print("inferon/letta: stream complete, reply {d} bytes, result_line={}\n", .{ reply.items.len, stream_ok });
|
||||
|
||||
if (reply.items.len == 0 and raw_tail.items.len > 0) {
|
||||
std.debug.print("inferon/letta: EMPTY reply, raw stream tail:\n{s}\n", .{raw_tail.items});
|
||||
}
|
||||
@@ -169,19 +198,22 @@ fn handleLine(
|
||||
|
||||
if (std.mem.eql(u8, mt, "tool_call_message")) {
|
||||
const name = dupeStr(allocator, line, "name") orelse allocator.dupe(u8, "?") catch return;
|
||||
defer allocator.free(name);
|
||||
// Different tools carry their primary arg under different keys;
|
||||
// try the common ones so the summary is never empty.
|
||||
const args = dupeStr(allocator, line, "command") orelse
|
||||
dupeStr(allocator, line, "description") orelse allocator.dupe(u8, "") catch return;
|
||||
defer allocator.free(args);
|
||||
const text = std.fmt.allocPrint(allocator, "> {s} {s}", .{ name, args }) catch return;
|
||||
defer allocator.free(text);
|
||||
onEvent(ctx, .{ .step = text });
|
||||
dupeStr(allocator, line, "file_path") orelse
|
||||
dupeStr(allocator, line, "path") orelse
|
||||
dupeStr(allocator, line, "query") orelse
|
||||
dupeStr(allocator, line, "url") orelse
|
||||
dupeStr(allocator, line, "pattern") orelse
|
||||
dupeStr(allocator, line, "description") orelse allocator.dupe(u8, "") catch {
|
||||
allocator.free(name);
|
||||
return;
|
||||
};
|
||||
onEvent(ctx, .{ .tool_call = .{ .name = name, .args = args } });
|
||||
} else if (std.mem.eql(u8, mt, "tool_return_message")) {
|
||||
const ret = dupeStr(allocator, line, "tool_return") orelse allocator.dupe(u8, "") catch return;
|
||||
defer allocator.free(ret);
|
||||
const text = std.fmt.allocPrint(allocator, " {s}", .{ret}) catch return;
|
||||
defer allocator.free(text);
|
||||
onEvent(ctx, .{ .step = text });
|
||||
onEvent(ctx, .{ .tool_return = ret });
|
||||
} else if (std.mem.eql(u8, mt, "assistant_message")) {
|
||||
const txt = dupeStr(allocator, line, "text") orelse return;
|
||||
defer allocator.free(txt);
|
||||
@@ -191,11 +223,35 @@ fn handleLine(
|
||||
}
|
||||
|
||||
/// Extract "key":"value" (with escape handling), allocated with `allocator`.
|
||||
/// Also matches the escaped variant `\"key\":\"` — nested JSON strings (e.g.
|
||||
/// tool_call arguments) carry their keys escaped on the wire.
|
||||
pub fn dupeStr(allocator: std.mem.Allocator, line: []const u8, key: []const u8) ?[]u8 {
|
||||
var pat_buf: [64]u8 = undefined;
|
||||
const pat = std.fmt.bufPrint(&pat_buf, "\"{s}\":\"", .{key}) catch return null;
|
||||
const start = std.mem.indexOf(u8, line, pat) orelse return null;
|
||||
var it = line[start + pat.len ..];
|
||||
const start = std.mem.indexOf(u8, line, pat) orelse {
|
||||
// Escaped variant: \"key\":\" inside a nested JSON string.
|
||||
var pat2_buf: [70]u8 = undefined;
|
||||
const pat2 = std.fmt.bufPrint(&pat2_buf, "\\\"{s}\\\":\\\"", .{key}) catch return null;
|
||||
const s2 = std.mem.indexOf(u8, line, pat2) orelse return null;
|
||||
const v = dupeValue(allocator, line[s2 + pat2.len ..]) orelse return null;
|
||||
// The escaped variant's closing quote is escaped too, so the value
|
||||
// decoder overruns into the following keys. Cut at the first
|
||||
// unescaped key boundary left in the decoded text.
|
||||
if (std.mem.indexOf(u8, v, "\",\"")) |cut| {
|
||||
const out = allocator.dupe(u8, v[0..cut]) catch {
|
||||
allocator.free(v);
|
||||
return null;
|
||||
};
|
||||
allocator.free(v);
|
||||
return out;
|
||||
}
|
||||
return v;
|
||||
};
|
||||
return dupeValue(allocator, line[start + pat.len ..]);
|
||||
}
|
||||
|
||||
fn dupeValue(allocator: std.mem.Allocator, from: []const u8) ?[]u8 {
|
||||
var it = from;
|
||||
|
||||
var out: std.ArrayList(u8) = .empty;
|
||||
errdefer out.deinit(allocator);
|
||||
|
||||
@@ -202,6 +202,17 @@ fn onLettaEvent(_: ?*anyopaque, ev: letta.Event) void {
|
||||
if (last == '.' or last == '!' or last == '?') flushSentence(false);
|
||||
}
|
||||
},
|
||||
.tool_call => |tc| {
|
||||
defer allocator.free(tc.name);
|
||||
defer allocator.free(tc.args);
|
||||
const line = std.fmt.allocPrint(allocator, "> {s} {s}", .{ tc.name, tc.args }) catch return;
|
||||
pushOverlay(.step, line);
|
||||
},
|
||||
.tool_return => |ret| {
|
||||
defer allocator.free(ret);
|
||||
const line = std.fmt.allocPrint(allocator, " {s}", .{ret}) catch return;
|
||||
pushOverlay(.step, line);
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+559
@@ -0,0 +1,559 @@
|
||||
//! Matrix client (scaffold) — conduwuit @ chaosmith.systems.
|
||||
//!
|
||||
//! HTTP via curl subprocess behind a narrow seam (native TLS swap later).
|
||||
//! Sync loop + send + login. No olm/megolm: rooms must be unencrypted
|
||||
//! (trusted_private_chat without encryption).
|
||||
|
||||
const std = @import("std");
|
||||
const Io = std.Io;
|
||||
const letta = @import("letta.zig");
|
||||
|
||||
pub const HOMESERVER = "https://matrix.chaosmith.systems";
|
||||
pub const DEVICE_NAME = "inferon-harness";
|
||||
|
||||
// ------------------------------------------------------------- transport ---
|
||||
|
||||
/// HTTP via curl subprocess. Returns body (allocated; caller frees).
|
||||
fn http(alloc: std.mem.Allocator, io: Io, method: []const u8, url: []const u8, token: ?[]const u8, json_body: ?[]const u8) ![]u8 {
|
||||
var argv: std.ArrayList([]const u8) = .empty;
|
||||
defer argv.deinit(alloc);
|
||||
try argv.appendSlice(alloc, &.{ "curl", "-sS", "-m", "25", "-X", method });
|
||||
if (token) |t| {
|
||||
// NOTE: no defer-free — h must outlive the spawn below (argv holds it).
|
||||
// Leaks one small string per call; fine for the harness, revisit with
|
||||
// the native-TLS rewrite.
|
||||
const h = try std.fmt.allocPrint(alloc, "Authorization: Bearer {s}", .{t});
|
||||
try argv.appendSlice(alloc, &.{ "-H", h });
|
||||
}
|
||||
if (json_body) |b| {
|
||||
try argv.appendSlice(alloc, &.{ "-H", "Content-Type: application/json", "-d", b });
|
||||
}
|
||||
try argv.append(alloc, url);
|
||||
|
||||
var child = try std.process.spawn(io, .{
|
||||
.argv = argv.items,
|
||||
.stdin = .ignore,
|
||||
.stdout = .pipe,
|
||||
.stderr = .ignore,
|
||||
});
|
||||
const out = child.stdout.?;
|
||||
var buf: std.ArrayList(u8) = .empty;
|
||||
errdefer buf.deinit(alloc);
|
||||
var tmp: [16 * 1024]u8 = undefined;
|
||||
while (true) {
|
||||
const got = out.readStreaming(io, &.{&tmp}) catch break;
|
||||
if (got == 0) break;
|
||||
try buf.appendSlice(alloc, tmp[0..got]);
|
||||
}
|
||||
_ = child.wait(io) catch {};
|
||||
return buf.toOwnedSlice(alloc);
|
||||
}
|
||||
|
||||
// ------------------------------------------------------------- api calls ---
|
||||
|
||||
pub const Session = struct { token: []u8, user_id: []u8 }; // caller frees both
|
||||
|
||||
/// HTTP with a raw (non-JSON) body + explicit content type — media upload.
|
||||
/// Returns body (allocated; caller frees).
|
||||
fn httpRaw(alloc: std.mem.Allocator, io: Io, method: []const u8, url: []const u8, token: ?[]const u8, content_type: []const u8, file_path: ?[]const u8) ![]u8 {
|
||||
// Arena for argv/header scratch — freed wholesale at scope exit (unlike
|
||||
// http(), which leaks its header string per call).
|
||||
var arena_state = std.heap.ArenaAllocator.init(alloc);
|
||||
defer arena_state.deinit();
|
||||
const a = arena_state.allocator();
|
||||
|
||||
var argv: std.ArrayList([]const u8) = .empty;
|
||||
defer argv.deinit(a); // buffer lives in the arena — no-op free
|
||||
try argv.appendSlice(a, &.{ "curl", "-sS", "-m", "120", "-X", method });
|
||||
if (token) |t| {
|
||||
const h = try std.fmt.allocPrint(a, "Authorization: Bearer {s}", .{t});
|
||||
try argv.appendSlice(a, &.{ "-H", h });
|
||||
}
|
||||
if (content_type.len > 0) {
|
||||
const ct = try std.fmt.allocPrint(a, "Content-Type: {s}", .{content_type});
|
||||
try argv.appendSlice(a, &.{ "-H", ct });
|
||||
}
|
||||
if (file_path) |p| {
|
||||
const d = try std.fmt.allocPrint(a, "@{s}", .{p});
|
||||
try argv.appendSlice(a, &.{ "--data-binary", d });
|
||||
}
|
||||
try argv.append(a, url);
|
||||
|
||||
var child = try std.process.spawn(io, .{
|
||||
.argv = argv.items,
|
||||
.stdin = .ignore,
|
||||
.stdout = .pipe,
|
||||
.stderr = .ignore,
|
||||
});
|
||||
const out = child.stdout.?;
|
||||
var buf: std.ArrayList(u8) = .empty;
|
||||
errdefer buf.deinit(alloc);
|
||||
var tmp: [16 * 1024]u8 = undefined;
|
||||
while (true) {
|
||||
const got = out.readStreaming(io, &.{&tmp}) catch break;
|
||||
if (got == 0) break;
|
||||
try buf.appendSlice(alloc, tmp[0..got]);
|
||||
}
|
||||
_ = child.wait(io) catch {};
|
||||
return buf.toOwnedSlice(alloc);
|
||||
}
|
||||
|
||||
/// Upload a local file to the homeserver media store.
|
||||
/// Returns the mxc:// URI (allocated; caller frees).
|
||||
pub fn uploadFile(alloc: std.mem.Allocator, io: Io, token: []const u8, path: []const u8, filename: []const u8, content_type: []const u8) ![]u8 {
|
||||
const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/media/v3/upload?filename={s}", .{ HOMESERVER, filename });
|
||||
defer alloc.free(url);
|
||||
const resp = try httpRaw(alloc, io, "POST", url, token, content_type, path);
|
||||
defer alloc.free(resp);
|
||||
const uri = letta.jsonStr(alloc, resp, 0, "content_uri") orelse {
|
||||
std.debug.print("matrix: upload failed, resp: {s}\n", .{resp[0..@min(resp.len, 400)]});
|
||||
return error.UploadFailed;
|
||||
};
|
||||
return uri.val;
|
||||
}
|
||||
|
||||
/// Download media by mxc:// URI. Returns raw bytes (allocated; caller frees).
|
||||
pub fn downloadFile(alloc: std.mem.Allocator, io: Io, token: []const u8, mxc: []const u8) ![]u8 {
|
||||
// mxc://server/mediaId
|
||||
if (!std.mem.startsWith(u8, mxc, "mxc://")) return error.BadMxc;
|
||||
const rest = mxc["mxc://".len..];
|
||||
const slash = std.mem.indexOfScalar(u8, rest, '/') orelse return error.BadMxc;
|
||||
const server = rest[0..slash];
|
||||
const media_id = rest[slash + 1 ..];
|
||||
// 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.
|
||||
if (resp.len > 0 and resp[0] == '{') {
|
||||
if (letta.jsonStr(alloc, resp, 0, "errcode") != null) {
|
||||
std.debug.print("matrix: download failed, resp: {s}\n", .{resp[0..@min(resp.len, 400)]});
|
||||
alloc.free(resp);
|
||||
return error.DownloadFailed;
|
||||
}
|
||||
}
|
||||
return resp;
|
||||
}
|
||||
|
||||
/// Send a file event (m.file / m.image / m.audio / m.video).
|
||||
/// Returns event id (allocated; caller frees).
|
||||
pub fn sendFile(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []const u8, msgtype: []const u8, mxc: []const u8, filename: []const u8, mimetype: []const u8, size: u64) ![]u8 {
|
||||
const ts = Io.Clock.now(.real, io).nanoseconds;
|
||||
const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/send/m.room.message/file-{d}", .{ HOMESERVER, room, ts });
|
||||
defer alloc.free(url);
|
||||
|
||||
var payload: std.ArrayList(u8) = .empty;
|
||||
defer payload.deinit(alloc);
|
||||
try payload.appendSlice(alloc, "{\"msgtype\":\"");
|
||||
try payload.appendSlice(alloc, msgtype);
|
||||
try payload.appendSlice(alloc, "\",\"body\":\"");
|
||||
var name_esc: std.ArrayList(u8) = .empty;
|
||||
defer name_esc.deinit(alloc);
|
||||
try jsonEscape(alloc, &name_esc, filename);
|
||||
try payload.appendSlice(alloc, name_esc.items);
|
||||
try payload.appendSlice(alloc, "\",\"url\":\"");
|
||||
try payload.appendSlice(alloc, mxc);
|
||||
try payload.appendSlice(alloc, "\",\"info\":{\"mimetype\":\"");
|
||||
try payload.appendSlice(alloc, mimetype);
|
||||
try payload.appendSlice(alloc, "\",\"size\":");
|
||||
var sz_buf: [24]u8 = undefined;
|
||||
const sz = std.fmt.bufPrint(&sz_buf, "{d}", .{size}) catch unreachable;
|
||||
try payload.appendSlice(alloc, sz);
|
||||
try payload.appendSlice(alloc, "}}");
|
||||
|
||||
const resp = try http(alloc, io, "PUT", url, token, payload.items);
|
||||
defer alloc.free(resp);
|
||||
const ev = letta.jsonStr(alloc, resp, 0, "event_id") orelse {
|
||||
std.debug.print("matrix: sendFile failed, resp: {s}\n", .{resp[0..@min(resp.len, 400)]});
|
||||
return error.SendFailed;
|
||||
};
|
||||
return ev.val;
|
||||
}
|
||||
|
||||
pub fn login(alloc: std.mem.Allocator, io: Io, user: []const u8, password: []const u8) !Session {
|
||||
const body = try std.fmt.allocPrint(alloc,
|
||||
\\{{"type":"m.login.password","identifier":{{"type":"m.id.user","user":"{s}"}},"password":"{s}","initial_device_display_name":"{s}"}}
|
||||
, .{ user, password, DEVICE_NAME });
|
||||
defer alloc.free(body);
|
||||
const resp = try http(alloc, io, "POST", HOMESERVER ++ "/_matrix/client/v3/login", null, body);
|
||||
defer alloc.free(resp);
|
||||
const tok = letta.jsonStr(alloc, resp, 0, "access_token") orelse return error.LoginFailed;
|
||||
errdefer alloc.free(tok.val);
|
||||
const uid = letta.jsonStr(alloc, resp, 0, "user_id") orelse return error.LoginFailed;
|
||||
return .{ .token = tok.val, .user_id = uid.val };
|
||||
}
|
||||
|
||||
/// Send a text message. Returns event id (allocated; caller frees).
|
||||
/// Escape text for embedding inside a JSON string literal.
|
||||
fn jsonEscape(alloc: std.mem.Allocator, out: *std.ArrayList(u8), text: []const u8) !void {
|
||||
for (text) |c| {
|
||||
switch (c) {
|
||||
'"' => try out.appendSlice(alloc, "\\\""),
|
||||
'\\' => try out.appendSlice(alloc, "\\\\"),
|
||||
'\n' => try out.appendSlice(alloc, "\\n"),
|
||||
'\r', '\t' => try out.append(alloc, ' '),
|
||||
// Any other control char is invalid raw in a JSON string —
|
||||
// escape as \u00XX so the server never rejects the send.
|
||||
// (Ranges disjoint from \t \n \r handled above.)
|
||||
0x00...0x08, 0x0b, 0x0c, 0x0e...0x1f => {
|
||||
var buf: [6]u8 = undefined;
|
||||
const s = std.fmt.bufPrint(&buf, "\\u{x:0>4}", .{c}) catch unreachable;
|
||||
try out.appendSlice(alloc, s);
|
||||
},
|
||||
else => try out.append(alloc, c),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Minimal markdown → HTML for formatted_body: **bold**, *italic*,
|
||||
/// `code`, ``` fences, paragraphs. Element only renders styling via
|
||||
/// org.matrix.custom.html — plain m.text bodies show raw asterisks.
|
||||
fn markdownToHtml(alloc: std.mem.Allocator, md: []const u8) ![]u8 {
|
||||
var out: std.ArrayList(u8) = .empty;
|
||||
errdefer out.deinit(alloc);
|
||||
var in_code = false;
|
||||
var i: usize = 0;
|
||||
while (i < md.len) {
|
||||
// Fenced code blocks
|
||||
if (std.mem.startsWith(u8, md[i..], "```")) {
|
||||
try out.appendSlice(alloc, if (in_code) "</pre>" else "<pre>");
|
||||
in_code = !in_code;
|
||||
i += 3;
|
||||
// Skip to end of line (language tag when opening).
|
||||
if (!in_code) {
|
||||
while (i < md.len and md[i] != '\n') i += 1;
|
||||
if (i < md.len) i += 1;
|
||||
} else {
|
||||
if (i < md.len and md[i] == '\n') i += 1;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (in_code) {
|
||||
try out.append(alloc, md[i]);
|
||||
i += 1;
|
||||
continue;
|
||||
}
|
||||
switch (md[i]) {
|
||||
'\n' => {
|
||||
try out.appendSlice(alloc, "<br>");
|
||||
i += 1;
|
||||
},
|
||||
'<' => {
|
||||
try out.appendSlice(alloc, "<");
|
||||
i += 1;
|
||||
},
|
||||
'>' => {
|
||||
try out.appendSlice(alloc, ">");
|
||||
i += 1;
|
||||
},
|
||||
'&' => {
|
||||
try out.appendSlice(alloc, "&");
|
||||
i += 1;
|
||||
},
|
||||
'`' => {
|
||||
// Inline code: `...`
|
||||
if (std.mem.indexOfScalarPos(u8, md, i + 1, '`')) |end| {
|
||||
try out.appendSlice(alloc, "<code>");
|
||||
try out.appendSlice(alloc, md[i + 1 .. end]);
|
||||
try out.appendSlice(alloc, "</code>");
|
||||
i = end + 1;
|
||||
} else {
|
||||
try out.append(alloc, '`');
|
||||
i += 1;
|
||||
}
|
||||
},
|
||||
'*' => {
|
||||
if (std.mem.startsWith(u8, md[i..], "**")) {
|
||||
if (std.mem.indexOfPos(u8, md, i + 2, "**")) |end| {
|
||||
try out.appendSlice(alloc, "<strong>");
|
||||
try out.appendSlice(alloc, md[i + 2 .. end]);
|
||||
try out.appendSlice(alloc, "</strong>");
|
||||
i = end + 2;
|
||||
} else {
|
||||
try out.appendSlice(alloc, "**");
|
||||
i += 2;
|
||||
}
|
||||
} else if (std.mem.indexOfScalarPos(u8, md, i + 1, '*')) |end| {
|
||||
try out.appendSlice(alloc, "<em>");
|
||||
try out.appendSlice(alloc, md[i + 1 .. end]);
|
||||
try out.appendSlice(alloc, "</em>");
|
||||
i = end + 1;
|
||||
} else {
|
||||
try out.append(alloc, '*');
|
||||
i += 1;
|
||||
}
|
||||
},
|
||||
else => {
|
||||
try out.append(alloc, md[i]);
|
||||
i += 1;
|
||||
},
|
||||
}
|
||||
}
|
||||
if (in_code) try out.appendSlice(alloc, "</pre>");
|
||||
return out.toOwnedSlice(alloc);
|
||||
}
|
||||
|
||||
|
||||
/// Typing notification. timeout_ms: how long the indicator lasts server-side
|
||||
/// (refresh periodically for long turns). typing=false clears it immediately.
|
||||
pub fn setTyping(alloc: std.mem.Allocator, io: Io, token: []const u8, user_id: []const u8, room: []const u8, typing: bool, timeout_ms: i64) void {
|
||||
var url_buf: std.ArrayListUnmanaged(u8) = .empty;
|
||||
defer url_buf.deinit(alloc);
|
||||
const user_esc = std.fmt.allocPrint(alloc, "{s}", .{user_id}) catch return;
|
||||
defer alloc.free(user_esc);
|
||||
// user ids contain chars that are fine unescaped in a path segment for
|
||||
// conduwuit, but encode the ':' minimally
|
||||
var user_path: std.ArrayListUnmanaged(u8) = .empty;
|
||||
defer user_path.deinit(alloc);
|
||||
for (user_esc) |ch| {
|
||||
if (ch == ':') user_path.appendSlice(alloc, "%3A") catch return else user_path.append(alloc, ch) catch return;
|
||||
}
|
||||
const url = std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/typing/{s}", .{ HOMESERVER, room, user_path.items }) catch return;
|
||||
defer alloc.free(url);
|
||||
var payload: std.ArrayListUnmanaged(u8) = .empty;
|
||||
defer payload.deinit(alloc);
|
||||
const body = if (typing)
|
||||
std.fmt.allocPrint(alloc, "{{\"typing\":true,\"timeout\":{d}}}", .{timeout_ms}) catch return
|
||||
else
|
||||
alloc.dupe(u8, "{\"typing\":false}") catch return;
|
||||
defer alloc.free(body);
|
||||
if (httpRaw(alloc, io, "PUT", url, token, "application/json", body)) |resp| {
|
||||
alloc.free(resp);
|
||||
} else |e| {
|
||||
std.debug.print("matrix_harness: setTyping failed: {s}\n", .{@errorName(e)});
|
||||
}
|
||||
}
|
||||
|
||||
pub fn sendText(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []const u8, body: []const u8) ![]u8 {
|
||||
const ts = Io.Clock.now(.real, io).nanoseconds;
|
||||
const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/send/m.room.message/infr-{d}", .{ HOMESERVER, room, ts });
|
||||
defer alloc.free(url);
|
||||
|
||||
const html = markdownToHtml(alloc, body) catch |e| switch (e) {
|
||||
error.OutOfMemory => return e,
|
||||
};
|
||||
defer alloc.free(html);
|
||||
|
||||
var payload: std.ArrayList(u8) = .empty;
|
||||
defer payload.deinit(alloc);
|
||||
try payload.appendSlice(alloc, "{\"msgtype\":\"m.text\",\"body\":\"");
|
||||
try jsonEscape(alloc, &payload, body);
|
||||
try payload.appendSlice(alloc, "\",\"format\":\"org.matrix.custom.html\",\"formatted_body\":\"");
|
||||
try jsonEscape(alloc, &payload, html);
|
||||
try payload.append(alloc, '"');
|
||||
try payload.append(alloc, '}');
|
||||
|
||||
const resp = try http(alloc, io, "PUT", url, token, payload.items);
|
||||
defer alloc.free(resp);
|
||||
const ev = letta.jsonStr(alloc, resp, 0, "event_id") orelse {
|
||||
// Diagnose silently-rejected sends: dump first 400 bytes of the
|
||||
// server response (errcode/message) and payload head to stderr.
|
||||
std.debug.print("matrix: send failed, resp: {s}\n", .{resp[0..@min(resp.len, 400)]});
|
||||
std.debug.print("matrix: payload head: {s}\n", .{payload.items[0..@min(payload.items.len, 400)]});
|
||||
return error.SendFailed;
|
||||
};
|
||||
return ev.val;
|
||||
}
|
||||
|
||||
/// Edit an existing m.text message in place (m.replace). `orig_event_id` is
|
||||
/// the event being replaced; `new_body` is the full replacement text
|
||||
/// (markdown — rendered via formatted_body, with plain-text fallback).
|
||||
pub fn editText(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []const u8, orig_event_id: []const u8, new_body: []const u8) ![]u8 {
|
||||
const ts = Io.Clock.now(.real, io).nanoseconds;
|
||||
const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/send/m.room.message/infr-edit-{d}", .{ HOMESERVER, room, ts });
|
||||
defer alloc.free(url);
|
||||
|
||||
const html = markdownToHtml(alloc, new_body) catch |e| switch (e) {
|
||||
error.OutOfMemory => return e,
|
||||
};
|
||||
defer alloc.free(html);
|
||||
|
||||
var payload: std.ArrayList(u8) = .empty;
|
||||
defer payload.deinit(alloc);
|
||||
try payload.appendSlice(alloc, "{\"msgtype\":\"m.text\",\"body\":\"* ");
|
||||
try jsonEscape(alloc, &payload, new_body);
|
||||
try payload.appendSlice(alloc, "\",\"m.new_content\":{\"msgtype\":\"m.text\",\"body\":\"");
|
||||
try jsonEscape(alloc, &payload, new_body);
|
||||
try payload.appendSlice(alloc, "\",\"format\":\"org.matrix.custom.html\",\"formatted_body\":\"");
|
||||
try jsonEscape(alloc, &payload, html);
|
||||
try payload.appendSlice(alloc, "\"},\"m.relates_to\":{\"rel_type\":\"m.replace\",\"event_id\":\"");
|
||||
try jsonEscape(alloc, &payload, orig_event_id);
|
||||
try payload.appendSlice(alloc, "\"}}");
|
||||
|
||||
const resp = try http(alloc, io, "PUT", url, token, payload.items);
|
||||
defer alloc.free(resp);
|
||||
const ev = letta.jsonStr(alloc, resp, 0, "event_id") orelse return error.SendFailed;
|
||||
return ev.val;
|
||||
}
|
||||
|
||||
/// Join a room by id (accepts invites).
|
||||
pub fn joinRoom(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []const u8) !void {
|
||||
const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/join", .{ HOMESERVER, room });
|
||||
defer alloc.free(url);
|
||||
const resp = try http(alloc, io, "POST", url, token, "{}");
|
||||
alloc.free(resp);
|
||||
}
|
||||
|
||||
pub const Event = struct {
|
||||
sender: []u8,
|
||||
body: []u8,
|
||||
ts: i64 = 0,
|
||||
/// mxc:// URI when the message carries an attachment (m.file/m.image/…).
|
||||
url: ?[]u8 = null, // allocated; caller frees
|
||||
mimetype: []u8 = "", // allocated; caller frees
|
||||
size: i64 = 0,
|
||||
/// Original attachment filename (distinct from body, which is the caption).
|
||||
filename: []u8 = &.{}, // allocated; caller frees
|
||||
};
|
||||
|
||||
/// One sync poll. Returns message events + next `since` token. Caller frees.
|
||||
/// (Sync first without a since token to establish one, then poll.)
|
||||
///
|
||||
/// Proper JSON parsing (std.json): the previous string-window scraper
|
||||
/// misattributed senders and duplicated events when multiple messages shared
|
||||
/// a response — sender/body could be scraped from different events.
|
||||
pub fn sync(alloc: std.mem.Allocator, io: Io, token: []const u8, since: []const u8, timeout_ms: u32) !struct { events: []Event, next: []u8 } {
|
||||
const url = if (since.len > 0)
|
||||
try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/sync?timeout={d}&since={s}", .{ HOMESERVER, timeout_ms, since })
|
||||
else
|
||||
try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/sync?timeout=0", .{HOMESERVER});
|
||||
defer alloc.free(url);
|
||||
const resp = try http(alloc, io, "GET", url, token, null);
|
||||
defer alloc.free(resp);
|
||||
|
||||
var events: std.ArrayList(Event) = .empty;
|
||||
errdefer events.deinit(alloc);
|
||||
|
||||
var arena_state = std.heap.ArenaAllocator.init(alloc);
|
||||
defer arena_state.deinit();
|
||||
const jalloc = arena_state.allocator();
|
||||
|
||||
var parsed = std.json.parseFromSlice(std.json.Value, jalloc, resp, .{}) catch return error.BadSyncJson;
|
||||
defer parsed.deinit();
|
||||
|
||||
const root = switch (parsed.value) {
|
||||
.object => |o| o,
|
||||
else => return error.BadSyncJson,
|
||||
};
|
||||
|
||||
var next: []u8 = try alloc.dupe(u8, since);
|
||||
if (root.get("next_batch")) |nb| switch (nb) {
|
||||
.string => |s| {
|
||||
alloc.free(next);
|
||||
next = try alloc.dupe(u8, s);
|
||||
},
|
||||
else => {},
|
||||
};
|
||||
|
||||
// Walk rooms.join.*.timeline.events[] — every room, in order.
|
||||
const join = if (root.get("rooms")) |rooms| switch (rooms) {
|
||||
.object => |o| o.get("join") orelse return .{ .events = try events.toOwnedSlice(alloc), .next = next },
|
||||
else => return .{ .events = try events.toOwnedSlice(alloc), .next = next },
|
||||
} else return .{ .events = try events.toOwnedSlice(alloc), .next = next };
|
||||
|
||||
const join_obj = switch (join) {
|
||||
.object => |o| o,
|
||||
else => return .{ .events = try events.toOwnedSlice(alloc), .next = next },
|
||||
};
|
||||
|
||||
var room_it = join_obj.iterator();
|
||||
while (room_it.next()) |room_entry| {
|
||||
const room = switch (room_entry.value_ptr.*) {
|
||||
.object => |o| o,
|
||||
else => continue,
|
||||
};
|
||||
const timeline = if (room.get("timeline")) |t| switch (t) {
|
||||
.object => |o| o,
|
||||
else => continue,
|
||||
} else continue;
|
||||
const evs = if (timeline.get("events")) |e| switch (e) {
|
||||
.array => |a| a,
|
||||
else => continue,
|
||||
} else continue;
|
||||
|
||||
for (evs.items) |ev| {
|
||||
const obj = switch (ev) {
|
||||
.object => |o| o,
|
||||
else => continue,
|
||||
};
|
||||
// Only m.room.message events.
|
||||
const ty = if (obj.get("type")) |t| switch (t) {
|
||||
.string => |s| s,
|
||||
else => continue,
|
||||
} else continue;
|
||||
if (!std.mem.eql(u8, ty, "m.room.message")) continue;
|
||||
|
||||
const sender_raw = if (obj.get("sender")) |s| switch (s) {
|
||||
.string => |v| v,
|
||||
else => continue,
|
||||
} else continue;
|
||||
|
||||
const content = if (obj.get("content")) |c| switch (c) {
|
||||
.object => |o| o,
|
||||
else => continue,
|
||||
} else continue;
|
||||
const body_raw = if (content.get("body")) |b| switch (b) {
|
||||
.string => |v| v,
|
||||
else => continue,
|
||||
} else continue;
|
||||
const ts: i64 = if (obj.get("origin_server_ts")) |t| switch (t) {
|
||||
.integer => |v| v,
|
||||
else => 0,
|
||||
} else 0;
|
||||
|
||||
// Attachment fields (m.file / m.image / m.audio / m.video).
|
||||
var url_copy: ?[]u8 = null;
|
||||
var fname_copy: []u8 = "";
|
||||
if (content.get("filename")) |f| switch (f) {
|
||||
.string => |v| fname_copy = alloc.dupe(u8, v) catch "",
|
||||
else => {},
|
||||
};
|
||||
if (content.get("url")) |u| switch (u) {
|
||||
.string => |v| url_copy = alloc.dupe(u8, v) catch null,
|
||||
else => {},
|
||||
};
|
||||
var mime: []u8 = "";
|
||||
var size: i64 = 0;
|
||||
if (content.get("info")) |inf| switch (inf) {
|
||||
.object => |io_| {
|
||||
if (io_.get("mimetype")) |m| switch (m) {
|
||||
.string => |v| mime = alloc.dupe(u8, v) catch "",
|
||||
else => {},
|
||||
};
|
||||
if (io_.get("size")) |s| switch (s) {
|
||||
.integer => |v| size = v,
|
||||
else => {},
|
||||
};
|
||||
},
|
||||
else => {},
|
||||
};
|
||||
|
||||
const sender_copy = alloc.dupe(u8, sender_raw) catch continue;
|
||||
const body_copy = alloc.dupe(u8, body_raw) catch {
|
||||
alloc.free(sender_copy);
|
||||
continue;
|
||||
};
|
||||
events.append(alloc, .{ .sender = sender_copy, .body = body_copy, .ts = ts, .url = url_copy, .filename = fname_copy, .mimetype = mime, .size = size }) catch {
|
||||
alloc.free(sender_copy);
|
||||
alloc.free(body_copy);
|
||||
};
|
||||
}
|
||||
}
|
||||
return .{ .events = try events.toOwnedSlice(alloc), .next = next };
|
||||
}
|
||||
|
||||
/// Create a direct chat room with `peer`; returns room id (allocated).
|
||||
pub fn createDirect(alloc: std.mem.Allocator, io: Io, token: []const u8, peer: []const u8) ![]u8 {
|
||||
const body = try std.fmt.allocPrint(alloc,
|
||||
\\{{"preset":"trusted_private_chat","invite":["{s}"],"is_direct":true}}
|
||||
, .{peer});
|
||||
defer alloc.free(body);
|
||||
const resp = try http(alloc, io, "POST", HOMESERVER ++ "/_matrix/client/v3/createRoom", token, body);
|
||||
defer alloc.free(resp);
|
||||
const room = letta.jsonStr(alloc, resp, 0, "room_id") orelse {
|
||||
std.debug.print("matrix: createRoom failed, resp: {s}\n", .{resp});
|
||||
return error.CreateFailed;
|
||||
};
|
||||
return room.val;
|
||||
}
|
||||
@@ -0,0 +1,663 @@
|
||||
//! matrix_harness — standalone Matrix bridge test harness.
|
||||
//!
|
||||
//! Logs in as the inferon Matrix account, opens/creates a DM with MATRIX_PEER
|
||||
//! (default @pierre:chaosmith.systems), and loops: incoming peer messages →
|
||||
//! letta CLI (chaos-prime) → text reply. Independent of the tray daemon so
|
||||
//! calling/voice work can be layered on top without touching it.
|
||||
//!
|
||||
//! Env: MATRIX_USER, MATRIX_PASSWORD, MATRIX_PEER, LETTA_MATRIX_AGENT.
|
||||
|
||||
const std = @import("std");
|
||||
const matrix = @import("matrix.zig");
|
||||
const letta = @import("letta.zig");
|
||||
|
||||
fn env_or(name: [:0]const u8, default: []const u8) []const u8 {
|
||||
if (std.c.getenv(name)) |v| {
|
||||
var len: usize = 0;
|
||||
while (v[len] != 0) len += 1;
|
||||
if (len > 0) return v[0..len];
|
||||
}
|
||||
return default;
|
||||
}
|
||||
|
||||
/// One-shot proactive send: `matrix_harness send "message body"`.
|
||||
/// Lets the agent reach Pierre unprompted (announcements, alerts).
|
||||
fn runSend(alloc: std.mem.Allocator, io: std.Io, text: []const u8) !void {
|
||||
letta.initConversationsDir();
|
||||
const user = env_or("MATRIX_USER", "@chaos-prime:chaosmith.systems");
|
||||
const password = env_or("MATRIX_PASSWORD", "");
|
||||
const room = env_or("MATRIX_ROOM", "");
|
||||
if (password.len == 0 or room.len == 0) {
|
||||
std.debug.print("matrix_harness send: MATRIX_PASSWORD and MATRIX_ROOM required\n", .{});
|
||||
return error.NoCredentials;
|
||||
}
|
||||
const session = try matrix.login(alloc, io, user, password);
|
||||
defer alloc.free(session.token);
|
||||
defer alloc.free(session.user_id);
|
||||
const ev_id = try matrix.sendText(alloc, io, session.token, room, text);
|
||||
std.debug.print("matrix_harness: sent {s}\n", .{ev_id});
|
||||
alloc.free(ev_id);
|
||||
}
|
||||
|
||||
/// Infer msgtype + mimetype from filename extension.
|
||||
const FileType = struct { msgtype: []const u8, mimetype: []const u8 };
|
||||
fn fileTypeOf(path: []const u8) FileType {
|
||||
const ext = if (std.mem.lastIndexOfScalar(u8, path, '.')) |i| path[i + 1 ..] else "";
|
||||
const map = .{
|
||||
.{ ".png", "m.image", "image/png" },
|
||||
.{ ".jpg", "m.image", "image/jpeg" },
|
||||
.{ ".jpeg", "m.image", "image/jpeg" },
|
||||
.{ ".webp", "m.image", "image/webp" },
|
||||
.{ ".gif", "m.image", "image/gif" },
|
||||
.{ ".svg", "m.image", "image/svg+xml" },
|
||||
.{ ".mp3", "m.audio", "audio/mpeg" },
|
||||
.{ ".ogg", "m.audio", "audio/ogg" },
|
||||
.{ ".wav", "m.audio", "audio/wav" },
|
||||
.{ ".mp4", "m.video", "video/mp4" },
|
||||
.{ ".webm", "m.video", "video/webm" },
|
||||
.{ ".pdf", "m.file", "application/pdf" },
|
||||
.{ ".txt", "m.file", "text/plain" },
|
||||
.{ ".md", "m.file", "text/markdown" },
|
||||
.{ ".json", "m.file", "application/json" },
|
||||
.{ ".zip", "m.file", "application/zip" },
|
||||
.{ ".gz", "m.file", "application/gzip" },
|
||||
.{ ".log", "m.file", "text/plain" },
|
||||
};
|
||||
inline for (map) |e| {
|
||||
if (std.ascii.eqlIgnoreCase(ext, e[0][1..])) return .{ .msgtype = e[1], .mimetype = e[2] };
|
||||
}
|
||||
return .{ .msgtype = "m.file", .mimetype = "application/octet-stream" };
|
||||
}
|
||||
|
||||
/// Basename of a path.
|
||||
fn basenameOf(path: []const u8) []const u8 {
|
||||
if (std.mem.lastIndexOfScalar(u8, path, '/')) |i| return path[i + 1 ..];
|
||||
return path;
|
||||
}
|
||||
|
||||
/// One-shot file send: `matrix_harness send-file /path/to/file`.
|
||||
/// Uploads to the homeserver media store, then sends as m.image/m.file.
|
||||
fn runSendFile(alloc: std.mem.Allocator, io: std.Io, path: []const u8) !void {
|
||||
letta.initConversationsDir();
|
||||
const user = env_or("MATRIX_USER", "@chaos-prime:chaosmith.systems");
|
||||
const password = env_or("MATRIX_PASSWORD", "");
|
||||
const room = env_or("MATRIX_ROOM", "");
|
||||
if (password.len == 0 or room.len == 0) {
|
||||
std.debug.print("matrix_harness send-file: MATRIX_PASSWORD and MATRIX_ROOM required\n", .{});
|
||||
return error.NoCredentials;
|
||||
}
|
||||
const dir = std.Io.Dir.cwd();
|
||||
const st = dir.statFile(io, path, .{}) catch {
|
||||
std.debug.print("matrix_harness send-file: cannot stat {s}\n", .{path});
|
||||
return error.NoFile;
|
||||
};
|
||||
const size: u64 = switch (st.kind) {
|
||||
.file => st.size,
|
||||
else => return error.NotAFile,
|
||||
};
|
||||
|
||||
const ft = fileTypeOf(path);
|
||||
const name = basenameOf(path);
|
||||
std.debug.print("matrix_harness: uploading {s} ({d} bytes, {s})...\n", .{ name, size, ft.mimetype });
|
||||
const session = try matrix.login(alloc, io, user, password);
|
||||
defer alloc.free(session.token);
|
||||
defer alloc.free(session.user_id);
|
||||
const mxc = try matrix.uploadFile(alloc, io, session.token, path, name, ft.mimetype);
|
||||
defer alloc.free(mxc);
|
||||
const ev_id = try matrix.sendFile(alloc, io, session.token, room, ft.msgtype, mxc, name, ft.mimetype, size);
|
||||
std.debug.print("matrix_harness: sent {s} ({s}) as {s}\n", .{ ev_id, name, ft.msgtype });
|
||||
alloc.free(ev_id);
|
||||
}
|
||||
|
||||
/// Where incoming peer attachments are saved. The harness passes this path
|
||||
/// to the agent in the forwarded message.
|
||||
const FILES_DIR = "/tmp/matrix-files";
|
||||
|
||||
/// Sanitize a filename for local saving (keep it recognizable, drop path
|
||||
/// separators and control chars).
|
||||
fn sanitize(alloc: std.mem.Allocator, name: []const u8) ![]u8 {
|
||||
var out: std.ArrayList(u8) = .empty;
|
||||
errdefer out.deinit(alloc);
|
||||
for (name) |c| {
|
||||
switch (c) {
|
||||
'/', '\\', 0...0x1f, 0x7f => try out.append(alloc, '_'),
|
||||
' ' => try out.append(alloc, '_'),
|
||||
else => try out.append(alloc, c),
|
||||
}
|
||||
}
|
||||
if (out.items.len == 0) try out.appendSlice(alloc, "attachment");
|
||||
return out.toOwnedSlice(alloc);
|
||||
}
|
||||
|
||||
/// Download an incoming attachment to FILES_DIR. Returns local path
|
||||
/// (allocated; caller frees).
|
||||
fn saveAttachment(alloc: std.mem.Allocator, io: std.Io, token: []const u8, ev_ts: i64, filename: []const u8, mxc: []const u8) ![]u8 {
|
||||
std.Io.Dir.cwd().createDirPath(io, FILES_DIR) catch {};
|
||||
const safe = try sanitize(alloc, filename);
|
||||
defer alloc.free(safe);
|
||||
const path = try std.fmt.allocPrint(alloc, "{s}/{d}-{s}", .{ FILES_DIR, if (ev_ts > 0) ev_ts else 0, safe });
|
||||
errdefer alloc.free(path);
|
||||
const data = try matrix.downloadFile(alloc, io, token, mxc);
|
||||
defer alloc.free(data);
|
||||
const f = try std.Io.Dir.createFileAbsolute(io, path, .{ .truncate = true });
|
||||
defer f.close(io);
|
||||
var wbuf: [16 * 1024]u8 = undefined;
|
||||
var w = f.writer(io, &wbuf);
|
||||
try w.interface.writeAll(data);
|
||||
try w.interface.flush();
|
||||
return path;
|
||||
}
|
||||
|
||||
/// Persistent seen-state: last homeserver event timestamp we already
|
||||
/// forwarded. Survives restarts so backlog never re-sends read messages.
|
||||
const STATE_FILE = "/tmp/matrix-harness.state";
|
||||
|
||||
/// "typing..." indicator: while turn_busy is true, a keepalive thread
|
||||
/// refreshes the typing notification so Pierre sees the agent working.
|
||||
var turn_busy: std.atomic.Value(bool) = std.atomic.Value(bool).init(false);
|
||||
var typing_room: []const u8 = "";
|
||||
var typing_token: []const u8 = "";
|
||||
var typing_uid: []const u8 = "";
|
||||
var typing_io: ?std.Io = null;
|
||||
|
||||
fn typingKeepalive() void {
|
||||
const io = typing_io orelse return;
|
||||
const pa = std.heap.page_allocator;
|
||||
while (turn_busy.load(.acquire)) {
|
||||
matrix.setTyping(pa, io, typing_token, typing_uid, typing_room, true, 10000);
|
||||
var slept: u32 = 0;
|
||||
while (slept < 2_000 and turn_busy.load(.acquire)) : (slept += 250) {
|
||||
io.sleep(.{ .nanoseconds = 250 * std.time.ns_per_ms }, .real) catch {};
|
||||
}
|
||||
}
|
||||
matrix.setTyping(pa, io, typing_token, typing_uid, typing_room, false, 0);
|
||||
}
|
||||
|
||||
fn beginTurn(io: std.Io, token: []const u8, uid: []const u8, room: []const u8) void {
|
||||
typing_io = io;
|
||||
typing_token = token;
|
||||
typing_uid = uid;
|
||||
typing_room = room;
|
||||
turn_busy.store(true, .release);
|
||||
if (std.Thread.spawn(.{}, typingKeepalive, .{})) |t| {
|
||||
t.detach();
|
||||
} else |_| {}
|
||||
}
|
||||
|
||||
fn endTurn() void {
|
||||
turn_busy.store(false, .release);
|
||||
}
|
||||
|
||||
// --------------------------------------------------- tool-call progress ---
|
||||
//
|
||||
// Live "working" bubble: during a turn we send ONE message and keep editing
|
||||
// it (m.replace) as the agent makes tool calls, so the DM shows what's
|
||||
// happening without spamming. The final reply is a separate fresh message.
|
||||
|
||||
/// Max tool-call blocks kept in the bubble before collapsing to "+N earlier".
|
||||
const PROGRESS_MAX_BLOCKS = 10;
|
||||
/// Char cap per tool arg summary / tool output block.
|
||||
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,
|
||||
io: std.Io,
|
||||
token: []const u8,
|
||||
room: []const u8,
|
||||
event_id: ?[]u8 = null, // the bubble message we keep editing
|
||||
blocks: std.ArrayListUnmanaged([]u8) = .empty, // markdown per tool call
|
||||
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);
|
||||
if (self.reply_event_id) |e| self.alloc.free(e);
|
||||
for (self.blocks.items) |b| self.alloc.free(b);
|
||||
self.blocks.deinit(self.alloc);
|
||||
self.reply.deinit(self.alloc);
|
||||
}
|
||||
|
||||
/// Send the initial bubble. Failure is non-fatal (progress is cosmetic).
|
||||
fn begin(self: *Progress) void {
|
||||
const id = matrix.sendText(self.alloc, self.io, self.token, self.room, "⚙️ working…") catch return;
|
||||
self.event_id = id;
|
||||
}
|
||||
|
||||
/// One-line summary: first line only, truncated, no backticks (they'd
|
||||
/// interact with markdown).
|
||||
fn summarize(self: *Progress, text: []const u8, cap: usize) ![]u8 {
|
||||
const nl = std.mem.indexOfScalar(u8, text, '\n') orelse text.len;
|
||||
var line = text[0..nl];
|
||||
// Neutralize backtick runs so they can't break our fences.
|
||||
var clean: std.ArrayListUnmanaged(u8) = .empty;
|
||||
var run: usize = 0;
|
||||
for (line) |c| {
|
||||
if (c == '`') {
|
||||
run += 1;
|
||||
if (run < 3) try clean.append(self.alloc, c);
|
||||
} else {
|
||||
run = 0;
|
||||
try clean.append(self.alloc, c);
|
||||
}
|
||||
if (clean.items.len >= cap) break;
|
||||
}
|
||||
line = clean.items;
|
||||
if (text.len > nl or line.len >= cap and text.len > line.len) {
|
||||
try clean.appendSlice(self.alloc, "…");
|
||||
}
|
||||
return clean.toOwnedSlice(self.alloc);
|
||||
}
|
||||
|
||||
/// Append a markdown block for one tool call and push an edit. Each
|
||||
/// block: bold name + arg summary, then the output in a fenced code
|
||||
/// block once the return arrives. Non-fatal on failure.
|
||||
fn call(self: *Progress, name: []const u8, args: []const u8) void {
|
||||
const arg_line = self.summarize(args, PROGRESS_ARG_CHARS) catch return;
|
||||
defer self.alloc.free(arg_line);
|
||||
const block = std.fmt.allocPrint(self.alloc, "**{s}**: `{s}`\n", .{ name, arg_line }) catch return;
|
||||
self.blocks.append(self.alloc, block) catch {
|
||||
self.alloc.free(block);
|
||||
return;
|
||||
};
|
||||
self.push();
|
||||
}
|
||||
|
||||
fn ret(self: *Progress, out: []const u8) void {
|
||||
const out_line = self.summarize(out, PROGRESS_RET_CHARS) catch return;
|
||||
defer self.alloc.free(out_line);
|
||||
if (self.blocks.items.len == 0) return;
|
||||
const last = self.blocks.items[self.blocks.items.len - 1];
|
||||
// Trailing newline is load-bearing: without it the next block's
|
||||
// header concatenates onto the closing fence and gets swallowed
|
||||
// by the <pre> in Element's renderer.
|
||||
const with_fence = std.fmt.allocPrint(self.alloc, "{s}---\n```\n{s}\n```\n", .{ last, out_line }) catch return;
|
||||
self.alloc.free(last);
|
||||
self.blocks.items[self.blocks.items.len - 1] = with_fence;
|
||||
self.push();
|
||||
}
|
||||
|
||||
/// Current bubble text (caller frees): hidden-count header + last blocks.
|
||||
fn body(self: *Progress) ![]u8 {
|
||||
var out: std.ArrayListUnmanaged(u8) = .empty;
|
||||
errdefer out.deinit(self.alloc);
|
||||
const visible: usize = @min(self.blocks.items.len, PROGRESS_MAX_BLOCKS);
|
||||
const hidden = self.blocks.items.len - visible;
|
||||
if (hidden > 0) {
|
||||
try out.appendSlice(self.alloc, "… +");
|
||||
var num_buf: [16]u8 = undefined;
|
||||
try out.appendSlice(self.alloc, std.fmt.bufPrint(&num_buf, "{d}", .{hidden}) catch "?");
|
||||
try out.appendSlice(self.alloc, " earlier calls\n");
|
||||
}
|
||||
// Walk backwards from the newest block, accumulating until the char
|
||||
// cap is exceeded; everything from there on is visible. Keep at
|
||||
// least the newest block, and at most PROGRESS_MAX_BLOCKS.
|
||||
var start: usize = self.blocks.items.len;
|
||||
var total: usize = 0;
|
||||
while (start > 0) {
|
||||
const l = self.blocks.items[start - 1].len;
|
||||
if (total > 0 and total + l > PROGRESS_MAX_CHARS) break;
|
||||
total += l;
|
||||
start -= 1;
|
||||
}
|
||||
if (self.blocks.items.len - start > visible) start = self.blocks.items.len - visible;
|
||||
for (self.blocks.items[start..]) |b| try out.appendSlice(self.alloc, b);
|
||||
return out.toOwnedSlice(self.alloc);
|
||||
}
|
||||
|
||||
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;
|
||||
};
|
||||
const text = self.body() catch |e| {
|
||||
std.debug.print("DBG push: body failed: {s}\n", .{@errorName(e)});
|
||||
return;
|
||||
};
|
||||
defer self.alloc.free(text);
|
||||
std.debug.print("DBG push: blocks={d} text.len={d}\n", .{ self.blocks.items.len, text.len });
|
||||
if (matrix.editText(self.alloc, self.io, self.token, self.room, id, text)) |new_id| {
|
||||
self.alloc.free(new_id);
|
||||
} else |_| {
|
||||
// Transient homeserver hiccups shouldn't kill the turn; drop the
|
||||
// bubble rather than retrying (next step may re-establish it).
|
||||
std.debug.print("matrix_harness: progress edit failed (dropping bubble)\n", .{});
|
||||
}
|
||||
}
|
||||
|
||||
/// Streamed reply text: on first chunk send a fresh message, then edit
|
||||
/// it as more text arrives. `final` stamps the complete reply (callers
|
||||
/// skip their own sendText when this returned true).
|
||||
fn replyChunk(self: *Progress, text: []const u8, final: bool) bool {
|
||||
self.reply.appendSlice(self.alloc, text) catch return false;
|
||||
if (self.reply_event_id == null) {
|
||||
const id = matrix.sendText(self.alloc, self.io, self.token, self.room, "…") catch return false;
|
||||
self.reply_event_id = id;
|
||||
}
|
||||
const id = self.reply_event_id.?;
|
||||
if (matrix.editText(self.alloc, self.io, self.token, self.room, id, self.reply.items)) |new_id| {
|
||||
self.alloc.free(new_id);
|
||||
self.replied = final;
|
||||
return final;
|
||||
} else |_| return false;
|
||||
}
|
||||
};
|
||||
|
||||
/// streamInfer callback: forward structured events into the Progress bubble.
|
||||
fn onProgressEvent(ctx: ?*anyopaque, ev: letta.Event) void {
|
||||
const p: *Progress = @ptrCast(@alignCast(ctx orelse return));
|
||||
switch (ev) {
|
||||
.tool_call => |tc| {
|
||||
defer p.alloc.free(tc.name);
|
||||
defer p.alloc.free(tc.args);
|
||||
p.call(tc.name, tc.args);
|
||||
},
|
||||
.tool_return => |out| {
|
||||
defer p.alloc.free(out);
|
||||
p.ret(out);
|
||||
},
|
||||
.chunk => |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
|
||||
}
|
||||
}
|
||||
|
||||
/// One full agent turn with a live tool-call bubble and a streamed reply
|
||||
/// message. Returns the reply text (allocated, caller frees). If the reply
|
||||
/// was already streamed and finalized in-message (p.replied), the caller
|
||||
/// should NOT send it again; the bubble stays as the tool-call record.
|
||||
fn runAgentTurn(p: *Progress, prompt: []const u8) ![]u8 {
|
||||
p.begin();
|
||||
const reply = try letta.streamInferConv(p.io, p.alloc, "", letta.CONVERSATION_ID, prompt, p, onProgressEvent);
|
||||
// Stamp the authoritative full reply into the streamed message: REPLACE
|
||||
// the accumulated chunks (appending would duplicate the whole text).
|
||||
if (p.reply_event_id != null) {
|
||||
p.reply.clearRetainingCapacity();
|
||||
_ = p.replyChunk(reply, true);
|
||||
}
|
||||
return reply;
|
||||
}
|
||||
|
||||
fn loadSeenTs(io: std.Io) i64 {
|
||||
var f = std.Io.Dir.openFileAbsolute(io, STATE_FILE, .{}) catch return 0;
|
||||
defer f.close(io);
|
||||
var buf: [64]u8 = undefined;
|
||||
var rbuf: [64]u8 = undefined;
|
||||
var r = f.reader(io, &rbuf);
|
||||
const n = r.interface.readSliceShort(&buf) catch return 0;
|
||||
return std.fmt.parseInt(i64, std.mem.trim(u8, buf[0..n], " \n"), 10) catch 0;
|
||||
}
|
||||
|
||||
fn saveSeenTs(io: std.Io, ts: i64) void {
|
||||
const f = std.Io.Dir.createFileAbsolute(io, STATE_FILE, .{ .truncate = true }) catch return;
|
||||
var buf: [24]u8 = undefined;
|
||||
const body = std.fmt.bufPrint(&buf, "{d}\n", .{ts}) catch return;
|
||||
var wbuf: [64]u8 = undefined;
|
||||
var w = f.writer(io, &wbuf);
|
||||
w.interface.writeAll(body) catch {};
|
||||
w.interface.flush() catch {};
|
||||
f.close(io);
|
||||
}
|
||||
|
||||
pub fn main(init: std.process.Init) !void {
|
||||
const alloc = init.gpa;
|
||||
const io = init.io;
|
||||
|
||||
// Proactive-send mode: `matrix_harness send "text"`.
|
||||
{
|
||||
var args_it = std.process.Args.Iterator.init(init.minimal.args);
|
||||
_ = args_it.next(); // program name
|
||||
const first_arg = args_it.next();
|
||||
if (first_arg != null and std.mem.eql(u8, first_arg.?, "send")) {
|
||||
if (args_it.next()) |msg| {
|
||||
return runSend(alloc, io, msg);
|
||||
}
|
||||
std.debug.print("usage: matrix_harness send \"message\"\n", .{});
|
||||
return error.MissingMessage;
|
||||
}
|
||||
if (first_arg != null and std.mem.eql(u8, first_arg.?, "send-file")) {
|
||||
if (args_it.next()) |p| {
|
||||
return runSendFile(alloc, io, p);
|
||||
}
|
||||
std.debug.print("usage: matrix_harness send-file /path/to/file\n", .{});
|
||||
return error.MissingFile;
|
||||
}
|
||||
}
|
||||
|
||||
letta.initConversationsDir();
|
||||
|
||||
const user = env_or("MATRIX_USER", "@chaos-prime:chaosmith.systems");
|
||||
const password = env_or("MATRIX_PASSWORD", "");
|
||||
const peer = env_or("MATRIX_PEER", "@pierre:chaosmith.systems");
|
||||
const agent = env_or("LETTA_MATRIX_AGENT", "agent-local-273c95cf-82a0-4711-bcd5-6139661e8488");
|
||||
|
||||
if (password.len == 0) {
|
||||
std.debug.print("matrix_harness: MATRIX_PASSWORD not set\n", .{});
|
||||
return error.NoCredentials;
|
||||
}
|
||||
|
||||
std.debug.print("matrix_harness: logging in as {s}...\n", .{user});
|
||||
const session = matrix.login(alloc, io, user, password) catch |e| {
|
||||
std.debug.print("matrix_harness: login failed: {s}\n", .{@errorName(e)});
|
||||
return e;
|
||||
};
|
||||
defer alloc.free(session.token);
|
||||
defer alloc.free(session.user_id);
|
||||
std.debug.print("matrix_harness: logged in ({s})\n", .{session.user_id});
|
||||
|
||||
// Single persistent room: MATRIX_ROOM env if set (reuse!); only create
|
||||
// a fresh DM when it's unset. All replies go to this one room.
|
||||
const room: []const u8 = env_or("MATRIX_ROOM", "");
|
||||
const owned_room: ?[]u8 = if (room.len > 0) null else matrix.createDirect(alloc, io, session.token, peer) catch |e| {
|
||||
std.debug.print("matrix_harness: createDirect failed: {s}\n", .{@errorName(e)});
|
||||
return e;
|
||||
};
|
||||
const room_id: []const u8 = owned_room orelse room;
|
||||
defer if (owned_room) |r| alloc.free(r);
|
||||
std.debug.print("matrix_harness: room {s}\n", .{room_id});
|
||||
|
||||
// No room banner on startup — restarts (deploys, watchdog) would spam
|
||||
// the DM. Startup is logged to journald instead.
|
||||
|
||||
// Establish a sync point. RECENT peer messages (last 30 min only, so
|
||||
// harness restarts don't re-forward ancient history every boot) are
|
||||
// forwarded to the agent as one [matrix backlog] context message
|
||||
// (reply suppressed).
|
||||
const first = try matrix.sync(alloc, io, session.token, "", 0);
|
||||
{
|
||||
const now_ms: i64 = @intCast(@divTrunc(std.Io.Clock.now(.real, io).nanoseconds, std.time.ns_per_ms));
|
||||
const seen_ts = loadSeenTs(io);
|
||||
var new_max = seen_ts;
|
||||
var backlog: std.ArrayListUnmanaged(u8) = .empty;
|
||||
defer backlog.deinit(alloc);
|
||||
for (first.events) |ev| {
|
||||
defer alloc.free(ev.sender);
|
||||
defer alloc.free(ev.body);
|
||||
defer if (ev.filename.len > 0) alloc.free(ev.filename);
|
||||
defer if (ev.url) |u| alloc.free(u);
|
||||
defer if (ev.mimetype.len > 0) alloc.free(ev.mimetype);
|
||||
if (ev.ts > new_max) new_max = ev.ts;
|
||||
if (!std.mem.eql(u8, ev.sender, peer)) continue;
|
||||
// Attachments in backlog are stale — note them but don't download.
|
||||
if (ev.url != null) {
|
||||
if (ev.ts > 0 and now_ms - ev.ts > 30 * 60 * 1000) continue;
|
||||
if (ev.ts > 0 and ev.ts <= seen_ts) continue;
|
||||
try backlog.appendSlice(alloc, ev.sender);
|
||||
try backlog.appendSlice(alloc, ": [sent a file: ");
|
||||
try backlog.appendSlice(alloc, ev.body);
|
||||
try backlog.appendSlice(alloc, " — not downloaded, ask if needed]\n");
|
||||
continue;
|
||||
}
|
||||
if (ev.body.len == 0) continue;
|
||||
if (ev.ts > 0 and now_ms - ev.ts > 30 * 60 * 1000) continue; // stale
|
||||
if (ev.ts > 0 and ev.ts <= seen_ts) continue; // already seen/answered
|
||||
try backlog.appendSlice(alloc, ev.sender);
|
||||
try backlog.appendSlice(alloc, ": ");
|
||||
try backlog.appendSlice(alloc, ev.body);
|
||||
try backlog.append(alloc, '\n');
|
||||
}
|
||||
if (backlog.items.len > 0) {
|
||||
std.debug.print("matrix_harness: forwarding backlog ({d} bytes) to agent\n", .{backlog.items.len});
|
||||
const hist_msg = try std.fmt.allocPrint(alloc, "[matrix backlog — messages received while harness was offline]\n{s}", .{backlog.items});
|
||||
defer alloc.free(hist_msg);
|
||||
// Do NOT exit on backlog failure: with systemd Restart=on-failure
|
||||
// this would re-forward the same backlog every 3s forever.
|
||||
// Log, skip, keep polling.
|
||||
beginTurn(io, session.token, session.user_id, room_id);
|
||||
defer endTurn();
|
||||
var prog = Progress{ .alloc = alloc, .io = io, .token = session.token, .room = room_id };
|
||||
defer prog.deinit();
|
||||
if (runAgentTurn(&prog, hist_msg)) |reply| {
|
||||
alloc.free(reply); // backlog replies are not sent back
|
||||
} else |e| {
|
||||
std.debug.print("matrix_harness: backlog letta failed: {s} (skipping)\n", .{@errorName(e)});
|
||||
}
|
||||
}
|
||||
// Advance persisted seen-state even when nothing was forwarded:
|
||||
// "read up to here" is authoritative from the homeserver ts.
|
||||
saveSeenTs(io, new_max);
|
||||
}
|
||||
alloc.free(first.events);
|
||||
var since = first.next;
|
||||
|
||||
std.debug.print("matrix_harness: polling sync (peer {s}, agent {s})\n", .{ peer, agent });
|
||||
|
||||
while (true) {
|
||||
const result = matrix.sync(alloc, io, session.token, since, 15000) catch |e| {
|
||||
std.debug.print("matrix_harness: sync error: {s}, retrying\n", .{@errorName(e)});
|
||||
io.sleep(.{ .nanoseconds = 3 * std.time.ns_per_s }, .real) catch {};
|
||||
continue;
|
||||
};
|
||||
alloc.free(since);
|
||||
since = result.next;
|
||||
|
||||
for (result.events) |ev| {
|
||||
defer alloc.free(ev.sender);
|
||||
defer alloc.free(ev.body);
|
||||
defer if (ev.filename.len > 0) alloc.free(ev.filename);
|
||||
defer if (ev.url) |u| alloc.free(u);
|
||||
defer if (ev.mimetype.len > 0) alloc.free(ev.mimetype);
|
||||
defer if (ev.ts > 0) saveSeenTs(io, ev.ts);
|
||||
// Only peer messages (skip own sends and server noise).
|
||||
if (!std.mem.eql(u8, ev.sender, peer)) continue;
|
||||
if (ev.body.len == 0) continue;
|
||||
|
||||
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>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.
|
||||
if (ev.url) |mxc| {
|
||||
// Use the attachment's REAL filename; the event body is just
|
||||
// the caption ("This us?" is not a filename). Fall back to
|
||||
// body, and ensure an extension via mimetype.
|
||||
const fname_src = if (ev.filename.len > 0) ev.filename else ev.body;
|
||||
var fname_buf: std.ArrayListUnmanaged(u8) = .empty;
|
||||
var fname: []const u8 = fname_src;
|
||||
if (std.mem.lastIndexOfScalar(u8, fname_src, '.') == null and ev.mimetype.len > 0) {
|
||||
const ext_map = [_]struct{ m: []const u8, e: []const u8 }{
|
||||
.{ .m = "image/jpeg", .e = ".jpg" }, .{ .m = "image/png", .e = ".png" },
|
||||
.{ .m = "image/webp", .e = ".webp" }, .{ .m = "application/pdf", .e = ".pdf" },
|
||||
.{ .m = "application/gzip", .e = ".gz" }, .{ .m = "application/zip", .e = ".zip" },
|
||||
.{ .m = "text/plain", .e = ".txt" },
|
||||
};
|
||||
for (ext_map) |pair| {
|
||||
if (std.mem.eql(u8, ev.mimetype, pair.m)) {
|
||||
fname_buf.appendSlice(alloc, fname_src) catch {};
|
||||
fname_buf.appendSlice(alloc, pair.e) catch {};
|
||||
fname = fname_buf.items;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
defer fname_buf.deinit(alloc);
|
||||
const local = saveAttachment(alloc, io, session.token, ev.ts, fname, mxc) catch |e| blk: {
|
||||
std.debug.print("matrix_harness: attachment download failed: {s}\n", .{@errorName(e)});
|
||||
break :blk null;
|
||||
};
|
||||
defer if (local) |p| alloc.free(p);
|
||||
var size_buf: [24]u8 = undefined;
|
||||
const size_str = std.fmt.bufPrint(&size_buf, "{d}", .{ev.size}) catch "?";
|
||||
const note = if (local) |p|
|
||||
try std.fmt.allocPrint(alloc, "[matrix{s}] sent a file: \"{s}\" ({s}, {s} bytes) — saved to {s}", .{ ts_str, ev.body, ev.mimetype, size_str, p })
|
||||
else
|
||||
try std.fmt.allocPrint(alloc, "[matrix{s}] sent a file: \"{s}\" ({s}, {s} bytes) — download failed, ask them to resend", .{ ts_str, ev.body, ev.mimetype, size_str });
|
||||
defer alloc.free(note);
|
||||
beginTurn(io, session.token, session.user_id, room_id);
|
||||
var prog = Progress{ .alloc = alloc, .io = io, .token = session.token, .room = room_id };
|
||||
defer prog.deinit();
|
||||
const reply = runAgentTurn(&prog, note) catch |e| {
|
||||
endTurn();
|
||||
std.debug.print("matrix_harness: letta failed: {s}\n", .{@errorName(e)});
|
||||
continue;
|
||||
};
|
||||
defer alloc.free(reply);
|
||||
if (matrix.sendText(alloc, io, session.token, room_id, reply)) |ev_id| alloc.free(ev_id) else |e| {
|
||||
std.debug.print("matrix_harness: send failed: {s}\n", .{@errorName(e)});
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
std.debug.print("matrix_harness: [{s}] {s}\n", .{ ev.sender, ev.body });
|
||||
|
||||
// Prefix with origin + event UTC timestamp so the agent can tell
|
||||
// Matrix-origin messages from CLI ones (and when they were sent).
|
||||
const prefixed = try std.fmt.allocPrint(alloc, "[matrix{s}] {s}", .{ ts_str, ev.body });
|
||||
defer alloc.free(prefixed);
|
||||
|
||||
beginTurn(io, session.token, session.user_id, room_id);
|
||||
var prog = Progress{ .alloc = alloc, .io = io, .token = session.token, .room = room_id };
|
||||
defer prog.deinit();
|
||||
const reply = runAgentTurn(&prog, prefixed) catch |e| {
|
||||
endTurn();
|
||||
std.debug.print("matrix_harness: letta failed: {s}\n", .{@errorName(e)});
|
||||
const err_txt = "agent failed — see harness logs";
|
||||
if (matrix.sendText(alloc, io, session.token, room_id, err_txt)) |ev_id| alloc.free(ev_id) else |_| {}
|
||||
continue;
|
||||
};
|
||||
defer alloc.free(reply);
|
||||
defer endTurn();
|
||||
|
||||
std.debug.print("matrix_harness: reply: {s}\n", .{reply});
|
||||
if (prog.replied) {
|
||||
// Reply was streamed into its own message; final edit already
|
||||
// carries the full authoritative text.
|
||||
} else if (matrix.sendText(alloc, io, session.token, room_id, reply)) |ev_id| {
|
||||
alloc.free(ev_id);
|
||||
} else |e| {
|
||||
std.debug.print("matrix_harness: send failed: {s}\n", .{@errorName(e)});
|
||||
}
|
||||
}
|
||||
alloc.free(result.events);
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user