Skip to content
Open
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
5 changes: 5 additions & 0 deletions .changeset/readbuffer-linear-append.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@modelcontextprotocol/core-internal': patch
---

Eliminate quadratic copying and rescanning in `ReadBuffer` when receiving messages across multiple chunks. Incoming chunks are retained without intermediate concatenation, and message delimiter scanning inspects newly appended bytes rather than repeatedly rescanning earlier buffered content.
90 changes: 81 additions & 9 deletions packages/core-internal/src/shared/stdio.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,31 +7,100 @@ export const STDIO_DEFAULT_MAX_BUFFER_SIZE = 10 * 1024 * 1024;
* Buffers a continuous stdio stream into discrete JSON-RPC messages.
*/
export class ReadBuffer {
private _buffer?: Buffer;
private _chunks: Buffer[] = [];
private _totalLength = 0;
private _newlineChunkIndex = -1;
private _newlineOffsetInChunk = -1;
private _maxBufferSize: number;

constructor(options?: { maxBufferSize?: number }) {
this._maxBufferSize = options?.maxBufferSize ?? STDIO_DEFAULT_MAX_BUFFER_SIZE;
}

append(chunk: Buffer): void {
const newSize = (this._buffer?.length ?? 0) + chunk.length;
if (chunk.length === 0) {
return;
}

const newSize = this._totalLength + chunk.length;
if (newSize > this._maxBufferSize) {
this.clear();
throw new Error(`ReadBuffer exceeded maximum size of ${this._maxBufferSize} bytes`);
}
this._buffer = this._buffer ? Buffer.concat([this._buffer, chunk]) : chunk;

if (this._newlineChunkIndex === -1) {
const index = chunk.indexOf(0x0a);
if (index !== -1) {
this._newlineChunkIndex = this._chunks.length;
this._newlineOffsetInChunk = index;
}
}

this._chunks.push(chunk);
this._totalLength = newSize;
}

readMessage(): JSONRPCMessage | null {
while (this._buffer) {
const index = this._buffer.indexOf('\n');
if (index === -1) {
while (this._newlineChunkIndex !== -1) {
const targetChunkIndex = this._newlineChunkIndex;
const newlineOffset = this._newlineOffsetInChunk;
const targetChunk = this._chunks[targetChunkIndex];

if (!targetChunk) {
this.clear();
return null;
}

const line = this._buffer.toString('utf8', 0, index).replace(/\r$/, '');
this._buffer = this._buffer.subarray(index + 1);
const messageParts: Buffer[] = [];
let consumedBytes = newlineOffset + 1;
for (let i = 0; i < targetChunkIndex; i++) {
const chunk = this._chunks[i];
if (chunk) {
messageParts.push(chunk);
consumedBytes += chunk.length;
}
}
if (newlineOffset > 0) {
messageParts.push(targetChunk.subarray(0, newlineOffset));
}

const messageBuffer =
messageParts.length === 1
? (messageParts[0] ?? Buffer.alloc(0))
: messageParts.length === 0
? Buffer.alloc(0)
: Buffer.concat(messageParts);

const remainder = targetChunk.subarray(newlineOffset + 1);
const remainingChunks: Buffer[] = [];
if (remainder.length > 0) {
remainingChunks.push(remainder);
}
for (let i = targetChunkIndex + 1; i < this._chunks.length; i++) {
const chunk = this._chunks[i];
if (chunk) {
remainingChunks.push(chunk);
}
}

this._chunks = remainingChunks;
this._totalLength -= consumedBytes;

this._newlineChunkIndex = -1;
this._newlineOffsetInChunk = -1;
for (let i = 0; i < this._chunks.length; i++) {
const chunk = this._chunks[i];
if (chunk) {
const index = chunk.indexOf(0x0a);
if (index !== -1) {
this._newlineChunkIndex = i;
this._newlineOffsetInChunk = index;
break;
}
}
}

const line = messageBuffer.toString('utf8').replace(/\r$/, '');

try {
return deserializeMessage(line);
Expand All @@ -49,7 +118,10 @@ export class ReadBuffer {
}

clear(): void {
this._buffer = undefined;
this._chunks = [];
this._totalLength = 0;
this._newlineChunkIndex = -1;
this._newlineOffsetInChunk = -1;
}
}

Expand Down
86 changes: 86 additions & 0 deletions packages/core-internal/test/shared/stdio.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -156,3 +156,89 @@ describe('buffer size limit', () => {
expect(readBuffer.readMessage()).not.toBeNull();
});
});

describe('chunked message reading', () => {
test('should handle large messages received across many small chunks without quadratic overhead', () => {
const readBuffer = new ReadBuffer();
const payloadText = 'a'.repeat(2 * 1024 * 1024); // 2 MB
const message: JSONRPCMessage = {
jsonrpc: '2.0',
id: 'large-msg',
result: { text: payloadText }
};
const raw = Buffer.from(JSON.stringify(message) + '\n');
const chunkSize = 1024; // 1 KB chunks (over 2000 chunks)

for (let offset = 0; offset < raw.length; offset += chunkSize) {
readBuffer.append(raw.subarray(offset, Math.min(offset + chunkSize, raw.length)));
// Read before final chunk should be null
if (offset + chunkSize < raw.length) {
expect(readBuffer.readMessage()).toBeNull();
}
}

const parsed = readBuffer.readMessage();
expect(parsed).toEqual(message);
expect(readBuffer.readMessage()).toBeNull();
});

test('should extract multiple messages delivered in a single chunk', () => {
const readBuffer = new ReadBuffer();
const msg1: JSONRPCMessage = { jsonrpc: '2.0', method: 'm1' };
const msg2: JSONRPCMessage = { jsonrpc: '2.0', method: 'm2' };
const msg3: JSONRPCMessage = { jsonrpc: '2.0', method: 'm3' };

const chunk = Buffer.from(JSON.stringify(msg1) + '\n' + JSON.stringify(msg2) + '\n' + JSON.stringify(msg3) + '\n');

readBuffer.append(chunk);
expect(readBuffer.readMessage()).toEqual(msg1);
expect(readBuffer.readMessage()).toEqual(msg2);
expect(readBuffer.readMessage()).toEqual(msg3);
expect(readBuffer.readMessage()).toBeNull();
});

test('should handle consecutive messages split unevenly across chunks', () => {
const readBuffer = new ReadBuffer();
const msg1: JSONRPCMessage = { jsonrpc: '2.0', method: 'first' };
const msg2: JSONRPCMessage = { jsonrpc: '2.0', method: 'second' };

const str1 = JSON.stringify(msg1) + '\n';
const str2 = JSON.stringify(msg2) + '\n';
const full = Buffer.from(str1 + str2);

// Split into 3 arbitrary chunks:
// Chunk 1: middle of msg1
// Chunk 2: end of msg1 + newline + start of msg2
// Chunk 3: rest of msg2 + newline
const split1 = Math.floor(str1.length / 2);
const split2 = str1.length + Math.floor(str2.length / 2);

readBuffer.append(full.subarray(0, split1));
expect(readBuffer.readMessage()).toBeNull();

readBuffer.append(full.subarray(split1, split2));
expect(readBuffer.readMessage()).toEqual(msg1);
expect(readBuffer.readMessage()).toBeNull();

readBuffer.append(full.subarray(split2));
expect(readBuffer.readMessage()).toEqual(msg2);
expect(readBuffer.readMessage()).toBeNull();
});

test('should handle empty chunks gracefully', () => {
const readBuffer = new ReadBuffer();
const msg: JSONRPCMessage = { jsonrpc: '2.0', method: 'test' };

readBuffer.append(Buffer.alloc(0));
expect(readBuffer.readMessage()).toBeNull();

readBuffer.append(Buffer.from(JSON.stringify(msg)));
readBuffer.append(Buffer.alloc(0));
expect(readBuffer.readMessage()).toBeNull();

readBuffer.append(Buffer.from('\n'));
readBuffer.append(Buffer.alloc(0));
expect(readBuffer.readMessage()).toEqual(msg);
expect(readBuffer.readMessage()).toBeNull();
});
});
Loading