diff --git a/README.md b/README.md index 9cb4d7345..5fdf94235 100644 --- a/README.md +++ b/README.md @@ -113,3 +113,9 @@ Visit our [contributing guide](CONTRIBUTING.md) to understand more about our dev ## License Testplane is [MIT licensed](LICENSE). + +### Общий срок WSDriver-команды + +Экспериментальный режим `TESTPLANE_WSDRIVER_DEADLINE_ENABLED=true` включает общий срок WSDriver-команды: соединение, reconnect, кодирование запроса и внутренние повторы входят в `httpTimeout` (или переданный транспортом `timeout.response`). По истечении срока транспорт закрывается, дальнейшая отправка прекращается, а браузер помечается сломанным и не возвращается в кэш сессий. Ошибка `WSDRIVER_REQUEST_DEADLINE` передаётся механизму повторов тестов; внутренний повтор WebdriverIO не начинает срок заново. + +Без флага сохраняются прежние таймауты и повторы команд. Закрытие незавершённого WS-соединения и отмена reconnect при закрытии работают в обоих режимах. diff --git a/src/browser/existing-browser.ts b/src/browser/existing-browser.ts index b9340f6ac..47ab59829 100644 --- a/src/browser/existing-browser.ts +++ b/src/browser/existing-browser.ts @@ -257,6 +257,8 @@ export class ExistingBrowser extends Browser { sessionCaps: sessionCaps as WebdriverIO.Capabilities, headers: sessionOpts.headers as Record, browserConfig: this._config, + // Сессия с отменённой WS-командой не должна возвращаться в кэш. + onRequestDeadline: () => this.markAsBroken({ stubBrowserCommands: true }), }); opts.customWdRequestAgent = this._wsDriver; diff --git a/src/browser/wsdriver/error.ts b/src/browser/wsdriver/error.ts index a5096af31..71685205c 100644 --- a/src/browser/wsdriver/error.ts +++ b/src/browser/wsdriver/error.ts @@ -27,3 +27,12 @@ export class WSDriverRequestError extends WsError { return true; } } + +export class WSDriverRequestDeadlineError extends Error { + readonly code = "WSDRIVER_REQUEST_DEADLINE"; + + constructor(timeout: number) { + super(`WSDriver command timed out after ${timeout}ms including connection and retries`); + this.name = "WSDriverRequestDeadlineError"; + } +} diff --git a/src/browser/wsdriver/index.ts b/src/browser/wsdriver/index.ts index 08758ef04..a3158d661 100644 --- a/src/browser/wsdriver/index.ts +++ b/src/browser/wsdriver/index.ts @@ -10,6 +10,7 @@ import { WSDriverError, WSDriverRequestError, WSDriverRequestTimeoutError, + WSDriverRequestDeadlineError, } from "./error"; import { WSD_ACCEPT_ENCODING_HEADER, @@ -36,7 +37,14 @@ import { BrowserConfig } from "../../config/browser-config"; import { constructWsDriverRequest } from "./request"; import { exponentiallyWait } from "../../ws-connection/utils"; +interface RequestDeadlineContext { + signal: AbortSignal; + remaining: () => number; + check: () => void; +} + interface WSDriverRequestAgentOptions { + onRequestDeadline?: () => void; sessionId: string; headers?: Record; requestTimeout: number; @@ -50,11 +58,27 @@ export class WSDriverRequestAgent { private _serverSupportedCompressionType?: WsDriverCompressionType; private _sessionId: string; private _sessionPrefix: string; + private readonly _requestTimeout: number; + private readonly _onRequestDeadline?: () => void; + private readonly _deadlineEnabled: boolean; + private readonly _requestAbort = new AbortController(); private constructor( wsdWsEndpoint: string, - { sessionId, headers, requestTimeout, clientSupportedCompressionTypes }: WSDriverRequestAgentOptions, + { + sessionId, + headers, + requestTimeout, + clientSupportedCompressionTypes, + onRequestDeadline, + }: WSDriverRequestAgentOptions, ) { + this._requestTimeout = requestTimeout; + this._onRequestDeadline = onRequestDeadline; + this._deadlineEnabled = + process.env.TESTPLANE_WSDRIVER_DEADLINE_ENABLED === "true" && + Number.isFinite(requestTimeout) && + requestTimeout > 0; headers ||= {}; headers[WSD_ACCEPT_ENCODING_HEADER] = clientSupportedCompressionTypes.join(", "); @@ -90,11 +114,13 @@ export class WSDriverRequestAgent { sessionCaps, headers = {}, browserConfig, + onRequestDeadline, }: { sessionId: string; sessionCaps: WebdriverIO.Capabilities; headers: Record; browserConfig: BrowserConfig; + onRequestDeadline?: () => void; }): WSDriverRequestAgent { if (!sessionCaps["se:wsdriver"]) { throw new WSDriverError({ message: "Couldn't determine wsdriver endpoint" }); @@ -122,11 +148,15 @@ export class WSDriverRequestAgent { requestTimeout, clientSupportedCompressionTypes, supportedVersions, + onRequestDeadline, }); } close(): void { - this._wsConnection.close(); + if (this._deadlineEnabled) { + this._requestAbort.abort(new WSDriverRequestAgentTerminatedError()); + } + this._wsConnection.close(this._deadlineEnabled); } private async _onMessage(data: RawData, isBinary: boolean): Promise { @@ -238,17 +268,64 @@ export class WSDriverRequestAgent { /** @description Performs high-level WSDriver request with timeout */ async request(url: URL, options: RequestWsDriverOptions): Promise { + if (!this._deadlineEnabled) return this._request(url, options); + const signal = this._requestAbort.signal; + signal.throwIfAborted(); + const responseTimeout = options.timeout?.response; + const requestTimeout = + Number.isFinite(responseTimeout) && responseTimeout! > 0 ? responseTimeout! : this._requestTimeout; + const deadline = performance.now() + requestTimeout; + let onAbort!: () => void; + const expire = (): void => { + if (signal.aborted) return; + // Ошибка не ETIMEDOUT: повтор всей команды в WebdriverIO обнулит общий срок. + const error = new WSDriverRequestDeadlineError(requestTimeout); + this._requestAbort.abort(error); + this._wsConnection.close(true); + this._onRequestDeadline?.(); + }; + const context: RequestDeadlineContext = { + signal, + remaining: () => Math.max(1, Math.ceil(deadline - performance.now())), + check: () => { + if (performance.now() >= deadline) expire(); + signal.throwIfAborted(); + }, + }; + const aborted = new Promise((_, reject) => { + onAbort = (): void => reject(signal.reason); + signal.addEventListener("abort", onAbort, { once: true }); + }); + const timer = setTimeout(expire, requestTimeout).unref(); + try { + // Отмена закрывает транспорт; проверки после await запрещают позднюю отправку. + return await Promise.race([this._request(url, options, context), aborted]); + } finally { + clearTimeout(timer); + signal.removeEventListener("abort", onAbort); + } + } + + private async _request( + url: URL, + options: RequestWsDriverOptions, + context?: RequestDeadlineContext, + ): Promise { let requestId!: number; let result!: IncomingWsDriverMessage | WsError; for (let retriesLeft = WSD_REQUEST_RETRIES; retriesLeft >= 0; retriesLeft--) { + context?.check(); requestId = this._wsConnection.getRequestId(); + const compressionType = await this._getRequestCompressionType(); + context?.check(); const requestMessage = await constructWsDriverRequest(url, options, { requestId, sessionPrefix: this._sessionPrefix, - compressionType: await this._getRequestCompressionType(), + compressionType, }); + context?.check(); if (debugWSDriver.enabled) { const header = requestMessage.readUint8(1); const commandEndIdx = requestMessage.indexOf(0, 8); @@ -273,10 +350,11 @@ export class WSDriverRequestAgent { ); } - result = (await this._wsConnection.makeRequest(requestId, requestMessage).catch((err: WsError) => err)) as - | IncomingWsDriverMessage - | WsError; + result = (await this._wsConnection + .makeRequest(requestId, requestMessage, context?.remaining()) + .catch((err: WsError) => err)) as IncomingWsDriverMessage | WsError; + context?.check(); if (result instanceof WSDriverRequestTimeoutError) { const requestError = new Error(result.message); requestError.stack = result.stack; @@ -289,6 +367,7 @@ export class WSDriverRequestAgent { break; } + context?.check(); if (debugWSDriver.enabled) { const header = requestMessage.readUint8(1); const commandEndIdx = requestMessage.indexOf(0, 8); @@ -310,6 +389,7 @@ export class WSDriverRequestAgent { await exponentiallyWait({ baseDelay: WSD_REQUEST_RETRY_BASE_DELAY, attempt: WSD_REQUEST_RETRIES - retriesLeft, + signal: context?.signal, }); } diff --git a/src/browser/wsdriver/types.ts b/src/browser/wsdriver/types.ts index 55e52a2a5..8ceee0b31 100644 --- a/src/browser/wsdriver/types.ts +++ b/src/browser/wsdriver/types.ts @@ -69,6 +69,7 @@ interface IncomingWsDriverStringMessage extends IncomingWsDriverGeneralMessage { export type IncomingWsDriverMessage = IncomingWsDriverJsonMessage | IncomingWsDriverStringMessage; export interface RequestWsDriverOptions { + timeout?: { response: number }; path?: string; method?: WsDriverRequestMethodString; json?: Record; diff --git a/src/ws-connection/index.ts b/src/ws-connection/index.ts index fa05149de..2c17f2d69 100644 --- a/src/ws-connection/index.ts +++ b/src/ws-connection/index.ts @@ -58,14 +58,14 @@ interface WsConnectionOptions { // Closing WS when its still not connected produces error: // https://github.com/websockets/ws/blob/86eac5b44ac2bff9087ec40c9bd06bc7b4f0da07/lib/websocket.js#L297-L301 -const closeWsConnection = (ws: WebSocket): void => { - if (ws.readyState !== ws.CONNECTING) { - ws.close(); - } else { - ws.once("open", () => { - ws.close(); - }); +const closeWsConnection = (ws: WebSocket, force = false): void => { + if (force || ws.readyState === ws.CONNECTING) { + // terminate прерывает и незавершённый HTTP upgrade; ws сообщает об этом error. + ws.once("error", () => {}); + ws.terminate(); + return; } + ws.close(); }; export class WsConnection< @@ -88,7 +88,9 @@ export class WsConnection< private _pingShouldSkip = false; private _pingInterval: ReturnType | null = null; private _pingSubsequentFails = 0; - private _onConnectionCloseFn: (() => void) | null = null; // Defined, if there is connection attempt at the moment + private _pongTimeout?: ReturnType; + private readonly _connectionAbort = new AbortController(); + private _onConnectionCloseFn: ((force?: boolean) => void) | null = null; // Defined, if there is connection attempt at the moment private _wsConnectionStatus: WsConnectionStatus = WsConnectionStatus.DISCONNECTED; private _wsConnection: WebSocket | null = null; private _wsConnectionPromise: Promise | null = null; @@ -119,23 +121,22 @@ export class WsConnection< /** @description Tries to establish ws connection with timeout */ private async _tryToEstablishWsConnection(endpoint: string): Promise { + if (this._wsConnectionStatus === WsConnectionStatus.CLOSED) { + return new this._errors.ConnectionTerminated(); + } return new Promise(resolve => { try { - const onConnectionCloseFn = (): void => done(new this._errors.ConnectionTerminated()); - - if (this._wsConnectionStatus === WsConnectionStatus.CLOSED) { - onConnectionCloseFn(); - } else { - this._onConnectionCloseFn = onConnectionCloseFn; - } - // eslint-disable-next-line const cdpConnectionInstance = this; const ws = new WebSocket(endpoint, { headers: this._requestHeaders }); + this._onConnectionCloseFn = (force): void => { + closeWsConnection(ws, force); + done(new this._errors.ConnectionTerminated()); + }; let isSettled = false; const timeoutId = setTimeout(() => { - closeWsConnection(ws); + closeWsConnection(ws, true); done( new this._errors.ConnectionTimeout({ message: `Couldn't establish WS connection to "${endpoint}" in ${this._timeouts.createSession}ms`, @@ -147,10 +148,9 @@ export class WsConnection< done(ws); }; - const onUnexpectedResponse = async (ws: WebSocket, res: IncomingMessage): Promise => { - closeWsConnection(ws); - + const onUnexpectedResponse = async (_request: unknown, res: IncomingMessage): Promise => { const reason = await consumeText(res).catch(() => "Unknown reason"); + closeWsConnection(ws, true); done( new this._errors.ConnectionEstablishment({ @@ -161,7 +161,7 @@ export class WsConnection< }; const onError = (error: unknown): void => { - closeWsConnection(ws); + closeWsConnection(ws, true); done( new this._errors.ConnectionEstablishment({ message: `Couldn't establish WS connection to "${endpoint}": ${error}`, @@ -300,6 +300,7 @@ export class WsConnection< baseDelay: this._retries.baseDelay, attempt: this._retries.count - retriesLeft, factor: this._retries.factor, + signal: this._connectionAbort.signal, }); } } @@ -336,7 +337,11 @@ export class WsConnection< } /** @description Performs WS request with timeout */ - async makeRequest(requestId: number, requestMessage: RequestMessageType): Promise { + async makeRequest( + requestId: number, + requestMessage: RequestMessageType, + requestTimeout = this._timeouts.request, + ): Promise { const ws = await this._getWsConnection(); if (this._wsConnectionStatus === WsConnectionStatus.CLOSED) { @@ -352,12 +357,12 @@ export class WsConnection< const onTimeout = setTimeout(() => { const err = new this._errors.RequestTimeout({ - message: `Timed out while waiting for request in ${this._timeouts.request}ms`, + message: `Timed out while waiting for request in ${requestTimeout}ms`, requestId, }); done(err); - }, this._timeouts.request).unref(); + }, requestTimeout).unref(); function done(response: ResponseMessageType | WsError): void { if (isSettled) { @@ -440,26 +445,27 @@ export class WsConnection< private _closeWsConnection( sessionAbortMessage: string, status: WsConnectionStatus.CLOSED | WsConnectionStatus.DISCONNECTED, + force = false, ): void { const ws = this._wsConnection; - if (!ws || this._wsConnectionStatus === WsConnectionStatus.CLOSED) { + const isClosing = status === WsConnectionStatus.CLOSED; + if (this._wsConnectionStatus === WsConnectionStatus.CLOSED || (!ws && !isClosing)) { this._wsConnection = null; return; } this._debugFn(`\u2718 ${sessionAbortMessage}; endpoint: "${this._endpoint}"`); - const isClosing = status === WsConnectionStatus.CLOSED; - + // Сначала запрещаем reconnect, в том числе во время handshake и backoff. + this._wsConnectionStatus = status; + if (isClosing) this._connectionAbort.abort(); if (isClosing && this._onConnectionCloseFn) { - this._onConnectionCloseFn(); + this._onConnectionCloseFn(force); } - + this._pingHealthCheckStop(); this._wsConnection = null; - this._wsConnectionStatus = status; this._abortPendingRequests(`Request was aborted because ${sessionAbortMessage}`, isClosing); - this._pingHealthCheckStop(); if (isClosing) { this._onClose?.(); @@ -467,7 +473,7 @@ export class WsConnection< this._onDisconnect?.(); } - closeWsConnection(ws); + if (ws) closeWsConnection(ws, force); } /** @@ -488,11 +494,12 @@ export class WsConnection< } /** @description Closes websocket connection, terminating all pending requests */ - close(): void { - this._closeWsConnection("Connection was closed manually", WsConnectionStatus.CLOSED); + close(force = false): void { + this._closeWsConnection("Connection was closed manually", WsConnectionStatus.CLOSED, force); } private _pingHealthCheckStop(): void { + clearTimeout(this._pongTimeout); this._pingSubsequentFails = 0; if (this._pingInterval) { @@ -547,7 +554,7 @@ export class WsConnection< return; } - pongTimeout = setTimeout(() => { + pongTimeout = this._pongTimeout = setTimeout(() => { if (isWaitingForPong && this._isWebSocketActive(ws)) { isWaitingForPong = false; diff --git a/src/ws-connection/utils.ts b/src/ws-connection/utils.ts index badcef9e9..0a0df160d 100644 --- a/src/ws-connection/utils.ts +++ b/src/ws-connection/utils.ts @@ -1,10 +1,21 @@ +import { setTimeout as delayWithSignal } from "node:timers/promises"; + export const exponentiallyWait = ({ baseDelay = 500, attempt = 0, factor = 2, jitter = 100, -}: { baseDelay?: number; attempt?: number; factor?: number; jitter?: number } = {}): Promise => { + signal, +}: { + baseDelay?: number; + attempt?: number; + factor?: number; + jitter?: number; + signal?: AbortSignal; +} = {}): Promise => { const delay = Math.round(baseDelay * factor ** attempt + Math.random() * jitter); + if (signal) return delayWithSignal(delay, undefined, { signal }); + return new Promise(resolve => setTimeout(resolve, delay)); }; diff --git a/test/src/browser/existing-browser.js b/test/src/browser/existing-browser.js index a6cbb730a..69c9c4469 100644 --- a/test/src/browser/existing-browser.js +++ b/test/src/browser/existing-browser.js @@ -836,7 +836,11 @@ describe("ExistingBrowser", () => { sessionCaps: { "se:wsdriver": "ws://grid.url/session/test", "se:wsdriverVersion": "1" }, headers: { "X-Custom-Header": "test" }, browserConfig: browser.config, + onRequestDeadline: sinon.match.func, }); + const markAsBroken = sandbox.stub(browser, "markAsBroken"); + WSDriverRequestAgentCreateStub.lastCall.args[0].onRequestDeadline(); + assert.calledOnceWith(markAsBroken, { stubBrowserCommands: true }); }); it("should set customWdRequestAgent if se:wsdriverVersion includes multiple versions with version 1", async () => { diff --git a/test/src/browser/wsdriver/deadline.ts b/test/src/browser/wsdriver/deadline.ts new file mode 100644 index 000000000..e7957692f --- /dev/null +++ b/test/src/browser/wsdriver/deadline.ts @@ -0,0 +1,167 @@ +import { URL } from "node:url"; +import { createServer, Server } from "node:http"; +import { AddressInfo, Socket } from "node:net"; +import { once } from "node:events"; +import { setTimeout as delay } from "node:timers/promises"; +import { WebSocketServer, WebSocket } from "ws"; +import sinon from "sinon"; +import { WSDriverRequestAgent } from "src/browser/wsdriver"; +import { WSDriverRequestDeadlineError } from "src/browser/wsdriver/error"; +import { BrowserConfig } from "src/config/browser-config"; + +const flag = "TESTPLANE_WSDRIVER_DEADLINE_ENABLED"; + +describe("WSDriver command deadline", () => { + let oldFlag: string | undefined; + let server: Server; + let wsServer: WebSocketServer; + let agent: WSDriverRequestAgent; + let requests: number; + let connections: number; + let sockets: Set; + let onDeadline: sinon.SinonSpy; + + beforeEach(() => { + oldFlag = process.env[flag]; + process.env[flag] = "true"; + requests = 0; + connections = 0; + sockets = new Set(); + onDeadline = sinon.spy(); + }); + + afterEach(async () => { + agent?.close(); + for (const socket of sockets) socket.destroy(); + if (wsServer) await new Promise(resolve => wsServer.close(() => resolve())); + if (server) await new Promise(resolve => server.close(() => resolve())); + if (oldFlag === undefined) delete process.env[flag]; + else process.env[flag] = oldFlag; + }); + + async function start( + handle?: (ws: WebSocket, request: Buffer) => void, + { handshake = true, httpTimeout = 100 }: { handshake?: boolean; httpTimeout?: number } = {}, + ): Promise { + server = createServer(); + server.on("connection", socket => { + sockets.add(socket); + socket.on("close", () => sockets.delete(socket)); + socket.on("end", () => socket.destroy()); + }); + wsServer = new WebSocketServer({ noServer: true }); + server.on("upgrade", (request, socket, head) => { + connections++; + socket.resume(); + if (handshake) { + wsServer.handleUpgrade(request, socket, head, ws => { + wsServer.emit("connection", ws); + }); + } + }); + wsServer.on("connection", ws => { + ws.on("message", message => { + requests++; + handle?.(ws, message as Buffer); + }); + }); + server.listen(0, "127.0.0.1"); + await once(server, "listening"); + const port = (server.address() as AddressInfo).port; + agent = WSDriverRequestAgent.create({ + sessionId: "session", + sessionCaps: { "se:wsdriver": `ws://127.0.0.1:${port}`, "se:wsdriverVersion": "1" }, + headers: {}, + browserConfig: { httpTimeout } as BrowserConfig, + onRequestDeadline: onDeadline, + }); + } + + function request(responseTimeout?: number): ReturnType { + return agent.request(new URL("http://localhost/session/session/title"), { + method: "GET", + ...(responseTimeout === undefined ? {} : { timeout: { response: responseTimeout } }), + }); + } + + function respond(ws: WebSocket, request: Buffer): void { + const response = Buffer.concat([request.subarray(0, 8), Buffer.from('title\0{"value":"ok"}')]); + response.writeUInt8(18, 1); + response.writeUInt16BE(200, 6); + ws.send(response); + } + + it("reuses a healthy connection and does not mark the session broken", async () => { + await start(respond); + assert.equal((await request()).statusCode, 200); + assert.equal((await request()).statusCode, 200); + assert.equal(connections, 1); + assert.equal(requests, 2); + assert.equal(onDeadline.callCount, 0); + }); + + it("aborts a pending handshake without sending commands or reconnecting", async () => { + await start(undefined, { handshake: false }); + await assert.isRejected(request(), WSDriverRequestDeadlineError); + await delay(50); + assert.equal(sockets.size, 0); + assert.equal(connections, 1); + assert.equal(requests, 0); + assert.equal(onDeadline.callCount, 1); + }); + + it("terminates a hung request and rejects subsequent commands on that session", async () => { + await start(); + const error = await request().catch(error => error); + assert.instanceOf(error, WSDriverRequestDeadlineError); + assert.equal(error.code, "WSDRIVER_REQUEST_DEADLINE"); + await assert.isRejected(request(), WSDriverRequestDeadlineError); + await delay(50); + assert.equal(sockets.size, 0); + assert.equal(requests, 1); + assert.equal(onDeadline.callCount, 1); + }); + + it("includes reconnect and request retry backoff in the same deadline", async () => { + await start(ws => ws.terminate()); + const started = performance.now(); + await assert.isRejected(request(), WSDriverRequestDeadlineError); + assert.isBelow(performance.now() - started, 1000); + const sentBeforeDeadline = requests; + await delay(150); + assert.equal(requests, sentBeforeDeadline); + assert.equal(sockets.size, 0); + assert.equal(onDeadline.callCount, 1); + }); + + it("uses the per-command response timeout", async () => { + await start(undefined, { httpTimeout: 5000 }); + await assert.isRejected(request(100), "timed out after 100ms"); + assert.equal(onDeadline.callCount, 1); + }); + + for (const value of [undefined, "false"]) { + it(`preserves the legacy timeout with flag=${value}`, async () => { + if (value === undefined) delete process.env[flag]; + else process.env[flag] = value; + await start(); + const error = await request().catch(error => error); + assert.equal(error.code, "ETIMEDOUT"); + assert.equal(onDeadline.callCount, 0); + }); + + it(`closes an unfinished handshake with flag=${value}`, async () => { + if (value === undefined) delete process.env[flag]; + else process.env[flag] = value; + await start(undefined, { handshake: false, httpTimeout: 5000 }); + const pending = request().catch(error => error); + while (connections === 0) await delay(5); + agent.close(); + assert.equal((await pending).name, "WSDriverRequestAgentTerminatedError"); + await delay(50); + assert.equal(sockets.size, 0); + assert.equal(requests, 0); + assert.equal(connections, 1); + }); + } +});