diff --git a/src/matrix.zig b/src/matrix.zig index becf4a1..bcecf6d 100644 --- a/src/matrix.zig +++ b/src/matrix.zig @@ -67,28 +67,143 @@ pub fn login(alloc: std.mem.Allocator, io: Io, user: []const u8, password: []con } /// 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) "" else "
");
+            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, "
"); + 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, ""); + try out.appendSlice(alloc, md[i + 1 .. end]); + try out.appendSlice(alloc, ""); + 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, ""); + try out.appendSlice(alloc, md[i + 2 .. end]); + try out.appendSlice(alloc, ""); + i = end + 2; + } else { + try out.appendSlice(alloc, "**"); + i += 2; + } + } else if (std.mem.indexOfScalarPos(u8, md, i + 1, '*')) |end| { + try out.appendSlice(alloc, ""); + try out.appendSlice(alloc, md[i + 1 .. end]); + try out.appendSlice(alloc, ""); + 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, "
"); + return out.toOwnedSlice(alloc); +} + 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 html = markdownToHtml(alloc, body) catch |e| switch (e) { + error.OutOfMemory => return e, + }; + defer alloc.free(html); - const resp = try http(alloc, io, "PUT", url, token, payload); + 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 return error.SendFailed; + 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; } @@ -100,10 +215,14 @@ 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 }; // both allocated +pub const Event = struct { sender: []u8, body: []u8, ts: i64 = 0 }; // sender/body allocated /// 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 }) @@ -113,28 +232,96 @@ pub fn sync(alloc: std.mem.Allocator, io: Io, token: []const u8, since: []const 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; - // Event layout: sender (and event_id) come BEFORE the type marker, - // body comes after it in content — window spans both sides. - const win_start = pos -| 600; - const win_end = @min(resp.len, pos + 1500); - const win = resp[win_start..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); + 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; + + 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 { + alloc.free(sender_copy); + alloc.free(body_copy); + }; + } } return .{ .events = try events.toOwnedSlice(alloc), .next = next }; } diff --git a/src/matrix_harness.zig b/src/matrix_harness.zig index 05eece3..1040230 100644 --- a/src/matrix_harness.zig +++ b/src/matrix_harness.zig @@ -55,13 +55,42 @@ pub fn main(init: std.process.Init) !void { defer if (owned_room) |r| alloc.free(r); std.debug.print("matrix_harness: room {s}\n", .{room_id}); - _ = try matrix.sendText(alloc, io, session.token, room_id, "matrix_harness online in persistent room. Messages reach the agent, replies come back here."); + // No room banner on startup — restarts (deploys, watchdog) would spam + // the DM. Startup is logged to journald instead. - // Establish a sync point (drop history before now). + // 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); - for (first.events) |ev| { - alloc.free(ev.sender); - alloc.free(ev.body); + { + const now_ms: i64 = @intCast(@divTrunc(std.Io.Clock.now(.real, io).nanoseconds, std.time.ns_per_ms)); + var backlog: std.ArrayListUnmanaged(u8) = .empty; + defer backlog.deinit(alloc); + for (first.events) |ev| { + defer alloc.free(ev.sender); + defer alloc.free(ev.body); + if (!std.mem.eql(u8, ev.sender, peer)) continue; + if (ev.body.len == 0) continue; + if (ev.ts > 0 and now_ms - ev.ts > 30 * 60 * 1000) continue; // stale + 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. + 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)}); + } + } } alloc.free(first.events); var since = first.next; @@ -86,7 +115,20 @@ pub fn main(init: std.process.Init) !void { std.debug.print("matrix_harness: [{s}] {s}\n", .{ ev.sender, ev.body }); - const reply = letta.infer(io, alloc, agent, ev.body) catch |e| { + // 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 ""; + 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 ""; + }; + const prefixed = try std.fmt.allocPrint(alloc, "[matrix{s}] {s}", .{ ts_str, ev.body }); + defer alloc.free(prefixed); + + const reply = letta.infer(io, alloc, agent, prefixed) 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_id, err_txt)) |ev_id| alloc.free(ev_id) else |_| {}