Files
stack/packages/discord/src/gateway.mjs
T
jason.woltjeandClaude Fable 5.1 786e379c49 feat(discord): connector pilot for the Sage seat, reviewed candidate (#1509)
Zero-dependency Discord connector under packages/discord: binding
validation, REST and gateway clients, pi engine adapter, journal with
append-only inbox, outbox, admissions and notices, and a run.lock
ownership record {pid, start, boot} whose identity is checked three ways
and whose cleanup is gated by STOP. CLI check|run|stop|unlock via
scripts/discord.sh; offline suite scripts/test-discord.sh (28 checks,
87 node tests).

Reviewed by rev-code-02 on #1509 over nine rounds; approved exact tree
4e0feb6758c0a7e4a71483912a8e0d3e3ec95aef at comment 26170. Corrections
(1) to (12) recorded in BUILD-LOG. No listener started, no token read,
no Discord write; the live pilot follows this commit per the brief.

Co-Authored-By: Claude Fable 5.1 <[email protected]>
2026-09-13 01:15:37 -05:00

270 lines
8.9 KiB
JavaScript

// Discord gateway v10 client on the built-in WebSocket. Identify, heartbeat,
// resume, and the four opcodes that steer reconnects. Nothing else. The
// token is sent in identify/resume and never appears in a log or event.
//
// Everything with a side effect is injectable so the offline suite can drive
// it with a fake socket and fake timers: WebSocketImpl, setTimeout,
// clearTimeout, random.
//
// Events (on(name, fn)):
// ready {sessionId, user, guilds} after READY
// resumed after RESUMED
// dispatch {t, d, s} every op 0 (including READY)
// closed {code, reason, willReconnect} every socket close
// fatal {code, reason} close code that must not reconnect
// log string
import { DiscordError } from "./errors.mjs";
export const OP = Object.freeze({
DISPATCH: 0, HEARTBEAT: 1, IDENTIFY: 2, RESUME: 6, RECONNECT: 7, INVALID_SESSION: 9, HELLO: 10, HEARTBEAT_ACK: 11,
});
export const INTENT = Object.freeze({ GUILDS: 1 << 0, GUILD_MESSAGES: 1 << 9, MESSAGE_CONTENT: 1 << 15 });
export const CONNECTOR_INTENTS = INTENT.GUILDS | INTENT.GUILD_MESSAGES | INTENT.MESSAGE_CONTENT;
// Close codes after which reconnecting is wrong. 4014 is the one `check`
// cares about: the message-content intent is not granted in the portal.
export const FATAL_CLOSE = Object.freeze({
4004: "authentication failed (bad token)",
4010: "invalid shard",
4011: "sharding required",
4012: "invalid API version",
4013: "invalid intents",
4014: "disallowed intent: a privileged intent (message content) is not enabled for this bot in the developer portal",
});
// Codes after which the session cannot be resumed; identify instead.
const NO_RESUME_CLOSE = new Set([1000, 1001, 4007, 4009]);
export const GATEWAY_QUERY = "/?v=10&encoding=json";
const MISSED_ACK_CLOSE = 4900; // app-private, used when a heartbeat ack never arrived
const MAX_BACKOFF_MS = 60000;
export function createGateway({
url, token, intents = CONNECTOR_INTENTS,
WebSocketImpl = globalThis.WebSocket,
setTimeoutImpl = globalThis.setTimeout, clearTimeoutImpl = globalThis.clearTimeout,
random = Math.random,
properties = { os: process.platform, browser: "mosaic-stack", device: "mosaic-stack" },
} = {}) {
if (typeof url !== "string" || !url.startsWith("wss://")) throw new DiscordError(`gateway: url must be wss://, got ${JSON.stringify(url)}`, 1);
if (typeof token !== "string" || token.length === 0) throw new DiscordError("gateway: token required", 1);
if (typeof WebSocketImpl !== "function") throw new DiscordError("gateway: WebSocket implementation required", 1);
const listeners = new Map();
const emit = (name, payload) => {
for (const fn of listeners.get(name) || []) fn(payload);
};
const state = {
ws: null, seq: null, sessionId: null, resumeUrl: null,
heartbeatTimer: null, ackPending: false, reconnectTimer: null,
attempts: 0, closedByUs: false, stopped: false, connected: false,
};
function log(msg) {
emit("log", msg);
}
function send(payload) {
if (!state.ws || state.ws.readyState !== 1) return false;
state.ws.send(JSON.stringify(payload));
return true;
}
function stopHeartbeat() {
if (state.heartbeatTimer !== null) {
clearTimeoutImpl(state.heartbeatTimer);
state.heartbeatTimer = null;
}
}
function beat(intervalMs) {
if (state.ackPending) {
log("heartbeat ack missed; reconnecting");
closeSocket(MISSED_ACK_CLOSE, "missed heartbeat ack");
return;
}
state.ackPending = true;
send({ op: OP.HEARTBEAT, d: state.seq });
state.heartbeatTimer = setTimeoutImpl(() => beat(intervalMs), intervalMs);
}
function startHeartbeat(intervalMs) {
stopHeartbeat();
state.ackPending = false;
const first = Math.floor(intervalMs * random());
state.heartbeatTimer = setTimeoutImpl(() => beat(intervalMs), first);
}
function identify() {
send({ op: OP.IDENTIFY, d: { token, intents, properties } });
}
function resume() {
send({ op: OP.RESUME, d: { token, session_id: state.sessionId, seq: state.seq } });
}
function closeSocket(code, reason) {
stopHeartbeat();
state.closedByUs = true;
const ws = state.ws;
if (ws && (ws.readyState === 0 || ws.readyState === 1)) {
try {
ws.close(code, reason);
} catch (err) {
log(`close failed: ${err.message}`);
}
}
}
function scheduleReconnect(canResume) {
if (state.stopped) return;
if (!canResume) {
state.sessionId = null;
state.resumeUrl = null;
state.seq = null;
}
state.attempts += 1;
const delay = Math.min(1000 * 2 ** Math.min(state.attempts - 1, 6), MAX_BACKOFF_MS) + Math.floor(random() * 1000);
log(`reconnect in ${delay} ms (${canResume ? "resume" : "identify"})`);
state.reconnectTimer = setTimeoutImpl(() => {
state.reconnectTimer = null;
open();
}, delay);
}
function handlePayload(payload) {
if (typeof payload.s === "number") state.seq = payload.s;
switch (payload.op) {
case OP.HELLO: {
const interval = payload.d && typeof payload.d.heartbeat_interval === "number" ? payload.d.heartbeat_interval : 41250;
startHeartbeat(interval);
if (state.sessionId) resume();
else identify();
return;
}
case OP.HEARTBEAT:
state.ackPending = false;
send({ op: OP.HEARTBEAT, d: state.seq });
return;
case OP.HEARTBEAT_ACK:
state.ackPending = false;
state.attempts = 0;
return;
case OP.RECONNECT:
log("gateway asked for reconnect");
closeSocket(4000, "reconnect requested");
return;
case OP.INVALID_SESSION: {
const resumable = payload.d === true;
log(`invalid session (resumable=${resumable})`);
if (!resumable) {
state.sessionId = null;
state.resumeUrl = null;
state.seq = null;
}
closeSocket(4000, "invalid session");
return;
}
case OP.DISPATCH: {
if (payload.t === "READY") {
state.sessionId = payload.d.session_id;
state.resumeUrl = payload.d.resume_gateway_url || null;
state.attempts = 0;
emit("ready", { sessionId: state.sessionId, user: payload.d.user, guilds: payload.d.guilds || [] });
} else if (payload.t === "RESUMED") {
state.attempts = 0;
emit("resumed", {});
}
emit("dispatch", { t: payload.t, d: payload.d, s: payload.s });
return;
}
default:
return;
}
}
function open() {
if (state.stopped) return;
const target = (state.sessionId && state.resumeUrl ? state.resumeUrl : url).replace(/\/+$/, "") + GATEWAY_QUERY;
state.closedByUs = false;
let ws;
try {
ws = new WebSocketImpl(target);
} catch (err) {
log(`socket open failed: ${err.message}`);
scheduleReconnect(Boolean(state.sessionId));
return;
}
state.ws = ws;
ws.onopen = () => {
state.connected = true;
};
ws.onmessage = (ev) => {
let payload;
try {
payload = JSON.parse(typeof ev.data === "string" ? ev.data : String(ev.data));
} catch {
log("unparseable gateway frame ignored");
return;
}
if (!payload || typeof payload !== "object") return;
handlePayload(payload);
};
ws.onerror = (ev) => {
log(`socket error: ${(ev && ev.message) || "unknown"}`);
};
ws.onclose = (ev) => {
if (state.ws !== ws) return;
state.ws = null;
state.connected = false;
stopHeartbeat();
const code = ev && typeof ev.code === "number" ? ev.code : 1006;
const reason = (ev && ev.reason) || "";
if (FATAL_CLOSE[code]) {
state.stopped = true;
emit("closed", { code, reason, willReconnect: false });
emit("fatal", { code, reason: FATAL_CLOSE[code] });
return;
}
if (state.stopped) {
emit("closed", { code, reason, willReconnect: false });
return;
}
const canResume = Boolean(state.sessionId) && !NO_RESUME_CLOSE.has(code);
emit("closed", { code, reason, willReconnect: true });
scheduleReconnect(canResume);
};
}
return {
on(name, fn) {
if (!listeners.has(name)) listeners.set(name, new Set());
listeners.get(name).add(fn);
return () => listeners.get(name).delete(fn);
},
connect() {
if (state.stopped) throw new DiscordError("gateway: already stopped", 1);
open();
},
// Final. No reconnect after this.
close(code = 1000, reason = "closing") {
state.stopped = true;
if (state.reconnectTimer !== null) {
clearTimeoutImpl(state.reconnectTimer);
state.reconnectTimer = null;
}
closeSocket(code, reason);
},
get connected() {
return state.connected;
},
get sessionId() {
return state.sessionId;
},
get seq() {
return state.seq;
},
};
}