Files
stack/packages/control-board/src/scan.mjs
T
jason.woltjeandClaude Opus 5.5 af4203ca92 feat(board): session attention, Discord rows, task attribution and relaunch activity (rows 18, 22, #1511, #1512)
One cumulative control-board, webui and seat state. The four rows edit the
same files (scan.mjs, page.html, README.md, app.js), so they land together,
each on its own receipt:

- Row 18, Discord connector rows on the board (#1509): R3 approved by
  Darkwing and Dewey, Gitea comment 26257, manifest 254403b8. Jason
  accepted the visual test.
- Row 22, board attention status (#1503): Filbert approved R1, comment
  26248, manifest e40b58ec; restart receipt 26249.
- #1511, task attribution (row 6 code phase): R2 approved by Filbert and
  Dewey, manifest d4c96395. docs/TOOLS.md carries the approved --by usage
  line (tools-usage.patch 86bcba3c).
- #1512, relaunch activity (row 6 pilot): R1 approved by Darkwing and
  Dewey, candidate manifest 47769fad. All seven source files match it.

Row 16, internal development bootstrap (#1510): the seven files outside
shared records match Filbert's R1 pins, receipt 26204 (agents/researcher/*,
scripts/test-darkwing-launch.mjs, the bootstrap plan).

packages/webui/src/public/app.js is committed at its #1512 R1 pin ce7d79a4.
The working copy holds Dewey's unreviewed return-flow candidate on top of
that, and it stays uncommitted.

Also: the four row briefs and Darkwing's evidence records under
agents/darkwing/work, including the 2026-09-26 tree manifest and the #1512
re-run against 21e3e908. Serial acceptance command: 397/397, three runs.
The failures that only show when tests run concurrently are in #1509 engine
tests, and they reproduce on clean HEAD.

Suites on the exact staged tree: config 24, task 90, foundation 43,
conductor 17, release 14, auth 15, discord 63; package union 397/397
(serial); test-darkwing-launch 5/5.

Shared records (BUILD-LOG, QUEUE, CURRENT, DEFERRED, SESSIONS, AGENTS.md,
agents/README.md) follow in Sage's records commit.

Co-Authored-By: Claude Opus 5.5 <[email protected]>
2026-09-26 14:54:18 -05:00

465 lines
22 KiB
JavaScript

// Control board step 1: read each agent's newest pi session log and tmux
// liveness, and write one small JSON status file per agent.
//
// States (plain words):
// working - the agent is in the middle of a turn (thinking or running tools)
// waiting - the agent explicitly requests human input in its completed reply
// error - the agent's last turn ended in an error, was aborted, or was cut off; look at it
// offline - no tmux session for this agent, or its session no longer runs pi
// idle - the agent is available, with no explicit input request
// unknown - liveness could not be checked (tmux missing or unresponsive); not a guess
//
// Board files are derived and rewritable. They are not run records. The one
// exception is <boardDir>/seen.json, which holds Jason's "seen" marks and is
// only changed when he clicks; a scan reads it and never rewrites it.
import { existsSync, readFileSync, readdirSync, statSync, mkdirSync, writeFileSync, renameSync } from "node:fs";
import { join, basename, isAbsolute, resolve } from "node:path";
import { homedir } from "node:os";
import { spawnSync } from "node:child_process";
import { readRegistration, samePath, SeatError, LAYOUTS } from "../../seat/src/seat.mjs";
import { discoverDiscordAgents, inspectDiscord, newestDiscordSession } from './discord.mjs';
export const STATES = Object.freeze(["working", "waiting", "error", "offline", "idle", "unknown"]);
const TEXT_LIMIT = 240;
export class ConfigError extends Error {}
export function defaultConfigPath() {
return join(homedir(), ".config", "mosaic-dev", "config.json");
}
// Fail closed: the config must exist, parse, and name an absolute dataRoot.
export function loadConfig(path = defaultConfigPath()) {
if (!existsSync(path)) throw new ConfigError(`config not found: ${path}`);
let raw;
try {
raw = JSON.parse(readFileSync(path, "utf8"));
} catch (err) {
throw new ConfigError(`config is not valid JSON: ${path} (${err.message})`);
}
if (!raw || typeof raw !== "object" || Array.isArray(raw)) throw new ConfigError(`config is not an object: ${path}`);
if (typeof raw.dataRoot !== "string" || !isAbsolute(raw.dataRoot)) {
throw new ConfigError(`config.dataRoot must be an absolute path: ${path}`);
}
return { dataRoot: raw.dataRoot };
}
// Newest *.jsonl under a directory tree, by mtime. Returns null when none.
export function findNewestSession(dir) {
if (!existsSync(dir)) return null;
let best = null;
const walk = (d) => {
for (const name of readdirSync(d)) {
const p = join(d, name);
const st = statSync(p);
if (st.isDirectory()) walk(p);
else if (name.endsWith(".jsonl") && (!best || st.mtimeMs > best.mtimeMs)) best = { path: p, mtimeMs: st.mtimeMs };
}
};
walk(dir);
return best ? best.path : null;
}
function collapse(text) {
const one = text.replace(/\s+/g, " ").trim();
return one.length > TEXT_LIMIT ? one.slice(0, TEXT_LIMIT - 1) + "…" : one;
}
// Read a pi session log. A partially written final line is skipped and counted,
// not treated as fatal, because pi may be appending while we read.
export function readSession(file) {
const lines = readFileSync(file, "utf8").split("\n");
let sessionId = null, cwd = null, lastTimestamp = null, lastMessage = null, lastAssistantText = null, lastError = null;
let firstUserText = null;
let model = null, provider = null;
let skippedLines = 0;
for (const line of lines) {
if (!line.trim()) continue;
let entry;
try {
entry = JSON.parse(line);
} catch {
skippedLines += 1;
continue;
}
if (entry.timestamp) lastTimestamp = entry.timestamp;
if (entry.type === "session") {
sessionId = entry.id ?? sessionId;
cwd = entry.cwd ?? cwd;
} else if (entry.type === "model_change") {
// pi logs a model switch (/model, or the launch default) as its own
// entry: { provider, modelId }. The latest one is the model in use.
if (typeof entry.modelId === "string" && entry.modelId) model = entry.modelId;
if (typeof entry.provider === "string" && entry.provider) provider = entry.provider;
} else if (entry.type === "message" && entry.message) {
lastMessage = entry.message;
if (entry.message.role === "assistant") {
// Each assistant turn also names the model that produced it.
if (typeof entry.message.model === "string" && entry.message.model) model = entry.message.model;
if (typeof entry.message.provider === "string" && entry.message.provider) provider = entry.message.provider;
}
if (entry.message.role === "user" && firstUserText === null) {
const text = userText(entry.message.content);
if (text.trim()) firstUserText = collapse(text);
}
if (entry.message.role === "assistant" && Array.isArray(entry.message.content)) {
const text = entry.message.content.filter((c) => c && c.type === "text" && typeof c.text === "string").map((c) => c.text).join("\n");
if (text.trim()) lastAssistantText = collapse(text);
lastError = entry.message.stopReason === "error" && typeof entry.message.errorMessage === "string" ? collapse(entry.message.errorMessage) : null;
}
}
}
return { file, sessionId, cwd, lastTimestamp, lastMessage, lastAssistantText, lastError, firstUserText, model, provider, skippedLines };
}
// A user message's text: pi writes either a plain string or a list of blocks.
function userText(content) {
if (typeof content === "string") return content;
if (!Array.isArray(content)) return "";
return content.filter((c) => c && c.type === "text" && typeof c.text === "string").map((c) => c.text).join("\n");
}
// Nearest ancestor (including dir itself) that holds a .git entry. A .git
// file counts too, because git worktrees use one. Null when there is none.
export function findRepoRoot(dir) {
if (typeof dir !== "string" || !isAbsolute(dir)) return null;
let cur = resolve(dir);
for (;;) {
if (existsSync(join(cur, ".git"))) return cur;
const parent = resolve(cur, "..");
if (parent === cur) return null;
cur = parent;
}
}
// True when an assistant message carries a tool call in its content.
export function hasToolCall(message) {
return Array.isArray(message?.content) && message.content.some((c) => c && c.type === "toolCall");
}
// Explicit display convention, not a language classifier or control permission.
// Only a plain first nonblank text line can request human input. Quoted/code
// examples and thinking blocks must not create attention events.
function requestsHumanInput(message) {
if (!Array.isArray(message.content)) return false;
const text = message.content.filter((c) => c?.type === "text" && typeof c.text === "string")
.map((c) => c.text).join("\n");
const first = text.split(/\r?\n/).find((line) => line.trim()) ?? "";
const prefix = "Input needed: ";
return first.startsWith(prefix) && first.slice(prefix.length).trim().length > 0;
}
// A missing liveness check stays unknown. Offline/errors/tool activity take
// precedence over attention text. Completion alone means idle, not waiting.
export function deriveState({ alive, session }) {
if (alive === false) return "offline";
if (alive !== true) return "unknown";
if (!session || !session.lastMessage) return "idle";
const m = session.lastMessage;
if (m.role === "assistant") {
if (m.stopReason === "error" || m.stopReason === "aborted" || m.stopReason === "length") return "error";
if (hasToolCall(m)) return "working";
if (m.stopReason === "stop") return requestsHumanInput(m) ? "waiting" : "idle";
return "working";
}
return "working";
}
// Programs that count as a live pi agent in a tmux pane. A tmux session that
// still exists but only runs a shell (or another harness) is not alive: its
// pi session log is history, not status.
export const PI_COMMANDS = Object.freeze(["pi"]);
export function panesRunPi(listPanesOutput) {
return parsePanes(listPanesOutput).some((p) => PI_COMMANDS.includes(p.command));
}
// One line per pane: "<command>\t<path>" (the path column is optional).
export function parsePanes(listPanesOutput) {
return String(listPanesOutput)
.split("\n")
.filter((l) => l.trim())
.map((l) => {
const [command = "", path = ""] = l.split("\t");
return { command: command.trim(), path: path.trim() || null };
});
}
// Ask tmux about one session. Returns { alive, workspace }:
// alive true when a pane runs pi; false when no such session or no pane
// runs pi; null when tmux could not be run at all (reported as
// "unknown", never assumed alive).
// workspace the current path of the first pane running pi, else null.
// `exec` is injectable for tests.
export function tmuxInspect({ socket, session }, { exec = spawnSync } = {}) {
const args = [];
if (socket) args.push("-L", socket);
args.push("list-panes", "-s", "-t", `=${session}`, "-F", "#{pane_current_command}\t#{pane_current_path}");
const r = exec("tmux", args, { encoding: "utf8", timeout: 5000 });
if (r.error) return { alive: null, workspace: null };
if (r.status !== 0) return { alive: false, workspace: null };
const pane = parsePanes(r.stdout ?? "").find((p) => PI_COMMANDS.includes(p.command));
return { alive: Boolean(pane), workspace: pane?.path ?? null };
}
// Liveness only, for callers that do not need the pane path.
export function tmuxIsAlive(tmux, opts) {
return tmuxInspect(tmux, opts).alive;
}
// "Seen" marks: { "<project>/<agent>": "<lastActivity ISO>" }. A mark only
// applies while the agent's newest message still has that timestamp; anything
// the agent writes afterwards clears it automatically.
export function seenKey(rec) {
return `${rec.project}/${rec.agent}`;
}
export function seenPath(boardDir) {
return join(boardDir, "seen.json");
}
// Fail closed: a present but unreadable seen.json refuses the scan rather than
// silently dropping every mark.
export function loadSeen(boardDir) {
const path = seenPath(boardDir);
if (!existsSync(path)) return {};
let parsed;
try {
parsed = JSON.parse(readFileSync(path, "utf8"));
} catch (err) {
throw new ConfigError(`seen marks file is not valid JSON: ${path} (${err.message})`);
}
if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) throw new ConfigError(`seen marks file must be a JSON object: ${path}`);
for (const [k, v] of Object.entries(parsed)) {
if (typeof v !== "string") throw new ConfigError(`seen marks file has a non-string value for ${JSON.stringify(k)}: ${path}`);
}
return parsed;
}
export function saveSeen(boardDir, marks) {
mkdirSync(boardDir, { recursive: true, mode: 0o700 });
writeAtomic(seenPath(boardDir), marks);
}
// Set or clear one mark. Returns the updated map.
export function markSeen(boardDir, { project, agent, lastActivity, seen = true }) {
for (const [name, v] of Object.entries({ project, agent, lastActivity })) {
if (typeof v !== "string" || v.length === 0 || v.length > 512) throw new ConfigError(`${name} must be a non-empty string`);
}
if (project.includes("/")) throw new ConfigError("project must not contain '/'");
if (typeof seen !== "boolean") throw new ConfigError("seen must be true or false");
const marks = loadSeen(boardDir);
const key = seenKey({ project, agent });
if (seen) marks[key] = lastActivity;
else delete marks[key];
saveSeen(boardDir, marks);
return marks;
}
// `isAlive` may return a bare liveness value (true/false/null) or the richer
// { alive, workspace } shape from tmuxInspect. Both are accepted.
function liveness(result) {
if (result && typeof result === "object") return { alive: result.alive ?? null, workspace: result.workspace ?? null };
return { alive: result ?? null, workspace: null };
}
// Registrations written by `mosaic launch <seat>` (packages/seat): one
// record per seat under <dataRoot>/seats/<layout>/<seat>/registration.json. Returns
// the readable records plus one error line per unreadable one; a bad record
// must not take the whole board down, but it is not silently dropped either.
export function loadRegistrations(seatsDir) {
const registrations = [];
const errors = [];
if (!seatsDir || !existsSync(seatsDir)) return { registrations, errors };
for (const layout of LAYOUTS) {
const layoutDir = join(seatsDir, layout);
if (!existsSync(layoutDir) || !statSync(layoutDir).isDirectory()) continue;
for (const seat of readdirSync(layoutDir).sort()) {
if (!statSync(join(layoutDir, seat)).isDirectory()) continue;
try {
const rec = readRegistration(seatsDir, seat, layout);
if (rec) registrations.push(rec);
} catch (err) {
if (!(err instanceof SeatError)) throw err;
errors.push(`seat ${layout}/${seat}: ${err.message}`);
}
}
}
return { registrations, errors };
}
// A registration belongs to a row when it names the same sessions directory.
// Seat names alone are not enough: the repo and fleet layouts both have a
// "darkwing", and they are different seats.
export function matchRegistration(spec, registrations) {
return registrations.find((r) => r.sessionsDir && samePath(r.sessionsDir, spec.sessionsDir)) ?? null;
}
// True when a pid is running (a signal-0 probe; EPERM still means running),
// false when it is gone, null when there is no pid to check.
export function pidAlive(pid) {
if (!(Number.isInteger(pid) && pid > 0)) return null;
try {
process.kill(pid, 0);
return true;
} catch (err) {
return err.code === "EPERM";
}
}
// One agent -> one status record.
//
// Three fields answer "what is this seat doing, and where" (Gate A ask,
// 2026-09-12). Each is derived, never guessed; null means "unknown". A seat
// started through `mosaic launch` has a registration, and a registered task,
// project or workspace wins over the derived value; the *Source field says
// which one the row shows.
// task registration.task, else the session's first user message
// (the log has no task envelope entry).
// workspace registration.workspace, else the live pane path from tmux,
// else the session log's cwd.
// activeProject registration.project, else the basename of the nearest git
// repo root above the workspace.
// A registration is written before the launch script's own checks run, so a
// refused launch leaves a record whose pid is gone. Such a record is stale:
// it is still reported under `registered` (with alive false) but the derived
// values win, because the record describes a launch that is not running.
export function scanAgent(spec, { isAlive = tmuxInspect, now = () => new Date(), seen = {}, registration = null, isPidAlive = pidAlive } = {}) {
const connector = spec.connector ? inspectDiscord(spec.connector) : null;
const live = connector ? { alive: connector.alive, workspace: null } : liveness(isAlive(spec.tmux));
const alive = live.alive;
const file = connector ? newestDiscordSession(spec.sessionsDir) : findNewestSession(spec.sessionsDir);
const session = file ? readSession(file) : null;
const state = deriveState({ alive, session });
const scannedAt = now();
const lastActivity = session?.lastTimestamp ?? null;
const ageSeconds = lastActivity ? Math.max(0, Math.round((scannedAt.getTime() - Date.parse(lastActivity)) / 1000)) : null;
const needsYou = state === "waiting" || state === "error";
const isSeen = needsYou && lastActivity !== null && seen[seenKey(spec)] === lastActivity;
const cwd = session?.cwd ?? null;
const record = !connector && registration && typeof registration === "object" ? registration : null;
const registeredAlive = record ? isPidAlive(record.pid) : null;
// Presentation only: retain history and attention/Seen semantics. Unknown
// liveness, unmatched records and unknown timestamps cannot assert relaunch.
const launchTime = typeof record?.startedAt === "string" ? Date.parse(record.startedAt) : NaN;
const activityTime = typeof lastActivity === "string" ? Date.parse(lastActivity) : NaN;
const relaunchedAt = alive === true && registeredAlive === true &&
samePath(record?.sessionsDir, spec.sessionsDir) && Number.isFinite(launchTime) &&
Number.isFinite(activityTime) && launchTime > activityTime
? new Date(launchTime).toISOString() : null;
const reg = record && registeredAlive !== false ? record : null;
const derivedWorkspace = live.workspace ?? cwd;
const workspace = reg?.workspace ?? derivedWorkspace;
const workspaceSource = reg?.workspace ? "registration" : live.workspace ? "tmux-pane" : cwd ? "session-cwd" : null;
const repoRoot = derivedWorkspace ? findRepoRoot(derivedWorkspace) : null;
const derivedProject = repoRoot ? basename(repoRoot) : null;
const activeProject = reg?.project ?? derivedProject;
const activeProjectSource = reg?.project ? "registration" : derivedProject ? "workspace-git-root" : null;
const firstUserText = session?.firstUserText ?? null;
// Discord user text starts with a routing envelope containing private IDs.
// Connector task metadata is fixed, never inferred from that transcript.
const task = connector ? "Discord connector" : reg?.task ? reg.task : firstUserText;
const taskSource = connector ? "connector" : reg?.task ? "registration" : firstUserText ? "first-user-message" : null;
// Who set the task (#1511), only when the task shown *is* the registered
// one: a record written before the field existed reads "unknown"; a
// transcript-derived task, a stale registration and a connector's fixed
// task carry no attribution at all (null), so a record can never lend its
// setter to a task it did not set. It is the caller's claim, bounded by
// the seat package on read, and it grants nothing.
const taskSetBy = taskSource === "registration" ? (reg.taskSetBy ?? "unknown") : null;
return {
agent: spec.agent,
project: spec.project,
...(connector ? { connector } : {}),
state,
waitingOnYou: needsYou && !isSeen,
seen: isSeen,
alive,
tmux: spec.tmux,
sessionFile: file,
sessionId: session?.sessionId ?? null,
cwd,
// The model the seat is running, from the log's latest model_change entry
// or assistant turn, whichever came last. null when the log has neither
// (a seat with no log, or one that has not answered yet).
model: session?.model ?? null,
provider: session?.provider ?? null,
task,
taskSource,
taskSetBy,
workspace,
workspaceSource,
activeProject,
activeProjectSource,
registered: record
? { startedAt: record.startedAt, updatedAt: record.updatedAt, harness: record.harness, pid: record.pid, alive: registeredAlive, tmux: record.tmux, layout: record.layout, launchScript: record.launchScript }
: null,
lastActivity,
relaunchedAt,
ageSeconds,
lastAssistantText: session?.lastAssistantText ?? null,
lastError: session?.lastError ?? null,
skippedLines: session?.skippedLines ?? 0,
scannedAt: scannedAt.toISOString(),
};
}
// Repo agents: <repo>/.pi/state/<agent>/sessions, tmux default socket, session = agent.
export function discoverRepoAgents(repoRoot) {
const stateDir = join(repoRoot, ".pi", "state");
if (!existsSync(stateDir)) return [];
const project = basename(resolve(repoRoot));
return readdirSync(stateDir)
.filter((n) => existsSync(join(stateDir, n, "sessions")))
.sort()
.map((agent) => ({ agent, project, sessionsDir: join(stateDir, agent, "sessions"), tmux: { socket: null, session: agent } }));
}
// Fleet agents: <fleet>/<agent>/.pi/agent/sessions, tmux socket mosaic-fleet, session = agent.
export function discoverFleetAgents(fleetRoot, { project = "fleet", socket = "mosaic-fleet" } = {}) {
if (!existsSync(fleetRoot)) return [];
return readdirSync(fleetRoot)
.filter((n) => existsSync(join(fleetRoot, n, ".pi", "agent", "sessions")))
.sort()
.map((agent) => ({ agent, project, sessionsDir: join(fleetRoot, agent, ".pi", "agent", "sessions"), tmux: { socket, session: agent } }));
}
function writeAtomic(path, data) {
const tmp = `${path}.tmp-${process.pid}`;
writeFileSync(tmp, JSON.stringify(data, null, 2) + "\n", { mode: 0o600 });
renameSync(tmp, path);
}
// Scan every spec and write <boardDir>/sessions/<project>/<agent>.json plus index.json.
// seatsDir (optional): where `mosaic launch` registrations live; read only.
export function scan(specs, { boardDir, isAlive, now, seatsDir = null, isPidAlive, discordDataRoot = null } = {}) {
if (!boardDir || !isAbsolute(boardDir)) throw new ConfigError("boardDir must be an absolute path");
if (seatsDir !== null && (typeof seatsDir !== "string" || !isAbsolute(seatsDir))) throw new ConfigError("seatsDir must be an absolute path or null");
const seen = loadSeen(boardDir);
const { registrations, errors: registrationErrors } = loadRegistrations(seatsDir);
const discord = discordDataRoot ? discoverDiscordAgents(discordDataRoot) : { specs: [], errors: [] };
const records = [...specs, ...discord.specs].map((spec) => scanAgent(spec, { isAlive, now, seen, registration: matchRegistration(spec, registrations), isPidAlive }));
for (const rec of records) {
const dir = join(boardDir, "sessions", rec.project);
mkdirSync(dir, { recursive: true, mode: 0o700 });
writeAtomic(join(dir, `${rec.agent}.json`), rec);
}
const generatedAt = (now ? now() : new Date()).toISOString();
const index = {
generatedAt,
counts: Object.fromEntries(STATES.map((s) => [s, records.filter((r) => r.state === s).length])),
waitingOnYou: records.filter((r) => r.waitingOnYou).map(seenKey),
seen: records.filter((r) => r.seen).map(seenKey),
registered: records.filter((r) => r.registered).map(seenKey),
registrationStale: records.filter((r) => r.registered && r.registered.alive === false).map(seenKey),
registrationErrors,
discoveryErrors: discord.errors,
sessions: records,
};
mkdirSync(boardDir, { recursive: true, mode: 0o700 });
writeAtomic(join(boardDir, "index.json"), index);
return index;
}