From 81fcbabefc4ca16eb51ce86bf0bb2fec55261e2e Mon Sep 17 00:00:00 2001 From: Pierre De Lancre Date: Sun, 20 Sep 2026 15:26:12 +0300 Subject: [PATCH] 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 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- src/letta.zig | 57 ++++++++++--- src/main.zig | 11 +++ src/matrix.zig | 11 ++- src/matrix_harness.zig | 184 ++++++++++++++++++++++++++++++++--------- 4 files changed, 209 insertions(+), 54 deletions(-) diff --git a/src/letta.zig b/src/letta.zig index 940a0eb..9b7a8b8 100644 --- a/src/letta.zig +++ b/src/letta.zig @@ -54,8 +54,12 @@ pub fn infer( // ------------------------------------------------------------- streaming --- 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 + /// 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 @@ -176,19 +180,22 @@ fn handleLine( if (std.mem.eql(u8, mt, "tool_call_message")) { 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 - dupeStr(allocator, line, "description") orelse allocator.dupe(u8, "") catch return; - defer allocator.free(args); - const text = std.fmt.allocPrint(allocator, "> {s} {s}", .{ name, args }) catch return; - defer allocator.free(text); - onEvent(ctx, .{ .step = text }); + dupeStr(allocator, line, "file_path") orelse + dupeStr(allocator, line, "path") orelse + dupeStr(allocator, line, "query") orelse + dupeStr(allocator, line, "url") orelse + 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")) { const ret = dupeStr(allocator, line, "tool_return") orelse allocator.dupe(u8, "") catch return; - defer allocator.free(ret); - const text = std.fmt.allocPrint(allocator, " {s}", .{ret}) catch return; - defer allocator.free(text); - onEvent(ctx, .{ .step = text }); + onEvent(ctx, .{ .tool_return = ret }); } else if (std.mem.eql(u8, mt, "assistant_message")) { const txt = dupeStr(allocator, line, "text") orelse return; defer allocator.free(txt); @@ -198,11 +205,35 @@ fn handleLine( } /// 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 { var pat_buf: [64]u8 = undefined; const pat = std.fmt.bufPrint(&pat_buf, "\"{s}\":\"", .{key}) catch return null; - const start = std.mem.indexOf(u8, line, pat) orelse return null; - var it = line[start + pat.len ..]; + const start = std.mem.indexOf(u8, line, pat) orelse { + // 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; errdefer out.deinit(allocator); diff --git a/src/main.zig b/src/main.zig index a18d9ac..e39647f 100644 --- a/src/main.zig +++ b/src/main.zig @@ -202,6 +202,17 @@ fn onLettaEvent(_: ?*anyopaque, ev: letta.Event) void { 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); + }, } } diff --git a/src/matrix.zig b/src/matrix.zig index a944df7..de4979b 100644 --- a/src/matrix.zig +++ b/src/matrix.zig @@ -355,19 +355,26 @@ pub fn sendText(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []con } /// 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. The -/// plain-text `body` carries the "* " legacy-edit prefix for old clients. +/// 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, "\"}}"); diff --git a/src/matrix_harness.zig b/src/matrix_harness.zig index 7477682..3cb7eca 100644 --- a/src/matrix_harness.zig +++ b/src/matrix_harness.zig @@ -194,8 +194,11 @@ fn endTurn() void { // 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 lines kept in the bubble before collapsing to "+N more". -const PROGRESS_MAX_LINES = 10; +/// 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; @@ -205,13 +208,17 @@ const Progress = struct { token: []const u8, room: []const u8, event_id: ?[]u8 = null, // the bubble message we keep editing - lines: std.ArrayListUnmanaged(u8) = .empty, // accumulated tool lines - line_count: usize = 0, // total steps seen (for "+N more") - body_chars: usize = 0, + 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); - self.lines.deinit(self.alloc); + 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). @@ -220,41 +227,98 @@ const Progress = struct { self.event_id = id; } - /// Current bubble text (caller frees). Collapses when over the caps. + /// 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
 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);
-        if (self.line_count > PROGRESS_MAX_LINES or self.body_chars > PROGRESS_MAX_CHARS) {
-            // Keep only the last PROGRESS_MAX_LINES lines.
-            var kept: usize = 0;
-            var start: usize = self.lines.items.len;
-            while (start > 0) {
-                const prev = std.mem.lastIndexOfScalar(u8, self.lines.items[0 .. start - 1], '\n') orelse break;
-                start = prev + 1;
-                kept += 1;
-                if (kept == PROGRESS_MAX_LINES) break;
-            }
-            const hidden = self.line_count - kept;
-            var num_buf: [16]u8 = undefined;
+        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, " more\n");
-            try out.appendSlice(self.alloc, self.lines.items[start..]);
-        } else {
-            try out.appendSlice(self.alloc, self.lines.items);
+            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);
     }
 
-    /// Append a step line and push an edit. Non-fatal on failure.
-    fn step(self: *Progress, line: []const u8) void {
-        self.lines.appendSlice(self.alloc, line) catch return;
-        self.lines.append(self.alloc, '\n') catch return;
-        self.line_count += 1;
-        self.body_chars += line.len + 1;
-        const id = self.event_id orelse return;
-        const text = self.body() catch return;
+    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 |_| {
@@ -263,23 +327,60 @@ const Progress = struct {
             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 step events into the Progress bubble.
+/// 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) {
-        .step => |line| p.step(line),
-        .chunk => {}, // streamed reply text: comes back as the final reply
+        .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. Returns the reply text
-/// (allocated, caller frees). The bubble is left in the room as a record of
-/// the turn; the reply is sent by the caller as a fresh message.
+/// 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();
-    return letta.streamInferConv(p.io, p.alloc, "", letta.CONVERSATION_ID, prompt, p, onProgressEvent);
+    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 {
@@ -528,7 +629,12 @@ pub fn main(init: std.process.Init) !void {
             defer endTurn();
 
             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)});
             }
         }