mirror of
https://github.com/stablyai/orca.git
synced 2026-09-22 16:02:32 +00:00
perf: count local RPC request bytes incrementally (#20243)
Co-authored-by: Orca Worker <orca-worker@localhost>
This commit is contained in:
@@ -1,4 +1,5 @@
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { StringDecoder } from 'node:string_decoder'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { Socket } from 'node:net'
|
||||
import { UnixSocketTransport } from './unix-socket-transport'
|
||||
@@ -11,6 +12,7 @@ class FakeSocket extends EventEmitter {
|
||||
setEncoding(): void {}
|
||||
setNoDelay(): void {}
|
||||
setTimeout(): void {}
|
||||
end(): void {}
|
||||
|
||||
write(data: string): boolean {
|
||||
this.writes.push(data)
|
||||
@@ -40,6 +42,47 @@ describe('UnixSocketTransport', () => {
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
function createReceiver() {
|
||||
const transport = new UnixSocketTransport({ endpoint: 'test-pipe', kind: 'named-pipe' })
|
||||
const socket = new FakeSocket()
|
||||
const received: string[] = []
|
||||
transport.onMessage((message, reply) => {
|
||||
received.push(message)
|
||||
reply('ok')
|
||||
})
|
||||
;(transport as unknown as UnixSocketTransportInternals).handleConnection(
|
||||
socket as unknown as Socket
|
||||
)
|
||||
return { socket, received }
|
||||
}
|
||||
|
||||
it.each([false, true])('preserves the UTF-8 byte boundary with oversized=%s', (oversized) => {
|
||||
const { socket, received } = createReceiver()
|
||||
const message = `${'é'.repeat(524287)}a${oversized ? 'x' : ''}`
|
||||
const wire = Buffer.from(`${message}\n`)
|
||||
const decoder = new StringDecoder('utf8')
|
||||
for (let offset = 0; offset < wire.length; offset += 4095) {
|
||||
socket.emit('data', decoder.write(wire.subarray(offset, offset + 4095)))
|
||||
}
|
||||
expect(received).toEqual([oversized ? '' : message])
|
||||
})
|
||||
|
||||
it('retains only the byte count of the partial tail between messages', () => {
|
||||
const { socket, received } = createReceiver()
|
||||
const large = 'a'.repeat(700000)
|
||||
socket.emit('data', `${large}\npart`)
|
||||
socket.emit('data', `ial\r\n\n${large}\n`)
|
||||
expect(received).toEqual([large, 'partial', large])
|
||||
})
|
||||
|
||||
it('checks the combined incoming buffer before dispatching any complete messages', () => {
|
||||
const { socket, received } = createReceiver()
|
||||
socket.emit('data', `${'a'.repeat(700000)}\n${'b'.repeat(700000)}\n`)
|
||||
expect(received).toEqual([''])
|
||||
socket.emit('data', 'later\n')
|
||||
expect(received).toEqual([''])
|
||||
})
|
||||
|
||||
it('clears request keepalive timers when the socket closes before a reply', () => {
|
||||
const transport = new UnixSocketTransport({
|
||||
endpoint: '/tmp/orca-runtime-rpc-test.sock',
|
||||
|
||||
@@ -105,6 +105,7 @@ export class UnixSocketTransport implements RpcTransport {
|
||||
private handleConnection(socket: Socket): void {
|
||||
this.activeSockets.add(socket)
|
||||
let buffer = ''
|
||||
let retainedBytes = 0
|
||||
let oversized = false
|
||||
// Why: each in-flight dispatch registers its own AbortController here so
|
||||
// `socket.on('close')` can abort them all at once. Keeping the set scoped
|
||||
@@ -134,10 +135,12 @@ export class UnixSocketTransport implements RpcTransport {
|
||||
return
|
||||
}
|
||||
buffer += chunk
|
||||
// setEncoding('utf8') keeps split codepoints intact, so chunk byte lengths add exactly.
|
||||
retainedBytes += Buffer.byteLength(chunk, 'utf8')
|
||||
// Why: the Orca runtime lives in Electron main, so it must reject
|
||||
// oversized local RPC frames instead of letting a local client grow an
|
||||
// unbounded buffer and stall the app.
|
||||
if (Buffer.byteLength(buffer, 'utf8') > MAX_RUNTIME_RPC_MESSAGE_BYTES) {
|
||||
if (retainedBytes > MAX_RUNTIME_RPC_MESSAGE_BYTES) {
|
||||
oversized = true
|
||||
this.messageHandler?.('', (response) => {
|
||||
socket.write(`${response}\n`)
|
||||
@@ -145,6 +148,9 @@ export class UnixSocketTransport implements RpcTransport {
|
||||
})
|
||||
return
|
||||
}
|
||||
if (!chunk.includes('\n')) {
|
||||
return
|
||||
}
|
||||
let newlineIndex = buffer.indexOf('\n')
|
||||
while (newlineIndex !== -1) {
|
||||
const rawMessage = buffer.slice(0, newlineIndex).trim()
|
||||
@@ -154,6 +160,7 @@ export class UnixSocketTransport implements RpcTransport {
|
||||
}
|
||||
newlineIndex = buffer.indexOf('\n')
|
||||
}
|
||||
retainedBytes = Buffer.byteLength(buffer, 'utf8')
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user