#!/usr/bin/env node // Usage: // mosaic-discord check [--config PATH] [--repo PATH] // mosaic-discord run [--config PATH] [--repo PATH] [--supervised] // mosaic-discord stop [--config PATH] // mosaic-discord unlock [--config PATH] // mosaic-discord recover [--config PATH] // mosaic-discord reload [--config PATH] // // names /discord/.json. The repository wrapper // is scripts/discord.sh. // // check: binding, token file, context files and pi are validated; then the // bot identity, guild and every listed channel are read over REST; then one // gateway connection is made and closed after READY. Nothing is sent to a // channel. A close code 4014 means the message-content intent is not // granted in the developer portal. // run: refuses when STOP exists or the outbox cannot be reconciled; otherwise // starts the engine and the gateway and serves turns until SIGTERM, SIGINT // or `stop`. // stop: writes STOP and sends SIGTERM to the owner in run.lock, only when that // process is alive and its start time and boot id match the record. // unlock: writes STOP, then removes a run.lock whose owner is gone (crash, // reboot, interrupted start). Refuses while the owner is live (use stop), alive // with unverifiable identity, or recorded in a file it cannot read. `run` never reclaims on its own, and a // claim that finds STOP after publishing releases itself, so unlock cannot // race a start. Remove STOP to run again. // recover: the supervised pre-start. Refuses while STOP is present or the // binding is held by a live or unverifiable process; clears a lock whose // owner is gone the way unlock does, then removes the STOP it wrote for // that, so the run that follows can claim. Never removes a STOP an operator // wrote. `run --supervised` does the same first thing itself; the service // unit (scripts/discord-service.sh) uses that form, because systemd only // honours a never-retry exit status from the main process, not from a // pre-start command. // // reload: validates the binding file, then sends SIGHUP to the live owner in // run.lock. The running process re-reads the file and applies guildName, // channels, users and limits in place; a new channel is read over REST and // must be in the bound guild. Any fixed key changed, an invalid file or a // failed lookup refuses the reload and keeps the old binding. Each attempt is // one line in reloads.jsonl. The service unit maps `systemctl --user reload` // to the same signal. // // Exit codes: 0 ok; 1 operation failed; 2 invalid data or configuration; // 3 refused by a brake (STOP present or the binding held; a supervisor must // not retry); 4 usage. import { existsSync, mkdirSync, mkdtempSync, writeFileSync, statSync, readdirSync } from "node:fs"; import { join, resolve } from "node:path"; import { createHash } from "node:crypto"; import { DiscordError } from "./errors.mjs"; import { defaultConfigPath, loadDataRoot, bindingPath, bindingDataDir, loadBinding, readToken, resolveContextFiles, resolveToolRoots, reloadDiff } from "./binding.mjs"; import { createRest } from "./rest.mjs"; import { createGateway, CONNECTOR_INTENTS } from "./gateway.mjs"; import { createEngine, buildPiArgs } from "./engine-pi.mjs"; import { TOOLS_ENV, enabledToolNames } from "./tools.mjs"; import { createSetsparkApi } from "./setspark.mjs"; import { assembleContext } from "./context.mjs"; import { createConnector } from "./connector.mjs"; import { ensureJournal, requestStop, stopRequested, readPid, stopTarget, writePid, clearPid, unlock, recover, appendReload, BRAKE_EXIT } from "./journal.mjs"; const USAGE = [ "usage: mosaic-discord check [--config PATH] [--repo PATH]", " mosaic-discord run [--config PATH] [--repo PATH] [--supervised]", " mosaic-discord stop [--config PATH]", " mosaic-discord unlock [--config PATH]", " mosaic-discord recover [--config PATH]", " mosaic-discord reload [--config PATH]", ].join("\n"); function parse(argv) { const opts = { command: null, binding: null, config: defaultConfigPath(), repo: process.cwd(), supervised: false }; for (let i = 0; i < argv.length; i++) { const a = argv[i]; if (a === "--config" || a === "--repo") { if (i + 1 >= argv.length) throw new DiscordError(`missing value for ${a}`, 4); opts[a.slice(2)] = resolve(argv[++i]); } else if (a === "--help" || a === "-h") { opts.command = "help"; } else if (a === "--supervised") { opts.supervised = true; } else if (a.startsWith("--")) throw new DiscordError(`unknown argument: ${a}\n${USAGE}`, 4); else if (opts.command === null) opts.command = a; else if (opts.binding === null) opts.binding = a; else throw new DiscordError(`unexpected argument: ${a}\n${USAGE}`, 4); } if (opts.command === "help") return opts; if (!["check", "run", "stop", "unlock", "recover", "reload"].includes(opts.command)) throw new DiscordError(USAGE, 4); if (opts.binding === null) throw new DiscordError(`${opts.command} needs a binding name\n${USAGE}`, 4); if (opts.supervised && opts.command !== "run") throw new DiscordError(`--supervised applies to run only\n${USAGE}`, 4); return opts; } function say(msg) { process.stdout.write(`${msg}\n`); } function warn(msg) { process.stderr.write(`discord: ${msg}\n`); } // Everything that can be checked without the network, shared by check and run. function prepare(opts) { const dataRoot = loadDataRoot(opts.config); const bindingFile = bindingPath(dataRoot, opts.binding); const binding = loadBinding(bindingFile); if (binding.name !== opts.binding) throw new DiscordError(`binding name ${JSON.stringify(binding.name)} does not match file name ${opts.binding}`); const contextFiles = resolveContextFiles(binding, opts.repo); const toolRoots = resolveToolRoots(binding, { dataRoot }); const pi = join(opts.repo, "node_modules", ".bin", "pi"); if (!existsSync(pi)) throw new DiscordError(`pi not found at ${pi}; run npm ci in the repository`); const journalDir = bindingDataDir(dataRoot, binding.name); const sessionDir = join(dataRoot, "sessions", `discord-${binding.name}`); return { dataRoot, bindingFile, binding, contextFiles, toolRoots, pi, journalDir, sessionDir }; } async function check(opts) { const { binding, contextFiles, toolRoots, pi, journalDir, sessionDir } = prepare(opts); const token = readToken(binding); say(`binding ${binding.name}: seat ${binding.seat}, guild ${binding.guildId} (${binding.guildName}), ${binding.channels.length} channel(s), ${binding.users.length} user(s)`); say(`engine ${binding.engine.provider}/${binding.engine.model}:${binding.engine.thinking}, limits ${JSON.stringify(binding.limits)}`); say(`context ${contextFiles.length} file(s); pi ${pi}; journal ${journalDir}; session ${sessionDir}`); say(toolRoots ? `tools: ${enabledToolNames(toolRoots).join(",")}; roots ${toolRoots.roots.map((r) => `${r.name}=${r.path}${r.write ? " (writable)" : ""}`).join(" ")}, ${toolRoots.maxCallsPerTurn} calls/message, ${toolRoots.maxFileBytes} bytes/file${toolRoots.web ? `, web via ${toolRoots.web.searxng} (${toolRoots.web.maxFetchBytes} bytes/page)` : ""}${toolRoots.setspark ? `, setspark via ${toolRoots.setspark.baseUrl} as ${toolRoots.setspark.principal}` : ""}` : "tools: none"); say(`token file mode 0600 ok; STOP ${stopRequested(journalDir) ? "PRESENT" : "absent"}`); const rest = createRest({ token, log: warn }); const me = await rest.getMe(); if (me.id !== binding.botUserId) throw new DiscordError(`token belongs to bot ${me.id}, binding says ${binding.botUserId}`); say(`rest: bot ${me.id} (${me.username}) matches binding`); const guild = await rest.getGuild(binding.guildId); say(`rest: guild ${guild.id} "${guild.name}" visible`); for (const c of binding.channels) { const ch = await rest.getChannel(c.id); if (ch.guild_id !== binding.guildId) throw new DiscordError(`channel ${c.id} is in guild ${ch.guild_id}, not ${binding.guildId}`); say(`rest: channel ${c.id} "#${ch.name}" type ${ch.type} (${c.mode}) visible`); } const gw = await rest.getGatewayBot(); say(`rest: gateway ${gw.url}, sessions remaining today ${gw.session_start_limit ? gw.session_start_limit.remaining : "?"}`); const gateway = createGateway({ url: gw.url, token, intents: CONNECTOR_INTENTS }); gateway.on("log", (m) => warn(`gateway: ${m}`)); const outcome = await new Promise((resolveOutcome) => { const timer = setTimeout(() => resolveOutcome({ error: "no READY within 30 s" }), 30000); gateway.on("ready", (r) => { clearTimeout(timer); resolveOutcome({ ready: r }); }); gateway.on("fatal", (f) => { clearTimeout(timer); resolveOutcome({ fatal: f }); }); gateway.connect(); }); gateway.close(1000, "check complete"); if (outcome.fatal) throw new DiscordError(`gateway close ${outcome.fatal.code}: ${outcome.fatal.reason}`); if (outcome.error) throw new DiscordError(`gateway: ${outcome.error}`, 1); const { ready } = outcome; if (!ready.user || ready.user.id !== binding.botUserId) throw new DiscordError(`gateway READY user ${ready.user && ready.user.id} does not match binding`); const inGuild = ready.guilds.some((g) => g.id === binding.guildId); if (!inGuild) throw new DiscordError(`gateway READY does not list guild ${binding.guildId}; is the bot in the server?`); say(`gateway: READY as ${ready.user.id}, intents ${CONNECTOR_INTENTS} accepted (message content granted), guild listed`); say("check passed; nothing was sent"); } async function run(opts) { const { bindingFile, binding, contextFiles, toolRoots, pi, journalDir, sessionDir } = prepare(opts); const token = readToken(binding); ensureJournal(journalDir); if (opts.supervised) { const outcome = recover(journalDir); if (outcome === "cleared") warn("supervised start: run.lock left by a process that is gone was removed"); } if (stopRequested(journalDir)) throw new DiscordError(`STOP is present in ${journalDir}; remove it to run`, BRAKE_EXIT); writePid(journalDir, process.pid); const cleanupPid = () => clearPid(journalDir, process.pid); try { await runClaimed(); } catch (err) { cleanupPid(); throw err; } async function runClaimed() { mkdirSync(sessionDir, { recursive: true, mode: 0o700 }); const launches = join(journalDir, "launches"); mkdirSync(launches, { recursive: true, mode: 0o700 }); const launch = mkdtempSync(join(launches, "launch.")); const snapshot = assembleContext(contextFiles, binding); const promptFile = join(launch, "context.md"); writeFileSync(promptFile, snapshot.text, { mode: 0o600 }); writeFileSync(join(launch, "context.sha256"), `${snapshot.sha256} context.md\n`, { mode: 0o600 }); const continueSession = existsSync(sessionDir) && statSync(sessionDir).isDirectory() && hasJsonl(sessionDir); const rest = createRest({ token, log: warn }); const gw = await rest.getGatewayBot(); const gateway = createGateway({ url: gw.url, token, intents: CONNECTOR_INTENTS }); gateway.on("log", (m) => warn(`gateway: ${m}`)); gateway.on("ready", (r) => warn(`gateway: READY as ${r.user && r.user.id}, session ${r.sessionId}`)); gateway.on("resumed", () => warn("gateway: RESUMED")); gateway.on("closed", (c) => warn(`gateway: closed ${c.code} ${c.reason} reconnect=${c.willReconnect}`)); const engine = createEngine({ command: pi, args: buildPiArgs({ ...binding.engine, sessionDir, appendSystemPromptFile: promptFile, continueSession, tools: toolRoots }), cwd: opts.repo, env: { ...process.env, MOSAIC_AGENT_NAME: binding.seat, ...(toolRoots ? { [TOOLS_ENV]: JSON.stringify(toolRoots) } : {}) }, log: warn, onExit: (e) => { warn(`engine exited: ${JSON.stringify(e)}; stopping`); shutdown(1); }, }); // The connector's own SetSpark client (bind and approvals) uses the same // key as the model's verbs; without a setspark key it has none, and an // approval request from the model is refused. const api = binding.tools && binding.tools.setspark ? createSetsparkApi(binding.tools.setspark) : null; const connector = createConnector({ binding, journalDir, rest, gateway, engine, api, log: warn }); let shuttingDown = false; let exitCode = 0; const shutdown = (code) => { if (shuttingDown) return; shuttingDown = true; exitCode = code; warn("stopping: finishing in-flight turns"); connector.stop().catch((err) => warn(`stop failed: ${err.message}`)).finally(() => { cleanupPid(); process.exit(exitCode); }); }; gateway.on("fatal", (f) => { warn(`gateway fatal close ${f.code}: ${f.reason}`); shutdown(2); }); process.on("SIGTERM", () => shutdown(0)); process.on("SIGINT", () => shutdown(0)); // Reload on SIGHUP: re-read the file, refuse anything that is not a // reloadable difference, verify new channels over REST, then swap. // Attempts are serialized; the outcome is journaled and logged, and a // refusal leaves the binding as it was. const reloadNow = async () => { const at = new Date().toISOString(); try { const next = loadBinding(bindingFile); if (next.name !== binding.name) throw new DiscordError(`reload: name ${JSON.stringify(next.name)} does not match ${binding.name}`); const diff = reloadDiff(connector.binding, next); for (const id of diff.channels.added) { const ch = await rest.getChannel(id); if (ch.guild_id !== next.guildId) throw new DiscordError(`reload: channel ${id} is in guild ${ch.guild_id}, not ${next.guildId}`); } connector.reload(next); appendReload(journalDir, { at, outcome: "applied", ...diff }); const n = (d) => `+${d.added.length} -${d.removed.length} ~${d.changed.length}`; warn(`reload applied: channels ${n(diff.channels)}, users ${n(diff.users)}, limits ${diff.limits.length ? diff.limits.join(",") : "unchanged"}${diff.guildName ? ", guildName" : ""}`); } catch (err) { appendReload(journalDir, { at, outcome: "refused", error: err.message }); warn(`reload refused, binding unchanged: ${err.message}`); } }; let reloadChain = Promise.resolve(); process.on("SIGHUP", () => { reloadChain = reloadChain.then(reloadNow, reloadNow); }); const started = await connector.start(); say(`run ${binding.name}: pid ${process.pid}, ${started.inbox} inbox id(s), ${started.reconciled.length} reconciled, context sha256 ${snapshot.sha256}, session ${continueSession ? "continued" : "new"}`); say(`stop with: scripts/discord.sh stop ${binding.name}`); } } function hasJsonl(dir) { try { return readdirSync(dir).some((f) => f.endsWith(".jsonl")); } catch { return false; } } function stop(opts) { const dataRoot = loadDataRoot(opts.config); const binding = loadBinding(bindingPath(dataRoot, opts.binding)); const journalDir = bindingDataDir(dataRoot, binding.name); ensureJournal(journalDir); const path = requestStop(journalDir, "cli stop"); say(`STOP written: ${path}`); const pid = stopTarget(journalDir); if (pid !== null) { process.kill(pid, "SIGTERM"); say(`SIGTERM sent to pid ${pid}; the current turn finishes or times out, then it exits`); } else if (readPid(journalDir) !== null) { say("run.lock exists but its process is gone, is a different process, or cannot be verified; nothing signaled. STOP stays in place; `unlock` clears a lock whose owner is gone"); } else { say("no running connector found for this binding"); } } function unlockCommand(opts) { const dataRoot = loadDataRoot(opts.config); const binding = loadBinding(bindingPath(dataRoot, opts.binding)); const journalDir = bindingDataDir(dataRoot, binding.name); ensureJournal(journalDir); const cleared = unlock(journalDir); say(`STOP written: ${join(journalDir, "STOP")}`); if (cleared === false) say("no run.lock for this binding"); else if (cleared === null) say("run.lock removed; it had no owner record (interrupted start)"); else say(`run.lock removed; its owner pid ${cleared.pid} is gone`); say("remove STOP to run again"); } function reloadCommand(opts) { const dataRoot = loadDataRoot(opts.config); const binding = loadBinding(bindingPath(dataRoot, opts.binding)); if (binding.name !== opts.binding) throw new DiscordError(`binding name ${JSON.stringify(binding.name)} does not match file name ${opts.binding}`); const journalDir = bindingDataDir(dataRoot, binding.name); ensureJournal(journalDir); const pid = stopTarget(journalDir); if (pid === null) throw new DiscordError("no running connector for this binding; the file is valid and applies at the next start", 1); process.kill(pid, "SIGHUP"); say(`SIGHUP sent to pid ${pid}; the outcome is the last line of ${join(journalDir, "reloads.jsonl")} and in its log`); } function recoverCommand(opts) { const dataRoot = loadDataRoot(opts.config); const binding = loadBinding(bindingPath(dataRoot, opts.binding)); const journalDir = bindingDataDir(dataRoot, binding.name); ensureJournal(journalDir); const outcome = recover(journalDir); if (outcome === "cleared") say("run.lock left by a process that is gone was removed; STOP is absent; ready to run"); else say("no lock and no STOP; ready to run"); } async function main() { const opts = parse(process.argv.slice(2)); if (opts.command === "help") { say(USAGE); return 0; } if (opts.command === "check") await check(opts); else if (opts.command === "run") await run(opts); else if (opts.command === "unlock") unlockCommand(opts); else if (opts.command === "recover") recoverCommand(opts); else if (opts.command === "reload") reloadCommand(opts); else stop(opts); return opts.command === "run" ? null : 0; } main().then((code) => { if (code !== null) process.exit(code); }).catch((err) => { const code = err instanceof DiscordError ? err.exitCode : 1; warn(err.message); if (err.details && err.details.nonces) warn(`nonces: ${err.details.nonces.join(", ")}`); process.exit(code); });