Files
orca/src/main/daemon/daemon-stream-backpressure-socket.test.ts
T
84d827a6ab fix(daemon): pause producers when stream backlogs grow (#20947)
* fix(daemon): pause producers when stream backlogs grow

* fix(daemon): reset stream backpressure on socket replacement

* docs(daemon): point retention audit at current reproducer

* test(daemon): validate stream retention audit outcomes

* fix(daemon): bound the stream producer stall and leave a visible gap

Stream backpressure pauses a session's PTY with no deadline: the only
un-pause comes from the consumer draining, so a half-open peer that stops
reading without closing freezes the shell for the rest of the session.

Arm a 60s watchdog on the false->true stream-pause transition (not on the
re-assertions refresh() makes for neighbouring sessions). On fire, mark the
session stall-released: it becomes keep-tail droppable, its backlog is
thinned behind a dataGap, and the producer runs again. The existing dataGap
path makes the renderer restore that pane from the daemon's snapshot, so the
user sees the terminal jump to current rather than sit frozen. The mark
clears once the session's last byte leaves the daemon, restoring ordinary
pausing. Nothing here reports a process exit - loss of contact with a
consumer is not evidence about the child.

Also enable TCP keepalive on the stream socket so a genuinely dead peer
closes and onStreamDisconnected clears the pause.

* test(daemon): put each casting SAFETY: directive on one line

`oxlint-disable-next-line` covers only the line directly after it, so a
rationale wrapped onto a second comment line suppressed nothing and the
casts failed the changed-code quality gate. Drop the remaining JSON.parse
cast for an annotated binding.

---------

Co-authored-by: m4air <m4air@Mac.localdomain>
Co-authored-by: Neil <neil@stably.ai>
2026-09-19 17:51:59 -07:00

146 lines
4.9 KiB
TypeScript

import { once } from 'node:events'
import { createServer, Socket } from 'node:net'
import { setImmediate as nextTurn, setTimeout as delay } from 'node:timers/promises'
import { describe, expect, it } from 'vitest'
import { DaemonStreamDataBatcher } from './daemon-stream-data-batcher'
import { createNdjsonParser } from './ndjson'
import { SessionProducerPause } from './session-producer-pause'
const MiB = 1024 * 1024
describe('daemon stream backpressure with a stalled Socket reader', () => {
it.each([
{ chunkSize: 64 * 1024, hidden: false },
{ chunkSize: 1024, hidden: false },
{ chunkSize: 64 * 1024, hidden: true },
{ chunkSize: 1024, hidden: true }
])(
'bounds $chunkSize-character writes (hidden=$hidden) and resumes',
async ({ chunkSize, hidden }) => {
const server = createServer()
const writer = new Socket()
let reader: Socket | undefined
let paused = false
let produced = 0
let dropped = 0
const received = new Map<string, number>()
const producer = new SessionProducerPause({
pause: () => {
paused = true
},
resume: () => {
paused = false
}
})
const batcher = new DaemonStreamDataBatcher(() => ({ streamSocket: writer }), {
isSessionDroppable: (sessionId) => hidden && sessionId === 'flood',
onProducerBackpressureChanged: (sessionId, value) => {
expect(sessionId).toBe('flood')
producer.setStreamBackpressured(value)
}
})
writer.on('drain', () => batcher.flush('client'))
try {
const accepted = once(server, 'connection')
server.listen(0, '127.0.0.1')
await once(server, 'listening')
const address = server.address()
if (!address || typeof address === 'string') {
throw new Error('Expected TCP address')
}
writer.connect(address.port, '127.0.0.1')
await once(writer, 'connect')
const [acceptedSocket] = await accepted
if (!(acceptedSocket instanceof Socket)) {
throw new Error('Expected accepted Socket')
}
reader = acceptedSocket
const parser = createNdjsonParser(
(event) => {
if (
!event ||
typeof event !== 'object' ||
!('sessionId' in event) ||
typeof event.sessionId !== 'string' ||
!('payload' in event)
) {
return
}
const payload = event.payload
if (
payload &&
typeof payload === 'object' &&
'droppedChars' in payload &&
typeof payload.droppedChars === 'number'
) {
dropped += payload.droppedChars
}
if (
!payload ||
typeof payload !== 'object' ||
!('data' in payload) ||
typeof payload.data !== 'string'
) {
return
}
received.set(
event.sessionId,
(received.get(event.sessionId) ?? 0) + payload.data.length
)
},
() => {
throw new Error('Invalid stream frame')
}
)
reader.on('data', (data) => parser.feed(data.toString('utf8')))
reader.pause()
const chunk = 'x'.repeat(chunkSize)
while (!paused && produced < 8 * MiB) {
batcher.enqueue('client', 'flood', chunk, {
flushImmediately: chunkSize <= 1024,
flushMaxChars: 1024
})
batcher.flush('client')
produced += chunkSize
if (produced % MiB === 0) {
await nextTurn()
}
}
expect(paused).toBe(!hidden)
expect(writer.writableLength + 2 * batcher.queuedCharsForClient('client')).toBeLessThan(
5 * MiB
)
if (!hidden) {
producer.pause()
producer.resumeClient()
expect(paused).toBe(true)
}
batcher.enqueue('client', 'typing', 'echo', { flushImmediately: true, flushMaxChars: 1024 })
reader.resume()
const startedAt = performance.now()
while (paused || writer.writableLength || batcher.queuedCharsForClient('client')) {
if (performance.now() - startedAt > 5_000) {
throw new Error('Stream failed to drain')
}
await delay(5)
}
const ended = once(reader, 'end')
writer.end()
await ended
expect((received.get('flood') ?? 0) + dropped).toBe(produced)
expect(dropped > 0).toBe(hidden)
expect(received.get('typing')).toBe(4)
expect(paused).toBe(false)
} finally {
batcher.clear()
producer.release({ resume: false })
writer.destroy()
reader?.destroy()
await new Promise<void>((resolve) => server.close(() => resolve()))
}
}
)
})