feat(system): implement graceful shutdown and expand websocket events
Improve system reliability and real-time capabilities by implementing a robust lifecycle management system and adding new broadcast events for attachments and voice recordings. - Implement asynchronous graceful shutdown in backend to close HTTP, WebSocket, Redis, and database connections. - Add new WebSocket broadcast events: `attachment_created`, `voice_recording_started`, `voice_recording_stopped`, `voice_recording_uploaded`, and `analysis_queue_status`. - Refactor media status handling to use boolean `playing` state instead of string-based status. - Centralize `PageResult` and `VoiceRecording` types to improve consistency between frontend and backend. - Update frontend API client to include `listRecordings` and handle new WebSocket event types. - Fix type mismatches in voice command handling and media status reporting.
This commit is contained in:
@@ -1,6 +1,10 @@
|
||||
import type { Server } from "node:http";
|
||||
import { createChildLogger } from "@bete/shared/logger";
|
||||
import { startHttpServer } from "./http/server.js";
|
||||
import { closeDatabase } from "./shared/database/index.js";
|
||||
import { stopCommandBridge } from "./shared/redis/index.js";
|
||||
import { stopRedisBridge as stopEventBridge } from "./ws/redis-bridge.js";
|
||||
import { closeWebSocketServer } from "./ws/server.js";
|
||||
|
||||
const logger = createChildLogger("backend");
|
||||
|
||||
@@ -17,22 +21,43 @@ async function main() {
|
||||
}
|
||||
}
|
||||
|
||||
function shutdown(signal: string) {
|
||||
async function shutdown(signal: string) {
|
||||
logger.info({ signal }, "Shutting down gracefully");
|
||||
|
||||
if (httpServer) {
|
||||
httpServer.close(() => {
|
||||
logger.info("HTTP server closed");
|
||||
process.exit(0);
|
||||
});
|
||||
try {
|
||||
// 1. Stop accepting new HTTP connections
|
||||
if (httpServer) {
|
||||
await new Promise<void>((resolve) => {
|
||||
httpServer!.close(() => {
|
||||
logger.info("HTTP server closed");
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
// Force exit after 10s if connections don't close
|
||||
setTimeout(() => {
|
||||
logger.error("Forced shutdown after timeout");
|
||||
process.exit(1);
|
||||
}, 10_000).unref();
|
||||
} else {
|
||||
// 2. Close WebSocket server
|
||||
closeWebSocketServer();
|
||||
|
||||
// 3. Stop Redis bridges (event subscriptions + command channel)
|
||||
await Promise.allSettled([
|
||||
stopEventBridge().catch((err) =>
|
||||
logger.warn({ err }, "Error stopping event bridge"),
|
||||
),
|
||||
stopCommandBridge().catch((err) =>
|
||||
logger.warn({ err }, "Error stopping command bridge"),
|
||||
),
|
||||
]);
|
||||
|
||||
// 4. Close database pool
|
||||
await closeDatabase().catch((err) =>
|
||||
logger.warn({ err }, "Error closing database"),
|
||||
);
|
||||
|
||||
logger.info("Graceful shutdown completed");
|
||||
process.exit(0);
|
||||
} catch (err) {
|
||||
logger.error({ err }, "Error during graceful shutdown");
|
||||
process.exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -41,12 +66,12 @@ process.on("SIGTERM", () => shutdown("SIGTERM"));
|
||||
|
||||
process.on("uncaughtException", (err) => {
|
||||
logger.error({ err }, "Uncaught exception");
|
||||
process.exit(1);
|
||||
shutdown("uncaughtException");
|
||||
});
|
||||
|
||||
process.on("unhandledRejection", (reason) => {
|
||||
logger.error({ reason }, "Unhandled rejection");
|
||||
process.exit(1);
|
||||
shutdown("unhandledRejection");
|
||||
});
|
||||
|
||||
main();
|
||||
|
||||
@@ -47,8 +47,14 @@ export async function getStatus(): Promise<MediaState> {
|
||||
const cached = await readRedisStatus("media:status");
|
||||
|
||||
if (cached) {
|
||||
const rawPlaying = cached.playing;
|
||||
// Handle both boolean (new) and string (legacy from String(discordPlayer.getStatus()))
|
||||
const playing =
|
||||
rawPlaying === true ||
|
||||
rawPlaying === "playing" ||
|
||||
rawPlaying === "buffering";
|
||||
return {
|
||||
playing: Boolean(cached.playing),
|
||||
playing,
|
||||
musicVolume: Number(cached.musicVolume ?? 1.0),
|
||||
current: (cached.current as MediaItem | null) ?? null,
|
||||
queue: (cached.queue as MediaItem[]) ?? [],
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import type { PageResult } from "@bete/shared";
|
||||
import { createChildLogger } from "@bete/shared/logger";
|
||||
import { getPool } from "../../shared/database/index.js";
|
||||
import type {
|
||||
@@ -8,11 +9,6 @@ import type {
|
||||
|
||||
const logger = createChildLogger("messages.repository");
|
||||
|
||||
export interface PageResult<T> {
|
||||
data: T[];
|
||||
nextCursor: string | null;
|
||||
}
|
||||
|
||||
export interface AttachmentResult {
|
||||
id: string;
|
||||
message_id: string;
|
||||
|
||||
@@ -12,13 +12,15 @@ export async function handleGetVoiceStatus(_req: Request, res: Response) {
|
||||
res.json(status);
|
||||
}
|
||||
|
||||
/** Safely extract a string value that may be a single string or string array. */
|
||||
function asString(val: unknown): string {
|
||||
if (Array.isArray(val)) return String(val[0] ?? "");
|
||||
return String(val ?? "");
|
||||
}
|
||||
|
||||
export async function handleConnectVoice(req: Request, res: Response) {
|
||||
const guildId = Array.isArray(req.body.guildId)
|
||||
? req.body.guildId[0]
|
||||
: req.body.guildId;
|
||||
const channelId = Array.isArray(req.body.channelId)
|
||||
? req.body.channelId[0]
|
||||
: req.body.channelId;
|
||||
const guildId = asString(req.body.guildId);
|
||||
const channelId = asString(req.body.channelId);
|
||||
if (!guildId || !channelId) {
|
||||
return res.status(400).json({
|
||||
error: "VALIDATION_ERROR",
|
||||
@@ -35,17 +37,13 @@ export async function handleDisconnectVoice(_req: Request, res: Response) {
|
||||
}
|
||||
|
||||
export async function handleGetVoiceChannels(req: Request, res: Response) {
|
||||
const guildId = Array.isArray(req.params.guildId)
|
||||
? req.params.guildId[0]
|
||||
: req.params.guildId;
|
||||
const guildId = asString(req.params.guildId);
|
||||
const channels = await getVoiceChannels(guildId);
|
||||
res.json(channels);
|
||||
}
|
||||
|
||||
export async function handleVoiceCommand(req: Request, res: Response) {
|
||||
const command = Array.isArray(req.body.command)
|
||||
? req.body.command[0]
|
||||
: req.body.command;
|
||||
const command = asString(req.body.command);
|
||||
|
||||
if (!command) {
|
||||
return res.status(400).json({
|
||||
@@ -55,7 +53,7 @@ export async function handleVoiceCommand(req: Request, res: Response) {
|
||||
}
|
||||
|
||||
try {
|
||||
await publishCommandNoReply(command as string);
|
||||
await publishCommandNoReply(command);
|
||||
res.json({ success: true, command });
|
||||
} catch (err) {
|
||||
res.status(500).json({
|
||||
|
||||
@@ -260,7 +260,7 @@ export async function writeRedisStatus(
|
||||
// Lifecycle
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export async function startRedisBridge(): Promise<void> {
|
||||
export async function startCommandBridge(): Promise<void> {
|
||||
if (!ensureRedisConfig()) {
|
||||
logger.info("Redis not configured, skipping command channel bridge");
|
||||
return;
|
||||
@@ -272,7 +272,7 @@ export async function startRedisBridge(): Promise<void> {
|
||||
logger.info("Redis command channel initialized");
|
||||
}
|
||||
|
||||
export async function stopRedisBridge(): Promise<void> {
|
||||
export async function stopCommandBridge(): Promise<void> {
|
||||
if (publisherClient) {
|
||||
await publisherClient.quit();
|
||||
publisherClient = null;
|
||||
|
||||
@@ -18,9 +18,14 @@ export interface BroadcastFunctions {
|
||||
messageUpdated: BroadcastFn;
|
||||
messageDeleted: BroadcastFn;
|
||||
messageAnalyzed: BroadcastFn;
|
||||
attachmentCreated: BroadcastFn;
|
||||
attachmentUploaded: BroadcastFn;
|
||||
voiceRecordingStarted: BroadcastFn;
|
||||
voiceRecordingStopped: BroadcastFn;
|
||||
voiceRecordingUploaded: BroadcastFn;
|
||||
voicePcmData: BroadcastFn;
|
||||
voiceActiveUser: BroadcastFn;
|
||||
analysisQueueStatus: BroadcastFn;
|
||||
raw: BroadcastRawFn;
|
||||
binary: BroadcastBinaryFn;
|
||||
}
|
||||
@@ -53,18 +58,33 @@ export const broadcastMessageUpdated: BroadcastFn = (data) =>
|
||||
export const broadcastMessageDeleted: BroadcastFn = (data) =>
|
||||
(_fns?.messageDeleted ?? noop)(data);
|
||||
|
||||
export const broadcastAttachmentCreated: BroadcastFn = (data) =>
|
||||
(_fns?.attachmentCreated ?? noop)(data);
|
||||
|
||||
export const broadcastAttachmentUploaded: BroadcastFn = (data) =>
|
||||
(_fns?.attachmentUploaded ?? noop)(data);
|
||||
|
||||
export const broadcastMessageAnalyzed: BroadcastFn = (data) =>
|
||||
(_fns?.messageAnalyzed ?? noop)(data);
|
||||
|
||||
export const broadcastVoiceRecordingStarted: BroadcastFn = (data) =>
|
||||
(_fns?.voiceRecordingStarted ?? noop)(data);
|
||||
|
||||
export const broadcastVoiceRecordingStopped: BroadcastFn = (data) =>
|
||||
(_fns?.voiceRecordingStopped ?? noop)(data);
|
||||
|
||||
export const broadcastVoiceRecordingUploaded: BroadcastFn = (data) =>
|
||||
(_fns?.voiceRecordingUploaded ?? noop)(data);
|
||||
|
||||
export const broadcastVoicePcmData: BroadcastFn = (data) =>
|
||||
(_fns?.voicePcmData ?? noop)(data);
|
||||
|
||||
export const broadcastVoiceActiveUser: BroadcastFn = (data) =>
|
||||
(_fns?.voiceActiveUser ?? noop)(data);
|
||||
|
||||
export const broadcastAnalysisQueueStatus: BroadcastFn = (data) =>
|
||||
(_fns?.analysisQueueStatus ?? noop)(data);
|
||||
|
||||
export const broadcastRaw: BroadcastRawFn = (type, data) =>
|
||||
(_fns?.raw ?? noopRaw)(type, data);
|
||||
|
||||
|
||||
@@ -11,6 +11,9 @@ interface BroadcastEvent {
|
||||
timestamp: string;
|
||||
}
|
||||
|
||||
// Track the active WebSocket server for lifecycle management
|
||||
let _wss: WebSocketServer | null = null;
|
||||
|
||||
async function sendInitialStates(ws: WebSocket): Promise<void> {
|
||||
// Send initial user state
|
||||
ws.send(
|
||||
@@ -51,10 +54,18 @@ async function sendInitialStates(ws: WebSocket): Promise<void> {
|
||||
}
|
||||
}
|
||||
|
||||
export function closeWebSocketServer(): void {
|
||||
if (!_wss) return;
|
||||
logger.info("Closing WebSocket server");
|
||||
_wss.close(() => logger.info("WebSocket server closed"));
|
||||
_wss = null;
|
||||
}
|
||||
|
||||
export function createWebSocketServer(server: Server): WebSocketServer {
|
||||
const clients = new Set<WebSocket>();
|
||||
|
||||
const wss = new WebSocketServer({ server, path: "/ws" });
|
||||
_wss = wss;
|
||||
|
||||
wss.on("connection", (ws: WebSocket) => {
|
||||
clients.add(ws);
|
||||
@@ -188,12 +199,22 @@ export function createWebSocketServer(server: Server): WebSocketServer {
|
||||
broadcast({ type: "message_deleted", data }),
|
||||
messageAnalyzed: (data: unknown) =>
|
||||
broadcast({ type: "message_analyzed", data }),
|
||||
attachmentCreated: (data: unknown) =>
|
||||
broadcast({ type: "attachment_created", data }),
|
||||
attachmentUploaded: (data: unknown) =>
|
||||
broadcast({ type: "attachment_uploaded", data }),
|
||||
voiceRecordingStarted: (data: unknown) =>
|
||||
broadcast({ type: "voice_recording_started", data }),
|
||||
voiceRecordingStopped: (data: unknown) =>
|
||||
broadcast({ type: "voice_recording_stopped", data }),
|
||||
voiceRecordingUploaded: (data: unknown) =>
|
||||
broadcast({ type: "voice_recording_uploaded", data }),
|
||||
voicePcmData: (data: unknown) =>
|
||||
broadcast({ type: "voice_pcm_data", data }),
|
||||
voiceActiveUser: (data: unknown) =>
|
||||
broadcast({ type: "voice_active_user", data }),
|
||||
analysisQueueStatus: (data: unknown) =>
|
||||
broadcast({ type: "analysis_queue_status", data }),
|
||||
raw: (type: string, data: unknown) => broadcast({ type, data }),
|
||||
binary: broadcastBinary,
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user