Skip to content
Open
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
7 changes: 6 additions & 1 deletion apps/server/src/graceful-shutdown.integration.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,11 @@
import { describe, expect, it } from "bun:test";
import { REPO_ROOT } from "./test-helpers/workspace";

const REPO_ROOT = "/Users/brain/Coding/snipeship/ccflare";
// REPO_ROOT was previously hardcoded to the original author's local machine path
// ("/Users/brain/Coding/snipeship/ccflare"), which does not exist on any other checkout and made
// this test non-portable (Bun.spawn's cwd pointing at a nonexistent directory fails with a
// misleading ENOENT attributed to the spawned binary, not the cwd). Reusing the existing portable
// helper (computed via import.meta.dir, not hardcoded) fixes this for any checkout location.

describe("graceful shutdown integration", () => {
it("awaits programmatic stop until pending request writes are flushed", async () => {
Expand Down
175 changes: 175 additions & 0 deletions packages/proxy/src/usage-worker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -287,6 +287,181 @@ describe("UsageWorkerController", () => {
expect(controller.getHealthSnapshot().state).toBe("stopped");
});

it("flushes a message queued just before shutdown once the worker becomes ready during the wait, instead of dropping it", async () => {
const workers: FakeWorker[] = [];
const logger = new TestLogger();
const controller = new UsageWorkerController({
createWorker() {
const worker = new FakeWorker();
workers.push(worker);
return worker;
},
readyTimeoutMs: 200,
ackTimeoutMs: 10_000,
shutdownDelayMs: 0,
logger,
});
controllers.push(controller);

// Worker is not ready yet, so this queues instead of sending immediately --
// this is the exact scenario that used to lose the message on shutdown.
controller.postMessage(createStartMessage());
expect(workers[0]?.postedMessages).toHaveLength(0);

const shutdownPromise = controller.terminateGracefully();

// Ready arrives DURING the shutdown wait window (before readyTimeoutMs elapses).
await Bun.sleep(5);
workers[0]?.emitMessage({ type: "ready" });

await waitFor(
() =>
workers[0]?.postedMessages.some(
(message) => (message as { type?: string }).type === "start",
) ?? false,
);

workers[0]?.emitMessage({
type: "shutdown-complete",
asyncWriter: { healthy: true, failureCount: 0, queuedJobs: 0 },
});

await shutdownPromise;
expect(
logger.warnings.some((warning) => warning.includes("Dropping")),
).toBe(false);
});

it("resolves shutdown within readyTimeoutMs + shutdownDelayMs (never hangs) and drops queued messages with a warning when the worker never becomes ready", async () => {
const workers: FakeWorker[] = [];
const logger = new TestLogger();
const readyTimeoutMs = 30;
const shutdownDelayMs = 20;
const controller = new UsageWorkerController({
createWorker() {
const worker = new FakeWorker();
workers.push(worker);
return worker;
},
readyTimeoutMs,
ackTimeoutMs: 10_000,
shutdownDelayMs,
logger,
});
controllers.push(controller);

// Worker never emits "ready" in this test.
controller.postMessage(createStartMessage());
expect(workers[0]?.postedMessages).toHaveLength(0);

const startedAt = Date.now();
await expect(controller.terminateGracefully()).rejects.toThrow();
const elapsedMs = Date.now() - startedAt;

// Generous slack for test-runner scheduling jitter -- the point is this is
// BOUNDED (readyTimeoutMs + shutdownDelayMs), not that it hangs indefinitely.
expect(elapsedMs).toBeLessThan(readyTimeoutMs + shutdownDelayMs + 1_000);
expect(
logger.warnings.some((warning) =>
warning.includes("Dropping 1 queued usage worker message"),
),
).toBe(true);
expect(workers[0]?.terminateCalls).toBe(1);
});

it("force termination flushes a queued message only if the worker is already ready, without waiting", async () => {
const workers: FakeWorker[] = [];
const logger = new TestLogger();
const controller = new UsageWorkerController({
createWorker() {
const worker = new FakeWorker();
workers.push(worker);
return worker;
},
readyTimeoutMs: 10_000,
ackTimeoutMs: 10_000,
shutdownDelayMs: 0,
logger,
});
controllers.push(controller);

workers[0]?.emitMessage({ type: "ready" });
controller.postMessage(createStartMessage());
expect(workers[0]?.postedMessages).toHaveLength(1);

controller.forceTerminate();

expect(workers[0]?.terminateCalls).toBe(1);
expect(
logger.warnings.some((warning) => warning.includes("Dropping")),
).toBe(false);
});

it("promptly rejects an in-flight terminateGracefully() wait when forceTerminate() races it, instead of idling out the full readyTimeoutMs", async () => {
const workers: FakeWorker[] = [];
const logger = new TestLogger();
const readyTimeoutMs = 300;
const controller = new UsageWorkerController({
createWorker() {
const worker = new FakeWorker();
workers.push(worker);
return worker;
},
readyTimeoutMs,
ackTimeoutMs: 10_000,
shutdownDelayMs: 0,
logger,
});
controllers.push(controller);

// Worker never becomes ready, so terminateGracefully() below enters its
// wait-for-ready phase and suspends for up to readyTimeoutMs.
controller.postMessage(createStartMessage());
const gracefulShutdown = controller.terminateGracefully();

// Race it with a forceTerminate() call (mirrors getUsageWorker()'s self-heal
// path in proxy.ts, which calls forceTerminate() on an instance for which
// isShuttingDown() is already true).
await Bun.sleep(5);
const startedAt = Date.now();
controller.forceTerminate();

await expect(gracefulShutdown).rejects.toThrow(
"Usage worker was force terminated",
);
const elapsedMs = Date.now() - startedAt;

// Must settle promptly (driven by forceTerminate unblocking the wait), not by
// idling out the full readyTimeoutMs.
expect(elapsedMs).toBeLessThan(readyTimeoutMs / 2);
});

it("does not let a standalone forceTerminate() poison a later, unrelated terminateGracefully() call on the same instance", async () => {
const workers: FakeWorker[] = [];
const logger = new TestLogger();
const controller = new UsageWorkerController({
createWorker() {
const worker = new FakeWorker();
workers.push(worker);
return worker;
},
readyTimeoutMs: 10_000,
ackTimeoutMs: 10_000,
shutdownDelayMs: 0,
logger,
});
controllers.push(controller);

// Force-terminate with nothing racing it -- this used to leave a stale
// earlyTerminationError sitting on the instance.
controller.forceTerminate();

// A later, wholly independent terminateGracefully() call must resolve
// normally (there is no worker left and nothing queued), not reject with the
// stale "Usage worker was force terminated" error from the earlier call.
await expect(controller.terminateGracefully()).resolves.toBeUndefined();
});

it("rejects graceful shutdown when the worker never confirms completion", async () => {
const workers: FakeWorker[] = [];
const controller = new UsageWorkerController({
Expand Down
Loading