Files
inferon/src/matrix_harness.zig
T
2026-09-05 11:43:30 +03:00

438 lines
20 KiB
Zig

//! 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);
}
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();
if (letta.infer(io, alloc, agent, 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: [16]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 sod = epoch.getDaySeconds();
break :blk std.fmt.bufPrint(&ts_buf, " {d:0>2}:{d:0>2}Z", .{ 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);
const reply = letta.infer(io, alloc, agent, 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);
const reply = letta.infer(io, alloc, agent, 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 (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);
}
}