mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 17:53:39 +00:00
fix(test): retain claimed listeners during initial allocation (#155209)
Compose the cooperative port claim and real loopback listener for initial TCP fixture and HTTP Gateway reservations. Select another automatic port only after a real bind collision and fully joined rollback; preserve explicit ports, pinned reacquisition, cancellation, and cleanup failures. Reproduce both adapters with real competing listeners. All 122 focused cases and the complete selected changed gate pass; independent P2 review is clean. The historical acquisition CI failures remain unattributed.
This commit is contained in:
parent
c69b8b9868
commit
485fa310fc
4 changed files with 175 additions and 20 deletions
|
|
@ -3,7 +3,11 @@ import { createServer } from "node:http";
|
|||
import type { Socket } from "node:net";
|
||||
import { collectNestedErrorCandidates } from "@openclaw/normalization-core/error-coercion";
|
||||
import { hasErrnoCode } from "../infra/errno.js";
|
||||
import { acquireTestPortBlock, type TestPortClaim } from "../test-utils/port-claims.js";
|
||||
import {
|
||||
acquireTestPortBlock,
|
||||
reserveTestPortListener,
|
||||
type TestPortClaim,
|
||||
} from "../test-utils/port-claims.js";
|
||||
import { getDeterministicFreePortBlock } from "../test-utils/ports.js";
|
||||
import { GatewayStartupCleanupError } from "./server-shutdown.js";
|
||||
import type { GatewayServer } from "./server.js";
|
||||
|
|
@ -81,29 +85,22 @@ let activeDispatcher:
|
|||
|
||||
/** Retain the exact loopback listener until the real Gateway transport adopts it. */
|
||||
export async function reserveGatewayTestListener(port = 0) {
|
||||
const claim = await acquireTestPortBlock({ offsets: [0, 1, 2, 3, 4], ...(port ? { port } : {}) });
|
||||
const listener = createServer();
|
||||
const rejectEarlyConnection = (socket: Socket) => socket.destroy();
|
||||
listener.on("connection", rejectEarlyConnection);
|
||||
const { claim, listener, releaseListener } = await reserveTestPortListener({
|
||||
offsets: [0, 1, 2, 3, 4],
|
||||
...(port ? { port } : {}),
|
||||
createListener: () => createServer().on("connection", rejectEarlyConnection),
|
||||
});
|
||||
let adopted = false;
|
||||
const closeUnadopted = async () => {
|
||||
if (!adopted && listener.listening) {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
listener.close((error) => (error ? reject(error) : resolve()));
|
||||
});
|
||||
await releaseListener();
|
||||
}
|
||||
if (!listener.listening) {
|
||||
await claim.release();
|
||||
}
|
||||
};
|
||||
try {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
listener.once("error", reject);
|
||||
listener.listen(claim.port, "127.0.0.1", () => {
|
||||
listener.off("error", reject);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
const address = listener.address();
|
||||
assert(address && typeof address !== "string");
|
||||
return {
|
||||
|
|
|
|||
|
|
@ -1,3 +1,5 @@
|
|||
import type { Server } from "node:net";
|
||||
import { runQaGatewayFixture } from "../../test/helpers/qa-gateway-cleanup.js";
|
||||
import { hasErrnoCode } from "../infra/errno.js";
|
||||
import { FILE_LOCK_TIMEOUT_ERROR_CODE } from "../infra/file-lock.js";
|
||||
import { claimTestPortBlock, type TestPortClaim } from "./port-claim-lock.js";
|
||||
|
|
@ -56,3 +58,70 @@ export async function acquireTestPortBlock(params: {
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Hold the real loopback listener as well as the cooperative port claim. */
|
||||
export async function reserveTestPortListener<T extends Server>(params: {
|
||||
offsets: number[];
|
||||
port?: number;
|
||||
signal?: AbortSignal;
|
||||
createListener: () => T;
|
||||
verifyCleanup?: (cleanup: () => Promise<void>) => Promise<void>;
|
||||
}) {
|
||||
const verifyCleanup = params.verifyCleanup ?? ((cleanup: () => Promise<void>) => cleanup());
|
||||
const seen = new Set<number>();
|
||||
while (true) {
|
||||
const claim = await acquireTestPortBlock(params);
|
||||
let reservation: { listener: T; releaseListener: () => Promise<void> } | undefined;
|
||||
let bindError: unknown;
|
||||
try {
|
||||
if (seen.has(claim.port)) {
|
||||
throw new Error("no unclaimed test Gateway port block available");
|
||||
}
|
||||
seen.add(claim.port);
|
||||
params.signal?.throwIfAborted();
|
||||
const listener = params.createListener();
|
||||
const releaseListener = () =>
|
||||
new Promise<void>((resolve, reject) => {
|
||||
listener.close((error) => (error ? reject(error) : resolve()));
|
||||
});
|
||||
reservation = { listener, releaseListener };
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
const failed = (error: Error) => {
|
||||
bindError = error;
|
||||
reject(error);
|
||||
};
|
||||
listener.once("error", failed);
|
||||
listener.listen(claim.port, "127.0.0.1", () => {
|
||||
listener.off("error", failed);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
params.signal?.throwIfAborted();
|
||||
return { claim, ...reservation };
|
||||
} catch (error) {
|
||||
try {
|
||||
await runQaGatewayFixture(
|
||||
async (): Promise<never> => {
|
||||
throw error;
|
||||
},
|
||||
() =>
|
||||
reservation?.listener.listening
|
||||
? verifyCleanup(reservation.releaseListener)
|
||||
: undefined,
|
||||
() => verifyCleanup(claim.release),
|
||||
);
|
||||
} catch (rollbackError) {
|
||||
// An unrelated listener can win after the free-port probe closes. Only
|
||||
// initial automatic selection may move, after both provisional owners drain.
|
||||
if (
|
||||
rollbackError !== error ||
|
||||
params.port !== undefined ||
|
||||
error !== bindError ||
|
||||
!hasErrnoCode(error, "EADDRINUSE")
|
||||
) {
|
||||
throw rollbackError;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -7,13 +7,98 @@ import { createVitestResourceOwner } from "../../scripts/lib/vitest-resource-own
|
|||
import { resolveGatewayPort } from "../../src/config/paths.js";
|
||||
import type { OpenClawConfig } from "../../src/config/types.openclaw.js";
|
||||
import { resolveGatewayUrlOverride } from "../../src/gateway/client-bootstrap.js";
|
||||
import { reserveGatewayTestListener } from "../../src/gateway/test-helpers.listener.js";
|
||||
import { captureFullEnv, withEnvAsync } from "../../src/test-utils/env.js";
|
||||
import { acquireTestPortBlock } from "../../src/test-utils/port-claims.js";
|
||||
import * as testPorts from "../../src/test-utils/ports.js";
|
||||
import { createFixtureLifetime } from "./fixture-lifetime.js";
|
||||
import { createOpenClawTestInstance } from "./openclaw-test-instance.js";
|
||||
import { createDeferred, withTestTimeout } from "./promise.js";
|
||||
import { runQaGatewayFixture } from "./qa-gateway-cleanup.js";
|
||||
|
||||
describe("createOpenClawTestInstance acquisition", () => {
|
||||
it.each([
|
||||
{
|
||||
name: "child-process",
|
||||
offsets: [0, 1],
|
||||
acquire: async () => {
|
||||
const instance = await createOpenClawTestInstance({ name: "initial-listener-race" });
|
||||
return { port: instance.port, cleanup: () => instance.cleanup() };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "in-process",
|
||||
offsets: [0, 1, 2, 3, 4],
|
||||
acquire: async () => {
|
||||
const reservation = await reserveGatewayTestListener();
|
||||
return { port: reservation.port, cleanup: reservation.closeUnadopted };
|
||||
},
|
||||
},
|
||||
])(
|
||||
"retains another $name reservation when an unclaimed listener wins the probe",
|
||||
async (adapter) => {
|
||||
const competitor = net.createServer((socket) => socket.destroy());
|
||||
const exclusiveProbe = net.createServer((socket) => socket.destroy());
|
||||
const allocate = testPorts.getDeterministicFreePortBlock;
|
||||
let competitorPort: number | undefined;
|
||||
let reserved: { port: number; cleanup: () => Promise<void> } | undefined;
|
||||
const listen = (server: net.Server, port: number) =>
|
||||
new Promise<void>((resolve, reject) => {
|
||||
server.once("error", reject);
|
||||
server.listen(port, "127.0.0.1", () => {
|
||||
server.off("error", reject);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
const close = async (server: net.Server) => {
|
||||
if (server.listening) {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
server.close((error) => (error ? reject(error) : resolve()));
|
||||
});
|
||||
}
|
||||
};
|
||||
const allocationSpy = vi
|
||||
.spyOn(testPorts, "getDeterministicFreePortBlock")
|
||||
.mockImplementationOnce(async (params) => {
|
||||
competitorPort = await allocate(params);
|
||||
// This listener does not borrow a file claim; the real probe has already closed.
|
||||
await listen(competitor, competitorPort);
|
||||
return competitorPort;
|
||||
});
|
||||
await runQaGatewayFixture(
|
||||
async () => {
|
||||
reserved = await adapter.acquire();
|
||||
expect(competitor.listening).toBe(true);
|
||||
if (competitorPort === undefined) {
|
||||
throw new Error("real allocator did not return a competitor port");
|
||||
}
|
||||
expect(reserved.port).not.toBe(competitorPort);
|
||||
await expect(listen(exclusiveProbe, reserved.port)).rejects.toMatchObject({
|
||||
code: "EADDRINUSE",
|
||||
});
|
||||
const abandoned = await acquireTestPortBlock({
|
||||
port: competitorPort,
|
||||
offsets: adapter.offsets,
|
||||
});
|
||||
await abandoned.release();
|
||||
// An explicitly requested port remains pinned, even when its socket is occupied.
|
||||
await expect(reserveGatewayTestListener(competitorPort)).rejects.toMatchObject({
|
||||
code: "EADDRINUSE",
|
||||
});
|
||||
},
|
||||
() => allocationSpy.mockRestore(),
|
||||
() => close(exclusiveProbe),
|
||||
() => reserved?.cleanup(),
|
||||
() => close(competitor),
|
||||
() => {
|
||||
expect(competitor.listening).toBe(false);
|
||||
expect(competitor.address()).toBeNull();
|
||||
expect(exclusiveProbe.listening).toBe(false);
|
||||
},
|
||||
);
|
||||
},
|
||||
);
|
||||
|
||||
it.skipIf(process.platform !== "linux")(
|
||||
"keeps Gateway and deferred sandbox listeners outside the kernel client-port range",
|
||||
async () => {
|
||||
|
|
|
|||
|
|
@ -30,7 +30,7 @@ import {
|
|||
createOpenClawTestState,
|
||||
type OpenClawTestState,
|
||||
} from "../../src/test-utils/openclaw-test-state.js";
|
||||
import { acquireTestPortBlock } from "../../src/test-utils/port-claims.js";
|
||||
import { reserveTestPortListener } from "../../src/test-utils/port-claims.js";
|
||||
import { cleanupSessionStateForTest } from "../../src/test-utils/session-state-cleanup.js";
|
||||
import { sleep } from "../../src/utils.js";
|
||||
import { decodeUtf8Tail } from "./bounded-child-output.js";
|
||||
|
|
@ -754,11 +754,15 @@ export async function createOpenClawTestInstance(
|
|||
if (options.port !== undefined) {
|
||||
port = options.port;
|
||||
} else {
|
||||
const claimed = await acquireTestPortBlock({ offsets: [0, 1], signal });
|
||||
port = claimed.port;
|
||||
releasePortClaims = claimed.release;
|
||||
signal?.throwIfAborted();
|
||||
reservation = await reserveGatewayPort(port, options.verifyCleanup);
|
||||
const reserved = await reserveTestPortListener({
|
||||
offsets: [0, 1],
|
||||
signal,
|
||||
createListener: () => net.createServer((socket) => socket.destroy()),
|
||||
verifyCleanup: options.verifyCleanup,
|
||||
});
|
||||
port = reserved.claim.port;
|
||||
releasePortClaims = reserved.claim.release;
|
||||
reservation = { release: reserved.releaseListener };
|
||||
}
|
||||
signal?.throwIfAborted();
|
||||
state = await createOpenClawTestState({
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue