harness: typing indicator during agent turns (keepalive refreshes 60s window, cleared on reply)
This commit is contained in:
1 parent
2f0890a61c
commit
bf40b50f52
2 files changed
+70
No files matched your search
@@ -292,6 +292,33 @@ fn markdownToHtml(alloc: std.mem.Allocator, md: []const u8) ![]u8 {
|
||||
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);
|
||||
_ = httpRaw(alloc, io, "PUT", url, token, "application/json", body) catch null;
|
||||
}
|
||||
|
||||
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 url = try std.fmt.allocPrint(alloc, "{s}/_matrix/client/v3/rooms/{s}/send/m.room.message/infr-{d}", .{ HOMESERVER, room, ts });
|
||||
|
||||
@@ -152,6 +152,42 @@ fn saveAttachment(alloc: std.mem.Allocator, io: std.Io, token: []const u8, ev_ts
|
||||
/// 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, 60000);
|
||||
var slept: u32 = 0;
|
||||
while (slept < 45_000 and turn_busy.load(.acquire)) : (slept += 500) {
|
||||
io.sleep(.{ .nanoseconds = 500 * 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);
|
||||
}
|
||||
|
||||
fn loadSeenTs(io: std.Io) i64 {
|
||||
var f = std.Io.Dir.openFileAbsolute(io, STATE_FILE, .{}) catch return 0;
|
||||
defer f.close(io);
|
||||
@@ -277,6 +313,8 @@ pub fn main(init: std.process.Init) !void {
|
||||
// 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();
|
||||
if (letta.infer(io, alloc, agent, hist_msg)) |reply| {
|
||||
alloc.free(reply); // backlog replies are not sent back
|
||||
} else |e| {
|
||||
@@ -358,7 +396,9 @@ pub fn main(init: std.process.Init) !void {
|
||||
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);
|
||||
const reply = letta.infer(io, alloc, agent, note) catch |e| {
|
||||
endTurn();
|
||||
std.debug.print("matrix_harness: letta failed: {s}\n", .{@errorName(e)});
|
||||
continue;
|
||||
};
|
||||
@@ -376,13 +416,16 @@ pub fn main(init: std.process.Init) !void {
|
||||
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);
|
||||
const reply = letta.infer(io, alloc, agent, prefixed) catch |e| {
|
||||
endTurn();
|
||||
std.debug.print("matrix_harness: letta failed: {s}\n", .{@errorName(e)});
|
||||
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 |_| {}
|
||||
continue;
|
||||
};
|
||||
defer alloc.free(reply);
|
||||
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| {
|
||||
|
||||
Reference in new issue
Block a user