From 679672ad15a1fce6adcc664854e19fb4104bf9f2 Mon Sep 17 00:00:00 2001 From: Pete Hampton Date: Wed, 20 May 2026 18:33:53 +0100 Subject: [PATCH] =?UTF-8?q?=F0=9F=AA=82=20feat:=20Graceful=20HTTP=20shutdo?= =?UTF-8?q?wn=20on=20SIGTERM/SIGINT=20(#13211)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * 🪂 feat: Graceful HTTP shutdown on SIGTERM/SIGINT * Address feedback * don't treat ERR_SERVER_NOT_RUNNING as fatal; route telemetry shutdown through coordinator --- api/cache/getLogStores.js | 22 +-- api/server/index.js | 5 +- packages/api/src/app/index.ts | 1 + packages/api/src/app/shutdown.spec.ts | 257 +++++++++++++++++++++++++ packages/api/src/app/shutdown.ts | 114 +++++++++++ packages/api/src/flow/manager.ts | 14 +- packages/api/src/telemetry/sdk.spec.ts | 67 ++----- packages/api/src/telemetry/sdk.ts | 56 ++---- 8 files changed, 428 insertions(+), 108 deletions(-) create mode 100644 packages/api/src/app/shutdown.spec.ts create mode 100644 packages/api/src/app/shutdown.ts diff --git a/api/cache/getLogStores.js b/api/cache/getLogStores.js index 70eb681e53..45a6a69994 100644 --- a/api/cache/getLogStores.js +++ b/api/cache/getLogStores.js @@ -7,6 +7,7 @@ const { sessionCache, standardCache, violationCache, + registerShutdownTask, } = require('@librechat/api'); const namespaces = { @@ -195,23 +196,16 @@ if (!cacheConfig.USE_REDIS && !cacheConfig.CI) { cleanupIntervals.add(monitor); } - const dispose = () => { + // Register cleanup with the centralized graceful-shutdown coordinator + // (see packages/api/src/app/shutdown.ts) rather than attaching a direct + // signal handler — multiple competing handlers race the HTTP drain. + registerShutdownTask('cache cleanup', async () => { cacheConfig.DEBUG_MEMORY_CACHE && console.log('[Cache] Cleaning up and shutting down...'); cleanupIntervals.forEach((interval) => clearInterval(interval)); cleanupIntervals.clear(); - - // One final cleanup before exit - clearAllExpiredFromCache().then(() => { - cacheConfig.DEBUG_MEMORY_CACHE && console.log('[Cache] Final cleanup completed'); - process.exit(0); - }); - }; - - // Handle various termination signals - process.on('SIGTERM', dispose); - process.on('SIGINT', dispose); - process.on('SIGQUIT', dispose); - process.on('SIGHUP', dispose); + await clearAllExpiredFromCache(); + cacheConfig.DEBUG_MEMORY_CACHE && console.log('[Cache] Final cleanup completed'); + }); } /** diff --git a/api/server/index.js b/api/server/index.js index 60b1a96b3d..e7ce041fa4 100644 --- a/api/server/index.js +++ b/api/server/index.js @@ -22,6 +22,7 @@ const { createStreamServices, initializeFileStorage, preAuthTenantMiddleware, + setupGracefulShutdown, updateInterfacePermissions, } = require('@librechat/api'); const { connectDb, indexSync } = require('~/db'); @@ -240,7 +241,7 @@ const startServer = async () => { /** Error handler (must be last - Express identifies error middleware by its 4-arg signature) */ app.use(ErrorController); - app.listen(port, host, async (err) => { + const server = app.listen(port, host, async (err) => { if (err) { logger.error('Failed to start server:', err); process.exit(1); @@ -283,6 +284,8 @@ const startServer = async () => { process.exit(1); } }); + + setupGracefulShutdown(server); }; /** diff --git a/packages/api/src/app/index.ts b/packages/api/src/app/index.ts index 3bb619ab13..a2eef88ae1 100644 --- a/packages/api/src/app/index.ts +++ b/packages/api/src/app/index.ts @@ -5,3 +5,4 @@ export * from './permissions'; export * from './cdn'; export * from './checks'; export * from './resolve'; +export * from './shutdown'; diff --git a/packages/api/src/app/shutdown.spec.ts b/packages/api/src/app/shutdown.spec.ts new file mode 100644 index 0000000000..baabd8aa8b --- /dev/null +++ b/packages/api/src/app/shutdown.spec.ts @@ -0,0 +1,257 @@ +/// +import http from 'http'; +import { + setupGracefulShutdown, + registerShutdownTask, + __resetShutdownStateForTests, +} from './shutdown'; + +jest.mock('@librechat/data-schemas', () => ({ + logger: { info: jest.fn(), warn: jest.fn(), error: jest.fn() }, +})); + +const triggerSignal = (signal: NodeJS.Signals): void => { + const listeners = process.listeners(signal) as NodeJS.SignalsListener[]; + listeners.forEach((listener) => listener(signal)); +}; + +const flush = (): Promise => new Promise((resolve) => setImmediate(resolve)); + +describe('setupGracefulShutdown', () => { + let server: http.Server; + let exitSpy: jest.SpyInstance; + let originalSigterm: NodeJS.SignalsListener[]; + let originalSigint: NodeJS.SignalsListener[]; + let originalSigquit: NodeJS.SignalsListener[]; + let originalSighup: NodeJS.SignalsListener[]; + + beforeEach(() => { + server = http.createServer(); + // Most tests exercise the close path, so we fake `listening` to true. + // Tests that want to exercise the not-listening short-circuit override + // this back to false. + Object.defineProperty(server, 'listening', { value: true, configurable: true }); + exitSpy = jest.spyOn(process, 'exit').mockImplementation((() => undefined) as never); + originalSigterm = process.listeners('SIGTERM').slice() as NodeJS.SignalsListener[]; + originalSigint = process.listeners('SIGINT').slice() as NodeJS.SignalsListener[]; + originalSigquit = process.listeners('SIGQUIT').slice() as NodeJS.SignalsListener[]; + originalSighup = process.listeners('SIGHUP').slice() as NodeJS.SignalsListener[]; + process.removeAllListeners('SIGTERM'); + process.removeAllListeners('SIGINT'); + process.removeAllListeners('SIGQUIT'); + process.removeAllListeners('SIGHUP'); + __resetShutdownStateForTests(); + }); + + afterEach(() => { + process.removeAllListeners('SIGTERM'); + process.removeAllListeners('SIGINT'); + process.removeAllListeners('SIGQUIT'); + process.removeAllListeners('SIGHUP'); + originalSigterm.forEach((listener) => process.on('SIGTERM', listener)); + originalSigint.forEach((listener) => process.on('SIGINT', listener)); + originalSigquit.forEach((listener) => process.on('SIGQUIT', listener)); + originalSighup.forEach((listener) => process.on('SIGHUP', listener)); + exitSpy.mockRestore(); + if (server.listening) { + server.close(); + } + __resetShutdownStateForTests(); + jest.useRealTimers(); + }); + + it('registers handlers for SIGTERM, SIGINT, SIGQUIT, and SIGHUP', () => { + for (const signal of ['SIGTERM', 'SIGINT', 'SIGQUIT', 'SIGHUP'] as NodeJS.Signals[]) { + expect(process.listenerCount(signal)).toBe(0); + } + setupGracefulShutdown(server); + for (const signal of ['SIGTERM', 'SIGINT', 'SIGQUIT', 'SIGHUP'] as NodeJS.Signals[]) { + expect(process.listenerCount(signal)).toBe(1); + } + }); + + it('closes the server and exits 0 on SIGTERM', async () => { + const closeSpy = jest.spyOn(server, 'close').mockImplementation((cb?: (err?: Error) => void) => { + if (cb) { + setImmediate(() => cb()); + } + return server; + }); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + await flush(); + await flush(); + expect(closeSpy).toHaveBeenCalled(); + expect(exitSpy).toHaveBeenCalledWith(0); + }); + + it('closes the server and exits 0 on SIGINT', async () => { + const closeSpy = jest.spyOn(server, 'close').mockImplementation((cb?: (err?: Error) => void) => { + if (cb) { + setImmediate(() => cb()); + } + return server; + }); + setupGracefulShutdown(server); + triggerSignal('SIGINT'); + await flush(); + await flush(); + expect(closeSpy).toHaveBeenCalled(); + expect(exitSpy).toHaveBeenCalledWith(0); + }); + + it('exits 1 if server.close yields an error', async () => { + const closeErr = new Error('close failed'); + jest.spyOn(server, 'close').mockImplementation((cb?: (err?: Error) => void) => { + if (cb) { + setImmediate(() => cb(closeErr)); + } + return server; + }); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + await flush(); + await flush(); + expect(exitSpy).toHaveBeenCalledWith(1); + }); + + it('treats ERR_SERVER_NOT_RUNNING from close() as a successful shutdown', async () => { + const notRunning: NodeJS.ErrnoException = Object.assign(new Error('Server is not running.'), { + code: 'ERR_SERVER_NOT_RUNNING', + }); + // Force listening=true so the pre-check doesn't short-circuit; + // we want to exercise the post-check that ignores ERR_SERVER_NOT_RUNNING. + Object.defineProperty(server, 'listening', { value: true, configurable: true }); + jest.spyOn(server, 'close').mockImplementation((cb?: (err?: Error) => void) => { + if (cb) { + setImmediate(() => cb(notRunning)); + } + return server; + }); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + await flush(); + await flush(); + expect(exitSpy).toHaveBeenCalledWith(0); + }); + + it('skips close() entirely when the server is not listening yet', async () => { + // SIGTERM during the startup window (before app.listen finishes binding + // the socket) must not be treated as a failed shutdown. + Object.defineProperty(server, 'listening', { value: false, configurable: true }); + const closeSpy = jest.spyOn(server, 'close'); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + await flush(); + expect(closeSpy).not.toHaveBeenCalled(); + expect(exitSpy).toHaveBeenCalledWith(0); + }); + + it('is idempotent — a second signal does not trigger another shutdown', async () => { + const closeSpy = jest.spyOn(server, 'close'); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + triggerSignal('SIGTERM'); + await flush(); + expect(closeSpy).toHaveBeenCalledTimes(1); + }); + + it('force-exits with code 1 if shutdown exceeds the timeout', () => { + jest.useFakeTimers(); + jest.spyOn(server, 'close').mockImplementation(() => server); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + jest.advanceTimersByTime(60_000); + expect(exitSpy).toHaveBeenCalledWith(1); + }); + + it('runs registered tasks after server.close and before exit', async () => { + const calls: string[] = []; + jest.spyOn(server, 'close').mockImplementation((cb?: (err?: Error) => void) => { + calls.push('server.close'); + if (cb) { + setImmediate(() => cb()); + } + return server; + }); + registerShutdownTask('task-a', () => { + calls.push('task-a'); + }); + registerShutdownTask('task-b', () => { + calls.push('task-b'); + }); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + await flush(); + await flush(); + expect(calls).toEqual(['server.close', 'task-a', 'task-b']); + expect(exitSpy).toHaveBeenCalledWith(0); + }); + + it('runs tasks in registration order', async () => { + const order: string[] = []; + jest.spyOn(server, 'close').mockImplementation((cb?: (err?: Error) => void) => { + if (cb) { + setImmediate(() => cb()); + } + return server; + }); + ['first', 'second', 'third'].forEach((name) => { + registerShutdownTask(name, () => { + order.push(name); + }); + }); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + await flush(); + await flush(); + expect(order).toEqual(['first', 'second', 'third']); + }); + + it('continues subsequent tasks and still exits if one task throws', async () => { + const calls: string[] = []; + jest.spyOn(server, 'close').mockImplementation((cb?: (err?: Error) => void) => { + if (cb) { + setImmediate(() => cb()); + } + return server; + }); + registerShutdownTask('ok-before', () => { + calls.push('ok-before'); + }); + registerShutdownTask('throws', () => { + calls.push('throws'); + throw new Error('boom'); + }); + registerShutdownTask('ok-after', () => { + calls.push('ok-after'); + }); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + await flush(); + await flush(); + expect(calls).toEqual(['ok-before', 'throws', 'ok-after']); + expect(exitSpy).toHaveBeenCalledWith(0); + }); + + it('awaits async tasks before exiting', async () => { + const calls: string[] = []; + jest.spyOn(server, 'close').mockImplementation((cb?: (err?: Error) => void) => { + if (cb) { + setImmediate(() => cb()); + } + return server; + }); + registerShutdownTask('async-task', async () => { + await new Promise((resolve) => setImmediate(resolve)); + calls.push('async-done'); + }); + setupGracefulShutdown(server); + triggerSignal('SIGTERM'); + await flush(); + await flush(); + await flush(); + expect(calls).toEqual(['async-done']); + expect(exitSpy).toHaveBeenCalledWith(0); + }); +}); diff --git a/packages/api/src/app/shutdown.ts b/packages/api/src/app/shutdown.ts new file mode 100644 index 0000000000..3fc0a9d105 --- /dev/null +++ b/packages/api/src/app/shutdown.ts @@ -0,0 +1,114 @@ +import { logger } from '@librechat/data-schemas'; +import type { Server } from 'http'; + +const SHUTDOWN_TIMEOUT_MS = 60_000; +const SIGNALS: NodeJS.Signals[] = ['SIGTERM', 'SIGINT', 'SIGQUIT', 'SIGHUP']; + +type ShutdownTask = { + name: string; + fn: () => void | Promise; +}; + +const tasks: ShutdownTask[] = []; +let isShuttingDown = false; +let httpServer: Server | null = null; + +/** + * Register a cleanup task to run after the HTTP server has closed. + * Tasks run in registration order; if one throws, subsequent tasks + * and the final exit are not blocked. Use this instead of attaching + * `process.on('SIGTERM', ...)` handlers directly — multiple competing + * signal handlers race with the HTTP drain because Node dispatches + * listeners in registration order and any one of them can call + * `process.exit` before the HTTP server has finished closing. + */ +export function registerShutdownTask( + name: string, + fn: () => void | Promise, +): void { + tasks.push({ name, fn }); +} + +/** + * Wires SIGTERM, SIGINT, SIGQUIT, and SIGHUP to a graceful shutdown + * sequence: close the HTTP server (stop accepting new connections, let + * in-flight requests finish), run any tasks registered via + * `registerShutdownTask`, then `process.exit(0)`. After + * SHUTDOWN_TIMEOUT_MS the process is force-exited with code 1 — a + * safety net for long-lived connections such as SSE streams that may + * not finish in time. + */ +export function setupGracefulShutdown(server: Server): void { + httpServer = server; + for (const signal of SIGNALS) { + process.on(signal, () => { + void shutdown(signal); + }); + } +} + +/** + * @internal Reset module state for tests. Not part of the public API. + */ +export function __resetShutdownStateForTests(): void { + tasks.length = 0; + isShuttingDown = false; + httpServer = null; +} + +async function shutdown(signal: NodeJS.Signals): Promise { + if (isShuttingDown) { + return; + } + isShuttingDown = true; + logger.info(`Received ${signal}, draining HTTP server...`); + + const forceExit = setTimeout(() => { + logger.warn(`Graceful shutdown exceeded ${SHUTDOWN_TIMEOUT_MS}ms, forcing exit`); + process.exit(1); + }, SHUTDOWN_TIMEOUT_MS); + forceExit.unref(); + + let exitCode = 0; + + try { + await closeHttpServer(); + } catch (err) { + logger.error('Error closing HTTP server during graceful shutdown:', err); + exitCode = 1; + } + + for (const task of tasks) { + try { + logger.info(`Running shutdown task: ${task.name}`); + await task.fn(); + } catch (err) { + logger.error(`Shutdown task "${task.name}" failed:`, err); + } + } + + clearTimeout(forceExit); + logger.info('Graceful shutdown complete, exiting'); + process.exit(exitCode); +} + +function closeHttpServer(): Promise { + return new Promise((resolve, reject) => { + if (!httpServer || !httpServer.listening) { + // SIGTERM can arrive during startup before the listen socket is open, + // in which case there is nothing to drain. Node also surfaces this as + // an ERR_SERVER_NOT_RUNNING error in the close callback — treated + // below as a successful close so a routine shutdown doesn't trip + // orchestrator restart/backoff with exit code 1. + resolve(); + return; + } + httpServer.close((err) => { + if (!err || (err as NodeJS.ErrnoException).code === 'ERR_SERVER_NOT_RUNNING') { + resolve(); + return; + } + reject(err); + }); + }); +} diff --git a/packages/api/src/flow/manager.ts b/packages/api/src/flow/manager.ts index 544cba9560..f775222282 100644 --- a/packages/api/src/flow/manager.ts +++ b/packages/api/src/flow/manager.ts @@ -2,6 +2,7 @@ import { Keyv } from 'keyv'; import { logger } from '@librechat/data-schemas'; import type { StoredDataNoRaw } from 'keyv'; import type { FlowState, FlowMetadata, FlowManagerOptions } from './types'; +import { registerShutdownTask } from '../app/shutdown'; export const PENDING_STALE_MS = 2 * 60 * 1000; @@ -40,17 +41,14 @@ export class FlowStateManager { } private setupCleanupHandlers() { - const cleanup = () => { + // Register cleanup with the centralized graceful-shutdown coordinator + // (see ../app/shutdown.ts) rather than attaching direct signal + // handlers — multiple competing handlers race the HTTP drain. + registerShutdownTask('flow manager cleanup', () => { logger.info('Cleaning up FlowStateManager intervals...'); this.intervals.forEach((interval) => clearInterval(interval)); this.intervals.clear(); - process.exit(0); - }; - - process.on('SIGTERM', cleanup); - process.on('SIGINT', cleanup); - process.on('SIGQUIT', cleanup); - process.on('SIGHUP', cleanup); + }); } /** diff --git a/packages/api/src/telemetry/sdk.spec.ts b/packages/api/src/telemetry/sdk.spec.ts index 951b1d91ef..32b33cb566 100644 --- a/packages/api/src/telemetry/sdk.spec.ts +++ b/packages/api/src/telemetry/sdk.spec.ts @@ -107,9 +107,9 @@ jest.mock( { virtual: true }, ); -async function flushSignalShutdown(): Promise { - await new Promise((resolve) => setImmediate(resolve)); -} +jest.mock('../app/shutdown', () => ({ + registerShutdownTask: jest.fn(), +})); describe('telemetry SDK lifecycle', () => { let emitWarningSpy: jest.SpyInstance; @@ -494,56 +494,27 @@ describe('telemetry SDK lifecycle', () => { expect(controller.status).toBe('stopped'); }); - it.each(['SIGTERM', 'SIGINT'])( - 'does not force process exit when another %s handler is registered', - async (signal) => { - const otherHandler = jest.fn(); - const killSpy = jest.spyOn(process, 'kill').mockImplementation(() => true); - initializeTelemetry({ OTEL_TRACING_ENABLED: 'true' }); - process.once(signal, otherHandler); - - process.emit(signal, signal); - await flushSignalShutdown(); - - expect(mockShutdown).toHaveBeenCalledTimes(1); - expect(otherHandler).toHaveBeenCalledTimes(1); - expect(killSpy).not.toHaveBeenCalled(); - - killSpy.mockRestore(); - }, - ); - - it.each(['SIGTERM', 'SIGINT'])( - 'reraises the shutdown %s signal when telemetry is the only signal handler', - async (signal) => { - const killSpy = jest.spyOn(process, 'kill').mockImplementation(() => true); - initializeTelemetry({ OTEL_TRACING_ENABLED: 'true' }); - - process.emit(signal, signal); - await flushSignalShutdown(); - - expect(mockShutdown).toHaveBeenCalledTimes(1); - expect(killSpy).toHaveBeenCalledWith(process.pid, signal); - - killSpy.mockRestore(); - }, - ); - - it('warns and reraises the signal when shutdown rejects', async () => { - mockShutdown.mockRejectedValueOnce(new Error('signal shutdown failed')); - const killSpy = jest.spyOn(process, 'kill').mockImplementation(() => true); + it('registers a shutdown task with the coordinator when initialized', async () => { + const { registerShutdownTask } = await import('../app/shutdown'); initializeTelemetry({ OTEL_TRACING_ENABLED: 'true' }); - process.emit('SIGTERM', 'SIGTERM'); - await flushSignalShutdown(); + expect(registerShutdownTask).toHaveBeenCalledWith('telemetry', expect.any(Function)); + }); + + it('the registered shutdown task warns when telemetry shutdown rejects', async () => { + const { registerShutdownTask } = await import('../app/shutdown'); + mockShutdown.mockRejectedValueOnce(new Error('flush failed')); + initializeTelemetry({ OTEL_TRACING_ENABLED: 'true' }); + + const taskFn = (registerShutdownTask as jest.Mock).mock.calls.at(-1)?.[1] as + | (() => Promise) + | undefined; + expect(taskFn).toBeDefined(); + await taskFn?.(); - expect(mockShutdown).toHaveBeenCalledTimes(1); expect(emitWarningSpy).toHaveBeenCalledWith( - 'OpenTelemetry shutdown failed: signal shutdown failed', + 'OpenTelemetry shutdown failed: flush failed', { code: 'LIBRECHAT_OTEL' }, ); - expect(killSpy).toHaveBeenCalledWith(process.pid, 'SIGTERM'); - - killSpy.mockRestore(); }); }); diff --git a/packages/api/src/telemetry/sdk.ts b/packages/api/src/telemetry/sdk.ts index 90c3b1977c..ea5e0a3de2 100644 --- a/packages/api/src/telemetry/sdk.ts +++ b/packages/api/src/telemetry/sdk.ts @@ -13,6 +13,7 @@ import type { Span, Attributes } from '@opentelemetry/api'; import type { RequestOptions } from 'node:http'; import type { TelemetryConfig, TelemetryStatus } from './config'; import { getTelemetryConfig } from './config'; +import { registerShutdownTask } from '../app/shutdown'; export interface TelemetryController { readonly enabled: boolean; @@ -24,11 +25,6 @@ const WARNING_CODE = 'LIBRECHAT_OTEL'; const REDACTED_QUERY_VALUE = '[REDACTED]'; const SIGNAL_SHUTDOWN_TIMEOUT_MS = 5_000; -interface RegisteredSignal { - signal: NodeJS.Signals; - listener: NodeJS.SignalsListener; -} - interface RequestUrlParts { href?: string; search?: string; @@ -49,7 +45,7 @@ let pendingSdk: NodeSDK | undefined; let startPromise: Promise | undefined; let shutdownPromise: Promise | undefined; let status: TelemetryStatus = 'stopped'; -let registeredSignals: RegisteredSignal[] = []; +let shutdownTaskRegistered = false; let requestSpans = new WeakMap(); function isBunRuntime(): boolean { @@ -370,35 +366,22 @@ function makeController(): TelemetryController { }; } -function unregisterShutdownHandlers(): void { - for (const { signal, listener } of registeredSignals) { - process.removeListener(signal, listener); - } - registeredSignals = []; -} - -function registerShutdownHandlers(): void { - if (registeredSignals.length > 0) { +function ensureShutdownTaskRegistered(): void { + if (shutdownTaskRegistered) { return; } - - const signals: NodeJS.Signals[] = ['SIGTERM', 'SIGINT']; - registeredSignals = signals.map((signal) => { - const listener: NodeJS.SignalsListener = () => { - const shouldReraiseSignal = process.listenerCount(signal) === 0; - withTimeout(shutdownTelemetry(), SIGNAL_SHUTDOWN_TIMEOUT_MS) - .catch((error) => { - emitWarning(`OpenTelemetry shutdown failed: ${getErrorMessage(error)}`); - }) - .finally(() => { - if (shouldReraiseSignal) { - process.kill(process.pid, signal); - } - }); - }; - process.once(signal, listener); - return { signal, listener }; - }); + shutdownTaskRegistered = true; + // Register with the centralized graceful-shutdown coordinator + // (see ../app/shutdown.ts) rather than attaching SIGTERM/SIGINT + // listeners directly — signal listener return values are ignored + // by Node, so a separate signal handler can let the coordinator + // exit before the async OpenTelemetry flush completes, dropping + // final spans during pod shutdowns. + registerShutdownTask('telemetry', () => + withTimeout(shutdownTelemetry(), SIGNAL_SHUTDOWN_TIMEOUT_MS).catch((error) => { + emitWarning(`OpenTelemetry shutdown failed: ${getErrorMessage(error)}`); + }), + ); } function withTimeout(promise: Promise, timeoutMs: number): Promise { @@ -440,7 +423,7 @@ export function initializeTelemetry(env: NodeJS.ProcessEnv = process.env): Telem pendingSdk = undefined; activeSdk = sdk; status = 'started'; - registerShutdownHandlers(); + ensureShutdownTaskRegistered(); } }) .catch((error) => { @@ -461,7 +444,7 @@ export function initializeTelemetry(env: NodeJS.ProcessEnv = process.env): Telem activeSdk = sdk; status = 'started'; - registerShutdownHandlers(); + ensureShutdownTaskRegistered(); return makeController(); } catch (error) { status = 'failed'; @@ -485,7 +468,6 @@ async function performShutdownTelemetry(): Promise { await sdk.shutdown(); activeSdk = undefined; status = 'stopped'; - unregisterShutdownHandlers(); } catch (error) { status = 'started'; throw error; @@ -520,6 +502,6 @@ export async function resetTelemetryForTests(): Promise { shutdownPromise = undefined; status = 'stopped'; requestSpans = new WeakMap(); - unregisterShutdownHandlers(); + shutdownTaskRegistered = false; } }