mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(relay): fail an over-budget RPC response, not the connection
The relay's control lane is a shared 1 MiB budget, and `sendResponse` admitted responses onto it with the fatal default: once the lane was full, admission closed the client. A ~900 KB `fs.listFiles` reply from a large remote workspace therefore took down the whole remote session -- every terminal on it -- rather than failing the one Quick Open request. The substitute `ResponseOverCapacity` frame already there only covered the `legacy-response` lane, because the fatal close beat it to the client. A JSON-RPC response is the droppable class of control frame: it carries an id, so one caller can be told and can retry. `pty.replay` and `notifyControl` keep the fatal default -- they are never re-sent, and a silent drop there desyncs the client with nothing to retry. Both response enqueues now pass `controlOverflow: 'reject'`, so the substitute error is what the caller sees; in the corner where even ~150 bytes will not fit, the caller's own 30s request timeout settles it and the session survives. Old clients are unaffected: they already decode this error code and message generically (`ssh-channel-multiplexer.handleResponse` rejects the pending promise with both), and the frame shape is unchanged. What changes is that a listing which used to drop the connection now returns an error on it.
This commit is contained in:
@@ -59,6 +59,14 @@ function makeBoundedClient(highWaterMark: number): BoundedClient {
|
||||
return client
|
||||
}
|
||||
|
||||
// Why: the sink accepts every write but never settles it, so control-lane bytes stay retained and the
|
||||
// queue fills, while the writer keeps pumping the other lanes — the shape of a peer whose socket is behind.
|
||||
function makeUnsettledWriteClient(highWaterMark: number): BoundedClient {
|
||||
const client = makeBoundedClient(highWaterMark)
|
||||
client.options = { ...client.options, supportsWriteCallback: true }
|
||||
return client
|
||||
}
|
||||
|
||||
function decodePayload(frame: Buffer): Record<string, unknown> {
|
||||
const length = frame.readUInt32BE(9)
|
||||
return JSON.parse(frame.subarray(13, 13 + length).toString('utf-8'))
|
||||
@@ -464,4 +472,92 @@ describe('RelayDispatcher bounded-capacity degradation', () => {
|
||||
bounded.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('answers an over-budget response with a capacity error instead of closing the connection', async () => {
|
||||
const primary = makeUnsettledWriteClient(65536)
|
||||
const bounded = new RelayDispatcher(primary.write, primary.options)
|
||||
try {
|
||||
const clientId = bounded.activeClientIds()[0]
|
||||
bounded.onRequest('fs.listFiles', async () => ({ paths: 'x'.repeat(700 * 1024) }))
|
||||
bounded.onRequest('workspace.get', async () => ({ name: 'workspace' }))
|
||||
|
||||
bounded.feed(encodeJsonRpcFrame({ jsonrpc: '2.0', id: 91, method: 'fs.listFiles' }, 1, 0))
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
expect(primary.frames).toHaveLength(1)
|
||||
|
||||
// The first reply still holds the shared control budget, so the second cannot fit under 1 MiB.
|
||||
bounded.feed(encodeJsonRpcFrame({ jsonrpc: '2.0', id: 92, method: 'fs.listFiles' }, 2, 0))
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
|
||||
expect(primary.closes).toBe(0)
|
||||
expect(bounded.isClientAttached(clientId)).toBe(true)
|
||||
expect(primary.frames).toHaveLength(2)
|
||||
const rejected = decodePayload(primary.frames[1]) as unknown as {
|
||||
id: number
|
||||
error: { code: number; message: string }
|
||||
}
|
||||
expect(rejected.id).toBe(92)
|
||||
expect(rejected.error.code).toBe(RelayErrorCode.ResponseOverCapacity)
|
||||
expect(rejected.error.message).toBe('Relay response exceeded the bounded transport capacity')
|
||||
|
||||
// Every other pane and request on this connection keeps working.
|
||||
bounded.feed(encodeJsonRpcFrame({ jsonrpc: '2.0', id: 93, method: 'workspace.get' }, 3, 0))
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
expect(decodePayload(primary.frames[2])).toMatchObject({
|
||||
id: 93,
|
||||
result: { name: 'workspace' }
|
||||
})
|
||||
|
||||
bounded.notify('pty.data', { paneId: 'pane-1', data: 'still-live' })
|
||||
expect(decodePayload(primary.frames[3])).toMatchObject({ method: 'pty.data' })
|
||||
expect(primary.closes).toBe(0)
|
||||
} finally {
|
||||
bounded.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('still closes the client when a protocol-critical control frame overflows', () => {
|
||||
const primary = makeUnsettledWriteClient(65536)
|
||||
const bounded = new RelayDispatcher(primary.write, primary.options)
|
||||
try {
|
||||
const clientId = bounded.activeClientIds()[0]
|
||||
bounded.notifyClient(clientId, 'workspace.stale', { blob: 'x'.repeat(700 * 1024) })
|
||||
expect(primary.closes).toBe(0)
|
||||
|
||||
// Replay is never re-sent, so an unnoticed drop strands the pane: overflow here stays fatal.
|
||||
bounded.notify('pty.replay', { paneKey: 'tab-1:pane-1', data: 'y'.repeat(700 * 1024) })
|
||||
expect(primary.closes).toBe(1)
|
||||
} finally {
|
||||
bounded.dispose()
|
||||
}
|
||||
})
|
||||
|
||||
it('drops an unsendable response without closing when even the capacity error will not fit', async () => {
|
||||
const primary = makeUnsettledWriteClient(65536)
|
||||
const bounded = new RelayDispatcher(primary.write, primary.options)
|
||||
try {
|
||||
const clientId = bounded.activeClientIds()[0]
|
||||
const settlements: SinkWriteSettlement[] = []
|
||||
bounded.onRequest('workspace.get', async (_params, context) => {
|
||||
context.onResponseSettled?.((result) => settlements.push(result))
|
||||
return { name: 'workspace' }
|
||||
})
|
||||
for (let index = 0; index < DISPATCHER_CONTROL_QUEUE_MAX_FRAMES; index += 1) {
|
||||
bounded.notifyClient(clientId, `control.${index}`)
|
||||
}
|
||||
const framesBefore = primary.frames.length
|
||||
|
||||
bounded.feed(encodeJsonRpcFrame({ jsonrpc: '2.0', id: 94, method: 'workspace.get' }, 1, 0))
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
|
||||
// Nothing goes out, but the connection lives and the caller's own request timeout settles it.
|
||||
expect(primary.closes).toBe(0)
|
||||
expect(primary.frames).toHaveLength(framesBefore)
|
||||
expect(settlements).toEqual([
|
||||
{ ok: false, error: new Error('Relay response was not admitted') }
|
||||
])
|
||||
} finally {
|
||||
bounded.dispose()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
@@ -203,12 +203,17 @@ export abstract class RelayDispatcherRpcRouting extends RelayDispatcherFrameCode
|
||||
const frame = this.prepareFrame(msg)
|
||||
const lane =
|
||||
frame.frameBytes > DISPATCHER_CONTROL_QUEUE_MAX_BYTES ? 'legacy-response' : 'control'
|
||||
const accepted = this.enqueuePreparedFrame(client, frame, lane, onSettled)
|
||||
// Why 'reject': the control lane is a shared budget, so a reply that fits the 1 MiB ceiling alone
|
||||
// still overflows it under concurrent traffic. Fatal admission would close the connection — every
|
||||
// pane on the host — over one listing. A response is the droppable class of control frame: it
|
||||
// carries an id, so the substitute below tells that one caller, and pty.replay/notifyControl keep
|
||||
// the fatal default because a silent drop there desyncs the client with nothing to retry.
|
||||
const accepted = this.enqueuePreparedFrame(client, frame, lane, onSettled, 'reject')
|
||||
if (accepted) {
|
||||
return true
|
||||
}
|
||||
// Why: an oversized response must fail its own request; closing would kill every pane on the host.
|
||||
// A rejected first enqueue either left onSettled untouched or closed the client, so exactly one settlement happens.
|
||||
// A rejected first enqueue leaves onSettled untouched, so exactly one settlement happens.
|
||||
return this.enqueuePreparedFrame(
|
||||
client,
|
||||
this.prepareFrame({
|
||||
@@ -227,7 +232,10 @@ export abstract class RelayDispatcherRpcRouting extends RelayDispatcherFrameCode
|
||||
settlement.ok
|
||||
? { ok: false, error: new Error(RESPONSE_OVER_CAPACITY_MESSAGE) }
|
||||
: settlement
|
||||
)
|
||||
),
|
||||
// Why 'reject': if even ~150 bytes will not fit, the caller's own request timeout settles it.
|
||||
// Closing to report that one request failed is the outcome this whole path exists to avoid.
|
||||
'reject'
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -105,10 +105,10 @@ describe('readRelayFileRange', () => {
|
||||
// which the writer refuses once the producer queue is busy -- so an over-wide
|
||||
// cap fails with ResponseOverCapacity depending on unrelated load.
|
||||
//
|
||||
// Fitting the lane once is not enough: the control queue is a SHARED budget
|
||||
// and overflowing it closes the client, so a full-cap frame has to leave room
|
||||
// for a second one. Anything wider lets two pipelined tail reads -- or one
|
||||
// read racing an unrelated response -- kill the connection.
|
||||
// Fitting the lane once is not enough: the control queue is a SHARED budget,
|
||||
// so a full-cap frame has to leave room for a second one. Anything wider lets
|
||||
// two pipelined tail reads -- or one read racing an unrelated response --
|
||||
// fail as ResponseOverCapacity on load that has nothing to do with them.
|
||||
it('leaves control-queue headroom for a second full-cap window', async () => {
|
||||
const contents = Buffer.allocUnsafe(MAX_FILE_RANGE_READ_BYTES)
|
||||
for (let i = 0; i < contents.length; i++) {
|
||||
|
||||
@@ -8,10 +8,11 @@
|
||||
* frames to ~350 KB, so a full-cap response takes the control lane instead.
|
||||
*
|
||||
* The control lane is a shared budget, not a per-frame one: two full-cap
|
||||
* responses fit alongside each other, and the third overflows -- which for a
|
||||
* response is fatal, it closes the client. That is the same exposure every
|
||||
* control-lane response already carries (`fs.readFile` frames any sub-1 MiB
|
||||
* file the same way), and the two-deep headroom is pinned by a test. Widening
|
||||
* responses fit alongside each other, and the third is refused. That refusal
|
||||
* costs the one request -- `sendResponse` admits responses with
|
||||
* `controlOverflow: 'reject'` and substitutes a `ResponseOverCapacity` error
|
||||
* rather than closing the connection -- but it still turns on unrelated load,
|
||||
* so the two-deep headroom that keeps it rare is pinned by a test. Widening
|
||||
* the cap spends that headroom, so bigger transfers belong on the ack-paced
|
||||
* bulk lane (`fs.readFileStream`) rather than on a wider window here.
|
||||
*
|
||||
|
||||
Reference in New Issue
Block a user