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>
This commit is contained in:
pierreandLetta Code committed 2026-09-20 14:20:13 +03:00
1 parent 58a66700c6
commit 8d89e3baeb
2 files changed
+127 -3

No files matched your search

+24
View File
@@ -354,6 +354,30 @@ pub fn sendText(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []con
return ev.val; 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). /// Join a room by id (accepts invites).
pub fn joinRoom(alloc: std.mem.Allocator, io: Io, token: []const u8, room: []const u8) !void { 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 }); const url = try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/join", .{ HOMESERVER, room });
+103 -3
View File
@@ -188,6 +188,100 @@ fn endTurn() void {
turn_busy.store(false, .release); 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 { fn loadSeenTs(io: std.Io) i64 {
var f = std.Io.Dir.openFileAbsolute(io, STATE_FILE, .{}) catch return 0; var f = std.Io.Dir.openFileAbsolute(io, STATE_FILE, .{}) catch return 0;
defer f.close(io); defer f.close(io);
@@ -315,7 +409,9 @@ pub fn main(init: std.process.Init) !void {
// Log, skip, keep polling. // Log, skip, keep polling.
beginTurn(io, session.token, session.user_id, room_id); beginTurn(io, session.token, session.user_id, room_id);
defer endTurn(); 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 alloc.free(reply); // backlog replies are not sent back
} else |e| { } else |e| {
std.debug.print("matrix_harness: backlog letta failed: {s} (skipping)\n", .{@errorName(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 }); 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); defer alloc.free(note);
beginTurn(io, session.token, session.user_id, room_id); 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(); endTurn();
std.debug.print("matrix_harness: letta failed: {s}\n", .{@errorName(e)}); std.debug.print("matrix_harness: letta failed: {s}\n", .{@errorName(e)});
continue; continue;
@@ -417,7 +515,9 @@ pub fn main(init: std.process.Init) !void {
defer alloc.free(prefixed); defer alloc.free(prefixed);
beginTurn(io, session.token, session.user_id, room_id); 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(); 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";