matrix bridge: proper JSON sync parsing, backlog forwarding, message provenance
- sync(): std.json walk of rooms.join.*.timeline (old string-window scraper misattributed senders across events); Event gains ts - startup backlog: peer messages from last 30min forwarded to agent as context (reply suppressed); failures logged, never fatal (systemd Restart=on-failure would otherwise re-forward every 3s) - messages to agent prefixed [matrix HH:MMZ] with UTC event timestamp - sendText: jsonEscape covers control chars (\u00XX), markdown -> HTML formatted_body (bold/italic/code/fences) for Element rendering - no startup banner in room; harness now runs as systemd user service matrix-harness (unit + .matrix-env, untracked)
This commit is contained in:
1 parent
77f3150af4
commit
9f1a839a22
2 files changed
+266
-37
No files matched your search
+219
-32
@@ -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).
|
/// 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) "</pre>" else "<pre>");
|
||||||
|
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, "<br>");
|
||||||
|
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, "<code>");
|
||||||
|
try out.appendSlice(alloc, md[i + 1 .. end]);
|
||||||
|
try out.appendSlice(alloc, "</code>");
|
||||||
|
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, "<strong>");
|
||||||
|
try out.appendSlice(alloc, md[i + 2 .. end]);
|
||||||
|
try out.appendSlice(alloc, "</strong>");
|
||||||
|
i = end + 2;
|
||||||
|
} else {
|
||||||
|
try out.appendSlice(alloc, "**");
|
||||||
|
i += 2;
|
||||||
|
}
|
||||||
|
} else if (std.mem.indexOfScalarPos(u8, md, i + 1, '*')) |end| {
|
||||||
|
try out.appendSlice(alloc, "<em>");
|
||||||
|
try out.appendSlice(alloc, md[i + 1 .. end]);
|
||||||
|
try out.appendSlice(alloc, "</em>");
|
||||||
|
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, "</pre>");
|
||||||
|
return out.toOwnedSlice(alloc);
|
||||||
|
}
|
||||||
|
|
||||||
pub fn sendText(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []const u8, body: []const u8) ![]u8 {
|
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 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 });
|
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);
|
defer alloc.free(url);
|
||||||
|
|
||||||
var esc: std.ArrayList(u8) = .empty;
|
const html = markdownToHtml(alloc, body) catch |e| switch (e) {
|
||||||
defer esc.deinit(alloc);
|
error.OutOfMemory => return e,
|
||||||
for (body) |c| {
|
};
|
||||||
switch (c) {
|
defer alloc.free(html);
|
||||||
'"' => 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);
|
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);
|
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;
|
return ev.val;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -100,10 +215,14 @@ 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 }; // 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.
|
/// 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.)
|
||||||
|
///
|
||||||
|
/// 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 } {
|
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)
|
const url = if (since.len > 0)
|
||||||
try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/sync?timeout={d}&since={s}", .{ HOMESERVER, timeout_ms, since })
|
try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/sync?timeout={d}&since={s}", .{ HOMESERVER, timeout_ms, since })
|
||||||
@@ -113,29 +232,97 @@ pub fn sync(alloc: std.mem.Allocator, io: Io, token: []const u8, since: []const
|
|||||||
const resp = try http(alloc, io, "GET", url, token, null);
|
const resp = try http(alloc, io, "GET", url, token, null);
|
||||||
defer alloc.free(resp);
|
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;
|
var events: std.ArrayList(Event) = .empty;
|
||||||
errdefer events.deinit(alloc);
|
errdefer events.deinit(alloc);
|
||||||
|
|
||||||
var i: usize = 0;
|
var arena_state = std.heap.ArenaAllocator.init(alloc);
|
||||||
while (std.mem.indexOfPos(u8, resp, i, "\"m.room.message\"")) |pos| {
|
defer arena_state.deinit();
|
||||||
defer i = pos + 15;
|
const jalloc = arena_state.allocator();
|
||||||
// Event layout: sender (and event_id) come BEFORE the type marker,
|
|
||||||
// body comes after it in content — window spans both sides.
|
var parsed = std.json.parseFromSlice(std.json.Value, jalloc, resp, .{}) catch return error.BadSyncJson;
|
||||||
const win_start = pos -| 600;
|
defer parsed.deinit();
|
||||||
const win_end = @min(resp.len, pos + 1500);
|
|
||||||
const win = resp[win_start..win_end];
|
const root = switch (parsed.value) {
|
||||||
const sender = letta.jsonStr(alloc, win, 0, "sender") orelse continue;
|
.object => |o| o,
|
||||||
const body = letta.jsonStr(alloc, win, 0, "body") orelse {
|
else => return error.BadSyncJson,
|
||||||
alloc.free(sender.val);
|
};
|
||||||
|
|
||||||
|
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;
|
continue;
|
||||||
};
|
};
|
||||||
events.append(alloc, .{ .sender = sender.val, .body = body.val }) catch {
|
events.append(alloc, .{ .sender = sender_copy, .body = body_copy, .ts = ts }) catch {
|
||||||
alloc.free(sender.val);
|
alloc.free(sender_copy);
|
||||||
alloc.free(body.val);
|
alloc.free(body_copy);
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
}
|
||||||
return .{ .events = try events.toOwnedSlice(alloc), .next = next };
|
return .{ .events = try events.toOwnedSlice(alloc), .next = next };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+47
-5
@@ -55,13 +55,42 @@ pub fn main(init: std.process.Init) !void {
|
|||||||
defer if (owned_room) |r| alloc.free(r);
|
defer if (owned_room) |r| alloc.free(r);
|
||||||
std.debug.print("matrix_harness: room {s}\n", .{room_id});
|
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);
|
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));
|
||||||
|
var backlog: std.ArrayListUnmanaged(u8) = .empty;
|
||||||
|
defer backlog.deinit(alloc);
|
||||||
for (first.events) |ev| {
|
for (first.events) |ev| {
|
||||||
alloc.free(ev.sender);
|
defer alloc.free(ev.sender);
|
||||||
alloc.free(ev.body);
|
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);
|
alloc.free(first.events);
|
||||||
var since = first.next;
|
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 });
|
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)});
|
std.debug.print("matrix_harness: letta failed: {s}\n", .{@errorName(e)});
|
||||||
const err_txt = "agent failed — see harness logs";
|
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 |_| {}
|
if (matrix.sendText(alloc, io, session.token, room_id, err_txt)) |ev_id| alloc.free(ev_id) else |_| {}
|
||||||
|
|||||||
Reference in new issue
Block a user