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>
This commit is contained in:
pierreandLetta Code committed 2026-09-20 15:26:12 +03:00
1 parent 8d89e3baeb
commit 81fcbabefc
4 files changed
+209 -54

No files matched your search

+44 -13
View File
@@ -54,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
@@ -176,19 +180,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);
@@ -198,11 +205,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);
},
} }
} }
+9 -2
View File
@@ -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 /// 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 /// the event being replaced; `new_body` is the full replacement text
/// plain-text `body` carries the "* " legacy-edit prefix for old clients. /// (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 { 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 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 }); 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); 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; var payload: std.ArrayList(u8) = .empty;
defer payload.deinit(alloc); defer payload.deinit(alloc);
try payload.appendSlice(alloc, "{\"msgtype\":\"m.text\",\"body\":\"* "); try payload.appendSlice(alloc, "{\"msgtype\":\"m.text\",\"body\":\"* ");
try jsonEscape(alloc, &payload, new_body); try jsonEscape(alloc, &payload, new_body);
try payload.appendSlice(alloc, "\",\"m.new_content\":{\"msgtype\":\"m.text\",\"body\":\""); try payload.appendSlice(alloc, "\",\"m.new_content\":{\"msgtype\":\"m.text\",\"body\":\"");
try jsonEscape(alloc, &payload, new_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 payload.appendSlice(alloc, "\"},\"m.relates_to\":{\"rel_type\":\"m.replace\",\"event_id\":\"");
try jsonEscape(alloc, &payload, orig_event_id); try jsonEscape(alloc, &payload, orig_event_id);
try payload.appendSlice(alloc, "\"}}"); try payload.appendSlice(alloc, "\"}}");
+145 -39
View File
@@ -194,8 +194,11 @@ fn endTurn() void {
// it (m.replace) as the agent makes tool calls, so the DM shows what's // 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. // happening without spamming. The final reply is a separate fresh message.
/// Max tool-call lines kept in the bubble before collapsing to "+N more". /// Max tool-call blocks kept in the bubble before collapsing to "+N earlier".
const PROGRESS_MAX_LINES = 10; 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). /// Char cap for the bubble body (edits re-send the full body every time).
const PROGRESS_MAX_CHARS = 3500; const PROGRESS_MAX_CHARS = 3500;
@@ -205,13 +208,17 @@ const Progress = struct {
token: []const u8, token: []const u8,
room: []const u8, room: []const u8,
event_id: ?[]u8 = null, // the bubble message we keep editing event_id: ?[]u8 = null, // the bubble message we keep editing
lines: std.ArrayListUnmanaged(u8) = .empty, // accumulated tool lines blocks: std.ArrayListUnmanaged([]u8) = .empty, // markdown per tool call
line_count: usize = 0, // total steps seen (for "+N more") reply_event_id: ?[]u8 = null, // the streamed-reply message (edited live)
body_chars: usize = 0, reply: std.ArrayListUnmanaged(u8) = .empty, // accumulated reply text
replied: bool = false, // true once the full reply was edited in
fn deinit(self: *Progress) void { fn deinit(self: *Progress) void {
if (self.event_id) |e| self.alloc.free(e); 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). /// Send the initial bubble. Failure is non-fatal (progress is cosmetic).
@@ -220,41 +227,98 @@ const Progress = struct {
self.event_id = id; 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 <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 { fn body(self: *Progress) ![]u8 {
var out: std.ArrayListUnmanaged(u8) = .empty; var out: std.ArrayListUnmanaged(u8) = .empty;
errdefer out.deinit(self.alloc); errdefer out.deinit(self.alloc);
if (self.line_count > PROGRESS_MAX_LINES or self.body_chars > PROGRESS_MAX_CHARS) { const visible: usize = @min(self.blocks.items.len, PROGRESS_MAX_BLOCKS);
// Keep only the last PROGRESS_MAX_LINES lines. const hidden = self.blocks.items.len - visible;
var kept: usize = 0; if (hidden > 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;
try out.appendSlice(self.alloc, "… +"); 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, std.fmt.bufPrint(&num_buf, "{d}", .{hidden}) catch "?");
try out.appendSlice(self.alloc, " more\n"); try out.appendSlice(self.alloc, " earlier calls\n");
try out.appendSlice(self.alloc, self.lines.items[start..]);
} else {
try out.appendSlice(self.alloc, self.lines.items);
} }
// 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); return out.toOwnedSlice(self.alloc);
} }
/// Append a step line and push an edit. Non-fatal on failure. fn push(self: *Progress) void {
fn step(self: *Progress, line: []const u8) void { const id = self.event_id orelse {
self.lines.appendSlice(self.alloc, line) catch return; std.debug.print("DBG push: no event_id, blocks={d}\n", .{self.blocks.items.len});
self.lines.append(self.alloc, '\n') catch return; return;
self.line_count += 1; };
self.body_chars += line.len + 1; const text = self.body() catch |e| {
const id = self.event_id orelse return; std.debug.print("DBG push: body failed: {s}\n", .{@errorName(e)});
const text = self.body() catch return; return;
};
defer self.alloc.free(text); 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| { if (matrix.editText(self.alloc, self.io, self.token, self.room, id, text)) |new_id| {
self.alloc.free(new_id); self.alloc.free(new_id);
} else |_| { } else |_| {
@@ -263,23 +327,60 @@ const Progress = struct {
std.debug.print("matrix_harness: progress edit failed (dropping bubble)\n", .{}); 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 { fn onProgressEvent(ctx: ?*anyopaque, ev: letta.Event) void {
const p: *Progress = @ptrCast(@alignCast(ctx orelse return)); const p: *Progress = @ptrCast(@alignCast(ctx orelse return));
switch (ev) { switch (ev) {
.step => |line| p.step(line), .tool_call => |tc| {
.chunk => {}, // streamed reply text: comes back as the final reply 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 /// One full agent turn with a live tool-call bubble and a streamed reply
/// (allocated, caller frees). The bubble is left in the room as a record of /// message. Returns the reply text (allocated, caller frees). If the reply
/// the turn; the reply is sent by the caller as a fresh message. /// 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 { fn runAgentTurn(p: *Progress, prompt: []const u8) ![]u8 {
p.begin(); 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 { fn loadSeenTs(io: std.Io) i64 {
@@ -528,7 +629,12 @@ pub fn main(init: std.process.Init) !void {
defer endTurn(); 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)});
} }
} }