Files
orca/src/main/daemon/daemon-stream-data-split.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

153 lines
4.2 KiB
TypeScript

/**
* Surrogate-safe splitting for daemon stream data events: NDJSON line-size
* chunking (the receiver's parser rejects oversized lines) and the safe-index
* clamp shared by the batcher's bulk write slicing and keep-tail dropping.
*/
import { encodeNdjson } from './ndjson'
export function encodeStreamDataEvent(
sessionId: string,
data: string,
rawLength?: number,
seq?: number,
transformed?: boolean
): string {
return encodeNdjson({
type: 'event',
event: 'data',
sessionId,
payload: {
data,
...(seq === undefined ? {} : { seq }),
...(rawLength === undefined ? {} : { rawLength }),
...(rawLength === undefined ? {} : { sequenceChars: rawLength }),
...(transformed ? { transformed: true } : {})
}
})
}
function streamDataEventLineBytes(sessionId: string, data: string, rawLength?: number): number {
return Buffer.byteLength(encodeStreamDataEvent(sessionId, data, rawLength), 'utf8')
}
function isHighSurrogate(value: number): boolean {
return value >= 0xd800 && value <= 0xdbff
}
function isLowSurrogate(value: number): boolean {
return value >= 0xdc00 && value <= 0xdfff
}
export function clampToSafeSplitIndex(value: string, start: number, end: number): number {
if (end <= start || end >= value.length) {
return end
}
const prev = value.charCodeAt(end - 1)
const next = value.charCodeAt(end)
return isHighSurrogate(prev) && isLowSurrogate(next) ? end - 1 : end
}
function nextSafeSplitIndex(value: string, start: number): number {
const next = Math.min(value.length, start + 1)
if (
next < value.length &&
isHighSurrogate(value.charCodeAt(start)) &&
isLowSurrogate(value.charCodeAt(next))
) {
return next + 1
}
return next
}
export function splitStreamDataForNdjson(
sessionId: string,
data: string,
maxLineBytes: number,
sequenceChars?: number
): string[] {
if (streamDataEventLineBytes(sessionId, data, sequenceChars) <= maxLineBytes) {
return [data]
}
return splitOversizedStreamDataForNdjson(sessionId, data, maxLineBytes, sequenceChars)
}
function splitOversizedStreamDataForNdjson(
sessionId: string,
data: string,
maxLineBytes: number,
sequenceChars?: number
): string[] {
const chunks: string[] = []
let start = 0
while (start < data.length) {
let low = start + 1
let high = data.length
let best = start
while (low <= high) {
const rawMid = Math.floor((low + high) / 2)
const mid = clampToSafeSplitIndex(data, start, rawMid)
if (mid <= start) {
low = rawMid + 1
continue
}
if (
streamDataEventLineBytes(sessionId, data.slice(start, mid), sequenceChars) <= maxLineBytes
) {
best = mid
low = rawMid + 1
} else {
high = rawMid - 1
}
}
const end = best > start ? best : nextSafeSplitIndex(data, start)
chunks.push(data.slice(start, end))
start = end
}
return chunks
}
export function writeStreamDataEvents(
streamSocket: { write(data: string): void },
sessionId: string,
data: string,
maxLineBytes: number,
rawLength = data.length,
seq?: number,
transformed = false
): void {
const explicitRawLength = rawLength === data.length ? undefined : rawLength
if (transformed) {
streamSocket.write(encodeStreamDataEvent(sessionId, data, rawLength, seq, true))
return
}
const carriesMetadata = explicitRawLength !== undefined || seq !== undefined
let chunks: string[]
if (!carriesMetadata) {
const line = encodeStreamDataEvent(sessionId, data)
if (Buffer.byteLength(line, 'utf8') <= maxLineBytes) {
streamSocket.write(line)
return
}
chunks = splitOversizedStreamDataForNdjson(sessionId, data, maxLineBytes)
} else {
chunks = splitStreamDataForNdjson(
sessionId,
data,
Math.max(1, maxLineBytes - 96),
explicitRawLength
)
}
let consumed = 0
for (const chunk of chunks) {
consumed += chunk.length
const chunkEndSeq = seq === undefined ? undefined : seq - (data.length - consumed)
const chunkRawLength = explicitRawLength === 0 ? 0 : carriesMetadata ? chunk.length : undefined
streamSocket.write(encodeStreamDataEvent(sessionId, chunk, chunkRawLength, chunkEndSeq))
}
}