Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 13 additions & 5 deletions apps/hub/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -354,6 +354,7 @@ import { createBootAssetWiring, REGISTRIES } from "./asset-service-factory";
import { createRoutineScheduler } from "./routine-scheduler";
import { createToolGrantsForPins } from "./tool-grants";
import { createMcpCredentialBindingsFor } from "./mcp-credential-bindings";
import { shutdownHub } from "./shutdown";

// Host policy constants, not configuration.
const MAX_TARBALL_BYTES = 10 * 1024 * 1024;
Expand Down Expand Up @@ -3454,11 +3455,18 @@ if (import.meta.main) {
});
const log = getLogger(["hub"]);
log.info`Hub serving on port ${port}`;
const shutdown = async () => {
await server.stop();
await hub.close();
process.exit(0);
};
const SHUTDOWN_DRAIN_MS = 10_000;
// `server.stop()` waits for open connections and websockets by default,
// so it sits inside the same bound as the hub's own closes.
const shutdown = () =>
shutdownHub({
drain: async () => {
await server.stop();
await hub.close();
},
timeoutMs: SHUTDOWN_DRAIN_MS,
exit: (code) => process.exit(code),
});
process.on("SIGINT", () => void shutdown());
process.on("SIGTERM", () => void shutdown());
}
78 changes: 78 additions & 0 deletions apps/hub/src/shutdown.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
import { expect, test } from "bun:test";
import { drainWithTimeout, shutdownHub } from "./shutdown";

test("drainWithTimeout resolves drained when the drain completes inside the bound", async () => {
const outcome = await drainWithTimeout(() => Promise.resolve(), 1_000);
expect(outcome).toEqual({ kind: "drained" });
});

test("drainWithTimeout resolves timed-out when the drain outlives the bound", async () => {
const outcome = await drainWithTimeout(
() => new Promise<void>(() => undefined),
10,
);
expect(outcome).toEqual({ kind: "timed-out" });
});

test("drainWithTimeout resolves failed with the thrown error when the drain throws", async () => {
const error = new Error("drain fault");
const outcome = await drainWithTimeout(() => Promise.reject(error), 1_000);
expect(outcome).toEqual({ kind: "failed", error });
});

test("shutdownHub exits 0 when the drain resolves inside the bound", async () => {
let exitCode: number | undefined;
const reported: unknown[] = [];
await shutdownHub({
drain: () => Promise.resolve(),
timeoutMs: 1_000,
exit: (code) => {
exitCode = code;
},
report: (error) => {
reported.push(error);
return "unused-ref-id";
},
});
expect(exitCode).toBe(0);
expect(reported).toEqual([]);
});

test("shutdownHub exits non-zero and reports the cause when the drain never settles", async () => {
let exitCode: number | undefined;
const reported: unknown[] = [];
await shutdownHub({
drain: () => new Promise<void>(() => undefined),
timeoutMs: 10,
exit: (code) => {
exitCode = code;
},
report: (error) => {
reported.push(error);
return "unused-ref-id";
},
});
expect(exitCode).toBe(1);
expect(reported).toHaveLength(1);
expect(reported[0]).toBeInstanceOf(Error);
expect((reported[0] as Error).message).toContain("exceeded 10ms");
});

test("shutdownHub exits non-zero and reports the cause when the drain rejects", async () => {
let exitCode: number | undefined;
const reported: unknown[] = [];
const error = new Error("close fault");
await shutdownHub({
drain: () => Promise.reject(error),
timeoutMs: 1_000,
exit: (code) => {
exitCode = code;
},
report: (reportedError) => {
reported.push(reportedError);
return "unused-ref-id";
},
});
expect(exitCode).toBe(1);
expect(reported).toEqual([error]);
});
71 changes: 71 additions & 0 deletions apps/hub/src/shutdown.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
// Bounded drain for process shutdown, mirroring apps/sidecar/src/shutdown.ts's
// shape. Kept as a local copy rather than a shared import or package: it's a
// ~15-line primitive, and the two apps already diverge on exit semantics
// (the sidecar treats a timeout as a clean-enough exit; the hub treats a
// timeout as a fault worth reporting), so sharing it would need a parameter
// immediately.
//
// The platform sends SIGTERM and expects a prompt exit; a drain that hangs
// would turn every deploy into an apparent crash, so the bound cuts it off
// and reports the cause instead of leaving an unhandled rejection.

import { reportError } from "@corbits/error-sink";

export type DrainOutcome =
| { kind: "drained" }
| { kind: "timed-out" }
| { kind: "failed"; error: unknown };

export async function drainWithTimeout(
drain: () => Promise<void>,
timeoutMs: number,
): Promise<DrainOutcome> {
let timer: ReturnType<typeof setTimeout> | undefined;
const timedOut = new Promise<DrainOutcome>((resolve) => {
timer = setTimeout(() => {
resolve({ kind: "timed-out" });
}, timeoutMs);
});
const drained = (async (): Promise<DrainOutcome> => {
try {
await drain();
return { kind: "drained" };
} catch (error) {
return { kind: "failed", error };
}
})();
const outcome = await Promise.race([drained, timedOut]);
clearTimeout(timer);
return outcome;
}

export type ShutdownHubDeps = {
drain: () => Promise<void>;
timeoutMs: number;
exit: (code: number) => void;
report?: typeof reportError;
};

/**
* Drains the hub within `timeoutMs` and always exits: 0 on a clean drain,
* non-zero with the cause reported through `reportError` on a throw or a
* timeout.
*/
export async function shutdownHub({
drain,
timeoutMs,
exit,
report = reportError,
}: ShutdownHubDeps): Promise<void> {
const outcome = await drainWithTimeout(drain, timeoutMs);
if (outcome.kind === "drained") {
exit(0);
return;
}
const error =
outcome.kind === "timed-out"
? new Error(`Hub shutdown drain exceeded ${timeoutMs}ms`)
: outcome.error;
report(error, { operation: "hub.shutdown" });
exit(1);
}
Loading