From da50dea0d111a1d89544d9bf93de380f898d303d Mon Sep 17 00:00:00 2001 From: Pierre De Lancre Date: Wed, 2 Sep 2026 02:09:42 +0300 Subject: [PATCH] Matrix bridge scaffold: test harness with texting E2E MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - src/matrix.zig: minimal Matrix client (conduwuit @ matrix.chaosmith.systems) via curl subprocess transport (native TLS swap later): login, createRoom (trusted_private_chat, unencrypted — no olm), sendText, sync polling with window-scan event extraction reusing letta.jsonStr - src/matrix_harness.zig: standalone entrypoint (independent of tray daemon): login -> DM with peer -> sync loop -> peer messages routed through letta CLI (chaos-prime) -> replies sent back to the room - build.zig: second exe matrix_harness; use_llvm fix promoted into build.zig for both executables (native lld vs GCC16 .sframe relocations) - calling/Element Call (livekit SFU already deployed on ass-host) layers on top of this harness next 👾 Generated with [Letta Code](https://letta.com) Co-Authored-By: Letta Code --- build.zig | 20 ++++++ src/matrix.zig | 152 +++++++++++++++++++++++++++++++++++++++++ src/matrix_harness.zig | 102 +++++++++++++++++++++++++++ 3 files changed, 274 insertions(+) create mode 100644 src/matrix.zig create mode 100644 src/matrix_harness.zig diff --git a/build.zig b/build.zig index bc48bcc..4478e99 100644 --- a/build.zig +++ b/build.zig @@ -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); diff --git a/src/matrix.zig b/src/matrix.zig new file mode 100644 index 0000000..d0a954f --- /dev/null +++ b/src/matrix.zig @@ -0,0 +1,152 @@ +//! 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 + +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). +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); + + var esc: std.ArrayList(u8) = .empty; + defer esc.deinit(alloc); + for (body) |c| { + switch (c) { + '"' => try esc.appendSlice(alloc, "\\\""), + '\\' => try esc.appendSlice(alloc, "\\\\"), + '\n' => try esc.appendSlice(alloc, "\\n"), + '\r', '\t' => try esc.appendSlice(alloc, " "), + else => try esc.append(alloc, c), + } + } + const payload = try std.fmt.allocPrint(alloc, "{{\"msgtype\":\"m.text\",\"body\":\"{s}\"}}", .{esc.items}); + defer alloc.free(payload); + + const resp = try http(alloc, io, "PUT", url, token, payload); + 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 }; // both allocated + +/// One sync poll. Returns message events + next `since` token. Caller frees. +/// (Sync first without a since token to establish one, then poll.) +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); + + const next = if (letta.jsonStr(alloc, resp, 0, "next_batch")) |nb| nb.val else try alloc.dupe(u8, since); + + var events: std.ArrayList(Event) = .empty; + errdefer events.deinit(alloc); + + var i: usize = 0; + while (std.mem.indexOfPos(u8, resp, i, "\"m.room.message\"")) |pos| { + defer i = pos + 15; + const win_end = @min(resp.len, pos + 1500); + const win = resp[pos..win_end]; + const sender = letta.jsonStr(alloc, win, 0, "sender") orelse continue; + const body = letta.jsonStr(alloc, win, 0, "body") orelse { + alloc.free(sender.val); + continue; + }; + events.append(alloc, .{ .sender = sender.val, .body = body.val }) catch { + alloc.free(sender.val); + alloc.free(body.val); + }; + } + 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; +} diff --git a/src/matrix_harness.zig b/src/matrix_harness.zig new file mode 100644 index 0000000..d6d8c95 --- /dev/null +++ b/src/matrix_harness.zig @@ -0,0 +1,102 @@ +//! 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 envOr(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; +} + +pub fn main(init: std.process.Init) !void { + const alloc = init.gpa; + const io = init.io; + letta.initConversationsDir(); + + const user = envOr("MATRIX_USER", "@chaos-prime:chaosmith.systems"); + const password = envOr("MATRIX_PASSWORD", ""); + const peer = envOr("MATRIX_PEER", "@pierre:chaosmith.systems"); + const agent = envOr("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}); + + // Create (or reuse) the DM. For scaffold simplicity we create fresh each + // run — Element dedupes DMs and the invite lands in Pierre's inbox. + const room = matrix.createDirect(alloc, io, session.token, peer) catch |e| { + std.debug.print("matrix_harness: createDirect failed: {s}\n", .{@errorName(e)}); + return e; + }; + defer alloc.free(room); + std.debug.print("matrix_harness: DM room {s}\n", .{room}); + + _ = try matrix.sendText(alloc, io, session.token, room, "matrix_harness online. Talk to me — messages go to the agent, replies come back here."); + + // Establish a sync point (drop history before now). + const first = try matrix.sync(alloc, io, session.token, "", 0); + for (first.events) |ev| { + alloc.free(ev.sender); + alloc.free(ev.body); + } + 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); + // Only peer messages (skip own sends and server noise). + if (!std.mem.eql(u8, ev.sender, peer)) continue; + if (ev.body.len == 0) continue; + + std.debug.print("matrix_harness: [{s}] {s}\n", .{ ev.sender, ev.body }); + + const reply = letta.infer(io, alloc, agent, ev.body) catch |e| { + 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, err_txt)) |ev_id| alloc.free(ev_id) else |_| {} + continue; + }; + defer alloc.free(reply); + + std.debug.print("matrix_harness: reply: {s}\n", .{reply}); + if (matrix.sendText(alloc, io, session.token, room, reply)) |ev_id| alloc.free(ev_id) else |e| { + std.debug.print("matrix_harness: send failed: {s}\n", .{@errorName(e)}); + } + } + alloc.free(result.events); + } +}