mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-19 00:52:10 -06:00
310 lines
10 KiB
TypeScript
310 lines
10 KiB
TypeScript
import type { Server as HttpServer } from "node:http";
|
|
import type { WebSocketServer } from "ws";
|
|
import { disposeRegisteredAgentHarnesses } from "../agents/harness/registry.js";
|
|
import { disposeAllSessionMcpRuntimes } from "../agents/pi-bundle-mcp-tools.js";
|
|
import type { CanvasHostHandler, CanvasHostServer } from "../canvas-host/server.js";
|
|
import { type ChannelId, listChannelPlugins } from "../channels/plugins/index.js";
|
|
import { stopGmailWatcher } from "../hooks/gmail-watcher.js";
|
|
import type { HeartbeatRunner } from "../infra/heartbeat-runner.js";
|
|
import { createSubsystemLogger } from "../logging/subsystem.js";
|
|
import type { PluginServicesHandle } from "../plugins/services.js";
|
|
import { normalizeOptionalString } from "../shared/string-coerce.js";
|
|
|
|
const shutdownLog = createSubsystemLogger("gateway/shutdown");
|
|
const WEBSOCKET_CLOSE_GRACE_MS = 1_000;
|
|
const WEBSOCKET_CLOSE_FORCE_CONTINUE_MS = 250;
|
|
const HTTP_CLOSE_GRACE_MS = 1_000;
|
|
const HTTP_CLOSE_FORCE_WAIT_MS = 5_000;
|
|
const MCP_RUNTIME_CLOSE_GRACE_MS = 5_000;
|
|
|
|
function createTimeoutRace<T>(timeoutMs: number, onTimeout: () => T) {
|
|
let timer: ReturnType<typeof setTimeout> | null = null;
|
|
timer = setTimeout(() => {
|
|
if (timer) {
|
|
clearTimeout(timer);
|
|
timer = null;
|
|
}
|
|
resolve(onTimeout());
|
|
}, timeoutMs);
|
|
timer.unref?.();
|
|
|
|
let resolve!: (value: T) => void;
|
|
const promise = new Promise<T>((innerResolve) => {
|
|
resolve = innerResolve;
|
|
});
|
|
|
|
return {
|
|
promise,
|
|
clear() {
|
|
if (timer) {
|
|
clearTimeout(timer);
|
|
timer = null;
|
|
}
|
|
},
|
|
};
|
|
}
|
|
|
|
export async function runGatewayClosePrelude(params: {
|
|
stopDiagnostics?: () => void;
|
|
clearSkillsRefreshTimer?: () => void;
|
|
skillsChangeUnsub?: () => void;
|
|
disposeAuthRateLimiter?: () => void;
|
|
disposeBrowserAuthRateLimiter: () => void;
|
|
stopModelPricingRefresh?: () => void;
|
|
stopChannelHealthMonitor?: () => void;
|
|
clearSecretsRuntimeSnapshot?: () => void;
|
|
closeMcpServer?: () => Promise<void>;
|
|
}): Promise<void> {
|
|
params.stopDiagnostics?.();
|
|
params.clearSkillsRefreshTimer?.();
|
|
params.skillsChangeUnsub?.();
|
|
params.disposeAuthRateLimiter?.();
|
|
params.disposeBrowserAuthRateLimiter();
|
|
params.stopModelPricingRefresh?.();
|
|
params.stopChannelHealthMonitor?.();
|
|
params.clearSecretsRuntimeSnapshot?.();
|
|
await params.closeMcpServer?.().catch(() => {});
|
|
}
|
|
|
|
function isServerNotRunningError(err: unknown): boolean {
|
|
return Boolean(
|
|
err &&
|
|
typeof err === "object" &&
|
|
"code" in err &&
|
|
(err as { code?: unknown }).code === "ERR_SERVER_NOT_RUNNING",
|
|
);
|
|
}
|
|
|
|
export function createGatewayCloseHandler(params: {
|
|
bonjourStop: (() => Promise<void>) | null;
|
|
tailscaleCleanup: (() => Promise<void>) | null;
|
|
canvasHost: CanvasHostHandler | null;
|
|
canvasHostServer: CanvasHostServer | null;
|
|
releasePluginRouteRegistry?: (() => void) | null;
|
|
stopChannel: (name: ChannelId, accountId?: string) => Promise<void>;
|
|
pluginServices: PluginServicesHandle | null;
|
|
disposeSessionMcpRuntimes?: () => Promise<void>;
|
|
cron: { stop: () => void };
|
|
heartbeatRunner: HeartbeatRunner;
|
|
updateCheckStop?: (() => void) | null;
|
|
stopTaskRegistryMaintenance?: (() => void) | null;
|
|
nodePresenceTimers: Map<string, ReturnType<typeof setInterval>>;
|
|
broadcast: (event: string, payload: unknown, opts?: { dropIfSlow?: boolean }) => void;
|
|
tickInterval: ReturnType<typeof setInterval>;
|
|
healthInterval: ReturnType<typeof setInterval>;
|
|
dedupeCleanup: ReturnType<typeof setInterval>;
|
|
mediaCleanup: ReturnType<typeof setInterval> | null;
|
|
agentUnsub: (() => void) | null;
|
|
heartbeatUnsub: (() => void) | null;
|
|
transcriptUnsub: (() => void) | null;
|
|
lifecycleUnsub: (() => void) | null;
|
|
chatRunState: { clear: () => void };
|
|
clients: Set<{ socket: { close: (code: number, reason: string) => void } }>;
|
|
configReloader: { stop: () => Promise<void> };
|
|
wss: WebSocketServer;
|
|
httpServer: HttpServer;
|
|
httpServers?: HttpServer[];
|
|
}) {
|
|
return async (opts?: { reason?: string; restartExpectedMs?: number | null }) => {
|
|
try {
|
|
const reasonRaw = normalizeOptionalString(opts?.reason) ?? "";
|
|
const reason = reasonRaw || "gateway stopping";
|
|
const restartExpectedMs =
|
|
typeof opts?.restartExpectedMs === "number" && Number.isFinite(opts.restartExpectedMs)
|
|
? Math.max(0, Math.floor(opts.restartExpectedMs))
|
|
: null;
|
|
if (params.bonjourStop) {
|
|
try {
|
|
await params.bonjourStop();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
if (params.tailscaleCleanup) {
|
|
await params.tailscaleCleanup();
|
|
}
|
|
if (params.canvasHost) {
|
|
try {
|
|
await params.canvasHost.close();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
if (params.canvasHostServer) {
|
|
try {
|
|
await params.canvasHostServer.close();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
for (const plugin of listChannelPlugins()) {
|
|
await params.stopChannel(plugin.id);
|
|
}
|
|
await disposeRegisteredAgentHarnesses();
|
|
const disposeMcpRuntimes = params.disposeSessionMcpRuntimes ?? disposeAllSessionMcpRuntimes;
|
|
const mcpDisposePromise = disposeMcpRuntimes().catch((err: unknown) => {
|
|
shutdownLog.warn(`bundle-mcp runtime disposal failed during shutdown: ${String(err)}`);
|
|
});
|
|
const mcpDisposeTimeout = createTimeoutRace(MCP_RUNTIME_CLOSE_GRACE_MS, () => {
|
|
shutdownLog.warn(
|
|
`bundle-mcp runtime disposal exceeded ${MCP_RUNTIME_CLOSE_GRACE_MS}ms; continuing shutdown`,
|
|
);
|
|
});
|
|
await Promise.race([mcpDisposePromise, mcpDisposeTimeout.promise]);
|
|
mcpDisposeTimeout.clear();
|
|
if (params.pluginServices) {
|
|
await params.pluginServices.stop().catch(() => {});
|
|
}
|
|
await stopGmailWatcher();
|
|
params.cron.stop();
|
|
params.heartbeatRunner.stop();
|
|
try {
|
|
params.stopTaskRegistryMaintenance?.();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
try {
|
|
params.updateCheckStop?.();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
for (const timer of params.nodePresenceTimers.values()) {
|
|
clearInterval(timer);
|
|
}
|
|
params.nodePresenceTimers.clear();
|
|
params.broadcast("shutdown", {
|
|
reason,
|
|
restartExpectedMs,
|
|
});
|
|
clearInterval(params.tickInterval);
|
|
clearInterval(params.healthInterval);
|
|
clearInterval(params.dedupeCleanup);
|
|
if (params.mediaCleanup) {
|
|
clearInterval(params.mediaCleanup);
|
|
}
|
|
if (params.agentUnsub) {
|
|
try {
|
|
params.agentUnsub();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
if (params.heartbeatUnsub) {
|
|
try {
|
|
params.heartbeatUnsub();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
if (params.transcriptUnsub) {
|
|
try {
|
|
params.transcriptUnsub();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
if (params.lifecycleUnsub) {
|
|
try {
|
|
params.lifecycleUnsub();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
params.chatRunState.clear();
|
|
for (const c of params.clients) {
|
|
try {
|
|
c.socket.close(1012, "service restart");
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
params.clients.clear();
|
|
await params.configReloader.stop().catch(() => {});
|
|
const wsClients = params.wss.clients ?? new Set();
|
|
const closePromise = new Promise<void>((resolve) => params.wss.close(() => resolve()));
|
|
const websocketGraceTimeout = createTimeoutRace(
|
|
WEBSOCKET_CLOSE_GRACE_MS,
|
|
() => false as const,
|
|
);
|
|
const closedWithinGrace = await Promise.race([
|
|
closePromise.then(() => true),
|
|
websocketGraceTimeout.promise,
|
|
]);
|
|
websocketGraceTimeout.clear();
|
|
if (!closedWithinGrace) {
|
|
shutdownLog.warn(
|
|
`websocket server close exceeded ${WEBSOCKET_CLOSE_GRACE_MS}ms; forcing shutdown continuation with ${wsClients.size} tracked client(s)`,
|
|
);
|
|
for (const client of wsClients) {
|
|
try {
|
|
client.terminate();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
const websocketForceTimeout = createTimeoutRace(WEBSOCKET_CLOSE_FORCE_CONTINUE_MS, () => {
|
|
shutdownLog.warn(
|
|
`websocket server close still pending after ${WEBSOCKET_CLOSE_FORCE_CONTINUE_MS}ms force window; continuing shutdown`,
|
|
);
|
|
});
|
|
await Promise.race([closePromise, websocketForceTimeout.promise]);
|
|
websocketForceTimeout.clear();
|
|
}
|
|
const servers =
|
|
params.httpServers && params.httpServers.length > 0
|
|
? params.httpServers
|
|
: [params.httpServer];
|
|
for (const server of servers) {
|
|
const httpServer = server as HttpServer & {
|
|
closeAllConnections?: () => void;
|
|
closeIdleConnections?: () => void;
|
|
};
|
|
if (typeof httpServer.closeIdleConnections === "function") {
|
|
httpServer.closeIdleConnections();
|
|
}
|
|
const closePromise = new Promise<void>((resolve, reject) =>
|
|
httpServer.close((err) => {
|
|
if (!err || isServerNotRunningError(err)) {
|
|
resolve();
|
|
return;
|
|
}
|
|
reject(err);
|
|
}),
|
|
);
|
|
const httpGraceTimeout = createTimeoutRace(HTTP_CLOSE_GRACE_MS, () => false as const);
|
|
const closedWithinGrace = await Promise.race([
|
|
closePromise.then(() => true),
|
|
httpGraceTimeout.promise,
|
|
]);
|
|
httpGraceTimeout.clear();
|
|
if (!closedWithinGrace) {
|
|
shutdownLog.warn(
|
|
`http server close exceeded ${HTTP_CLOSE_GRACE_MS}ms; forcing connection shutdown and waiting for close`,
|
|
);
|
|
httpServer.closeAllConnections?.();
|
|
const httpForceTimeout = createTimeoutRace(
|
|
HTTP_CLOSE_FORCE_WAIT_MS,
|
|
() => false as const,
|
|
);
|
|
const closedAfterForce = await Promise.race([
|
|
closePromise.then(() => true),
|
|
httpForceTimeout.promise,
|
|
]);
|
|
httpForceTimeout.clear();
|
|
if (!closedAfterForce) {
|
|
throw new Error(
|
|
`http server close still pending after forced connection shutdown (${HTTP_CLOSE_FORCE_WAIT_MS}ms)`,
|
|
);
|
|
}
|
|
}
|
|
}
|
|
} finally {
|
|
try {
|
|
params.releasePluginRouteRegistry?.();
|
|
} catch {
|
|
/* ignore */
|
|
}
|
|
}
|
|
};
|
|
}
|