From e44d2423cde89db6c1db4bc9d12448981230f3dd Mon Sep 17 00:00:00 2001 From: Rory& Date: Sat, 16 May 2026 21:06:58 +0200 Subject: [PATCH] Demote default rabbitmq behavior, implement init for all BaseEventWriter/Listener abstracted implementations --- src/util/util/ipc/Event.ts | 80 +++++++++++++++++++++----------------- 1 file changed, 45 insertions(+), 35 deletions(-) diff --git a/src/util/util/ipc/Event.ts b/src/util/util/ipc/Event.ts index 8a2c55d88..9e523ce6c 100644 --- a/src/util/util/ipc/Event.ts +++ b/src/util/util/ipc/Event.ts @@ -28,41 +28,47 @@ import { BaseEventWriter } from "./writer/BaseEventWriter"; import { UnixSocketWriter } from "./writer/UnixSocketWriter"; import { UnixSocketListener } from "./listener/UnixSocketListener"; import { RabbitMqSingleListener } from "./listener/RabbitMqSingleListener"; +import { RabbitMqSingleWriter } from "./writer/RabbitMqSingleWriter"; export const events = new EventEmitter(); let listener: BaseEventListener | null = null; let writer: BaseEventWriter | null = null; -let unixSocketListener: UnixSocketListener | null = null; -let unixSocketWriter: UnixSocketWriter | null = null; export async function emitEvent(payload: Omit) { const id = (payload.guild_id || payload.channel_id || payload.user_id || payload.session_id) as string; if (!id) return console.error("event doesn't contain any id", payload); - if (RabbitMQ.connection) { - await RabbitMQ.publishEvent(id, payload); - } else if (process.env.EVENT_TRANSMISSION === "unix" && process.env.EVENT_SOCKET_PATH) { - if (!unixSocketWriter) { - console.error("[Event] Unix socket writer not initialized, cannot emit event!"); - throw new Error("Unix socket writer not initialized"); - } - await unixSocketWriter.emit(payload); + if (writer) { + await writer.emit(payload); } else if (process.env.EVENT_TRANSMISSION === "process") { process.send?.({ type: "event", event: payload, id } as ProcessEvent); + } else if (RabbitMQ.connection) { + await RabbitMQ.publishEvent(id, payload); } else { events.emit(id, payload); } } export async function initEvent() { - await RabbitMQ.init(); // does nothing if rabbitmq is not setup - - if (process.env.EVENT_TRANSMISSION === "unix" && process.env.EVENT_SOCKET_PATH) { - if (!unixSocketWriter) { - writer = unixSocketWriter = new UnixSocketWriter(process.env.EVENT_SOCKET_PATH); - await unixSocketWriter.init(); + if (process.env.EVENT_TRANSMISSION === "rabbitmq-single") { + if (!Config.get().rabbitmq.host!) { + throw new Error("[Events] RabbitMQ is not configured."); } - } + + if (!writer) { + writer = new RabbitMqSingleWriter(Config.get().rabbitmq.host!); + await writer.init(); + } + } else if (process.env.EVENT_TRANSMISSION === "unix" && process.env.EVENT_SOCKET_PATH) { + if (!process.env.EVENT_SOCKET_PATH) { + throw new Error("[Events] EVENT_SOCKET_PATH is not configured."); + } + + if (!writer) { + writer = new UnixSocketWriter(process.env.EVENT_SOCKET_PATH); + await writer.init(); + } + } else await RabbitMQ.init(); // does nothing if rabbitmq is not setup // Set up the spacebar event listener (used for config reload, etc.) const setupSpacebarListener = async () => { @@ -104,32 +110,25 @@ export interface ProcessEvent { } export async function listenEvent(event: string, callback: (event: EventOpts) => unknown, opts?: ListenEventOpts): Promise<() => Promise> { - if (RabbitMQ.connection) { - if (process.env.EVENT_TRANSMISSION !== "rabbitmq-legacy") { - console.warn(yellow("[Events] Warning:"), "RabbitMQ replication without configuring EVENT_TRANSMISSION is deprecated."); - console.warn(yellow("[Events] Warning:"), "Set EVENT_TRANSMISSION to 'rabbitmq-legacy' in environment variables to silence this warning."); - } - const rabbitMQChannel = await RabbitMQ.getSafeChannel(); - const channel = opts?.channel || rabbitMQChannel; - if (!channel) throw new Error("[Events] An event was sent without an associated channel"); - return await rabbitListen(channel, event, callback, { - acknowledge: opts?.acknowledge, - }); - } else if (process.env.EVENT_TRANSMISSION === "rabbitmq-single") { + if (process.env.EVENT_TRANSMISSION === "rabbitmq-single") { if (!Config.get().rabbitmq.host) { - console.error("[Events] RabbitMQ is not configured."); + throw new Error("[Events] EVENT_SOCKET_PATH is not configured."); } if (!listener) { listener = new RabbitMqSingleListener(Config.get().rabbitmq.host!); await listener.init(); } return await listener.listen(event, callback); - } else if (process.env.EVENT_TRANSMISSION === "unix" && process.env.EVENT_SOCKET_PATH) { - if (!unixSocketListener) { - listener = unixSocketListener = new UnixSocketListener(path.join(process.env.EVENT_SOCKET_PATH, `${process.pid}.sock`)); - await unixSocketListener.init(); + } else if (process.env.EVENT_TRANSMISSION === "unix") { + if (!process.env.EVENT_SOCKET_PATH) { + throw new Error("[Events] EVENT_SOCKET_PATH is not configured."); } - return await unixSocketListener.listen(event, callback); + + if (!listener) { + listener = listener = new UnixSocketListener(path.join(process.env.EVENT_SOCKET_PATH, `${process.pid}.sock`)); + await listener.init(); + } + return await listener.listen(event, callback); } else if (process.env.EVENT_TRANSMISSION === "process") { const cancel = async () => { process.removeListener("message", listener); @@ -146,6 +145,17 @@ export async function listenEvent(event: string, callback: (event: EventOpts) => process.setMaxListeners(process.getMaxListeners() + 1); return cancel; + } else if (RabbitMQ.connection) { + if (process.env.EVENT_TRANSMISSION !== "rabbitmq-legacy") { + console.warn(yellow("[Events] Warning:"), "RabbitMQ replication without configuring EVENT_TRANSMISSION is deprecated."); + console.warn(yellow("[Events] Warning:"), "Set EVENT_TRANSMISSION to 'rabbitmq-legacy' in environment variables to silence this warning."); + } + const rabbitMQChannel = await RabbitMQ.getSafeChannel(); + const channel = opts?.channel || rabbitMQChannel; + if (!channel) throw new Error("[Events] An event was sent without an associated channel"); + return await rabbitListen(channel, event, callback, { + acknowledge: opts?.acknowledge, + }); } else { const listener = (opts: EventOpts) => callback({ ...opts, cancel }); const cancel = async () => {