From e212e3e12026b3c7354c4312a2c74b9127dbadd4 Mon Sep 17 00:00:00 2001 From: Pierre De Lancre Date: Thu, 3 Sep 2026 23:51:09 +0300 Subject: [PATCH] =?UTF-8?q?matrix:=20file=20transfer=20both=20ways=20?= =?UTF-8?q?=E2=80=94=20send-file=20one-shot=20(upload=20+=20m.image/m.file?= =?UTF-8?q?=20by=20ext),=20incoming=20attachments=20downloaded=20to=20/tmp?= =?UTF-8?q?/matrix-files=20and=20path=20forwarded=20to=20agent?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 👾 Generated with [Letta Code](https://letta.com) Co-Authored-By: Letta Code --- src/matrix.zig | 150 ++++++++++++++++++++++++++++++++++++- src/matrix_harness.zig | 164 ++++++++++++++++++++++++++++++++++++++++- 2 files changed, 308 insertions(+), 6 deletions(-) diff --git a/src/matrix.zig b/src/matrix.zig index bcecf6d..1b8ed05 100644 --- a/src/matrix.zig +++ b/src/matrix.zig @@ -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 +/// 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 { 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}"}} @@ -215,7 +331,15 @@ pub fn joinRoom(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []con 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. /// (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; + // 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 body_copy = alloc.dupe(u8, body_raw) catch { alloc.free(sender_copy); 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(body_copy); }; diff --git a/src/matrix_harness.zig b/src/matrix_harness.zig index 402c717..38eade7 100644 --- a/src/matrix_harness.zig +++ b/src/matrix_harness.zig @@ -39,6 +39,115 @@ fn runSend(alloc: std.mem.Allocator, io: std.Io, text: []const u8) !void { 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"; @@ -80,6 +189,13 @@ pub fn main(init: std.process.Init) !void { 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(); @@ -131,8 +247,20 @@ pub fn main(init: std.process.Init) !void { for (first.events) |ev| { defer alloc.free(ev.sender); 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 (!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 @@ -175,15 +303,13 @@ pub fn main(init: std.process.Init) !void { for (result.events) |ev| { defer alloc.free(ev.sender); 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); // 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 }); - - // 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; const ts_str: []const u8 = blk: { if (ev.ts <= 0) break :blk ""; @@ -192,6 +318,36 @@ pub fn main(init: std.process.Init) !void { 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| { + 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 }); defer alloc.free(prefixed);