From 6179a9ae74f311bed665adf0f7308d8438e51bde Mon Sep 17 00:00:00 2001 From: Evgen Byelozorov Date: Tue, 21 Jul 2026 09:04:55 +0200 Subject: [PATCH] fix(ipc): serialize concurrent execute() calls through a request queue The wire protocol carries no request id: a RESP is correlated to a SYNC request purely by being the next thing to arrive on the socket after it was sent, so at most one SYNC request may be outstanding at a time. But RayforceIpcClient.execute() had a single pendingResolve/pendingReject slot with no such guard, so a second concurrent call (a double Enter in the REPL, a pagination click landing while refreshEnv() still has a request in flight) silently overwrote the first call's callbacks: whichever response arrived first got handed to the wrong caller, and the caller whose callbacks got overwritten just hung until its own timeout. Reworked execute() around a FIFO request queue: each call enqueues and calls pumpQueue(), which sends the next request only once the current one has settled (response, timeout, or write error, via finishCurrent()). finishCurrent() checks that the request it's asked to settle is still the current one, so a stale timeout firing after rejectAll() already cleared it (disconnect mid-flight) is a no-op rather than a double-settle. disconnect() and the socket 'close' handler now drain the whole queue through rejectAll(), not just the one in-flight request, so queued callers fail immediately instead of hanging until their own timeout on a dead socket. scripts/test-ipc.js: added a concurrency test that fires 10 execute() calls back-to-back without awaiting and asserts each resolves to its own result. Confirmed this reproduces the bug against the pre-fix client (request #9's promise settles with request #0's result "0n", requests #0-8 hang to timeout) and passes against the fix. Full suite: 33/33. Co-authored-by: Jack Bell <6161232+belljack@users.noreply.github.com> --- scripts/test-ipc.js | 36 ++++++++++ src/rayforceIpc.ts | 168 ++++++++++++++++++++++++++++++++------------ 2 files changed, 161 insertions(+), 43 deletions(-) diff --git a/scripts/test-ipc.js b/scripts/test-ipc.js index 4cd0fa9..e7b141d 100644 --- a/scripts/test-ipc.js +++ b/scripts/test-ipc.js @@ -153,6 +153,41 @@ async function testValues(client) { wrapped[2] && wrapped[2]._type === 'table' && wrapped[2].values[0].length === 10, show(wrapped && wrapped[0])); } +async function testConcurrency(client) { + const N = 10; + const promises = []; + for (let i = 0; i < N; i++) { + // Deliberately not awaited here: firing all N execute() calls back + // to back (before any response comes back) is exactly the pattern + // that hit the single pendingResolve/pendingReject slot — a double + // Enter in the REPL, or a pagination click while refreshEnv() is + // still awaiting its own request. A short per-call timeout keeps a + // regression (misattributed/dropped responses hang until timeout) + // from stalling the whole suite. + promises.push(client.execute(`(* ${i} ${i})`, 3000)); + } + + const results = await Promise.allSettled(promises); + + let allMatched = true; + const detail = []; + for (let i = 0; i < N; i++) { + const r = results[i]; + const want = BigInt(i * i); + const got = r.status === 'fulfilled' ? r.value : `rejected: ${r.reason && r.reason.message}`; + if (!(r.status === 'fulfilled' && got === want)) { + allMatched = false; + } + detail.push(`${i}->${show(got)}`); + } + ok(`${N} concurrent execute() calls each resolve to their own result`, allMatched, detail.join(' ')); + + // The queue must also drain fully afterward — a stuck currentRequest + // would silently swallow every call made after this test + const after = await client.execute('(+ 100 1)'); + ok('client usable after concurrent burst', after === 101n, show(after)); +} + async function testAuth(port) { const noCreds = new RayforceIpcClient(HOST, port); try { @@ -201,6 +236,7 @@ async function main() { console.log('Connected (handshake v3).'); await testValues(client); + await testConcurrency(client); client.disconnect(); console.log('Auth handshake:'); diff --git a/src/rayforceIpc.ts b/src/rayforceIpc.ts index 92445c1..56d2f75 100644 --- a/src/rayforceIpc.ts +++ b/src/rayforceIpc.ts @@ -687,6 +687,13 @@ class Deserializer { // IPC Client // ============================================================================ +interface PendingRequest { + statement: string; + timeout: number; + resolve: (value: RayforceValue) => void; + reject: (reason: Error) => void; +} + export class RayforceIpcClient { private socket: net.Socket | null = null; private host: string; @@ -696,6 +703,17 @@ export class RayforceIpcClient { private pendingResolve: ((value: RayforceValue) => void) | null = null; private pendingReject: ((reason: Error) => void) | null = null; + // The wire protocol carries no request id: a RESP is correlated to a + // request purely by being the next thing to arrive on the socket after + // it was sent. So at most one SYNC request may be in flight at a time; + // concurrent execute() callers (e.g. a second command submitted before + // the first one's result comes back) are queued and sent one by one, + // instead of clobbering the single pendingResolve/pendingReject slot + // and misattributing the next response to the wrong caller. + private requestQueue: PendingRequest[] = []; + private currentRequest: PendingRequest | null = null; + private currentTimeoutHandle: ReturnType | null = null; + constructor(host: string, port: number) { this.host = host; this.port = port; @@ -739,12 +757,8 @@ export class RayforceIpcClient { settle(() => { reject(new Error('Connection closed')); }); - // Also handle pending execute operations - if (this.pendingReject) { - this.pendingReject(new Error('Connection closed')); - this.pendingResolve = null; - this.pendingReject = null; - } + // Also settle the in-flight request and drain the queue + this.rejectAll(new Error('Connection closed')); }); // Handshake: send [version, 0]; server replies [version, auth_required]. @@ -832,13 +846,11 @@ export class RayforceIpcClient { } disconnect(): void { - // Clear pending operations first - if (this.pendingReject) { - this.pendingReject(new Error('Disconnected')); - } - this.pendingResolve = null; - this.pendingReject = null; - + // Clear the in-flight request and drain anything still queued — + // otherwise a queued caller would hang until its own timeout fires + // on a socket that will never receive a reply. + this.rejectAll(new Error('Disconnected')); + if (this.socket) { this.socket.removeAllListeners(); this.socket.destroy(); @@ -852,42 +864,111 @@ export class RayforceIpcClient { return this.connected && this.socket !== null; } + /** + * Reject the in-flight request (if any) and every queued request. + * Clears state directly rather than going through the pendingReject + * wire-closure, so this never re-enters pumpQueue() and tries to write + * to a socket that's mid-teardown — after a disconnect there is nothing + * left to pump until a future execute() call starts a new request. + */ + private rejectAll(err: Error): void { + if (this.currentTimeoutHandle) { + clearTimeout(this.currentTimeoutHandle); + this.currentTimeoutHandle = null; + } + const current = this.currentRequest; + this.currentRequest = null; + this.pendingResolve = null; + this.pendingReject = null; + + const queued = this.requestQueue; + this.requestQueue = []; + + if (current) { + current.reject(err); + } + for (const req of queued) { + req.reject(err); + } + } + async execute(statement: string, timeout: number = 30000): Promise { if (!this.isConnected()) { throw new Error('Not connected to Rayforce instance'); } return new Promise((resolve, reject) => { - const timeoutHandle = setTimeout(() => { - this.pendingResolve = null; - this.pendingReject = null; - reject(new Error(`Execution timeout after ${timeout}ms`)); - }, timeout); + this.requestQueue.push({ statement, timeout, resolve, reject }); + this.pumpQueue(); + }); + } - this.pendingResolve = (value) => { - clearTimeout(timeoutHandle); - resolve(value); - }; - - this.pendingReject = (err) => { - clearTimeout(timeoutHandle); - reject(err); - }; + /** + * Send the next queued request, if the wire is free. Only one SYNC + * request may be outstanding at a time (see the requestQueue comment + * above), so this is a no-op while currentRequest is set; whichever + * settles the current request (response, timeout, or write error) calls + * finishCurrent(), which advances the queue. + */ + private pumpQueue(): void { + if (this.currentRequest || this.requestQueue.length === 0) { + return; + } + if (!this.isConnected()) { + // Socket died between enqueue and pump. rejectAll() already + // drains the queue on disconnect/'close', so this is normally + // unreachable — kept as a guard since this.socket!.write below + // would otherwise throw on a null socket. + this.rejectAll(new Error('Not connected to Rayforce instance')); + return; + } - const payload = Serializer.serializeString(statement); - const message = Serializer.createMessage(payload, MSG_TYPE_SYNC); + const request = this.requestQueue.shift()!; + this.currentRequest = request; - this.socket!.write(message, (err) => { - if (err) { - clearTimeout(timeoutHandle); - this.pendingResolve = null; - this.pendingReject = null; - reject(err); - } - }); + this.currentTimeoutHandle = setTimeout(() => { + this.finishCurrent(request, () => request.reject(new Error(`Execution timeout after ${request.timeout}ms`))); + }, request.timeout); + + this.pendingResolve = (value) => { + this.finishCurrent(request, () => request.resolve(value)); + }; + + this.pendingReject = (err) => { + this.finishCurrent(request, () => request.reject(err)); + }; + + const payload = Serializer.serializeString(request.statement); + const message = Serializer.createMessage(payload, MSG_TYPE_SYNC); + + this.socket!.write(message, (err) => { + if (err) { + this.finishCurrent(request, () => request.reject(err)); + } }); } + /** + * Settle `request` and advance the queue. A no-op if `request` is no + * longer the current one — guards against a stale timeout or write-error + * callback firing after rejectAll() already cleared it (e.g. the socket + * closed and a new request has since started). + */ + private finishCurrent(request: PendingRequest, fn: () => void): void { + if (this.currentRequest !== request) { + return; + } + if (this.currentTimeoutHandle) { + clearTimeout(this.currentTimeoutHandle); + this.currentTimeoutHandle = null; + } + this.currentRequest = null; + this.pendingResolve = null; + this.pendingReject = null; + fn(); + this.pumpQueue(); + } + async executeAsync(statement: string): Promise { if (!this.isConnected()) { throw new Error('Not connected to Rayforce instance'); @@ -925,10 +1006,13 @@ export class RayforceIpcClient { if (!header) break; if (header.prefix !== SERDE_PREFIX) { + // pendingReject (if set) is the finishCurrent-wrapped closure + // installed by pumpQueue(): invoking it clears the fields and + // advances to the next queued request itself, so this must + // not also null them afterward — that would stomp on the + // next request's just-installed pendingResolve/pendingReject. if (this.pendingReject) { this.pendingReject(new Error('Invalid response prefix')); - this.pendingResolve = null; - this.pendingReject = null; } this.responseBuffer = Buffer.alloc(0); break; @@ -947,16 +1031,14 @@ export class RayforceIpcClient { const deserializer = new Deserializer(payload); const value = deserializer.deserialize(); + // See the comment on the invalid-prefix branch above: these + // closures own clearing/advancing the queue themselves. if (header.msgtype === MSG_TYPE_RESP && this.pendingResolve) { this.pendingResolve(value); - this.pendingResolve = null; - this.pendingReject = null; } } catch (err) { if (this.pendingReject) { this.pendingReject(err instanceof Error ? err : new Error(String(err))); - this.pendingResolve = null; - this.pendingReject = null; } } }