246 lines
6.9 KiB
JavaScript
246 lines
6.9 KiB
JavaScript
/**
|
|||
|
|
* Codex Office live watcher.
|
||
|
|
*
|
||
|
|
* Tails the rollout-*.jsonl session transcripts under ~/.codex/sessions/
|
||
|
|
* (organised as YYYY/MM/DD/) and translates appended lines into office
|
||
|
|
* events, streamed over Server-Sent Events at http://localhost:5181/events.
|
||
|
|
*
|
||
|
|
* Sessions are grouped into office agents by their working directory (cwd),
|
||
|
|
* so each project appears as one pawn.
|
||
|
|
*
|
||
|
|
* Zero dependencies — plain Node 18+.
|
||
|
|
*/
|
||
|
|
import { createServer } from "node:http";
|
||
|
|
import { createReadStream, promises as fs } from "node:fs";
|
||
|
|
import { homedir } from "node:os";
|
||
|
|
import path from "node:path";
|
||
|
|
|
||
|
|
const PORT = Number(process.env.PORT || 5181);
|
||
|
|
const POLL_MS = 1500;
|
||
|
|
const ACTIVE_WINDOW_DAYS = 7;
|
||
|
|
|
||
|
|
const SESSIONS_ROOT = path.join(homedir(), ".codex", "sessions");
|
||
|
|
|
||
|
|
/** slug → { slug, name, cwd, lastActivity } */
|
||
|
|
const projects = new Map();
|
||
|
|
/** absolute file path → { offset, remainder, cwd, slug } */
|
||
|
|
const tails = new Map();
|
||
|
|
|
||
|
|
const clients = new Set();
|
||
|
|
|
||
|
|
function broadcast(event) {
|
||
|
|
const data = `data: ${JSON.stringify(event)}\n\n`;
|
||
|
|
for (const res of clients) res.write(data);
|
||
|
|
}
|
||
|
|
|
||
|
|
function slugOf(cwd) {
|
||
|
|
return cwd.replace(/[^a-zA-Z0-9]+/g, "-");
|
||
|
|
}
|
||
|
|
|
||
|
|
function nameOf(cwd) {
|
||
|
|
const seg = cwd.split(/[\\/]/).filter(Boolean);
|
||
|
|
return seg[seg.length - 1] || cwd;
|
||
|
|
}
|
||
|
|
|
||
|
|
function registerProject(cwd, lastActivity) {
|
||
|
|
const slug = slugOf(cwd);
|
||
|
|
let proj = projects.get(slug);
|
||
|
|
if (!proj) {
|
||
|
|
proj = { slug, name: nameOf(cwd), cwd, lastActivity };
|
||
|
|
projects.set(slug, proj);
|
||
|
|
broadcast({ type: "project", slug, name: proj.name, lastActivity });
|
||
|
|
} else if (lastActivity > proj.lastActivity) {
|
||
|
|
proj.lastActivity = lastActivity;
|
||
|
|
}
|
||
|
|
return proj;
|
||
|
|
}
|
||
|
|
|
||
|
|
function readRange(file, start, end) {
|
||
|
|
return new Promise((resolve, reject) => {
|
||
|
|
const chunks = [];
|
||
|
|
createReadStream(file, { start, end: Math.max(start, end - 1) })
|
||
|
|
.on("data", (c) => chunks.push(c))
|
||
|
|
.on("end", () => resolve(Buffer.concat(chunks)))
|
||
|
|
.on("error", reject);
|
||
|
|
});
|
||
|
|
}
|
||
|
|
|
||
|
|
/** Rollouts put session_meta (with cwd) in the first lines — sniff the head. */
|
||
|
|
async function sniffCwd(file, size) {
|
||
|
|
try {
|
||
|
|
const buf = await readRange(file, 0, Math.min(size, 8192));
|
||
|
|
const m = buf.toString("utf8").match(/"cwd":"((?:[^"\\]|\\.)*)"/);
|
||
|
|
if (m) return JSON.parse(`"${m[1]}"`);
|
||
|
|
} catch {
|
||
|
|
/* ignore */
|
||
|
|
}
|
||
|
|
return null;
|
||
|
|
}
|
||
|
|
|
||
|
|
function handleLine(tail, line) {
|
||
|
|
let o;
|
||
|
|
try {
|
||
|
|
o = JSON.parse(line);
|
||
|
|
} catch {
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
const p = o.payload ?? {};
|
||
|
|
|
||
|
|
// cwd can appear in session_meta and change per turn_context
|
||
|
|
if ((o.type === "session_meta" || o.type === "turn_context") && typeof p.cwd === "string") {
|
||
|
|
tail.cwd = p.cwd;
|
||
|
|
tail.slug = slugOf(p.cwd);
|
||
|
|
registerProject(p.cwd, Date.now());
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
|
||
|
|
if (!tail.slug) return;
|
||
|
|
const slug = tail.slug;
|
||
|
|
const proj = projects.get(slug);
|
||
|
|
if (proj) proj.lastActivity = Date.now();
|
||
|
|
|
||
|
|
if (o.type === "event_msg") {
|
||
|
|
switch (p.type) {
|
||
|
|
case "user_message":
|
||
|
|
broadcast({ type: "prompt", slug });
|
||
|
|
break;
|
||
|
|
case "task_started":
|
||
|
|
broadcast({ type: "thinking", slug });
|
||
|
|
break;
|
||
|
|
case "agent_message":
|
||
|
|
if (p.message) broadcast({ type: "speech", slug, text: String(p.message).slice(0, 80) });
|
||
|
|
break;
|
||
|
|
case "task_complete":
|
||
|
|
if (p.last_agent_message)
|
||
|
|
broadcast({ type: "speech", slug, text: String(p.last_agent_message).slice(0, 80) });
|
||
|
|
break;
|
||
|
|
case "mcp_tool_call_begin":
|
||
|
|
case "mcp_tool_call_end": {
|
||
|
|
const inv = p.invocation ?? {};
|
||
|
|
const name = [inv.server, inv.tool].filter(Boolean).join(":") || "mcp";
|
||
|
|
broadcast({ type: "tool", slug, tool: name });
|
||
|
|
break;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
|
||
|
|
if (o.type === "response_item") {
|
||
|
|
switch (p.type) {
|
||
|
|
case "custom_tool_call":
|
||
|
|
case "function_call":
|
||
|
|
broadcast({ type: "tool", slug, tool: p.name || "tool" });
|
||
|
|
break;
|
||
|
|
case "local_shell_call":
|
||
|
|
broadcast({ type: "tool", slug, tool: "shell" });
|
||
|
|
break;
|
||
|
|
case "reasoning":
|
||
|
|
broadcast({ type: "thinking", slug });
|
||
|
|
break;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
/** Recursively collect .jsonl files under the YYYY/MM/DD tree (depth ≤ 3). */
|
||
|
|
async function collectFiles(dir, depth) {
|
||
|
|
let entries;
|
||
|
|
try {
|
||
|
|
entries = await fs.readdir(dir, { withFileTypes: true });
|
||
|
|
} catch {
|
||
|
|
return [];
|
||
|
|
}
|
||
|
|
const files = [];
|
||
|
|
for (const e of entries) {
|
||
|
|
const full = path.join(dir, e.name);
|
||
|
|
if (e.isDirectory() && depth < 3) {
|
||
|
|
files.push(...(await collectFiles(full, depth + 1)));
|
||
|
|
} else if (e.isFile() && e.name.endsWith(".jsonl")) {
|
||
|
|
files.push(full);
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return files;
|
||
|
|
}
|
||
|
|
|
||
|
|
async function pollOnce(initial) {
|
||
|
|
const files = await collectFiles(SESSIONS_ROOT, 0);
|
||
|
|
const cutoff = Date.now() - ACTIVE_WINDOW_DAYS * 86400000;
|
||
|
|
|
||
|
|
for (const full of files) {
|
||
|
|
let stat;
|
||
|
|
try {
|
||
|
|
stat = await fs.stat(full);
|
||
|
|
} catch {
|
||
|
|
continue;
|
||
|
|
}
|
||
|
|
|
||
|
|
let tail = tails.get(full);
|
||
|
|
if (!tail) {
|
||
|
|
// First sighting: skip history, only follow new appends.
|
||
|
|
tail = { offset: stat.size, remainder: "", cwd: null, slug: null };
|
||
|
|
tails.set(full, tail);
|
||
|
|
const cwd = await sniffCwd(full, stat.size);
|
||
|
|
if (cwd) {
|
||
|
|
tail.cwd = cwd;
|
||
|
|
tail.slug = slugOf(cwd);
|
||
|
|
if (stat.mtimeMs >= cutoff) {
|
||
|
|
const proj = registerProject(cwd, stat.mtimeMs);
|
||
|
|
if (initial) proj.lastActivity = stat.mtimeMs;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
continue;
|
||
|
|
}
|
||
|
|
|
||
|
|
if (stat.size < tail.offset) {
|
||
|
|
tail.offset = stat.size;
|
||
|
|
tail.remainder = "";
|
||
|
|
}
|
||
|
|
if (stat.size > tail.offset) {
|
||
|
|
const buf = await readRange(full, tail.offset, stat.size);
|
||
|
|
tail.offset = stat.size;
|
||
|
|
const chunk = tail.remainder + buf.toString("utf8");
|
||
|
|
const lines = chunk.split("\n");
|
||
|
|
tail.remainder = lines.pop() ?? "";
|
||
|
|
for (const line of lines) {
|
||
|
|
if (line.trim()) handleLine(tail, line);
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
function snapshot() {
|
||
|
|
return {
|
||
|
|
type: "snapshot",
|
||
|
|
projects: [...projects.values()]
|
||
|
|
.sort((a, b) => b.lastActivity - a.lastActivity)
|
||
|
|
.map((p) => ({ slug: p.slug, name: p.name, lastActivity: p.lastActivity })),
|
||
|
|
subs: [],
|
||
|
|
};
|
||
|
|
}
|
||
|
|
|
||
|
|
const server = createServer((req, res) => {
|
||
|
|
if (req.url === "/events") {
|
||
|
|
res.writeHead(200, {
|
||
|
|
"Content-Type": "text/event-stream",
|
||
|
|
"Cache-Control": "no-cache",
|
||
|
|
Connection: "keep-alive",
|
||
|
|
"Access-Control-Allow-Origin": "*",
|
||
|
|
});
|
||
|
|
res.write(`data: ${JSON.stringify(snapshot())}\n\n`);
|
||
|
|
clients.add(res);
|
||
|
|
req.on("close", () => clients.delete(res));
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
res.writeHead(404, { "Access-Control-Allow-Origin": "*" });
|
||
|
|
res.end("not found");
|
||
|
|
});
|
||
|
|
|
||
|
|
await pollOnce(true);
|
||
|
|
// Transcripts are sensitive — never expose beyond this machine.
|
||
|
|
server.listen(PORT, "127.0.0.1", () => {
|
||
|
|
console.log(">_ Codex Office watcher");
|
||
|
|
console.log(` watching ${SESSIONS_ROOT}`);
|
||
|
|
console.log(` projects ${projects.size} active in last ${ACTIVE_WINDOW_DAYS} days`);
|
||
|
|
console.log(` SSE http://localhost:${PORT}/events`);
|
||
|
|
});
|
||
|
|
setInterval(() => pollOnce(false).catch(() => {}), POLL_MS);
|