import * as plugins from './plugins.js'; import { flexIpcMaximumFrameBytes } from './interfaces.flexipc.js'; export interface IFlexFramedTransportOptions { readable: plugins.stream.Readable; writable: plugins.stream.Writable; parseIncoming: (valueArg: unknown) => TIncoming; validateOutgoing?: (valueArg: unknown) => TOutgoing; maximumFrameBytes?: number; } type TMessageListener = (messageArg: TMessage) => void; type TErrorListener = (errorArg: Error) => void; type TCloseListener = () => void; export const flexIpcMaximumQueuedFrames = 128; export const flexIpcMaximumQueuedFrameBytes = 64 * 1024 * 1024; const flexIpcCloseDrainTimeoutMs = 2_000; const normalizeMaximumFrameBytes = (valueArg: number | undefined): number => { const selected = valueArg ?? flexIpcMaximumFrameBytes; if (!Number.isSafeInteger(selected) || selected < 1 || selected > flexIpcMaximumFrameBytes) { throw new Error(`maximumFrameBytes must be between 1 and ${flexIpcMaximumFrameBytes}.`); } return selected; }; export const encodeFlexFrame = ( valueArg: unknown, maximumFrameBytesArg = flexIpcMaximumFrameBytes, ): Buffer => { const serialized = Buffer.from(JSON.stringify(valueArg), 'utf8'); try { if (serialized.byteLength > maximumFrameBytesArg) { throw new Error('Flex IPC frame exceeds its size limit.'); } const frame = Buffer.allocUnsafe(4 + serialized.byteLength); frame.writeUInt32BE(serialized.byteLength, 0); serialized.copy(frame, 4); return frame; } finally { serialized.fill(0); } }; /** * A single-reader/single-writer JSON transport. The body buffer is allocated * only after the fixed header has passed the configured length limit. */ export class FlexFramedTransport { private readonly readable: plugins.stream.Readable; private readonly writable: plugins.stream.Writable; private readonly parseIncoming: (valueArg: unknown) => TIncoming; private readonly validateOutgoing?: (valueArg: unknown) => TOutgoing; private readonly maximumFrameBytes: number; private readonly header = Buffer.allocUnsafe(4); private headerOffset = 0; private body?: Buffer; private bodyOffset = 0; private messageListeners = new Set>(); private errorListeners = new Set(); private closeListeners = new Set(); private writeTail = Promise.resolve(); private queuedFrames = 0; private queuedFrameBytes = 0; private failed = false; private closed = false; private closePromise?: Promise; constructor(optionsArg: IFlexFramedTransportOptions) { this.readable = optionsArg.readable; this.writable = optionsArg.writable; this.parseIncoming = optionsArg.parseIncoming; this.validateOutgoing = optionsArg.validateOutgoing; this.maximumFrameBytes = normalizeMaximumFrameBytes(optionsArg.maximumFrameBytes); this.readable.on('data', this.onData); this.readable.once('error', this.onReadError); this.readable.once('end', this.onReadableClosed); this.readable.once('close', this.onReadableClosed); this.writable.once('error', this.onWriteError); this.writable.once('finish', this.onWritableClosed); this.writable.once('close', this.onWritableClosed); } public onMessage(listenerArg: TMessageListener): () => void { this.messageListeners.add(listenerArg); return () => this.messageListeners.delete(listenerArg); } public onError(listenerArg: TErrorListener): () => void { this.errorListeners.add(listenerArg); return () => this.errorListeners.delete(listenerArg); } public onClose(listenerArg: TCloseListener): () => void { this.closeListeners.add(listenerArg); return () => this.closeListeners.delete(listenerArg); } public async send(messageArg: TOutgoing): Promise { if (this.closed || this.failed) throw new Error('Flex IPC transport is closed.'); this.validateOutgoing?.(messageArg); const frame = encodeFlexFrame(messageArg, this.maximumFrameBytes); if ( this.queuedFrames >= flexIpcMaximumQueuedFrames || this.queuedFrameBytes + frame.byteLength > flexIpcMaximumQueuedFrameBytes ) throw new Error('Flex IPC outbound queue limit exceeded.'); this.queuedFrames++; this.queuedFrameBytes += frame.byteLength; const write = this.writeTail.then(() => this.writeFrame(frame)); this.writeTail = write.catch(() => undefined); try { await write; } finally { this.queuedFrames--; this.queuedFrameBytes -= frame.byteLength; frame.fill(0); } } public close(): Promise { if (this.closePromise) return this.closePromise; if (this.closed) return Promise.resolve(); this.closePromise = this.closeTransport(); return this.closePromise; } private async closeTransport(): Promise { this.closed = true; this.detachStreams(); let drainTimer: ReturnType | undefined; const drainTimeout = new Promise<'timeout'>((resolve) => { drainTimer = setTimeout(() => resolve('timeout'), flexIpcCloseDrainTimeoutMs); drainTimer.unref?.(); }); const gracefulClose = this.writeTail.then(async () => { if (!this.writable.destroyed && !this.writable.writableEnded) { await new Promise((resolve) => { this.writable.end(resolve); }); } return 'drained' as const; }); const drainResult = await Promise.race([ gracefulClose, drainTimeout, ]); if (drainTimer) clearTimeout(drainTimer); if (drainResult === 'timeout') { if (!this.writable.destroyed) this.writable.destroy(); if (!this.readable.destroyed) this.readable.destroy(); this.emitClose(); return; } if (!this.readable.destroyed) this.readable.destroy(); this.emitClose(); } private readonly onData = (chunkArg: unknown): void => { if (this.closed || this.failed) return; const chunk = typeof chunkArg === 'string' ? Buffer.from(chunkArg, 'utf8') : chunkArg instanceof Uint8Array ? Buffer.from(chunkArg.buffer, chunkArg.byteOffset, chunkArg.byteLength) : undefined; if (!chunk) { this.fail(new Error('Flex IPC received a non-byte chunk.')); return; } let chunkOffset = 0; try { while (chunkOffset < chunk.byteLength) { if (!this.body) { const headerBytes = Math.min(4 - this.headerOffset, chunk.byteLength - chunkOffset); chunk.copy(this.header, this.headerOffset, chunkOffset, chunkOffset + headerBytes); this.headerOffset += headerBytes; chunkOffset += headerBytes; if (this.headerOffset < 4) continue; const bodyLength = this.header.readUInt32BE(0); if (bodyLength > this.maximumFrameBytes) { throw new Error('Flex IPC frame exceeds its size limit.'); } this.body = Buffer.allocUnsafe(bodyLength); this.bodyOffset = 0; if (bodyLength === 0) this.finishBody(); continue; } const bodyBytes = Math.min( this.body.byteLength - this.bodyOffset, chunk.byteLength - chunkOffset, ); chunk.copy(this.body, this.bodyOffset, chunkOffset, chunkOffset + bodyBytes); this.bodyOffset += bodyBytes; chunkOffset += bodyBytes; if (this.bodyOffset === this.body.byteLength) this.finishBody(); } } catch (errorArg) { this.fail(errorArg instanceof Error ? errorArg : new Error('Flex IPC framing failed.')); } }; private finishBody(): void { const completedBody = this.body; this.body = undefined; this.bodyOffset = 0; this.headerOffset = 0; if (!completedBody) throw new Error('Flex IPC body state was lost.'); let parsed: unknown; try { parsed = JSON.parse(completedBody.toString('utf8')) as unknown; } catch { throw new Error('Flex IPC frame contains invalid JSON.'); } finally { completedBody.fill(0); } const message = this.parseIncoming(parsed); for (const listener of [...this.messageListeners]) listener(message); } private writeFrame(frameArg: Buffer): Promise { if (this.closed || this.failed || this.writable.destroyed || this.writable.writableEnded) { throw new Error('Flex IPC transport is closed.'); } return new Promise((resolve, reject) => { this.writable.write(frameArg, (errorArg?: Error | null) => { if (errorArg) reject(errorArg); else resolve(); }); }); } private readonly onReadError = (errorArg: Error): void => this.fail(errorArg); private readonly onWriteError = (errorArg: Error): void => this.fail(errorArg); private readonly onReadableClosed = (): void => { if (this.closed) return; if (this.headerOffset !== 0 || this.body) { this.fail(new Error('Flex IPC stream ended within a frame.')); return; } this.closed = true; this.detachStreams(); if (!this.readable.destroyed) this.readable.destroy(); if (!this.writable.destroyed) this.writable.destroy(); this.emitClose(); }; private readonly onWritableClosed = (): void => { if (this.closed) return; this.closed = true; this.detachStreams(); if (!this.readable.destroyed) this.readable.destroy(); if (!this.writable.destroyed) this.writable.destroy(); this.emitClose(); }; private fail(errorArg: Error): void { if (this.failed || this.closed) return; this.failed = true; this.detachStreams(); if (!this.readable.destroyed) this.readable.destroy(); if (!this.writable.destroyed) this.writable.destroy(); for (const listener of [...this.errorListeners]) listener(errorArg); this.emitClose(); } private detachStreams(): void { this.readable.removeListener('data', this.onData); this.readable.removeListener('error', this.onReadError); this.readable.removeListener('end', this.onReadableClosed); this.readable.removeListener('close', this.onReadableClosed); this.writable.removeListener('error', this.onWriteError); this.writable.removeListener('finish', this.onWritableClosed); this.writable.removeListener('close', this.onWritableClosed); } private emitClose(): void { const listeners = [...this.closeListeners]; this.closeListeners.clear(); for (const listener of listeners) listener(); } }