Files
inferon/src/letta.zig
T

378 lines
14 KiB
Zig

const std = @import("std");
const Io = std.Io;
const LETTA_PATH = "/home/pierre/.bun/bin/letta";
/// Conversation the matrix bridge routes messages into. Pinned to a specific
/// conversation ID (not "default") so a tainted/default conversation can be
/// swapped without touching the bridge.
pub const CONVERSATION_ID = "local-conv-487";
/// Runs `letta -n <agent> -p <prompt>`, captures everything it writes to
/// stdout, and returns it as an owned slice (caller frees it).
pub fn infer(
io: std.Io,
allocator: std.mem.Allocator,
agent: []const u8,
prompt: []const u8,
) ![]u8 {
_ = agent; // unused: pinned conversation implies the agent on the CLI side
var child = try std.process.spawn(io, .{
// Pinned conversation: --agent must be omitted for non-default conversations.
.argv = &.{ LETTA_PATH, "-p", prompt, "--conversation", CONVERSATION_ID, "--toolset", "default" },
.stdout = .pipe,
.stderr = .inherit, // letta's error output goes straight to your terminal
});
// Non-null because we asked for a pipe.
const stdout_file = child.stdout.?;
var output: std.ArrayList(u8) = .empty;
errdefer output.deinit(allocator);
var buf: [16 * 1024]u8 = undefined;
while (true) {
const n = stdout_file.readStreaming(io, &.{&buf}) catch |err| switch (err) {
// EOF in the new Io API: letta closed its stdout. That's success.
error.EndOfStream => break,
else => |e| return e,
};
if (n == 0) break; // defensive: never spin if a 0 ever comes back
try output.appendSlice(allocator, buf[0..n]);
}
const term = try child.wait(io);
const exited_cleanly = switch (term) {
.exited => |code| code == 0,
else => false,
};
if (!exited_cleanly) return error.LettaFailed;
return output.toOwnedSlice(allocator);
}
// ------------------------------------------------------------- streaming ---
pub const Event = union(enum) {
step: []const u8, // tool call / return, pre-formatted line
chunk: []const u8, // streamed reply text
};
/// Like infer(), but spawns letta with --output-format stream-json and calls
/// `onEvent` (on THIS thread — caller marshals to Qt) for every step/chunk.
/// Returns the final full reply text (allocated, caller frees).
pub fn streamInfer(
io: std.Io,
allocator: std.mem.Allocator,
agent: []const u8,
prompt: []const u8,
ctx: ?*anyopaque,
onEvent: *const fn (ctx: ?*anyopaque, ev: Event) void,
) ![]u8 {
return streamInferConv(io, allocator, agent, "default", prompt, ctx, onEvent);
}
pub fn streamInferConv(
io: std.Io,
allocator: std.mem.Allocator,
agent: []const u8,
conversation: []const u8,
prompt: []const u8,
ctx: ?*anyopaque,
onEvent: *const fn (ctx: ?*anyopaque, ev: Event) void,
) ![]u8 {
// CLI rule: --agent only pairs with --conversation default; for named
// conversations the agent is implied, so we omit it.
const is_default = std.mem.eql(u8, conversation, "default");
var argv: [11][]const u8 = undefined;
var ai: usize = 0;
argv[ai] = LETTA_PATH; ai += 1;
if (is_default) {
argv[ai] = "--agent"; ai += 1;
argv[ai] = agent; ai += 1;
}
argv[ai] = "-p"; ai += 1;
argv[ai] = prompt; ai += 1;
argv[ai] = "--conversation"; ai += 1;
argv[ai] = conversation; ai += 1;
argv[ai] = "--toolset"; ai += 1;
argv[ai] = "default"; ai += 1;
argv[ai] = "--output-format"; ai += 1;
argv[ai] = "stream-json"; ai += 1;
var child = try std.process.spawn(io, .{
.argv = argv[0..ai],
.stdout = .pipe,
.stderr = .inherit,
});
const stdout_file = child.stdout.?;
var reply: std.ArrayList(u8) = .empty;
errdefer reply.deinit(allocator);
var raw_tail: std.ArrayList(u8) = .empty; // last bytes, debug on empty reply
defer raw_tail.deinit(allocator);
var line_buf: [16 * 1024]u8 = undefined;
var line_len: usize = 0;
var buf: [16 * 1024]u8 = undefined;
while (true) {
const n = stdout_file.readStreaming(io, &.{&buf}) catch break;
if (n == 0) break;
raw_tail.appendSlice(allocator, buf[0..n]) catch {};
if (raw_tail.items.len > 4096) {
const excess = raw_tail.items.len - 4096;
std.mem.copyForwards(u8, raw_tail.items[0..4096], raw_tail.items[excess..]);
raw_tail.shrinkRetainingCapacity(4096);
}
for (buf[0..n]) |ch| {
if (ch == '\n') {
if (line_len > 0) {
handleLine(allocator, line_buf[0..line_len], ctx, onEvent, &reply);
line_len = 0;
}
} else if (line_len < line_buf.len) {
line_buf[line_len] = ch;
line_len += 1;
}
}
}
_ = child.wait(io) catch {};
if (reply.items.len == 0 and raw_tail.items.len > 0) {
std.debug.print("inferon/letta: EMPTY reply, raw stream tail:\n{s}\n", .{raw_tail.items});
}
return reply.toOwnedSlice(allocator);
}
/// Minimal hand parse of one stream-json line (we control the producer).
fn handleLine(
allocator: std.mem.Allocator,
line: []const u8,
ctx: ?*anyopaque,
onEvent: *const fn (ctx: ?*anyopaque, ev: Event) void,
reply: *std.ArrayList(u8),
) void {
const mt = dupeStr(allocator, line, "message_type") orelse {
// The "result" line holds the full text as backup.
if (std.mem.indexOf(u8, line, "\"type\":\"result\"") != null and reply.items.len == 0) {
if (dupeStr(allocator, line, "result")) |r| {
defer allocator.free(r);
reply.appendSlice(allocator, r) catch {};
}
return;
}
// Conversation-mode (no --agent) ignores --output-format and emits
// PLAIN TEXT — treat non-JSON lines as the reply itself.
if (line.len > 0 and line[0] != '{') {
reply.appendSlice(allocator, line) catch {};
reply.append(allocator, ' ') catch {};
const dup = allocator.dupe(u8, line) catch return;
onEvent(ctx, .{ .chunk = dup });
}
return;
};
defer allocator.free(mt);
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);
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 });
} 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 });
} else if (std.mem.eql(u8, mt, "assistant_message")) {
const txt = dupeStr(allocator, line, "text") orelse return;
defer allocator.free(txt);
reply.appendSlice(allocator, txt) catch {};
onEvent(ctx, .{ .chunk = txt });
}
}
/// Extract "key":"value" (with escape handling), allocated with `allocator`.
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 ..];
var out: std.ArrayList(u8) = .empty;
errdefer out.deinit(allocator);
while (it.len > 0 and it[0] != '"') {
if (it[0] == '\\' and it.len > 1) {
const esc: u8 = switch (it[1]) {
'n' => '\n',
't' => '\t',
else => it[1],
};
out.append(allocator, esc) catch return null;
it = it[2..];
} else {
out.append(allocator, it[0]) catch return null;
it = it[1..];
}
}
return out.toOwnedSlice(allocator) catch null;
}
// -------------------------------------------------------------- listings ---
pub const Agent = struct { id: []u8, name: []u8 }; // caller frees both
fn readAllStdout(io: std.Io, alloc: std.mem.Allocator, argv: []const []const u8) ![]u8 {
var child = try std.process.spawn(io, .{
.argv = argv,
.stdin = .ignore,
.stdout = .pipe,
.stderr = .ignore,
});
const out = child.stdout.?;
var buf: std.ArrayList(u8) = .empty;
errdefer buf.deinit(alloc);
var tmp: [16 * 1024]u8 = undefined;
while (true) {
const n = out.readStreaming(io, &.{&tmp}) catch break;
if (n == 0) break;
try buf.appendSlice(alloc, tmp[0..n]);
}
_ = child.wait(io) catch {};
return buf.toOwnedSlice(alloc);
}
/// Parse `letta agents list` JSON into id/name pairs.
pub fn listAgents(io: std.Io, alloc: std.mem.Allocator) ![]Agent {
const json = try readAllStdout(io, alloc, &.{ LETTA_PATH, "agents", "list" });
defer alloc.free(json);
var out: std.ArrayList(Agent) = .empty;
errdefer out.deinit(alloc);
var i: usize = 0;
while (jsonStr(alloc, json, i, "id")) |r| {
defer alloc.free(r.val);
if (!std.mem.startsWith(u8, r.val, "agent-")) {
i = r.end;
continue;
}
const id = try alloc.dupe(u8, r.val);
errdefer alloc.free(id);
// name is the next field after id in each item — no fallback leak.
const name: []u8 = if (jsonStr(alloc, json, r.end, "name")) |nr| nr.val else try alloc.dupe(u8, "?");
try out.append(alloc, .{ .id = id, .name = name });
i = r.end;
}
return out.toOwnedSlice(alloc);
}
pub const Convo = struct { id: []u8, preview: []u8 }; // caller frees both
/// Conversations for an agent, straight from the local backend's storage
/// (the CLI has no list-conversations surface; messages list only exports
/// "default"). Reads <backend>/conversations/*/conversation.json.
var conv_dir_buf: [512]u8 = undefined;
pub var CONVERSATIONS_DIR: []const u8 = "";
pub fn initConversationsDir() void {
const home = std.c.getenv("HOME") orelse return;
CONVERSATIONS_DIR = std.fmt.bufPrint(&conv_dir_buf, "{s}/.letta/lc-local-backend/conversations", .{std.mem.span(home)}) catch return;
}
pub fn listConversations(io: std.Io, alloc: std.mem.Allocator, agent: []const u8) ![]Convo {
var out: std.ArrayList(Convo) = .empty;
errdefer out.deinit(alloc);
try out.append(alloc, .{ .id = try alloc.dupe(u8, "default"), .preview = try alloc.dupe(u8, "") });
var dir = try Io.Dir.openDirAbsolute(io, CONVERSATIONS_DIR, .{ .iterate = true });
defer dir.close(io);
var it = dir.iterate();
while (try it.next(io)) |entry| {
if (entry.kind != .directory) continue;
var buf: [4096]u8 = undefined;
const sub = try std.fmt.bufPrint(&buf, "{s}/conversation.json", .{entry.name});
const json = dir.readFileAlloc(io, sub, alloc, .limited(64 * 1024)) catch continue;
defer alloc.free(json);
const aid = jsonStr(alloc, json, 0, "agent_id") orelse continue;
defer alloc.free(aid.val);
if (!std.mem.eql(u8, aid.val, agent)) continue;
// skip archived
if (std.mem.indexOf(u8, json, "\"archived\": true") != null) continue;
const cid = jsonStr(alloc, json, 0, "id") orelse continue;
var dup = false;
for (out.items) |existing| {
if (std.mem.eql(u8, existing.id, cid.val)) dup = true;
}
if (dup) {
alloc.free(cid.val);
continue;
}
// Preview: summary if set, else the last message text.
var preview: []u8 = try alloc.dupe(u8, "");
if (jsonStr(alloc, json, 0, "summary")) |s| {
if (s.val.len > 0) {
alloc.free(preview);
preview = s.val;
} else {
alloc.free(s.val);
}
}
if (preview.len == 0) {
var mbuf: [4096]u8 = undefined;
const msub = try std.fmt.bufPrint(&mbuf, "{s}/messages.jsonl", .{entry.name});
const msgs = dir.readFileAlloc(io, msub, alloc, .limited(16 * 1024 * 1024)) catch null;
if (msgs) |m| {
defer alloc.free(m);
if (std.mem.lastIndexOfScalar(u8, m, '\n')) |nl| {
const last = m[nl + 1 ..];
if (jsonStr(alloc, last, 0, "text")) |t| {
alloc.free(preview);
preview = t.val;
}
}
}
}
try out.append(alloc, .{ .id = cid.val, .preview = preview });
}
return out.toOwnedSlice(alloc);
}
/// Whitespace-tolerant JSON string-field extractor: finds "key" (with or
/// without space after the colon), returns the unescaped value.
pub fn jsonStr(alloc: std.mem.Allocator, json: []const u8, from: usize, key: []const u8) ?struct { val: []u8, end: usize } {
var pat_buf: [64]u8 = undefined;
const pat = std.fmt.bufPrint(&pat_buf, "\"{s}\"", .{key}) catch return null;
const ks = std.mem.indexOfPos(u8, json, from, pat) orelse return null;
var i = ks + pat.len;
while (i < json.len and (json[i] == ' ' or json[i] == ':')) i += 1;
if (i >= json.len or json[i] != '"') return null;
i += 1;
var out: std.ArrayList(u8) = .empty;
while (i < json.len and json[i] != '"') {
if (json[i] == '\\' and i + 1 < json.len) {
const esc: u8 = switch (json[i + 1]) {
'n' => '\n',
't' => '\t',
else => json[i + 1],
};
out.append(alloc, esc) catch return null;
i += 2;
} else {
out.append(alloc, json[i]) catch return null;
i += 1;
}
}
if (i >= json.len) return null;
return .{ .val = out.toOwnedSlice(alloc) catch return null, .end = i + 1 };
}