openclaw / src /process /supervisor /supervisor.test.ts
SaylorTwift's picture
SaylorTwift HF Staff
Add files using upload-large-folder tool
5cb63c1 verified
Raw History Blame Contribute Delete
37.5 kB
// Process supervisor tests cover lifecycle, restart, and termination behavior.
import { performance } from "node:perf_hooks";
import { expectDefined } from "@openclaw/normalization-core";
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../../test/helpers/promise.js";
import { mockProcessPlatform } from "../../test-utils/vitest-spies.js";
import {
createSilentIdleArgv,
createStubChildAdapter,
createWriteStdoutArgv,
spawnChild,
type StubChildAdapter,
} from "./supervisor.test-support.js";
import type { ManagedRun } from "./types.js";
const { createChildAdapterMock, createPtyAdapterMock } = vi.hoisted(() => ({
createChildAdapterMock: vi.fn(),
createPtyAdapterMock: vi.fn(),
}));
vi.mock("./adapters/child.js", () => ({
createChildAdapter: async (
...args: Parameters<typeof import("./adapters/child.js").createChildAdapter>
) => ({
adapter: await createChildAdapterMock(...args),
ready: Promise.resolve(),
}),
}));
vi.mock("./adapters/pty.js", () => ({
createPtyAdapter: createPtyAdapterMock,
}));
let createProcessSupervisor: typeof import("./supervisor.js").createProcessSupervisor;
describe("process supervisor", () => {
beforeAll(async () => {
vi.resetModules();
({ createProcessSupervisor } = await import("./supervisor.js"));
});
beforeEach(() => {
createChildAdapterMock.mockReset();
createPtyAdapterMock.mockReset();
vi.useRealTimers();
});
afterEach(() => {
vi.useRealTimers();
vi.restoreAllMocks();
});
it("passes private secret input and exact environment to the child adapter", async () => {
const adapter = createStubChildAdapter();
createChildAdapterMock.mockResolvedValue(adapter);
const secretInput = { fd: 3, createData: () => Buffer.from("secret") };
const supervisor = createProcessSupervisor();
const run = await spawnChild(supervisor, {
argv: createWriteStdoutArgv("ok"),
exactEnv: true,
secretInput,
});
adapter.settle(0);
await run.wait();
expect(createChildAdapterMock).toHaveBeenCalledWith(
expect.objectContaining({ exactEnv: true, secretInput }),
);
});
it("enforces no-output timeout for silent processes", async () => {
vi.useFakeTimers();
const adapter = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGKILL");
},
});
createChildAdapterMock.mockResolvedValue(adapter);
const supervisor = createProcessSupervisor();
const run = await spawnChild(supervisor, {
argv: createSilentIdleArgv(),
timeoutMs: 300,
noOutputTimeoutMs: 5,
stdinMode: "pipe-closed",
});
const exitPromise = run.wait();
await vi.advanceTimersByTimeAsync(5);
const exit = await exitPromise;
const expectedTimeoutSignal = process.platform === "win32" ? "SIGKILL" : "SIGTERM";
expect(adapter.killMock).toHaveBeenCalledWith(expectedTimeoutSignal);
if (process.platform !== "win32") {
await vi.advanceTimersByTimeAsync(5_000);
expect(adapter.killMock).not.toHaveBeenCalledWith("SIGKILL");
}
expect(exit.reason).toBe("no-output-timeout");
expect(exit.noOutputTimedOut).toBe(true);
expect(exit.timedOut).toBe(true);
});
it("coalesces overlapping Windows deadline cancellation while hard kill is pending", async () => {
vi.useFakeTimers();
mockProcessPlatform("win32");
const adapter = createStubChildAdapter();
createChildAdapterMock.mockResolvedValue(adapter);
const supervisor = createProcessSupervisor();
const run = await spawnChild(supervisor, {
argv: createSilentIdleArgv(),
timeoutMs: 20,
noOutputTimeoutMs: 5,
stdinMode: "pipe-closed",
});
const exitPromise = run.wait();
await vi.advanceTimersByTimeAsync(5);
expect(adapter.killMock).toHaveBeenCalledTimes(1);
expect(adapter.killMock).toHaveBeenCalledWith("SIGKILL");
await vi.advanceTimersByTimeAsync(15);
expect(adapter.killMock).toHaveBeenCalledTimes(1);
adapter.settle(null, "SIGKILL");
const exit = await exitPromise;
expect(exit.reason).toBe("no-output-timeout");
});
it("escalates cancellation to SIGKILL when graceful shutdown does not settle", async () => {
vi.useFakeTimers();
const adapter = createStubChildAdapter({
onKill: (signal, current) => {
if (signal === "SIGKILL") {
current.settle(null, signal);
}
},
});
createChildAdapterMock.mockResolvedValue(adapter);
const supervisor = createProcessSupervisor();
const run = await spawnChild(supervisor, {
argv: createSilentIdleArgv(),
timeoutMs: 1_000,
stdinMode: "pipe-closed",
});
const exitPromise = run.wait();
run.cancel("manual-cancel");
expect(adapter.killMock).toHaveBeenCalledWith("SIGTERM");
expect(adapter.killMock).not.toHaveBeenCalledWith("SIGKILL");
await vi.advanceTimersByTimeAsync(4_999);
expect(adapter.killMock).not.toHaveBeenCalledWith("SIGKILL");
await vi.advanceTimersByTimeAsync(1);
const exit = await exitPromise;
expect(adapter.killMock).toHaveBeenCalledWith("SIGKILL");
expect(exit.reason).toBe("manual-cancel");
expect(exit.exitSignal).toBe("SIGKILL");
});
it.each(["child", "pty"] as const)(
"cancels a %s process by run ID while its adapter is starting",
async (mode) => {
const adapter = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
});
const startup = createDeferred<StubChildAdapter>();
if (mode === "pty") {
createPtyAdapterMock.mockReturnValueOnce(startup.promise);
} else {
createChildAdapterMock.mockReturnValueOnce(startup.promise);
}
const supervisor = createProcessSupervisor();
const runId = `cancel-starting-${mode}`;
const pendingRun =
mode === "pty"
? supervisor.spawn({
runId,
mode: "pty",
argv: ["/bin/sh", "-c", "printf cancelled"],
scopeKey: "scope:cancel-starting",
})
: spawnChild(supervisor, {
runId,
scopeKey: "scope:cancel-starting",
argv: createSilentIdleArgv(),
});
expect(mode === "pty" ? createPtyAdapterMock : createChildAdapterMock).toHaveBeenCalledOnce();
supervisor.cancel(runId, "manual-cancel");
expect(adapter.killMock).not.toHaveBeenCalled();
startup.resolve(adapter);
const run = await pendingRun;
expect(adapter.killMock).toHaveBeenCalledWith("SIGKILL");
await run.waitForExtinction?.();
expect(adapter.disposeMock).toHaveBeenCalled();
await expect(run.wait()).resolves.toMatchObject({ reason: "manual-cancel" });
expect(run.activity.resultSettled).toBe(true);
},
);
it.each([
{
timeoutField: "timeoutMs" as const,
reason: "overall-timeout" as const,
},
{
timeoutField: "noOutputTimeoutMs" as const,
reason: "no-output-timeout" as const,
},
])("bounds a hung child adapter construction with $reason", async ({ timeoutField, reason }) => {
vi.useFakeTimers();
const startup = createDeferred<StubChildAdapter>();
createChildAdapterMock.mockReturnValueOnce(startup.promise);
const supervisor = createProcessSupervisor();
const runId = `hung-adapter-${reason}`;
const pendingRun = spawnChild(supervisor, {
runId,
argv: createSilentIdleArgv(),
[timeoutField]: 25,
stdinMode: "pipe-closed",
});
expect(createChildAdapterMock).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(25);
const constructionState = await Promise.race([
pendingRun.then(() => "settled" as const),
Promise.resolve().then(() => "pending" as const),
]);
expect(constructionState).toBe("settled");
const run = await pendingRun;
await expect(run.wait()).resolves.toMatchObject({
reason,
timedOut: true,
noOutputTimedOut: reason === "no-output-timeout",
});
expect(run.activity.resultSettled).toBe(true);
const killed = createDeferred();
const lateAdapter = createStubChildAdapter({ onKill: () => killed.resolve() });
startup.resolve(lateAdapter);
await killed.promise;
expect(lateAdapter.killMock).toHaveBeenCalledWith("SIGKILL");
expect(lateAdapter.disposeMock).not.toHaveBeenCalled();
lateAdapter.settle(null, "SIGKILL");
await run.waitForExtinction?.();
expect(lateAdapter.disposeMock).toHaveBeenCalled();
});
it("fences new runs and drains an unscoped startup during shutdown", async () => {
const adapter = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
});
const startup = createDeferred<StubChildAdapter>();
createChildAdapterMock.mockReturnValueOnce(startup.promise);
const supervisor = createProcessSupervisor();
const pendingRun = spawnChild(supervisor, {
runId: "shutdown-starting",
argv: createSilentIdleArgv(),
});
const shutdown = supervisor.shutdown();
await expect(
spawnChild(supervisor, {
runId: "shutdown-late",
argv: createSilentIdleArgv(),
}),
).rejects.toThrow("process supervisor is shut down");
expect(createChildAdapterMock).toHaveBeenCalledTimes(1);
startup.resolve(adapter);
const run = await pendingRun;
await expect(shutdown).resolves.toBeUndefined();
expect(adapter.killMock).toHaveBeenCalledWith("SIGKILL");
expect(adapter.disposeMock).toHaveBeenCalled();
await expect(run.wait()).resolves.toMatchObject({ reason: "manual-cancel" });
});
it("keeps shutdown fenced when live ownership extinction fails", async () => {
const extinction = createDeferred();
const adapter = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
});
adapter.waitForExtinction = () => extinction.promise;
createChildAdapterMock.mockResolvedValueOnce(adapter);
const supervisor = createProcessSupervisor();
await spawnChild(supervisor, {
runId: "shutdown-failed-extinction",
argv: createSilentIdleArgv(),
});
const shutdown = supervisor.shutdown();
extinction.reject(new Error("owner extinction failed"));
await expect(shutdown).rejects.toThrow("owner extinction failed");
await expect(
spawnChild(supervisor, {
runId: "shutdown-after-failed-extinction",
argv: createSilentIdleArgv(),
}),
).rejects.toThrow("process supervisor is shut down");
expect(createChildAdapterMock).toHaveBeenCalledOnce();
});
it("cancels every starting scoped process without canceling a later arrival", async () => {
const runCount = 16;
const startups = Array.from({ length: runCount }, () => createDeferred<StubChildAdapter>());
const adapters = Array.from({ length: runCount }, (_unused, index) =>
createStubChildAdapter({
pid: 7_000 + index,
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
}),
);
const laterAdapter = createStubChildAdapter({ pid: 8_000 });
let adapterIndex = 0;
createChildAdapterMock.mockImplementation(() => {
const index = adapterIndex++;
if (index === runCount) {
return Promise.resolve(laterAdapter);
}
const startup = startups[index];
if (!startup) {
throw new Error(`unexpected scope cancellation startup ${index}`);
}
return startup.promise;
});
const supervisor = createProcessSupervisor();
const pendingRuns = Array.from({ length: runCount }, (_unused, index) =>
spawnChild(supervisor, {
runId: `cancel-scope-starting-${index}`,
scopeKey: "scope:cancel-every-start",
argv: createSilentIdleArgv(),
}),
);
expect(createChildAdapterMock).toHaveBeenCalledTimes(runCount);
supervisor.cancelScope("scope:cancel-every-start", "manual-cancel");
for (const adapter of adapters) {
expect(adapter.killMock).not.toHaveBeenCalled();
}
const laterRun = await spawnChild(supervisor, {
runId: "cancel-scope-later-arrival",
scopeKey: "scope:cancel-every-start",
argv: createSilentIdleArgv(),
});
expect(laterAdapter.killMock).not.toHaveBeenCalled();
for (let index = runCount - 1; index >= 0; index -= 1) {
const startup = startups[index];
const adapter = adapters[index];
if (!startup || !adapter) {
throw new Error(`missing scope cancellation startup ${index}`);
}
startup.resolve(adapter);
}
const runs = await Promise.all(pendingRuns);
await Promise.all(
runs.map((run) => expectDefined(run.waitForExtinction, "cancelled construction cleanup")()),
);
for (const adapter of adapters) {
expect(adapter.killMock).toHaveBeenCalledWith("SIGKILL");
expect(adapter.disposeMock).toHaveBeenCalled();
}
await expect(Promise.all(runs.map((run) => run.wait()))).resolves.toEqual(
Array.from({ length: runCount }, () => expect.objectContaining({ reason: "manual-cancel" })),
);
expect(laterAdapter.killMock).not.toHaveBeenCalled();
laterAdapter.settle(0);
await expect(laterRun.wait()).resolves.toMatchObject({ reason: "exit" });
});
it("cancels a replacement that is fenced behind an earlier startup", async () => {
const first = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
});
const later = createStubChildAdapter();
const firstStartup = createDeferred<StubChildAdapter>();
createChildAdapterMock.mockReturnValueOnce(firstStartup.promise).mockResolvedValueOnce(later);
const supervisor = createProcessSupervisor();
const firstRunPromise = spawnChild(supervisor, {
runId: "cancel-fenced-first",
scopeKey: "scope:cancel-fenced",
argv: createSilentIdleArgv(),
});
let replacementCurrent = true;
const replacementPromise = spawnChild(supervisor, {
runId: "cancel-fenced-replacement",
scopeKey: "scope:cancel-fenced",
replaceExistingScope: true,
argv: createSilentIdleArgv(),
onCancel: () => {
replacementCurrent = false;
},
});
expect(createChildAdapterMock).toHaveBeenCalledTimes(1);
supervisor.cancelScope("scope:cancel-fenced", "manual-cancel");
expect(replacementCurrent).toBe(false);
const laterPromise = spawnChild(supervisor, {
runId: "cancel-fenced-later",
scopeKey: "scope:cancel-fenced",
argv: createSilentIdleArgv(),
});
firstStartup.resolve(first);
const [firstRun, replacementRun, laterRun] = await Promise.all([
firstRunPromise,
replacementPromise,
laterPromise,
]);
expect(createChildAdapterMock).toHaveBeenCalledTimes(2);
expect(first.killMock).toHaveBeenCalledWith("SIGKILL");
expect(first.disposeMock).toHaveBeenCalled();
expect(replacementRun.pid).toBeUndefined();
expect(replacementRun.waitForExtinction).toBeUndefined();
expect(later.killMock).not.toHaveBeenCalled();
later.settle(0);
await expect(
Promise.all([firstRun.wait(), replacementRun.wait(), laterRun.wait()]),
).resolves.toEqual([
expect.objectContaining({ reason: "manual-cancel" }),
expect.objectContaining({ reason: "manual-cancel" }),
expect.objectContaining({ reason: "exit" }),
]);
});
it("cancels prior scoped run when replaceExistingScope is enabled", async () => {
const first = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGKILL");
},
});
const second = createStubChildAdapter();
createChildAdapterMock.mockResolvedValueOnce(first).mockResolvedValueOnce(second);
const supervisor = createProcessSupervisor();
const firstRun = await spawnChild(supervisor, {
scopeKey: "scope:a",
argv: [process.execPath, "-e", "setTimeout(() => {}, 80)"],
timeoutMs: 1_000,
stdinMode: "pipe-open",
});
const secondRun = await spawnChild(supervisor, {
scopeKey: "scope:a",
replaceExistingScope: true,
argv: createWriteStdoutArgv("new"),
timeoutMs: 1_000,
stdinMode: "pipe-closed",
});
second.emitStdout("new");
second.settle(0);
const firstExit = await firstRun.wait();
const secondExit = await secondRun.wait();
expect(first.killMock).toHaveBeenCalledWith("SIGTERM");
expect(["manual-cancel", "signal"]).toContain(firstExit.reason);
expect(secondExit.reason).toBe("exit");
expect(secondExit.stdout).toBe("new");
});
it("waits for an in-flight scoped startup before replacing its run", async () => {
const first = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
});
const second = createStubChildAdapter();
const firstStartup = createDeferred<StubChildAdapter>();
createChildAdapterMock.mockReturnValueOnce(firstStartup.promise).mockResolvedValueOnce(second);
const supervisor = createProcessSupervisor();
const firstRunPromise = spawnChild(supervisor, {
runId: "scoped-start-first",
scopeKey: "scope:overlap",
argv: createSilentIdleArgv(),
});
const replacementPromise = spawnChild(supervisor, {
runId: "scoped-start-replacement",
scopeKey: "scope:overlap",
replaceExistingScope: true,
argv: createWriteStdoutArgv("replacement"),
});
expect(createChildAdapterMock).toHaveBeenCalledTimes(1);
firstStartup.resolve(first);
const [firstRun, replacement] = await Promise.all([firstRunPromise, replacementPromise]);
expect(createChildAdapterMock).toHaveBeenCalledTimes(2);
expect(first.killMock).toHaveBeenCalledWith("SIGTERM");
await expect(firstRun.wait()).resolves.toMatchObject({ reason: "manual-cancel" });
second.emitStdout("replacement");
second.settle(0);
await expect(replacement.wait()).resolves.toMatchObject({
reason: "exit",
stdout: "replacement",
});
});
it("starts same-scope runs concurrently when replacement is not requested", async () => {
const first = createStubChildAdapter();
const second = createStubChildAdapter();
const firstStartup = createDeferred<StubChildAdapter>();
const secondStartup = createDeferred<StubChildAdapter>();
createChildAdapterMock
.mockReturnValueOnce(firstStartup.promise)
.mockReturnValueOnce(secondStartup.promise);
const supervisor = createProcessSupervisor();
const firstRunPromise = spawnChild(supervisor, {
scopeKey: "scope:shared-concurrent",
argv: createSilentIdleArgv(),
});
const secondRunPromise = spawnChild(supervisor, {
scopeKey: "scope:shared-concurrent",
argv: createSilentIdleArgv(),
});
expect(createChildAdapterMock).toHaveBeenCalledTimes(2);
secondStartup.resolve(second);
const secondRun = await secondRunPromise;
firstStartup.resolve(first);
const firstRun = await firstRunPromise;
expect(first.killMock).not.toHaveBeenCalled();
expect(second.killMock).not.toHaveBeenCalled();
first.settle(0);
second.settle(0);
await expect(Promise.all([firstRun.wait(), secondRun.wait()])).resolves.toEqual([
expect.objectContaining({ reason: "exit" }),
expect.objectContaining({ reason: "exit" }),
]);
});
it("does not cancel a newer run while an earlier scoped replacement is pending", async () => {
const first = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
});
const replacementAdapter = createStubChildAdapter();
const newerAdapter = createStubChildAdapter();
const firstStartup = createDeferred<StubChildAdapter>();
const replacementStartup = createDeferred<StubChildAdapter>();
createChildAdapterMock
.mockReturnValueOnce(firstStartup.promise)
.mockReturnValueOnce(replacementStartup.promise)
.mockResolvedValueOnce(newerAdapter);
const supervisor = createProcessSupervisor();
const firstRunPromise = spawnChild(supervisor, {
runId: "replacement-fence-first",
scopeKey: "scope:replacement-fence",
argv: createSilentIdleArgv(),
});
const replacementPromise = spawnChild(supervisor, {
runId: "replacement-fence-replacement",
scopeKey: "scope:replacement-fence",
replaceExistingScope: true,
argv: createSilentIdleArgv(),
});
const newerRunPromise = spawnChild(supervisor, {
runId: "replacement-fence-newer",
scopeKey: "scope:replacement-fence",
argv: createSilentIdleArgv(),
});
expect(createChildAdapterMock).toHaveBeenCalledTimes(1);
firstStartup.resolve(first);
const firstRun = await firstRunPromise;
await vi.waitFor(() => {
expect(createChildAdapterMock).toHaveBeenCalledTimes(2);
});
expect(first.killMock).toHaveBeenCalledWith("SIGTERM");
expect(newerAdapter.killMock).not.toHaveBeenCalled();
replacementStartup.resolve(replacementAdapter);
const [replacement, newerRun] = await Promise.all([replacementPromise, newerRunPromise]);
expect(createChildAdapterMock).toHaveBeenCalledTimes(3);
expect(replacementAdapter.killMock).not.toHaveBeenCalled();
expect(newerAdapter.killMock).not.toHaveBeenCalled();
replacementAdapter.settle(0);
newerAdapter.settle(0);
await expect(
Promise.all([firstRun.wait(), replacement.wait(), newerRun.wait()]),
).resolves.toEqual([
expect.objectContaining({ reason: "manual-cancel" }),
expect.objectContaining({ reason: "exit" }),
expect.objectContaining({ reason: "exit" }),
]);
});
it("waits for every concurrent same-scope startup before replacing the scope", async () => {
const runCount = 8;
const startups = Array.from({ length: runCount }, () => createDeferred<StubChildAdapter>());
const adapters = Array.from({ length: runCount }, (_unused, index) =>
createStubChildAdapter({
pid: 5_000 + index,
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
}),
);
const replacementAdapter = createStubChildAdapter({ pid: 6_000 });
let adapterIndex = 0;
createChildAdapterMock.mockImplementation(() => {
const index = adapterIndex++;
if (index === runCount) {
return Promise.resolve(replacementAdapter);
}
const startup = startups[index];
if (!startup) {
throw new Error(`unexpected shared-scope startup ${index}`);
}
return startup.promise;
});
const supervisor = createProcessSupervisor();
const pendingRuns = Array.from({ length: runCount }, (_unused, index) =>
spawnChild(supervisor, {
runId: `shared-scope-${index}`,
scopeKey: "scope:shared-replacement",
argv: createSilentIdleArgv(),
}),
);
const replacementPromise = spawnChild(supervisor, {
runId: "shared-scope-replacement",
scopeKey: "scope:shared-replacement",
replaceExistingScope: true,
argv: createWriteStdoutArgv("replacement"),
});
expect(createChildAdapterMock).toHaveBeenCalledTimes(runCount);
for (let index = runCount - 1; index >= 0; index -= 1) {
const startup = startups[index];
const adapter = adapters[index];
if (!startup || !adapter) {
throw new Error(`missing shared-scope startup ${index}`);
}
startup.resolve(adapter);
}
const runs = await Promise.all(pendingRuns);
const replacement = await replacementPromise;
expect(createChildAdapterMock).toHaveBeenCalledTimes(runCount + 1);
for (const adapter of adapters) {
expect(adapter.killMock).toHaveBeenCalledWith("SIGTERM");
}
await expect(Promise.all(runs.map((run) => run.wait()))).resolves.toEqual(
Array.from({ length: runCount }, () => expect.objectContaining({ reason: "manual-cancel" })),
);
replacementAdapter.emitStdout("replacement");
replacementAdapter.settle(0);
await expect(replacement.wait()).resolves.toMatchObject({
reason: "exit",
stdout: "replacement",
});
});
it("keeps startup independent for different replacement scopes", async () => {
const first = createStubChildAdapter();
const second = createStubChildAdapter();
const firstStartup = createDeferred<StubChildAdapter>();
const secondStartup = createDeferred<StubChildAdapter>();
createChildAdapterMock
.mockReturnValueOnce(firstStartup.promise)
.mockReturnValueOnce(secondStartup.promise);
const supervisor = createProcessSupervisor();
const firstRunPromise = spawnChild(supervisor, {
scopeKey: "scope:independent-a",
replaceExistingScope: true,
argv: createSilentIdleArgv(),
});
const secondRunPromise = spawnChild(supervisor, {
scopeKey: "scope:independent-b",
replaceExistingScope: true,
argv: createSilentIdleArgv(),
});
expect(createChildAdapterMock).toHaveBeenCalledTimes(2);
secondStartup.resolve(second);
const secondRun = await secondRunPromise;
firstStartup.resolve(first);
const firstRun = await firstRunPromise;
expect(first.killMock).not.toHaveBeenCalled();
expect(second.killMock).not.toHaveBeenCalled();
first.settle(0);
second.settle(0);
await expect(Promise.all([firstRun.wait(), secondRun.wait()])).resolves.toEqual([
expect.objectContaining({ reason: "exit" }),
expect.objectContaining({ reason: "exit" }),
]);
});
it("starts a scoped replacement after the previous adapter fails to start", async () => {
const second = createStubChildAdapter();
createChildAdapterMock
.mockRejectedValueOnce(new Error("first adapter could not start"))
.mockResolvedValueOnce(second);
const supervisor = createProcessSupervisor();
const firstRunPromise = spawnChild(supervisor, {
runId: "failed-scoped-start",
scopeKey: "scope:recover-start",
argv: createSilentIdleArgv(),
});
const replacementPromise = spawnChild(supervisor, {
runId: "recovered-scoped-start",
scopeKey: "scope:recover-start",
replaceExistingScope: true,
argv: createWriteStdoutArgv("recovered"),
});
await expect(firstRunPromise).rejects.toThrow("first adapter could not start");
const replacement = await replacementPromise;
second.emitStdout("recovered");
second.settle(0);
await expect(replacement.wait()).resolves.toMatchObject({
reason: "exit",
stdout: "recovered",
});
});
it("keeps only the newest run across 64 concurrent same-scope replacements", async () => {
const runCount = 64;
const adapters = Array.from({ length: runCount }, (_unused, index) =>
createStubChildAdapter({
pid: 2_000 + index,
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
}),
);
let adapterIndex = 0;
createChildAdapterMock.mockImplementation(async () => {
const adapter = adapters[adapterIndex++];
if (!adapter) {
throw new Error("unexpected concurrent supervisor startup");
}
return adapter;
});
const supervisor = createProcessSupervisor();
const pendingRuns = Array.from({ length: runCount }, (_unused, index) =>
spawnChild(supervisor, {
runId: `scope-stress-${index}`,
scopeKey: "scope:stress",
replaceExistingScope: true,
argv: createSilentIdleArgv(),
}),
);
const runs = await Promise.all(pendingRuns);
for (const [index, adapter] of adapters.entries()) {
if (index === runCount - 1) {
expect(adapter.killMock, `newest scope owner ${index}`).not.toHaveBeenCalled();
} else {
expect(adapter.killMock, `superseded scope owner ${index}`).toHaveBeenCalledWith("SIGTERM");
}
}
const newestAdapter = adapters[runCount - 1];
const newestRun = runs[runCount - 1];
if (!newestAdapter || !newestRun) {
throw new Error("expected the newest scoped process");
}
newestAdapter.settle(0);
const exits = await Promise.all(runs.map((run) => run.wait()));
for (const [index, exit] of exits.entries()) {
expect(exit.reason, `scoped run ${index}`).toBe(
index === runCount - 1 ? "exit" : "manual-cancel",
);
}
});
it("continues replacing scoped runs across interleaved startup failures", async () => {
const runCount = 48;
const adapters = new Map<number, StubChildAdapter>();
let adapterIndex = 0;
createChildAdapterMock.mockImplementation(async () => {
const index = adapterIndex++;
if (index % 7 === 3) {
throw new Error(`adapter ${index} could not start`);
}
const adapter = createStubChildAdapter({
pid: 3_000 + index,
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGTERM");
},
});
adapters.set(index, adapter);
return adapter;
});
const supervisor = createProcessSupervisor();
const pendingRuns = Array.from({ length: runCount }, (_unused, index) =>
spawnChild(supervisor, {
runId: `scope-recovery-${index}`,
scopeKey: "scope:interleaved-recovery",
replaceExistingScope: true,
argv: createSilentIdleArgv(),
}),
);
const results = await Promise.allSettled(pendingRuns);
const running: ManagedRun[] = [];
for (const [index, result] of results.entries()) {
if (index % 7 === 3) {
expect(result.status, `failed adapter ${index}`).toBe("rejected");
if (result.status === "rejected") {
expect(result.reason).toMatchObject({ message: `adapter ${index} could not start` });
}
continue;
}
expect(result.status, `started adapter ${index}`).toBe("fulfilled");
if (result.status === "fulfilled") {
running.push(result.value);
}
}
const newestAdapter = adapters.get(runCount - 1);
if (!newestAdapter) {
throw new Error("expected the final replacement adapter");
}
for (const [index, adapter] of adapters) {
if (index === runCount - 1) {
expect(adapter.killMock, `newest adapter ${index}`).not.toHaveBeenCalled();
} else {
expect(adapter.killMock, `superseded adapter ${index}`).toHaveBeenCalledWith("SIGTERM");
}
}
newestAdapter.settle(0);
const exits = await Promise.all(running.map((run) => run.wait()));
expect(exits.filter((exit) => exit.reason === "exit")).toHaveLength(1);
expect(exits.filter((exit) => exit.reason === "manual-cancel")).toHaveLength(
running.length - 1,
);
});
it("starts 24 independent scoped processes without a global startup bottleneck", async () => {
const scopeCount = 24;
const startups = Array.from({ length: scopeCount }, () => createDeferred<StubChildAdapter>());
const adapters = Array.from({ length: scopeCount }, (_unused, index) =>
createStubChildAdapter({ pid: 4_000 + index }),
);
let adapterIndex = 0;
createChildAdapterMock.mockImplementation(() => {
const startup = startups[adapterIndex++];
if (!startup) {
throw new Error("unexpected independent supervisor startup");
}
return startup.promise;
});
const supervisor = createProcessSupervisor();
const pendingRuns = Array.from({ length: scopeCount }, (_unused, index) =>
spawnChild(supervisor, {
runId: `independent-scope-${index}`,
scopeKey: `scope:independent-${index}`,
replaceExistingScope: true,
argv: createSilentIdleArgv(),
}),
);
expect(createChildAdapterMock).toHaveBeenCalledTimes(scopeCount);
for (let index = scopeCount - 1; index >= 0; index -= 1) {
const startup = startups[index];
const adapter = adapters[index];
if (!startup || !adapter) {
throw new Error(`missing independent scope ${index}`);
}
startup.resolve(adapter);
}
const runs = await Promise.all(pendingRuns);
for (const adapter of adapters) {
expect(adapter.killMock).not.toHaveBeenCalled();
adapter.settle(0);
}
const exits = await Promise.all(runs.map((run) => run.wait()));
expect(exits.every((exit) => exit.reason === "exit")).toBe(true);
});
it("applies overall timeout even for near-immediate timer firing", async () => {
vi.useFakeTimers();
const adapter = createStubChildAdapter({
onKill: (signal, current) => {
current.settle(null, signal ?? "SIGKILL");
},
});
createChildAdapterMock.mockResolvedValue(adapter);
const supervisor = createProcessSupervisor();
const run = await spawnChild(supervisor, {
argv: createSilentIdleArgv(),
timeoutMs: 1,
stdinMode: "pipe-closed",
});
const exitPromise = run.wait();
await vi.advanceTimersByTimeAsync(1);
const exit = await exitPromise;
expect(adapter.killMock).toHaveBeenCalledWith(
process.platform === "win32" ? "SIGKILL" : "SIGTERM",
);
expect(exit.reason).toBe("overall-timeout");
expect(exit.timedOut).toBe(true);
});
it("classifies a natural close after a missed overall deadline as timed out", async () => {
vi.useFakeTimers();
const nowSpy = vi.spyOn(performance, "now").mockReturnValue(1_000);
const adapter = createStubChildAdapter();
createChildAdapterMock.mockResolvedValue(adapter);
const supervisor = createProcessSupervisor();
const run = await spawnChild(supervisor, {
argv: createSilentIdleArgv(),
timeoutMs: 10,
stdinMode: "pipe-closed",
});
const exitPromise = run.wait();
nowSpy.mockReturnValue(1_011);
adapter.settle(0);
const exit = await exitPromise;
expect(adapter.killMock).not.toHaveBeenCalled();
expect(exit.reason).toBe("overall-timeout");
expect(exit.timedOut).toBe(true);
});
it("uses the refreshed no-output deadline when a missed timer races natural close", async () => {
vi.useFakeTimers();
const nowSpy = vi.spyOn(performance, "now").mockReturnValue(1_000);
const adapter = createStubChildAdapter();
createChildAdapterMock.mockResolvedValue(adapter);
const supervisor = createProcessSupervisor();
const run = await spawnChild(supervisor, {
argv: createSilentIdleArgv(),
timeoutMs: 100,
noOutputTimeoutMs: 10,
stdinMode: "pipe-closed",
});
const exitPromise = run.wait();
nowSpy.mockReturnValue(1_005);
adapter.emitStdout("progress");
nowSpy.mockReturnValue(1_016);
adapter.settle(0);
const exit = await exitPromise;
expect(adapter.killMock).not.toHaveBeenCalled();
expect(exit.reason).toBe("no-output-timeout");
expect(exit.noOutputTimedOut).toBe(true);
expect(exit.timedOut).toBe(true);
});
it("can stream output without retaining it in RunExit payload", async () => {
const adapter = createStubChildAdapter();
createChildAdapterMock.mockResolvedValue(adapter);
const supervisor = createProcessSupervisor();
let streamed = "";
const run = await spawnChild(supervisor, {
argv: createWriteStdoutArgv("streamed"),
timeoutMs: 1_000,
stdinMode: "pipe-closed",
captureOutput: false,
onStdout: (chunk) => {
streamed += chunk;
},
});
adapter.emitStdout("streamed");
adapter.settle(0);
const exit = await run.wait();
expect(streamed).toBe("streamed");
expect(exit.stdout).toBe("");
});
it("bounds retained output on UTF-16 boundaries while streaming full chunks", async () => {
const adapter = createStubChildAdapter();
createChildAdapterMock.mockResolvedValue(adapter);
const supervisor = createProcessSupervisor();
let streamedStdout = "";
let streamedStderr = "";
const maxCapturedOutputChars = 256;
const stdoutMarker = `[openclaw: captured stdout truncated to last ${maxCapturedOutputChars} chars]\n`;
const stderrMarker = `[openclaw: captured stderr truncated to last ${maxCapturedOutputChars} chars]\n`;
const retainedChars = maxCapturedOutputChars - stdoutMarker.length - 1;
const stdoutChunk = `${"a".repeat(stdoutMarker.length)}😀${"s".repeat(retainedChars)}`;
const stderrChunk = `${"b".repeat(stderrMarker.length)}😀${"e".repeat(retainedChars)}`;
const run = await spawnChild(supervisor, {
argv: createWriteStdoutArgv(stdoutChunk),
timeoutMs: 1_000,
stdinMode: "pipe-closed",
maxCapturedOutputChars,
onStdout: (chunk) => {
streamedStdout += chunk;
},
onStderr: (chunk) => {
streamedStderr += chunk;
},
});
adapter.emitStdout(stdoutChunk);
adapter.emitStderr(stderrChunk);
adapter.settle(0);
const exit = await run.wait();
expect(streamedStdout).toBe(stdoutChunk);
expect(streamedStderr).toBe(stderrChunk);
expect(exit.stdout).toBe(`${stdoutMarker}${"s".repeat(retainedChars)}`);
expect(exit.stderr).toBe(`${stderrMarker}${"e".repeat(retainedChars)}`);
});
});