matrix: file transfer both ways — send-file one-shot (upload + m.image/m.file by ext), incoming attachments downloaded to /tmp/matrix-files and path forwarded to agent

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

Co-Authored-By: Letta Code <noreply@letta.com>
This commit is contained in:
pierreandLetta Code committed 2026-09-03 23:51:09 +03:00
1 parent dcf6d8df75
commit e212e3e120
2 files changed
+308 -6

No files matched your search

+148 -2
View File
@@ -53,6 +53,122 @@ fn http(alloc: std.mem.Allocator, io: Io, method: []const u8, url: []const u8, t
pub const Session = struct { token: []u8, user_id: []u8 }; // caller frees both 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 ..];
const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/media/v3/download/{s}/{s}?allow_redirect=true", .{ 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 { pub fn login(alloc: std.mem.Allocator, io: Io, user: []const u8, password: []const u8) !Session {
const body = try std.fmt.allocPrint(alloc, 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}"}} \\{{"type":"m.login.password","identifier":{{"type":"m.id.user","user":"{s}"}},"password":"{s}","initial_device_display_name":"{s}"}}
@@ -215,7 +331,15 @@ pub fn joinRoom(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []con
alloc.free(resp); alloc.free(resp);
} }
pub const Event = struct { sender: []u8, body: []u8, ts: i64 = 0 }; // sender/body allocated 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,
};
/// One sync poll. Returns message events + next `since` token. Caller frees. /// One sync poll. Returns message events + next `since` token. Caller frees.
/// (Sync first without a since token to establish one, then poll.) /// (Sync first without a since token to establish one, then poll.)
@@ -312,12 +436,34 @@ pub fn sync(alloc: std.mem.Allocator, io: Io, token: []const u8, since: []const
else => 0, else => 0,
} else 0; } else 0;
// Attachment fields (m.file / m.image / m.audio / m.video).
var url_copy: ?[]u8 = null;
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 sender_copy = alloc.dupe(u8, sender_raw) catch continue;
const body_copy = alloc.dupe(u8, body_raw) catch { const body_copy = alloc.dupe(u8, body_raw) catch {
alloc.free(sender_copy); alloc.free(sender_copy);
continue; continue;
}; };
events.append(alloc, .{ .sender = sender_copy, .body = body_copy, .ts = ts }) catch { events.append(alloc, .{ .sender = sender_copy, .body = body_copy, .ts = ts, .url = url_copy, .mimetype = mime, .size = size }) catch {
alloc.free(sender_copy); alloc.free(sender_copy);
alloc.free(body_copy); alloc.free(body_copy);
}; };
+160 -4
View File
@@ -39,6 +39,115 @@ fn runSend(alloc: std.mem.Allocator, io: std.Io, text: []const u8) !void {
alloc.free(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 /// Persistent seen-state: last homeserver event timestamp we already
/// forwarded. Survives restarts so backlog never re-sends read messages. /// forwarded. Survives restarts so backlog never re-sends read messages.
const STATE_FILE = "/tmp/matrix-harness.state"; const STATE_FILE = "/tmp/matrix-harness.state";
@@ -80,6 +189,13 @@ pub fn main(init: std.process.Init) !void {
std.debug.print("usage: matrix_harness send \"message\"\n", .{}); std.debug.print("usage: matrix_harness send \"message\"\n", .{});
return error.MissingMessage; 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(); letta.initConversationsDir();
@@ -131,8 +247,20 @@ pub fn main(init: std.process.Init) !void {
for (first.events) |ev| { for (first.events) |ev| {
defer alloc.free(ev.sender); defer alloc.free(ev.sender);
defer alloc.free(ev.body); defer alloc.free(ev.body);
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 (ev.ts > new_max) new_max = ev.ts;
if (!std.mem.eql(u8, ev.sender, peer)) continue; 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.body.len == 0) continue;
if (ev.ts > 0 and now_ms - ev.ts > 30 * 60 * 1000) continue; // stale 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 if (ev.ts > 0 and ev.ts <= seen_ts) continue; // already seen/answered
@@ -175,15 +303,13 @@ pub fn main(init: std.process.Init) !void {
for (result.events) |ev| { for (result.events) |ev| {
defer alloc.free(ev.sender); defer alloc.free(ev.sender);
defer alloc.free(ev.body); defer alloc.free(ev.body);
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); defer if (ev.ts > 0) saveSeenTs(io, ev.ts);
// Only peer messages (skip own sends and server noise). // Only peer messages (skip own sends and server noise).
if (!std.mem.eql(u8, ev.sender, peer)) continue; if (!std.mem.eql(u8, ev.sender, peer)) continue;
if (ev.body.len == 0) continue; if (ev.body.len == 0) 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).
var ts_buf: [16]u8 = undefined; var ts_buf: [16]u8 = undefined;
const ts_str: []const u8 = blk: { const ts_str: []const u8 = blk: {
if (ev.ts <= 0) break :blk ""; if (ev.ts <= 0) break :blk "";
@@ -192,6 +318,36 @@ pub fn main(init: std.process.Init) !void {
const sod = epoch.getDaySeconds(); const sod = epoch.getDaySeconds();
break :blk std.fmt.bufPrint(&ts_buf, " {d:0>2}:{d:0>2}Z", .{ sod.getHoursIntoDay(), sod.getMinutesIntoHour() }) catch ""; break :blk std.fmt.bufPrint(&ts_buf, " {d:0>2}:{d:0>2}Z", .{ sod.getHoursIntoDay(), sod.getMinutesIntoHour() }) catch "";
}; };
// Attachment: download locally and forward the path to the agent.
if (ev.url) |mxc| {
const local = saveAttachment(alloc, io, session.token, ev.ts, ev.body, 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);
const reply = letta.infer(io, alloc, agent, note) catch |e| {
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 }); const prefixed = try std.fmt.allocPrint(alloc, "[matrix{s}] {s}", .{ ts_str, ev.body });
defer alloc.free(prefixed); defer alloc.free(prefixed);