feat(discord): binding reload without a restart, and a per-user channel allowlist (#1509)
`reload` validates the binding file and sends SIGHUP to the live owner; the running connector re-reads it and swaps guildName, channels, users and limits in place. name, seat, guildId, botUserId, tokenFile, engine and context are fixed for the life of the process; a change there, an invalid file or a channel outside the guild refuses the reload and keeps the old binding. Every attempt is one line in reloads.jsonl. The service unit maps `systemctl --user reload` to the same signal. A user entry may carry `channels`, an allowlist of listed channel ids; absent means every listed channel. Outside the list the message is dropped as channel-not-for-user; threads count as their parent. Suite 41/41, 101 node tests. QUEUE rows 19 and 20 opened. Co-Authored-By: Claude Fable 5.1 <[email protected]>
This commit is contained in:
@@ -28,6 +28,7 @@ export const DROP = Object.freeze({
|
||||
SELF: "author-is-self",
|
||||
USER: "user-unlisted",
|
||||
CHANNEL: "channel-unlisted",
|
||||
USER_CHANNEL: "channel-not-for-user",
|
||||
THREAD_PARENT: "thread-parent-unlisted",
|
||||
MENTION: "no-mention",
|
||||
});
|
||||
@@ -47,7 +48,8 @@ export function authorize(binding, message, channelInfo = () => undefined) {
|
||||
if (author.id === binding.botUserId) return { ok: false, reason: DROP.SELF };
|
||||
if (message.webhook_id !== undefined && message.webhook_id !== null) return { ok: false, reason: DROP.WEBHOOK };
|
||||
if (author.bot === true || author.system === true) return { ok: false, reason: DROP.BOT };
|
||||
if (!binding.users.some((u) => u.id === author.id)) return { ok: false, reason: DROP.USER };
|
||||
const user = binding.users.find((u) => u.id === author.id);
|
||||
if (!user) return { ok: false, reason: DROP.USER };
|
||||
|
||||
let channel = binding.channels.find((c) => c.id === message.channel_id);
|
||||
let thread = null;
|
||||
@@ -59,6 +61,9 @@ export function authorize(binding, message, channelInfo = () => undefined) {
|
||||
if (!channel) return { ok: false, reason: DROP.THREAD_PARENT };
|
||||
thread = { id: message.channel_id, name: typeof info.name === "string" ? info.name : null };
|
||||
}
|
||||
// A user's channel allowlist, when present, is checked against the listed
|
||||
// channel, so a thread counts as its parent.
|
||||
if (user.channels && !user.channels.includes(channel.id)) return { ok: false, reason: DROP.USER_CHANNEL };
|
||||
// A thread inherits the parent's mode (Q4). @everyone is not a mention of
|
||||
// the bot; only an entry in `mentions` with the bot's id counts.
|
||||
if (channel.mode === "mention" && !mentionsBot(message, binding.botUserId)) return { ok: false, reason: DROP.MENTION };
|
||||
|
||||
@@ -7,6 +7,11 @@
|
||||
// Loading fails closed: unknown key, missing field, wrong type, bad mode,
|
||||
// symlink, empty allowlist. Nothing is defaulted silently except the limits'
|
||||
// documented defaults below, which the fixture spells out anyway.
|
||||
//
|
||||
// A user entry may carry `channels`, an allowlist of listed channel ids;
|
||||
// absent means every listed channel. A running connector may re-read the
|
||||
// file (`reload`): `reloadDiff` says which keys may change in place and
|
||||
// refuses the rest.
|
||||
|
||||
import { existsSync, lstatSync, readFileSync, realpathSync, statSync } from "node:fs";
|
||||
import { isAbsolute, join, resolve, sep } from "node:path";
|
||||
@@ -28,7 +33,7 @@ export const LIMIT_DEFAULTS = Object.freeze({
|
||||
|
||||
const TOP_KEYS = ["bindingVersion", "name", "seat", "guildId", "guildName", "botUserId", "tokenFile", "channels", "users", "engine", "limits", "context"];
|
||||
const CHANNEL_KEYS = ["id", "name", "mode"];
|
||||
const USER_KEYS = ["id", "name"];
|
||||
const USER_KEYS = ["id", "name", "channels"];
|
||||
const ENGINE_KEYS = ["provider", "model", "thinking"];
|
||||
const LIMIT_KEYS = Object.keys(LIMIT_DEFAULTS);
|
||||
const CONTEXT_KEYS = ["files"];
|
||||
@@ -123,7 +128,19 @@ export function validateBinding(raw, where = "binding") {
|
||||
const w = `${where}.users[${i}]`;
|
||||
if (!isObject(u)) throw new DiscordError(`${w}: not an object`);
|
||||
onlyKeys(u, USER_KEYS, w);
|
||||
return Object.freeze({ id: requireSnowflake(u, "id", w), name: requireString(u, "name", w) });
|
||||
const id = requireSnowflake(u, "id", w);
|
||||
const uname = requireString(u, "name", w);
|
||||
let allowed = null;
|
||||
if (u.channels !== undefined) {
|
||||
if (!Array.isArray(u.channels) || u.channels.length === 0) throw new DiscordError(`${w}: channels must be a non-empty array of listed channel ids`);
|
||||
allowed = u.channels.map((cid, j) => {
|
||||
if (typeof cid !== "string" || !SNOWFLAKE.test(cid)) throw new DiscordError(`${w}.channels[${j}]: not a Discord snowflake id`);
|
||||
if (!channels.some((c) => c.id === cid)) throw new DiscordError(`${w}.channels[${j}]: ${cid} is not a listed channel`);
|
||||
return cid;
|
||||
});
|
||||
if (new Set(allowed).size !== allowed.length) throw new DiscordError(`${w}: duplicate channel id`);
|
||||
}
|
||||
return Object.freeze({ id, name: uname, channels: allowed === null ? null : Object.freeze(allowed) });
|
||||
});
|
||||
if (new Set(users.map((u) => u.id)).size !== users.length) throw new DiscordError(`${where}: duplicate user id`);
|
||||
if (users.some((u) => u.id === botUserId)) throw new DiscordError(`${where}: the bot cannot be an authorized user`);
|
||||
@@ -166,6 +183,36 @@ export function validateBinding(raw, where = "binding") {
|
||||
});
|
||||
}
|
||||
|
||||
// What a running connector may take from a re-read binding, and what it may
|
||||
// not: the engine and its prompt are launched once, the token is read once,
|
||||
// and the journal directory is named after the binding. A change to a fixed
|
||||
// key needs a stop and a start. Returns a summary of the reloadable
|
||||
// differences or throws with exit 2.
|
||||
export const RELOADABLE_KEYS = Object.freeze(["guildName", "channels", "users", "limits"]);
|
||||
export const FIXED_KEYS = Object.freeze(["bindingVersion", "name", "seat", "guildId", "botUserId", "tokenFile", "engine", "context"]);
|
||||
|
||||
export function reloadDiff(current, next) {
|
||||
for (const k of FIXED_KEYS) {
|
||||
if (JSON.stringify(current[k]) !== JSON.stringify(next[k])) throw new DiscordError(`reload: ${k} cannot change while running; stop and start instead`);
|
||||
}
|
||||
const byId = (xs) => new Map(xs.map((x) => [x.id, JSON.stringify(x)]));
|
||||
const listDiff = (a, b) => {
|
||||
const A = byId(a);
|
||||
const B = byId(b);
|
||||
return Object.freeze({
|
||||
added: Object.freeze([...B.keys()].filter((id) => !A.has(id))),
|
||||
removed: Object.freeze([...A.keys()].filter((id) => !B.has(id))),
|
||||
changed: Object.freeze([...B.keys()].filter((id) => A.has(id) && A.get(id) !== B.get(id))),
|
||||
});
|
||||
};
|
||||
return Object.freeze({
|
||||
channels: listDiff(current.channels, next.channels),
|
||||
users: listDiff(current.users, next.users),
|
||||
limits: Object.freeze(LIMIT_KEYS.filter((k) => current.limits[k] !== next.limits[k])),
|
||||
guildName: current.guildName !== next.guildName,
|
||||
});
|
||||
}
|
||||
|
||||
// A private file: regular, not a symlink, owner-only (0600), non-empty.
|
||||
export function checkPrivateFile(path, what) {
|
||||
let st;
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
// mosaic-discord stop <binding> [--config PATH]
|
||||
// mosaic-discord unlock <binding> [--config PATH]
|
||||
// mosaic-discord recover <binding> [--config PATH]
|
||||
// mosaic-discord reload <binding> [--config PATH]
|
||||
//
|
||||
// <binding> names <dataRoot>/discord/<binding>.json. The repository wrapper
|
||||
// is scripts/discord.sh.
|
||||
@@ -33,6 +34,14 @@
|
||||
// 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.
|
||||
@@ -41,13 +50,13 @@ import { existsSync, mkdirSync, mkdtempSync, writeFileSync, statSync, readdirSyn
|
||||
import { join, resolve } from "node:path";
|
||||
import { createHash } from "node:crypto";
|
||||
import { DiscordError } from "./errors.mjs";
|
||||
import { defaultConfigPath, loadDataRoot, bindingPath, bindingDataDir, loadBinding, readToken, resolveContextFiles } from "./binding.mjs";
|
||||
import { defaultConfigPath, loadDataRoot, bindingPath, bindingDataDir, loadBinding, readToken, resolveContextFiles, reloadDiff } from "./binding.mjs";
|
||||
import { createRest } from "./rest.mjs";
|
||||
import { createGateway, CONNECTOR_INTENTS } from "./gateway.mjs";
|
||||
import { createEngine, buildPiArgs } from "./engine-pi.mjs";
|
||||
import { assembleContext } from "./context.mjs";
|
||||
import { createConnector } from "./connector.mjs";
|
||||
import { ensureJournal, requestStop, stopRequested, readPid, stopTarget, writePid, clearPid, unlock, recover, BRAKE_EXIT } from "./journal.mjs";
|
||||
import { ensureJournal, requestStop, stopRequested, readPid, stopTarget, writePid, clearPid, unlock, recover, appendReload, BRAKE_EXIT } from "./journal.mjs";
|
||||
|
||||
const USAGE = [
|
||||
"usage: mosaic-discord check <binding> [--config PATH] [--repo PATH]",
|
||||
@@ -55,6 +64,7 @@ const USAGE = [
|
||||
" mosaic-discord stop <binding> [--config PATH]",
|
||||
" mosaic-discord unlock <binding> [--config PATH]",
|
||||
" mosaic-discord recover <binding> [--config PATH]",
|
||||
" mosaic-discord reload <binding> [--config PATH]",
|
||||
].join("\n");
|
||||
|
||||
function parse(argv) {
|
||||
@@ -74,7 +84,7 @@ function parse(argv) {
|
||||
else throw new DiscordError(`unexpected argument: ${a}\n${USAGE}`, 4);
|
||||
}
|
||||
if (opts.command === "help") return opts;
|
||||
if (!["check", "run", "stop", "unlock", "recover"].includes(opts.command)) throw new DiscordError(USAGE, 4);
|
||||
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;
|
||||
@@ -90,14 +100,15 @@ function warn(msg) {
|
||||
// Everything that can be checked without the network, shared by check and run.
|
||||
function prepare(opts) {
|
||||
const dataRoot = loadDataRoot(opts.config);
|
||||
const binding = loadBinding(bindingPath(dataRoot, opts.binding));
|
||||
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 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, binding, contextFiles, pi, journalDir, sessionDir };
|
||||
return { dataRoot, bindingFile, binding, contextFiles, pi, journalDir, sessionDir };
|
||||
}
|
||||
|
||||
async function check(opts) {
|
||||
@@ -148,7 +159,7 @@ async function check(opts) {
|
||||
}
|
||||
|
||||
async function run(opts) {
|
||||
const { binding, contextFiles, pi, journalDir, sessionDir } = prepare(opts);
|
||||
const { bindingFile, binding, contextFiles, pi, journalDir, sessionDir } = prepare(opts);
|
||||
const token = readToken(binding);
|
||||
ensureJournal(journalDir);
|
||||
if (opts.supervised) {
|
||||
@@ -216,6 +227,34 @@ async function run(opts) {
|
||||
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}`);
|
||||
@@ -261,6 +300,18 @@ function unlockCommand(opts) {
|
||||
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));
|
||||
@@ -281,6 +332,7 @@ async function main() {
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -11,7 +11,10 @@
|
||||
// window is marked refused, not re-sent, because a re-send could post a
|
||||
// second reply. If any intent is still unknown after that, start refuses.
|
||||
// 4. A STOP file refuses new turns; the current one finishes or times out.
|
||||
// 5. The daily ceiling counts admissions (admissions.jsonl) on the current
|
||||
// 5. A reload swaps the binding between turns: authorization and limits
|
||||
// read the current binding at admission; a turn in flight keeps the
|
||||
// values it started with. Fixed keys (see binding.mjs) are refused.
|
||||
// 6. The daily ceiling counts admissions (admissions.jsonl) on the current
|
||||
// UTC date; an admission is written before the engine runs, so turns in
|
||||
// flight and turns cut short by a crash count too.
|
||||
// Over the ceiling, inbound messages are journaled as drops, one fixed
|
||||
@@ -21,6 +24,7 @@
|
||||
// timers. The offline suite drives this with fakes.
|
||||
|
||||
import { authorize, isThreadType } from "./authorize.mjs";
|
||||
import { reloadDiff } from "./binding.mjs";
|
||||
import { envelope, splitReply } from "./context.mjs";
|
||||
import {
|
||||
ensureJournal, appendInbox, readInboxIds, appendOutbox, unresolvedOutbox, appendDrop,
|
||||
@@ -41,17 +45,18 @@ export const FIXED_LINES = Object.freeze({
|
||||
export const RECONCILE_WINDOW_MS = 5 * 60 * 1000;
|
||||
|
||||
export function createConnector({
|
||||
binding, journalDir, rest, gateway, engine,
|
||||
binding: initialBinding, journalDir, rest, gateway, engine,
|
||||
now = () => Date.now(),
|
||||
setTimeoutImpl = globalThis.setTimeout, clearTimeoutImpl = globalThis.clearTimeout,
|
||||
typingIntervalMs = 8000,
|
||||
readReceipt = READ_RECEIPT,
|
||||
log = () => {},
|
||||
} = {}) {
|
||||
for (const [k, v] of Object.entries({ binding, journalDir, rest, gateway, engine })) {
|
||||
for (const [k, v] of Object.entries({ binding: initialBinding, journalDir, rest, gateway, engine })) {
|
||||
if (!v) throw new DiscordError(`connector: ${k} required`, 1);
|
||||
}
|
||||
ensureJournal(journalDir);
|
||||
let binding = initialBinding;
|
||||
|
||||
const state = {
|
||||
inbox: new Set(), channels: new Map(), inFlight: 0, typingTimer: null, typingChannel: null,
|
||||
@@ -308,6 +313,16 @@ export function createConnector({
|
||||
onDispatch,
|
||||
reconcile,
|
||||
rememberChannel,
|
||||
get binding() {
|
||||
return binding;
|
||||
},
|
||||
// Swap the binding in place. Throws (exit 2) and changes nothing when a
|
||||
// fixed key differs. Returns the summary of what changed.
|
||||
reload(next) {
|
||||
const diff = reloadDiff(binding, next);
|
||||
binding = next;
|
||||
return diff;
|
||||
},
|
||||
get inFlight() {
|
||||
return state.inFlight;
|
||||
},
|
||||
|
||||
@@ -10,6 +10,9 @@
|
||||
// ({at, reason}), appended, so the last line names who
|
||||
// braked; `recover` removes only a STOP it wrote itself
|
||||
// notices.jsonl once-per-day fixed lines already attempted (ceiling)
|
||||
// reloads.jsonl one line per binding reload attempt: applied (with the
|
||||
// differences) or refused (with the reason); the binding
|
||||
// in memory only changes on an applied line
|
||||
// run.lock/ ownership directory (mkdir is atomic) holding owner.json
|
||||
// {pid, start, boot}; `stop` signals only a live pid whose
|
||||
// start time and boot id match; a stale lock refuses `run`
|
||||
@@ -435,6 +438,17 @@ export function clearPid(dir, pid) {
|
||||
if (rec !== null && !rec.invalid && rec.pid === pid) rmSync(lockPath(dir), { recursive: true, force: true });
|
||||
}
|
||||
|
||||
// --- reloads ---
|
||||
|
||||
export function appendReload(dir, entry) {
|
||||
if (!["applied", "refused"].includes(entry.outcome)) throw new DiscordError("reload entry needs an outcome", 1);
|
||||
appendLine(join(dir, "reloads.jsonl"), entry);
|
||||
}
|
||||
|
||||
export function readReloads(dir) {
|
||||
return readLines(join(dir, "reloads.jsonl"));
|
||||
}
|
||||
|
||||
// --- notices ---
|
||||
// Fixed lines that must go out at most once per UTC day (the ceiling
|
||||
// notice). The line is appended before the delivery attempt, so a crash
|
||||
|
||||
Reference in New Issue
Block a user