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:
Peter Steinberger 2026-09-21 19:58:04 -07:00 • committed by GitHub
parent c69b8b9868
commit 485fa310fc
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 175 additions and 20 deletions

View file

@ -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 {

View file

@ -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;
}
}
}
}
}

View file

@ -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 () => {

View file

@ -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({