harness: proactive send mode + persistent backlog seen-state
- 'matrix_harness send "text"' — one-shot unprompted messaging to the pinned room (agent announcements/alerts) - backlog no longer re-forwards read messages: last homeserver origin_server_ts persisted to /tmp/matrix-harness.state, advanced on every processed event; events at/below it are skipped - 30-min staleness window retained as outer bound
This commit is contained in:
1 parent
9f1a839a22
commit
dcf6d8df75
1 file changed
+67
@@ -20,9 +20,68 @@ fn env_or(name: [:0]const u8, default: []const u8) []const u8 {
|
|||||||
return default;
|
return default;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// One-shot proactive send: `matrix_harness send "message body"`.
|
||||||
|
/// Lets the agent reach Pierre unprompted (announcements, alerts).
|
||||||
|
fn runSend(alloc: std.mem.Allocator, io: std.Io, text: []const u8) !void {
|
||||||
|
letta.initConversationsDir();
|
||||||
|
const user = env_or("MATRIX_USER", "@chaos-prime:chaosmith.systems");
|
||||||
|
const password = env_or("MATRIX_PASSWORD", "");
|
||||||
|
const room = env_or("MATRIX_ROOM", "");
|
||||||
|
if (password.len == 0 or room.len == 0) {
|
||||||
|
std.debug.print("matrix_harness send: MATRIX_PASSWORD and MATRIX_ROOM required\n", .{});
|
||||||
|
return error.NoCredentials;
|
||||||
|
}
|
||||||
|
const session = try matrix.login(alloc, io, user, password);
|
||||||
|
defer alloc.free(session.token);
|
||||||
|
defer alloc.free(session.user_id);
|
||||||
|
const ev_id = try matrix.sendText(alloc, io, session.token, room, text);
|
||||||
|
std.debug.print("matrix_harness: sent {s}\n", .{ev_id});
|
||||||
|
alloc.free(ev_id);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Persistent seen-state: last homeserver event timestamp we already
|
||||||
|
/// forwarded. Survives restarts so backlog never re-sends read messages.
|
||||||
|
const STATE_FILE = "/tmp/matrix-harness.state";
|
||||||
|
|
||||||
|
fn loadSeenTs(io: std.Io) i64 {
|
||||||
|
var f = std.Io.Dir.openFileAbsolute(io, STATE_FILE, .{}) catch return 0;
|
||||||
|
defer f.close(io);
|
||||||
|
var buf: [64]u8 = undefined;
|
||||||
|
var rbuf: [64]u8 = undefined;
|
||||||
|
var r = f.reader(io, &rbuf);
|
||||||
|
const n = r.interface.readSliceShort(&buf) catch return 0;
|
||||||
|
return std.fmt.parseInt(i64, std.mem.trim(u8, buf[0..n], " \n"), 10) catch 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
fn saveSeenTs(io: std.Io, ts: i64) void {
|
||||||
|
const f = std.Io.Dir.createFileAbsolute(io, STATE_FILE, .{ .truncate = true }) catch return;
|
||||||
|
var buf: [24]u8 = undefined;
|
||||||
|
const body = std.fmt.bufPrint(&buf, "{d}\n", .{ts}) catch return;
|
||||||
|
var wbuf: [64]u8 = undefined;
|
||||||
|
var w = f.writer(io, &wbuf);
|
||||||
|
w.interface.writeAll(body) catch {};
|
||||||
|
w.interface.flush() catch {};
|
||||||
|
f.close(io);
|
||||||
|
}
|
||||||
|
|
||||||
pub fn main(init: std.process.Init) !void {
|
pub fn main(init: std.process.Init) !void {
|
||||||
const alloc = init.gpa;
|
const alloc = init.gpa;
|
||||||
const io = init.io;
|
const io = init.io;
|
||||||
|
|
||||||
|
// Proactive-send mode: `matrix_harness send "text"`.
|
||||||
|
{
|
||||||
|
var args_it = std.process.Args.Iterator.init(init.minimal.args);
|
||||||
|
_ = args_it.next(); // program name
|
||||||
|
const first_arg = args_it.next();
|
||||||
|
if (first_arg != null and std.mem.eql(u8, first_arg.?, "send")) {
|
||||||
|
if (args_it.next()) |msg| {
|
||||||
|
return runSend(alloc, io, msg);
|
||||||
|
}
|
||||||
|
std.debug.print("usage: matrix_harness send \"message\"\n", .{});
|
||||||
|
return error.MissingMessage;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
letta.initConversationsDir();
|
letta.initConversationsDir();
|
||||||
|
|
||||||
const user = env_or("MATRIX_USER", "@chaos-prime:chaosmith.systems");
|
const user = env_or("MATRIX_USER", "@chaos-prime:chaosmith.systems");
|
||||||
@@ -65,14 +124,18 @@ pub fn main(init: std.process.Init) !void {
|
|||||||
const first = try matrix.sync(alloc, io, session.token, "", 0);
|
const first = try matrix.sync(alloc, io, session.token, "", 0);
|
||||||
{
|
{
|
||||||
const now_ms: i64 = @intCast(@divTrunc(std.Io.Clock.now(.real, io).nanoseconds, std.time.ns_per_ms));
|
const now_ms: i64 = @intCast(@divTrunc(std.Io.Clock.now(.real, io).nanoseconds, std.time.ns_per_ms));
|
||||||
|
const seen_ts = loadSeenTs(io);
|
||||||
|
var new_max = seen_ts;
|
||||||
var backlog: std.ArrayListUnmanaged(u8) = .empty;
|
var backlog: std.ArrayListUnmanaged(u8) = .empty;
|
||||||
defer backlog.deinit(alloc);
|
defer backlog.deinit(alloc);
|
||||||
for (first.events) |ev| {
|
for (first.events) |ev| {
|
||||||
defer alloc.free(ev.sender);
|
defer alloc.free(ev.sender);
|
||||||
defer alloc.free(ev.body);
|
defer alloc.free(ev.body);
|
||||||
|
if (ev.ts > new_max) new_max = ev.ts;
|
||||||
if (!std.mem.eql(u8, ev.sender, peer)) continue;
|
if (!std.mem.eql(u8, ev.sender, peer)) continue;
|
||||||
if (ev.body.len == 0) continue;
|
if (ev.body.len == 0) continue;
|
||||||
if (ev.ts > 0 and now_ms - ev.ts > 30 * 60 * 1000) continue; // stale
|
if (ev.ts > 0 and now_ms - ev.ts > 30 * 60 * 1000) continue; // stale
|
||||||
|
if (ev.ts > 0 and ev.ts <= seen_ts) continue; // already seen/answered
|
||||||
try backlog.appendSlice(alloc, ev.sender);
|
try backlog.appendSlice(alloc, ev.sender);
|
||||||
try backlog.appendSlice(alloc, ": ");
|
try backlog.appendSlice(alloc, ": ");
|
||||||
try backlog.appendSlice(alloc, ev.body);
|
try backlog.appendSlice(alloc, ev.body);
|
||||||
@@ -91,6 +154,9 @@ pub fn main(init: std.process.Init) !void {
|
|||||||
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)});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
// Advance persisted seen-state even when nothing was forwarded:
|
||||||
|
// "read up to here" is authoritative from the homeserver ts.
|
||||||
|
saveSeenTs(io, new_max);
|
||||||
}
|
}
|
||||||
alloc.free(first.events);
|
alloc.free(first.events);
|
||||||
var since = first.next;
|
var since = first.next;
|
||||||
@@ -109,6 +175,7 @@ pub fn main(init: std.process.Init) !void {
|
|||||||
for (result.events) |ev| {
|
for (result.events) |ev| {
|
||||||
defer alloc.free(ev.sender);
|
defer alloc.free(ev.sender);
|
||||||
defer alloc.free(ev.body);
|
defer alloc.free(ev.body);
|
||||||
|
defer if (ev.ts > 0) saveSeenTs(io, ev.ts);
|
||||||
// Only peer messages (skip own sends and server noise).
|
// Only peer messages (skip own sends and server noise).
|
||||||
if (!std.mem.eql(u8, ev.sender, peer)) continue;
|
if (!std.mem.eql(u8, ev.sender, peer)) continue;
|
||||||
if (ev.body.len == 0) continue;
|
if (ev.body.len == 0) continue;
|
||||||
|
|||||||
Reference in new issue
Block a user