From 8d89e3baeb72463203b522a9743bc4562e0b5a6b Mon Sep 17 00:00:00 2001 From: Pierre De Lancre Date: Sun, 20 Sep 2026 14:20:13 +0300 Subject: [PATCH] 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) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 👾 Generated with [Letta Code](https://letta.com) Co-Authored-By: Letta Code --- src/matrix.zig | 24 ++++++++++ src/matrix_harness.zig | 106 +++++++++++++++++++++++++++++++++++++++-- 2 files changed, 127 insertions(+), 3 deletions(-) diff --git a/src/matrix.zig b/src/matrix.zig index 2e54aea..a944df7 100644 --- a/src/matrix.zig +++ b/src/matrix.zig @@ -354,6 +354,30 @@ pub fn sendText(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []con 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. The +/// plain-text `body` carries the "* " legacy-edit prefix for old clients. +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); + + 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, "\"},\"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); + const ev = letta.jsonStr(alloc, resp, 0, "event_id") orelse return error.SendFailed; + return ev.val; +} + /// Join a room by id (accepts invites). pub fn joinRoom(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []const u8) !void { const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/join", .{ HOMESERVER, room }); diff --git a/src/matrix_harness.zig b/src/matrix_harness.zig index 2f05c75..7477682 100644 --- a/src/matrix_harness.zig +++ b/src/matrix_harness.zig @@ -188,6 +188,100 @@ 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 lines kept in the bubble before collapsing to "+N more". +const PROGRESS_MAX_LINES = 10; +/// 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 + lines: std.ArrayListUnmanaged(u8) = .empty, // accumulated tool lines + line_count: usize = 0, // total steps seen (for "+N more") + body_chars: usize = 0, + + fn deinit(self: *Progress) void { + if (self.event_id) |e| self.alloc.free(e); + self.lines.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; + } + + /// Current bubble text (caller frees). Collapses when over the caps. + 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; + try out.appendSlice(self.alloc, "… +"); + 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); + } + 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; + defer self.alloc.free(text); + 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", .{}); + } + } +}; + +/// streamInfer callback: forward step 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 + } +} + +/// 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. +fn runAgentTurn(p: *Progress, prompt: []const u8) ![]u8 { + p.begin(); + return letta.streamInferConv(p.io, p.alloc, "", letta.CONVERSATION_ID, prompt, p, onProgressEvent); +} + fn loadSeenTs(io: std.Io) i64 { var f = std.Io.Dir.openFileAbsolute(io, STATE_FILE, .{}) catch return 0; defer f.close(io); @@ -315,7 +409,9 @@ pub fn main(init: std.process.Init) !void { // Log, skip, keep polling. beginTurn(io, session.token, session.user_id, room_id); defer endTurn(); - if (letta.infer(io, alloc, agent, hist_msg)) |reply| { + 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)}); @@ -397,7 +493,9 @@ pub fn main(init: std.process.Init) !void { 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); - const reply = letta.infer(io, alloc, agent, note) catch |e| { + 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; @@ -417,7 +515,9 @@ pub fn main(init: std.process.Init) !void { defer alloc.free(prefixed); beginTurn(io, session.token, session.user_id, room_id); - const reply = letta.infer(io, alloc, agent, prefixed) catch |e| { + 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)}); const err_txt = "agent failed — see harness logs";