tools/core/coordinators/board.coordinator.ts
tools/core/coordinators/board.coordinator.ts is a file in Coordination Surface. 176 lines of code and 31 definitions.
import { argumentValue, finish } from "../readers/invocation.reader.ts";
import { currentWaiters, releaseWaiter, writeWaiters } from "../registries/board.registry.ts";
import { existsSync, readFileSync, statSync } from "node:fs";
import { isResolved, projectRoot, slotCount, surfacePath } from "../../../config/surface.config.ts";
import {
waitAllParked,
waitSolitary,
watchAbsent,
watchChanged,
watchQuiet,
watchRemoved,
watching,
} from "../strings/board.strings.ts";
import type { Admission } from "../types/board.types.ts";
import type { Invocation } from "../types/invocation.types.ts";
import { activeSeats } from "../analyzers/board.analyzer.ts";
import { barrierState } from "../runners/board.runner.ts";
import { changesSince } from "../reporters/board.reporter.ts";
import { heldClaims } from "../registries/claim.registry.ts";
import { lastInteraction } from "../registries/snapshot.registry.ts";
import { resolve } from "node:path";
const NO_WAIT_FLAG = "--no-wait";
const SLOT_EXPIRY_MS = 3_600_000;
const POLL_MS = 1000;
interface Watch {
readonly id: string;
readonly absolute: string;
readonly target: string;
readonly baseline: number;
readonly deadline: number;
readonly timeout: number;
}
const livenessWindow = function livenessWindow(): number | null {
if (!isResolved("convention", "run_live_window_ms")) {
return null;
}
return slotCount("convention", "run_live_window_ms");
};
export const writingSeats = function writingSeats(
repoRoot: string,
seats: readonly string[],
parked: readonly string[],
now: number,
): string[] {
const window = livenessWindow();
if (window === null) {
return [...seats];
}
const held = new Set(parked);
const running = new Set(heldClaims(repoRoot).map((claim) => claim.agent));
return seats.filter((seat) => {
if (held.has(seat)) {
return true;
}
if (running.has(seat)) {
return true;
}
const stamped = lastInteraction(repoRoot, seat);
return stamped > 0 && now - stamped < window;
});
};
const rosterOf = function rosterOf(repoRoot: string, board: string, parked: readonly string[], now: number): number {
if (!existsSync(board)) {
return 0;
}
const index = resolve(repoRoot, surfacePath("agent_index"));
const seated = activeSeats(readFileSync(board, "utf8"), existsSync(index) ? readFileSync(index, "utf8") : "");
return writingSeats(repoRoot, seated, parked, now).length;
};
export const runBarrier = function runBarrier(
repoRoot: string,
board: string,
now: number,
): { message: string; code: number } {
const waiting = currentWaiters(repoRoot, now);
return barrierState(
rosterOf(
repoRoot,
board,
waiting.map((entry) => entry.agent),
now,
),
waiting.length,
);
};
export const admitWaiter = function admitWaiter(
repoRoot: string,
board: string,
request: { agent: string; target: string; timeout: number; now: number; pid: number },
): Admission {
const waiting = currentWaiters(repoRoot, request.now);
const agents = rosterOf(
repoRoot,
board,
waiting.map((entry) => entry.agent),
request.now,
);
if (agents > 0 && waiting.length >= agents - 1) {
const solitary = waiting.length === 0;
return {
code: 3,
id: "",
message: solitary ? waitSolitary(agents, waiting.length) : waitAllParked(agents, waiting.length),
};
}
const id = `${String(request.pid)}-${String(request.now)}`;
writeWaiters(repoRoot, [...waiting, { agent: request.agent, expiresAt: request.now + request.timeout, id }]);
return {
code: 0,
id,
message: watching(request.target, Math.round(request.timeout / 1000), waiting.length + 1, agents),
};
};
const modifiedAt = function modifiedAt(path: string): number | null {
if (!existsSync(path)) {
return null;
}
return statSync(path).mtimeMs;
};
const sleep = async function sleep(ms: number): Promise<void> {
await new Promise<void>((done) => {
setTimeout(done, ms);
});
};
const deliverSince = function deliverSince(absolute: string, named: string): void {
process.stdout.write(changesSince(projectRoot(), absolute, argumentValue("--agent") ?? "", named));
};
const release = function release(id: string): void {
releaseWaiter(projectRoot(), id, Date.now());
};
const poll = async function poll(watch: Watch): Promise<never> {
if (Date.now() >= watch.deadline) {
release(watch.id);
process.stdout.write(watchQuiet(watch.target));
deliverSince(watch.absolute, watch.target);
process.exit(1);
}
await sleep(POLL_MS);
const current = modifiedAt(watch.absolute);
if (current === null) {
release(watch.id);
finish(watchRemoved(watch.target), 2);
}
if (current !== watch.baseline) {
release(watch.id);
const waited = Math.round((watch.timeout - (watch.deadline - Date.now())) / 1000);
process.stdout.write(watchChanged(watch.target, waited));
deliverSince(watch.absolute, watch.target);
process.exit(0);
}
return poll(watch);
};
export const awaitChange = async function awaitChange({ absolute, caller, target }: Invocation): Promise<never> {
if (process.argv.includes(NO_WAIT_FLAG)) {
deliverSince(absolute, target);
process.exit(0);
}
const baseline = modifiedAt(absolute);
if (baseline === null) {
finish(watchAbsent(target), 2);
}
const now = Date.now();
const timeout = SLOT_EXPIRY_MS;
const admission = admitWaiter(projectRoot(), resolve(projectRoot(), surfacePath("board")), {
agent: caller,
now,
pid: process.pid,
target,
timeout,
});
process.stdout.write(admission.message);
if (admission.code !== 0) {
process.exit(admission.code);
}
return poll({ absolute, baseline, deadline: now + timeout, id: admission.id, target, timeout });
};