Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 36 additions & 0 deletions scripts/test-ipc.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -201,6 +236,7 @@ async function main() {
console.log('Connected (handshake v3).');

await testValues(client);
await testConcurrency(client);
client.disconnect();

console.log('Auth handshake:');
Expand Down
168 changes: 125 additions & 43 deletions src/rayforceIpc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<typeof setTimeout> | null = null;

constructor(host: string, port: number) {
this.host = host;
this.port = port;
Expand Down Expand Up @@ -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].
Expand Down Expand Up @@ -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();
Expand All @@ -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<RayforceValue> {
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<void> {
if (!this.isConnected()) {
throw new Error('Not connected to Rayforce instance');
Expand Down Expand Up @@ -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;
Expand All @@ -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;
}
}
}
Expand Down