diff --git a/src/shared/stdio.ts b/src/shared/stdio.ts index 8f5478dfa8..905c9474dd 100644 --- a/src/shared/stdio.ts +++ b/src/shared/stdio.ts @@ -33,7 +33,8 @@ export class ReadBuffer { } const line = this._buffer.toString('utf8', 0, index).replace(/\r$/, ''); - this._buffer = this._buffer.subarray(index + 1); + const remainder = this._buffer.subarray(index + 1); + this._buffer = remainder.length === 0 ? undefined : remainder; return deserializeMessage(line); } diff --git a/test/shared/stdio.test.ts b/test/shared/stdio.test.ts index 87e666f7d0..4f4df52b37 100644 --- a/test/shared/stdio.test.ts +++ b/test/shared/stdio.test.ts @@ -1,3 +1,4 @@ +import { vi } from 'vitest'; import { JSONRPCMessage } from '../../src/types.js'; import { STDIO_DEFAULT_MAX_BUFFER_SIZE, ReadBuffer } from '../../src/shared/stdio.js'; @@ -34,6 +35,53 @@ test('should be reusable after clearing', () => { expect(readBuffer.readMessage()).toEqual(testMessage); }); +describe('consumed buffer', () => { + afterEach(() => { + vi.restoreAllMocks(); + }); + + test('should release the backing allocation after the final newline', () => { + const readBuffer = new ReadBuffer(); + const message = Buffer.from(JSON.stringify(testMessage) + '\n'); + const allocation = Buffer.alloc(64 * 1024); + message.copy(allocation); + readBuffer.append(allocation.subarray(0, message.length)); + + expect(readBuffer.readMessage()).toEqual(testMessage); + expect(readBuffer).toHaveProperty('_buffer', undefined); + expect(readBuffer.readMessage()).toBeNull(); + }); + + test('should append the next chunk without copying after draining', () => { + const readBuffer = new ReadBuffer(); + readBuffer.append(Buffer.from(JSON.stringify(testMessage) + '\n')); + expect(readBuffer.readMessage()).toEqual(testMessage); + + const nextChunk = Buffer.from(JSON.stringify(testMessage) + '\n'); + const concatSpy = vi.spyOn(Buffer, 'concat'); + readBuffer.append(nextChunk); + + expect(concatSpy).not.toHaveBeenCalled(); + expect(readBuffer.readMessage()).toEqual(testMessage); + expect(readBuffer.readMessage()).toBeNull(); + }); + + test('should preserve a partial message after a complete message', () => { + const readBuffer = new ReadBuffer(); + const nextMessage: JSONRPCMessage = { jsonrpc: '2.0', method: 'second' }; + const message = JSON.stringify(nextMessage); + const split = Math.floor(message.length / 2); + readBuffer.append(Buffer.from(JSON.stringify(testMessage) + '\n' + message.slice(0, split))); + + expect(readBuffer.readMessage()).toEqual(testMessage); + expect(readBuffer.readMessage()).toBeNull(); + + readBuffer.append(Buffer.from(message.slice(split) + '\r\n')); + expect(readBuffer.readMessage()).toEqual(nextMessage); + expect(readBuffer.readMessage()).toBeNull(); + }); +}); + describe('buffer size limit', () => { test('should throw when buffer exceeds default max size', () => { const readBuffer = new ReadBuffer();