Compare commits

..
Author SHA1 Message Date
pierreandLetta Code 27327e94b0 letta: fix truncated replies — process final unterminated line at EOF
The stream reader only handled lines terminated by newline. The letta
CLI final JSON line arrives without trailing newline at EOF, so the last
assistant message (or result line) sat in the line buffer and was
silently dropped — replies ended mid-sentence. Also: read errors no
longer silently end collection (logged), and stream completeness is
tracked via the result line with a warning when absent.

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

Co-Authored-By: Letta Code <noreply@letta.com>
2026-09-21 06:02:59 +03:00
pierreandLetta Code 81fcbabefc harness: markdown tool bubble (bold names, fenced outputs, --- separator) + streamed reply message; letta: structured tool_call/tool_return events, escaped-key arg extraction; matrix: editText with formatted_body
Fixes included: bubble body off-by-one, reply chunk duplication, use-after-free in escaped-key truncation, fence/header concatenation swallowing headers, Read file_path arg extraction.

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

Co-Authored-By: Letta Code <noreply@letta.com>
2026-09-20 15:26:12 +03:00
pierreandLetta Code 8d89e3baeb harness: live tool-call bubble during agent turns (one message, m.replace edits, capped at 10 lines/3.5k chars, final reply as fresh message)
👾 Generated with [Letta Code](https://letta.com)

Co-Authored-By: Letta Code <noreply@letta.com>
2026-09-20 14:20:13 +03:00
pierreandLetta Code 58a66700c6 harness: pin bridge conversation to local-conv-487 (CONVERSATION_ID constant, --agent omitted for pinned convs)
👾 Generated with [Letta Code](https://letta.com)

Co-Authored-By: Letta Code <noreply@letta.com>
2026-09-20 13:54:15 +03:00
pierre 65e2ee083e setTyping: failure diagnostics (journald log on error) — API confirmed working via curl (200); if Element shows no indicator during turns, it's display-side 2026-09-06 13:12:38 +03:00
pierre b958eaab4c typing: refresh every 2s 2026-09-05 11:43:30 +03:00
pierre bf40b50f52 harness: typing indicator during agent turns (keepalive refreshes 60s window, cleared on reply) 2026-09-05 03:15:58 +03:00
pierre 2f0890a61c gitignore build artifacts + local letta dir 2026-09-05 01:21:31 +03:00
pierre 861ac201cb untrack .matrix-env (contains bridge password) — never commit secrets 2026-09-05 01:11:54 +03:00
pierre cbe99e89d0 harness: attachments named by real m.filename, not the caption; mimetype-derived extension fallback 2026-09-05 01:11:29 +03:00
pierreandLetta Code e212e3e120 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>
2026-09-03 23:51:09 +03:00
pierre dcf6d8df75 harness: proactive send mode + persistent backlog seen-state
- 'matrix_harness send "text"' — one-shot unprompted messaging to the
  pinned room (agent announcements/alerts)
- backlog no longer re-forwards read messages: last homeserver
  origin_server_ts persisted to /tmp/matrix-harness.state, advanced on
  every processed event; events at/below it are skipped
- 30-min staleness window retained as outer bound
2026-09-03 00:58:42 +03:00
pierre 9f1a839a22 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)
2026-09-02 08:57:43 +03:00
pierre 77f3150af4 matrix: widen sync event window backwards — sender precedes type marker 2026-09-02 02:22:38 +03:00
5 changed files with 1068 additions and 56 deletions

No files matched your search

+5 -4
View File
@@ -1,4 +1,5 @@
zig-pkg .matrix-env
zig-out .letta/
.letta .zig-cache/
.zig-cache zig-out/
zig-pkg/
+71 -15
View File
@@ -3,6 +3,11 @@ const Io = std.Io;
const LETTA_PATH = "/home/pierre/.bun/bin/letta"; const LETTA_PATH = "/home/pierre/.bun/bin/letta";
/// Conversation the matrix bridge routes messages into. Pinned to a specific
/// conversation ID (not "default") so a tainted/default conversation can be
/// swapped without touching the bridge.
pub const CONVERSATION_ID = "local-conv-487";
/// Runs `letta -n <agent> -p <prompt>`, captures everything it writes to /// Runs `letta -n <agent> -p <prompt>`, captures everything it writes to
/// stdout, and returns it as an owned slice (caller frees it). /// stdout, and returns it as an owned slice (caller frees it).
pub fn infer( pub fn infer(
@@ -11,8 +16,10 @@ pub fn infer(
agent: []const u8, agent: []const u8,
prompt: []const u8, prompt: []const u8,
) ![]u8 { ) ![]u8 {
_ = agent; // unused: pinned conversation implies the agent on the CLI side
var child = try std.process.spawn(io, .{ var child = try std.process.spawn(io, .{
.argv = &.{ LETTA_PATH, "--agent", agent, "-p", prompt, "--conversation", "default", "--toolset", "default" }, // Pinned conversation: --agent must be omitted for non-default conversations.
.argv = &.{ LETTA_PATH, "-p", prompt, "--conversation", CONVERSATION_ID, "--toolset", "default" },
.stdout = .pipe, .stdout = .pipe,
.stderr = .inherit, // letta's error output goes straight to your terminal .stderr = .inherit, // letta's error output goes straight to your terminal
}); });
@@ -47,8 +54,12 @@ pub fn infer(
// ------------------------------------------------------------- streaming --- // ------------------------------------------------------------- streaming ---
pub const Event = union(enum) { pub const Event = union(enum) {
step: []const u8, // tool call / return, pre-formatted line step: []const u8, // tool call / return, pre-formatted line (legacy)
chunk: []const u8, // streamed reply text chunk: []const u8, // streamed reply text
/// Structured tool event (caller frees both fields).
tool_call: struct { name: []u8, args: []u8 },
/// Structured tool output (caller frees).
tool_return: []u8,
}; };
/// Like infer(), but spawns letta with --output-format stream-json and calls /// Like infer(), but spawns letta with --output-format stream-json and calls
@@ -109,8 +120,12 @@ pub fn streamInferConv(
var line_len: usize = 0; var line_len: usize = 0;
var buf: [16 * 1024]u8 = undefined; var buf: [16 * 1024]u8 = undefined;
var stream_ok = false; // did we see a proper result/completion line?
while (true) { while (true) {
const n = stdout_file.readStreaming(io, &.{&buf}) catch break; const n = stdout_file.readStreaming(io, &.{&buf}) catch |e| {
std.debug.print("inferon/letta: stdout read error {s} — reply may be TRUNCATED\n", .{@errorName(e)});
break;
};
if (n == 0) break; if (n == 0) break;
raw_tail.appendSlice(allocator, buf[0..n]) catch {}; raw_tail.appendSlice(allocator, buf[0..n]) catch {};
if (raw_tail.items.len > 4096) { if (raw_tail.items.len > 4096) {
@@ -122,6 +137,7 @@ pub fn streamInferConv(
if (ch == '\n') { if (ch == '\n') {
if (line_len > 0) { if (line_len > 0) {
handleLine(allocator, line_buf[0..line_len], ctx, onEvent, &reply); handleLine(allocator, line_buf[0..line_len], ctx, onEvent, &reply);
if (std.mem.indexOf(u8, line_buf[0..line_len], "\"type\":\"result\"") != null) stream_ok = true;
line_len = 0; line_len = 0;
} }
} else if (line_len < line_buf.len) { } else if (line_len < line_buf.len) {
@@ -130,8 +146,21 @@ pub fn streamInferConv(
} }
} }
} }
// CRITICAL: the last line may arrive WITHOUT a trailing newline (EOF
// right after the final JSON). It is still a complete line — dropping
// it silently truncated replies (the text of the last assistant message
// or the result line was lost).
if (line_len > 0) {
handleLine(allocator, line_buf[0..line_len], ctx, onEvent, &reply);
if (std.mem.indexOf(u8, line_buf[0..line_len], "\"type\":\"result\"") != null) stream_ok = true;
}
_ = child.wait(io) catch {}; _ = child.wait(io) catch {};
if (!stream_ok) {
std.debug.print("inferon/letta: stream ended WITHOUT result line — reply likely TRUNCATED ({d} bytes collected). tail:\n{s}\n", .{ reply.items.len, raw_tail.items });
}
std.debug.print("inferon/letta: stream complete, reply {d} bytes, result_line={}\n", .{ reply.items.len, stream_ok });
if (reply.items.len == 0 and raw_tail.items.len > 0) { if (reply.items.len == 0 and raw_tail.items.len > 0) {
std.debug.print("inferon/letta: EMPTY reply, raw stream tail:\n{s}\n", .{raw_tail.items}); std.debug.print("inferon/letta: EMPTY reply, raw stream tail:\n{s}\n", .{raw_tail.items});
} }
@@ -169,19 +198,22 @@ fn handleLine(
if (std.mem.eql(u8, mt, "tool_call_message")) { if (std.mem.eql(u8, mt, "tool_call_message")) {
const name = dupeStr(allocator, line, "name") orelse allocator.dupe(u8, "?") catch return; const name = dupeStr(allocator, line, "name") orelse allocator.dupe(u8, "?") catch return;
defer allocator.free(name); // Different tools carry their primary arg under different keys;
// try the common ones so the summary is never empty.
const args = dupeStr(allocator, line, "command") orelse const args = dupeStr(allocator, line, "command") orelse
dupeStr(allocator, line, "description") orelse allocator.dupe(u8, "") catch return; dupeStr(allocator, line, "file_path") orelse
defer allocator.free(args); dupeStr(allocator, line, "path") orelse
const text = std.fmt.allocPrint(allocator, "> {s} {s}", .{ name, args }) catch return; dupeStr(allocator, line, "query") orelse
defer allocator.free(text); dupeStr(allocator, line, "url") orelse
onEvent(ctx, .{ .step = text }); dupeStr(allocator, line, "pattern") orelse
dupeStr(allocator, line, "description") orelse allocator.dupe(u8, "") catch {
allocator.free(name);
return;
};
onEvent(ctx, .{ .tool_call = .{ .name = name, .args = args } });
} else if (std.mem.eql(u8, mt, "tool_return_message")) { } else if (std.mem.eql(u8, mt, "tool_return_message")) {
const ret = dupeStr(allocator, line, "tool_return") orelse allocator.dupe(u8, "") catch return; const ret = dupeStr(allocator, line, "tool_return") orelse allocator.dupe(u8, "") catch return;
defer allocator.free(ret); onEvent(ctx, .{ .tool_return = ret });
const text = std.fmt.allocPrint(allocator, " {s}", .{ret}) catch return;
defer allocator.free(text);
onEvent(ctx, .{ .step = text });
} else if (std.mem.eql(u8, mt, "assistant_message")) { } else if (std.mem.eql(u8, mt, "assistant_message")) {
const txt = dupeStr(allocator, line, "text") orelse return; const txt = dupeStr(allocator, line, "text") orelse return;
defer allocator.free(txt); defer allocator.free(txt);
@@ -191,11 +223,35 @@ fn handleLine(
} }
/// Extract "key":"value" (with escape handling), allocated with `allocator`. /// Extract "key":"value" (with escape handling), allocated with `allocator`.
/// Also matches the escaped variant `\"key\":\"` — nested JSON strings (e.g.
/// tool_call arguments) carry their keys escaped on the wire.
pub fn dupeStr(allocator: std.mem.Allocator, line: []const u8, key: []const u8) ?[]u8 { pub fn dupeStr(allocator: std.mem.Allocator, line: []const u8, key: []const u8) ?[]u8 {
var pat_buf: [64]u8 = undefined; var pat_buf: [64]u8 = undefined;
const pat = std.fmt.bufPrint(&pat_buf, "\"{s}\":\"", .{key}) catch return null; const pat = std.fmt.bufPrint(&pat_buf, "\"{s}\":\"", .{key}) catch return null;
const start = std.mem.indexOf(u8, line, pat) orelse return null; const start = std.mem.indexOf(u8, line, pat) orelse {
var it = line[start + pat.len ..]; // Escaped variant: \"key\":\" inside a nested JSON string.
var pat2_buf: [70]u8 = undefined;
const pat2 = std.fmt.bufPrint(&pat2_buf, "\\\"{s}\\\":\\\"", .{key}) catch return null;
const s2 = std.mem.indexOf(u8, line, pat2) orelse return null;
const v = dupeValue(allocator, line[s2 + pat2.len ..]) orelse return null;
// The escaped variant's closing quote is escaped too, so the value
// decoder overruns into the following keys. Cut at the first
// unescaped key boundary left in the decoded text.
if (std.mem.indexOf(u8, v, "\",\"")) |cut| {
const out = allocator.dupe(u8, v[0..cut]) catch {
allocator.free(v);
return null;
};
allocator.free(v);
return out;
}
return v;
};
return dupeValue(allocator, line[start + pat.len ..]);
}
fn dupeValue(allocator: std.mem.Allocator, from: []const u8) ?[]u8 {
var it = from;
var out: std.ArrayList(u8) = .empty; var out: std.ArrayList(u8) = .empty;
errdefer out.deinit(allocator); errdefer out.deinit(allocator);
+11
View File
@@ -202,6 +202,17 @@ fn onLettaEvent(_: ?*anyopaque, ev: letta.Event) void {
if (last == '.' or last == '!' or last == '?') flushSentence(false); if (last == '.' or last == '!' or last == '?') flushSentence(false);
} }
}, },
.tool_call => |tc| {
defer allocator.free(tc.name);
defer allocator.free(tc.args);
const line = std.fmt.allocPrint(allocator, "> {s} {s}", .{ tc.name, tc.args }) catch return;
pushOverlay(.step, line);
},
.tool_return => |ret| {
defer allocator.free(ret);
const line = std.fmt.allocPrint(allocator, " {s}", .{ret}) catch return;
pushOverlay(.step, line);
},
} }
} }
+435 -30
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}"}}
@@ -67,26 +183,203 @@ 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, "&lt;");
i += 1;
},
'>' => {
try out.appendSlice(alloc, "&gt;");
i += 1;
},
'&' => {
try out.appendSlice(alloc, "&amp;");
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);
}
/// Typing notification. timeout_ms: how long the indicator lasts server-side
/// (refresh periodically for long turns). typing=false clears it immediately.
pub fn setTyping(alloc: std.mem.Allocator, io: Io, token: []const u8, user_id: []const u8, room: []const u8, typing: bool, timeout_ms: i64) void {
var url_buf: std.ArrayListUnmanaged(u8) = .empty;
defer url_buf.deinit(alloc);
const user_esc = std.fmt.allocPrint(alloc, "{s}", .{user_id}) catch return;
defer alloc.free(user_esc);
// user ids contain chars that are fine unescaped in a path segment for
// conduwuit, but encode the ':' minimally
var user_path: std.ArrayListUnmanaged(u8) = .empty;
defer user_path.deinit(alloc);
for (user_esc) |ch| {
if (ch == ':') user_path.appendSlice(alloc, "%3A") catch return else user_path.append(alloc, ch) catch return;
}
const url = std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/typing/{s}", .{ HOMESERVER, room, user_path.items }) catch return;
defer alloc.free(url);
var payload: std.ArrayListUnmanaged(u8) = .empty;
defer payload.deinit(alloc);
const body = if (typing)
std.fmt.allocPrint(alloc, "{{\"typing\":true,\"timeout\":{d}}}", .{timeout_ms}) catch return
else
alloc.dupe(u8, "{\"typing\":false}") catch return;
defer alloc.free(body);
if (httpRaw(alloc, io, "PUT", url, token, "application/json", body)) |resp| {
alloc.free(resp);
} else |e| {
std.debug.print("matrix_harness: setTyping failed: {s}\n", .{@errorName(e)});
}
}
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);
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;
}
/// Edit an existing m.text message in place (m.replace). `orig_event_id` is
/// the event being replaced; `new_body` is the full replacement text
/// (markdown — rendered via formatted_body, with plain-text fallback).
pub fn editText(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []const u8, orig_event_id: []const u8, new_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-edit-{d}", .{ HOMESERVER, room, ts });
defer alloc.free(url);
const html = markdownToHtml(alloc, new_body) catch |e| switch (e) {
error.OutOfMemory => return e,
};
defer alloc.free(html);
var payload: std.ArrayList(u8) = .empty;
defer payload.deinit(alloc);
try payload.appendSlice(alloc, "{\"msgtype\":\"m.text\",\"body\":\"* ");
try jsonEscape(alloc, &payload, new_body);
try payload.appendSlice(alloc, "\",\"m.new_content\":{\"msgtype\":\"m.text\",\"body\":\"");
try jsonEscape(alloc, &payload, new_body);
try payload.appendSlice(alloc, "\",\"format\":\"org.matrix.custom.html\",\"formatted_body\":\"");
try jsonEscape(alloc, &payload, html);
try payload.appendSlice(alloc, "\"},\"m.relates_to\":{\"rel_type\":\"m.replace\",\"event_id\":\"");
try jsonEscape(alloc, &payload, orig_event_id);
try payload.appendSlice(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 return error.SendFailed;
return ev.val; return ev.val;
@@ -100,10 +393,24 @@ 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,
/// 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,
/// Original attachment filename (distinct from body, which is the caption).
filename: []u8 = &.{}, // allocated; caller frees
};
/// 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,25 +420,123 @@ 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();
const win_end = @min(resp.len, pos + 1500);
const win = resp[pos..win_end]; var parsed = std.json.parseFromSlice(std.json.Value, jalloc, resp, .{}) catch return error.BadSyncJson;
const sender = letta.jsonStr(alloc, win, 0, "sender") orelse continue; defer parsed.deinit();
const body = letta.jsonStr(alloc, win, 0, "body") orelse {
alloc.free(sender.val); const root = switch (parsed.value) {
continue; .object => |o| o,
}; else => return error.BadSyncJson,
events.append(alloc, .{ .sender = sender.val, .body = body.val }) catch { };
alloc.free(sender.val);
alloc.free(body.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;
// Attachment fields (m.file / m.image / m.audio / m.video).
var url_copy: ?[]u8 = null;
var fname_copy: []u8 = "";
if (content.get("filename")) |f| switch (f) {
.string => |v| fname_copy = alloc.dupe(u8, v) catch "",
else => {},
};
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, .url = url_copy, .filename = fname_copy, .mimetype = mime, .size = size }) catch {
alloc.free(sender_copy);
alloc.free(body_copy);
};
}
} }
return .{ .events = try events.toOwnedSlice(alloc), .next = next }; return .{ .events = try events.toOwnedSlice(alloc), .next = next };
} }
+546 -7
View File
@@ -20,9 +20,415 @@ fn env_or(name: [:0]const u8, default: []const u8) []const u8 {
return default; 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);
}
// --------------------------------------------------- tool-call progress ---
//
// Live "working" bubble: during a turn we send ONE message and keep editing
// it (m.replace) as the agent makes tool calls, so the DM shows what's
// happening without spamming. The final reply is a separate fresh message.
/// Max tool-call blocks kept in the bubble before collapsing to "+N earlier".
const PROGRESS_MAX_BLOCKS = 10;
/// Char cap per tool arg summary / tool output block.
const PROGRESS_ARG_CHARS = 120;
const PROGRESS_RET_CHARS = 400;
/// Char cap for the bubble body (edits re-send the full body every time).
const PROGRESS_MAX_CHARS = 3500;
const Progress = struct {
alloc: std.mem.Allocator,
io: std.Io,
token: []const u8,
room: []const u8,
event_id: ?[]u8 = null, // the bubble message we keep editing
blocks: std.ArrayListUnmanaged([]u8) = .empty, // markdown per tool call
reply_event_id: ?[]u8 = null, // the streamed-reply message (edited live)
reply: std.ArrayListUnmanaged(u8) = .empty, // accumulated reply text
replied: bool = false, // true once the full reply was edited in
fn deinit(self: *Progress) void {
if (self.event_id) |e| self.alloc.free(e);
if (self.reply_event_id) |e| self.alloc.free(e);
for (self.blocks.items) |b| self.alloc.free(b);
self.blocks.deinit(self.alloc);
self.reply.deinit(self.alloc);
}
/// Send the initial bubble. Failure is non-fatal (progress is cosmetic).
fn begin(self: *Progress) void {
const id = matrix.sendText(self.alloc, self.io, self.token, self.room, "⚙️ working…") catch return;
self.event_id = id;
}
/// One-line summary: first line only, truncated, no backticks (they'd
/// interact with markdown).
fn summarize(self: *Progress, text: []const u8, cap: usize) ![]u8 {
const nl = std.mem.indexOfScalar(u8, text, '\n') orelse text.len;
var line = text[0..nl];
// Neutralize backtick runs so they can't break our fences.
var clean: std.ArrayListUnmanaged(u8) = .empty;
var run: usize = 0;
for (line) |c| {
if (c == '`') {
run += 1;
if (run < 3) try clean.append(self.alloc, c);
} else {
run = 0;
try clean.append(self.alloc, c);
}
if (clean.items.len >= cap) break;
}
line = clean.items;
if (text.len > nl or line.len >= cap and text.len > line.len) {
try clean.appendSlice(self.alloc, "…");
}
return clean.toOwnedSlice(self.alloc);
}
/// Append a markdown block for one tool call and push an edit. Each
/// block: bold name + arg summary, then the output in a fenced code
/// block once the return arrives. Non-fatal on failure.
fn call(self: *Progress, name: []const u8, args: []const u8) void {
const arg_line = self.summarize(args, PROGRESS_ARG_CHARS) catch return;
defer self.alloc.free(arg_line);
const block = std.fmt.allocPrint(self.alloc, "**{s}**: `{s}`\n", .{ name, arg_line }) catch return;
self.blocks.append(self.alloc, block) catch {
self.alloc.free(block);
return;
};
self.push();
}
fn ret(self: *Progress, out: []const u8) void {
const out_line = self.summarize(out, PROGRESS_RET_CHARS) catch return;
defer self.alloc.free(out_line);
if (self.blocks.items.len == 0) return;
const last = self.blocks.items[self.blocks.items.len - 1];
// Trailing newline is load-bearing: without it the next block's
// header concatenates onto the closing fence and gets swallowed
// by the <pre> in Element's renderer.
const with_fence = std.fmt.allocPrint(self.alloc, "{s}---\n```\n{s}\n```\n", .{ last, out_line }) catch return;
self.alloc.free(last);
self.blocks.items[self.blocks.items.len - 1] = with_fence;
self.push();
}
/// Current bubble text (caller frees): hidden-count header + last blocks.
fn body(self: *Progress) ![]u8 {
var out: std.ArrayListUnmanaged(u8) = .empty;
errdefer out.deinit(self.alloc);
const visible: usize = @min(self.blocks.items.len, PROGRESS_MAX_BLOCKS);
const hidden = self.blocks.items.len - visible;
if (hidden > 0) {
try out.appendSlice(self.alloc, "… +");
var num_buf: [16]u8 = undefined;
try out.appendSlice(self.alloc, std.fmt.bufPrint(&num_buf, "{d}", .{hidden}) catch "?");
try out.appendSlice(self.alloc, " earlier calls\n");
}
// Walk backwards from the newest block, accumulating until the char
// cap is exceeded; everything from there on is visible. Keep at
// least the newest block, and at most PROGRESS_MAX_BLOCKS.
var start: usize = self.blocks.items.len;
var total: usize = 0;
while (start > 0) {
const l = self.blocks.items[start - 1].len;
if (total > 0 and total + l > PROGRESS_MAX_CHARS) break;
total += l;
start -= 1;
}
if (self.blocks.items.len - start > visible) start = self.blocks.items.len - visible;
for (self.blocks.items[start..]) |b| try out.appendSlice(self.alloc, b);
return out.toOwnedSlice(self.alloc);
}
fn push(self: *Progress) void {
const id = self.event_id orelse {
std.debug.print("DBG push: no event_id, blocks={d}\n", .{self.blocks.items.len});
return;
};
const text = self.body() catch |e| {
std.debug.print("DBG push: body failed: {s}\n", .{@errorName(e)});
return;
};
defer self.alloc.free(text);
std.debug.print("DBG push: blocks={d} text.len={d}\n", .{ self.blocks.items.len, text.len });
if (matrix.editText(self.alloc, self.io, self.token, self.room, id, text)) |new_id| {
self.alloc.free(new_id);
} else |_| {
// Transient homeserver hiccups shouldn't kill the turn; drop the
// bubble rather than retrying (next step may re-establish it).
std.debug.print("matrix_harness: progress edit failed (dropping bubble)\n", .{});
}
}
/// Streamed reply text: on first chunk send a fresh message, then edit
/// it as more text arrives. `final` stamps the complete reply (callers
/// skip their own sendText when this returned true).
fn replyChunk(self: *Progress, text: []const u8, final: bool) bool {
self.reply.appendSlice(self.alloc, text) catch return false;
if (self.reply_event_id == null) {
const id = matrix.sendText(self.alloc, self.io, self.token, self.room, "…") catch return false;
self.reply_event_id = id;
}
const id = self.reply_event_id.?;
if (matrix.editText(self.alloc, self.io, self.token, self.room, id, self.reply.items)) |new_id| {
self.alloc.free(new_id);
self.replied = final;
return final;
} else |_| return false;
}
};
/// streamInfer callback: forward structured events into the Progress bubble.
fn onProgressEvent(ctx: ?*anyopaque, ev: letta.Event) void {
const p: *Progress = @ptrCast(@alignCast(ctx orelse return));
switch (ev) {
.tool_call => |tc| {
defer p.alloc.free(tc.name);
defer p.alloc.free(tc.args);
p.call(tc.name, tc.args);
},
.tool_return => |out| {
defer p.alloc.free(out);
p.ret(out);
},
.chunk => |c| {
defer p.alloc.free(c);
_ = p.replyChunk(c, false);
},
.step => {}, // legacy pre-formatted lines: superseded by structured
}
}
/// One full agent turn with a live tool-call bubble and a streamed reply
/// message. Returns the reply text (allocated, caller frees). If the reply
/// was already streamed and finalized in-message (p.replied), the caller
/// should NOT send it again; the bubble stays as the tool-call record.
fn runAgentTurn(p: *Progress, prompt: []const u8) ![]u8 {
p.begin();
const reply = try letta.streamInferConv(p.io, p.alloc, "", letta.CONVERSATION_ID, prompt, p, onProgressEvent);
// Stamp the authoritative full reply into the streamed message: REPLACE
// the accumulated chunks (appending would duplicate the whole text).
if (p.reply_event_id != null) {
p.reply.clearRetainingCapacity();
_ = p.replyChunk(reply, true);
}
return reply;
}
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 { pub fn main(init: std.process.Init) !void {
const alloc = init.gpa; const alloc = init.gpa;
const io = init.io; 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(); letta.initConversationsDir();
const user = env_or("MATRIX_USER", "@chaos-prime:chaosmith.systems"); const user = env_or("MATRIX_USER", "@chaos-prime:chaosmith.systems");
@@ -55,13 +461,66 @@ 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);
for (first.events) |ev| { {
alloc.free(ev.sender); const now_ms: i64 = @intCast(@divTrunc(std.Io.Clock.now(.real, io).nanoseconds, std.time.ns_per_ms));
alloc.free(ev.body); 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();
var prog = Progress{ .alloc = alloc, .io = io, .token = session.token, .room = room_id };
defer prog.deinit();
if (runAgentTurn(&prog, 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); alloc.free(first.events);
var since = first.next; var since = first.next;
@@ -80,22 +539,102 @@ 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.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). // 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;
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);
var prog = Progress{ .alloc = alloc, .io = io, .token = session.token, .room = room_id };
defer prog.deinit();
const reply = runAgentTurn(&prog, 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 }); 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).
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);
var prog = Progress{ .alloc = alloc, .io = io, .token = session.token, .room = room_id };
defer prog.deinit();
const reply = runAgentTurn(&prog, prefixed) catch |e| {
endTurn();
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 |_| {}
continue; continue;
}; };
defer alloc.free(reply); defer alloc.free(reply);
defer endTurn();
std.debug.print("matrix_harness: reply: {s}\n", .{reply}); 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| { if (prog.replied) {
// Reply was streamed into its own message; final edit already
// carries the full authoritative text.
} else 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)}); std.debug.print("matrix_harness: send failed: {s}\n", .{@errorName(e)});
} }
} }