| |
| import { createServer } from "node:http"; |
| import { describe, expect, it, vi } from "vitest"; |
| import { resolveGatewayRuntimeConfig } from "../../gateway/server-runtime-config.js"; |
| import { GatewayLockError } from "../../infra/gateway-lock.js"; |
| import { StateDatabaseCoordinatorContentionError } from "../../infra/state-database-coordinator.js"; |
| import { TailscaleRouteOwnershipConflictError } from "../../infra/tailscale-route-ownership-error.js"; |
| import { OpenClawAgentDatabaseMediaMigrationRequiredError } from "../../state/openclaw-agent-db-migration-required.js"; |
| import { testing } from "./run.test-support.js"; |
|
|
| const inspectGatewayTlsCertificateMock = vi.hoisted(() => |
| vi.fn<typeof import("../../infra/tls/gateway.js").inspectGatewayTlsCertificate>(async () => ({ |
| ok: false, |
| error: "public certificate missing", |
| })), |
| ); |
|
|
| vi.mock("../../infra/tls/gateway.js", () => ({ |
| inspectGatewayTlsCertificate: inspectGatewayTlsCertificateMock, |
| })); |
|
|
| function createLogger() { |
| return { |
| info: vi.fn(), |
| warn: vi.fn(), |
| }; |
| } |
|
|
| describe("supervised gateway lock recovery", () => { |
| it("retries lifecycle contention without treating a healthy port as ownership", async () => { |
| const error = new GatewayLockError( |
| "failed to acquire gateway state ownership", |
| new StateDatabaseCoordinatorContentionError("gateway-lifecycle"), |
| ); |
| const startLoop = vi |
| .fn<() => Promise<void>>() |
| .mockRejectedValueOnce(error) |
| .mockResolvedValueOnce(); |
| const probeHealth = vi.fn(async () => true); |
| let elapsedMs = 0; |
| await testing.runGatewayLoopWithSupervisedLockRecovery({ |
| startLoop, |
| supervisor: "systemd", |
| port: 18789, |
| healthHost: "127.0.0.1", |
| log: createLogger(), |
| probeHealth, |
| now: () => elapsedMs, |
| sleep: async (ms) => { |
| elapsedMs += ms; |
| }, |
| }); |
| expect(startLoop).toHaveBeenCalledTimes(2); |
| expect(probeHealth).not.toHaveBeenCalled(); |
| }); |
|
|
| it("spends one budget across supervised retries and lifecycle acquisition", async () => { |
| let elapsedMs = 0; |
| const error = new GatewayLockError( |
| "failed to acquire gateway state ownership; waited 295000ms for gateway-lifecycle ownership", |
| new StateDatabaseCoordinatorContentionError("gateway-lifecycle"), |
| ); |
| const budgets: Array<number | undefined> = []; |
| const startLoop = vi.fn(async (deadlineMs?: number) => { |
| budgets.push(deadlineMs === undefined ? undefined : deadlineMs - elapsedMs); |
| if (budgets.length > 1) { |
| elapsedMs = deadlineMs ?? elapsedMs; |
| } |
| throw error; |
| }); |
| const sleep = vi.fn(async (ms: number) => { |
| elapsedMs += ms; |
| }); |
| await expect( |
| testing.runGatewayLoopWithSupervisedLockRecovery({ |
| startLoop, |
| supervisor: "systemd", |
| port: 18789, |
| healthHost: "127.0.0.1", |
| log: createLogger(), |
| now: () => elapsedMs, |
| sleep, |
| }), |
| ).rejects.toBe(error); |
| expect(budgets).toEqual([300_000, 295_000]); |
| expect(elapsedMs).toBe(300_000); |
| expect(sleep).toHaveBeenCalledTimes(1); |
| }); |
|
|
| it("uses exit 78 for an ambiguous persistent Tailscale route", () => { |
| expect( |
| testing.resolveGatewayStartupFailureExitCode(new TailscaleRouteOwnershipConflictError()), |
| ).toBe(78); |
| }); |
|
|
| it("uses exit 78 for an effective bind/Tailscale config conflict", async () => { |
| |
| |
| |
| |
| const conflict = await resolveGatewayRuntimeConfig({ |
| cfg: { |
| gateway: { |
| bind: "lan", |
| auth: { mode: "token", token: "test-token-123" }, |
| tailscale: { mode: "serve" }, |
| }, |
| }, |
| port: 18789, |
| }).catch((err: unknown) => err); |
|
|
| expect(conflict).toBeInstanceOf(Error); |
| expect((conflict as Error).message).toBe( |
| "tailscale serve/funnel requires gateway bind=loopback (127.0.0.1)", |
| ); |
| expect(testing.resolveGatewayStartupFailureExitCode(conflict)).toBe(78); |
| }); |
|
|
| it("uses exit 78 for offline agent database migration requirements", () => { |
| expect( |
| testing.resolveGatewayStartupFailureExitCode( |
| new OpenClawAgentDatabaseMediaMigrationRequiredError("/tmp/openclaw-agent.sqlite", 14), |
| ), |
| ).toBe(78); |
| }); |
|
|
| it("does not retry gateway lock errors outside a supervisor", async () => { |
| const err = new GatewayLockError("gateway already running"); |
| const startLoop = vi.fn(async () => { |
| throw err; |
| }); |
|
|
| await expect( |
| testing.runGatewayLoopWithSupervisedLockRecovery({ |
| startLoop, |
| supervisor: null, |
| port: 18789, |
| healthHost: "127.0.0.1", |
| log: createLogger(), |
| }), |
| ).rejects.toBe(err); |
|
|
| expect(startLoop).toHaveBeenCalledTimes(1); |
| }); |
|
|
| it("leaves a healthy launchd-supervised gateway in control", async () => { |
| const startLoop = vi.fn(async () => { |
| throw new GatewayLockError("gateway already running"); |
| }); |
| const probeHealth = vi.fn(async () => true); |
| const log = createLogger(); |
|
|
| await testing.runGatewayLoopWithSupervisedLockRecovery({ |
| startLoop, |
| supervisor: "launchd", |
| port: 18789, |
| healthHost: "0.0.0.0", |
| log, |
| probeHealth, |
| }); |
|
|
| expect(startLoop).toHaveBeenCalledTimes(1); |
| expect(probeHealth).toHaveBeenCalledWith({ host: "0.0.0.0", port: 18789 }); |
| expect(log.info).toHaveBeenCalledWith( |
| "gateway already running under launchd; existing gateway is healthy, leaving it in control", |
| ); |
| expect(log.warn).not.toHaveBeenCalled(); |
| }); |
|
|
| it("uses exit 78 semantics for healthy systemd-supervised lock conflicts", async () => { |
| const startLoop = vi.fn(async () => { |
| throw new GatewayLockError("another gateway instance is already listening"); |
| }); |
| const probeHealth = vi.fn(async () => true); |
|
|
| let failure: unknown; |
| try { |
| await testing.runGatewayLoopWithSupervisedLockRecovery({ |
| startLoop, |
| supervisor: "systemd", |
| port: 18789, |
| healthHost: "127.0.0.1", |
| log: createLogger(), |
| probeHealth, |
| }); |
| } catch (err) { |
| failure = err; |
| } |
|
|
| expect(failure).toMatchObject({ |
| message: expect.stringContaining( |
| "exiting with code 78 to prevent a systemd Restart=always loop", |
| ), |
| }); |
| expect(startLoop).toHaveBeenCalledTimes(1); |
| expect(probeHealth).toHaveBeenCalledWith({ host: "127.0.0.1", port: 18789 }); |
| expect(testing.resolveGatewayLockErrorExitCode(failure)).toBe(78); |
| }); |
|
|
| it("preserves an agent-embedded owner error under a supervisor", async () => { |
| const err = new GatewayLockError( |
| "another embedded OpenClaw state writer is active (pid 123); lock timeout after 5000ms", |
| ); |
| const startLoop = vi.fn(async () => { |
| throw err; |
| }); |
| const probeHealth = vi.fn(async () => true); |
|
|
| await expect( |
| testing.runGatewayLoopWithSupervisedLockRecovery({ |
| startLoop, |
| supervisor: "systemd", |
| port: 18789, |
| healthHost: "127.0.0.1", |
| log: createLogger(), |
| probeHealth, |
| }), |
| ).rejects.toBe(err); |
|
|
| expect(startLoop).toHaveBeenCalledTimes(1); |
| expect(probeHealth).not.toHaveBeenCalled(); |
| }); |
|
|
| it("bounds supervised retries when the existing gateway stays unhealthy", async () => { |
| let now = 0; |
| const startLoop = vi.fn(async () => { |
| throw new GatewayLockError("gateway already running"); |
| }); |
| const sleep = vi.fn(async (ms: number) => { |
| now += ms; |
| }); |
|
|
| let failure: unknown; |
| try { |
| await testing.runGatewayLoopWithSupervisedLockRecovery({ |
| startLoop, |
| supervisor: "systemd", |
| port: 18789, |
| healthHost: "127.0.0.1", |
| log: createLogger(), |
| probeHealth: vi.fn(async () => false), |
| now: () => now, |
| sleep, |
| retryMs: 5, |
| timeoutMs: 12, |
| }); |
| } catch (err) { |
| failure = err; |
| } |
|
|
| expect(failure).toMatchObject({ |
| message: |
| "gateway already running under systemd; existing gateway did not become healthy after 12ms", |
| }); |
| expect(testing.resolveGatewayLockErrorExitCode(failure)).toBe(1); |
| expect(startLoop).toHaveBeenCalledTimes(4); |
| expect(sleep).toHaveBeenNthCalledWith(1, 5); |
| expect(sleep).toHaveBeenNthCalledWith(2, 5); |
| expect(sleep).toHaveBeenNthCalledWith(3, 2); |
| }); |
|
|
| it("bounds supervised retries for EADDRINUSE lock errors", async () => { |
| let now = 0; |
| const startLoop = vi.fn(async () => { |
| throw new GatewayLockError( |
| "another gateway instance is already listening on ws://127.0.0.1:18789", |
| ); |
| }); |
| const sleep = vi.fn(async (ms: number) => { |
| now += ms; |
| }); |
|
|
| await expect( |
| testing.runGatewayLoopWithSupervisedLockRecovery({ |
| startLoop, |
| supervisor: "systemd", |
| port: 18789, |
| healthHost: "127.0.0.1", |
| log: createLogger(), |
| probeHealth: vi.fn(async () => false), |
| now: () => now, |
| sleep, |
| retryMs: 5, |
| timeoutMs: 12, |
| }), |
| ).rejects.toThrow( |
| "gateway already running under systemd; existing gateway did not become healthy after 12ms", |
| ); |
|
|
| expect(startLoop).toHaveBeenCalledTimes(4); |
| expect(sleep).toHaveBeenNthCalledWith(1, 5); |
| expect(sleep).toHaveBeenNthCalledWith(2, 5); |
| expect(sleep).toHaveBeenNthCalledWith(3, 2); |
| }); |
|
|
| it.each(["gateway already running", "another gateway instance is already listening"])( |
| "uses exit 1 for unmanaged lock errors: %s", |
| (message) => { |
| expect(testing.resolveGatewayLockErrorExitCode(new GatewayLockError(message))).toBe(1); |
| }, |
| ); |
|
|
| it("retries public certificate inspection while TLS material is unavailable", async () => { |
| inspectGatewayTlsCertificateMock.mockClear(); |
| const probeHealth = testing.createConfiguredGatewayHealthProbe({ |
| gateway: { tls: { enabled: true, autoGenerate: true } }, |
| }); |
|
|
| await expect(probeHealth({ host: "127.0.0.1", port: 18789 })).resolves.toBe(false); |
| await expect(probeHealth({ host: "127.0.0.1", port: 18789 })).resolves.toBe(false); |
|
|
| expect(inspectGatewayTlsCertificateMock).toHaveBeenCalledTimes(2); |
| expect(inspectGatewayTlsCertificateMock).toHaveBeenCalledWith({ |
| enabled: true, |
| autoGenerate: true, |
| }); |
| }); |
|
|
| it("recognizes only the OpenClaw health response", () => { |
| expect( |
| testing.isGatewayHealthzResponse(200, JSON.stringify({ ok: true, status: "live" })), |
| ).toBe(true); |
| expect( |
| testing.isGatewayHealthzResponse(200, JSON.stringify({ ok: true, status: "ready" })), |
| ).toBe(false); |
| expect(testing.isGatewayHealthzResponse(404, "not found")).toBe(false); |
| expect(testing.isGatewayHealthzResponse(200, "not json")).toBe(false); |
| }); |
|
|
| it("bounds slow health responses with an absolute deadline", async () => { |
| const server = createServer((_req, res) => { |
| res.writeHead(200, { "content-type": "application/json" }); |
| const interval = setInterval(() => { |
| res.write(" "); |
| }, 10); |
| res.once("close", () => clearInterval(interval)); |
| }); |
| await new Promise<void>((resolve) => { |
| server.listen(0, "127.0.0.1", () => resolve()); |
| }); |
|
|
| try { |
| const address = server.address(); |
| if (!address || typeof address === "string") { |
| throw new Error("expected TCP server address"); |
| } |
| const startedAt = Date.now(); |
| await expect( |
| testing.probeGatewayHealthz({ |
| host: "127.0.0.1", |
| port: address.port, |
| timeoutMs: 50, |
| }), |
| ).resolves.toBe(false); |
| expect(Date.now() - startedAt).toBeLessThan(500); |
| } finally { |
| await new Promise<void>((resolve, reject) => { |
| server.close((err) => (err ? reject(err) : resolve())); |
| }); |
| } |
| }); |
|
|
| it("normalizes wildcard bind hosts for local health probes", () => { |
| expect(testing.normalizeGatewayHealthProbeHost("0.0.0.0")).toBe("127.0.0.1"); |
| expect(testing.normalizeGatewayHealthProbeHost("::")).toBe("127.0.0.1"); |
| expect(testing.normalizeGatewayHealthProbeHost("127.0.0.1")).toBe("127.0.0.1"); |
| }); |
| }); |
|
|