mirror of
https://github.com/stablyai/orca.git
synced 2026-10-08 16:02:37 +00:00
Detach retained CI and terminal tails from oversized strings (#20960)
* fix(memory): detach retained CI and terminal tails from oversized strings * fix(terminal): detach retained error and reattach string slices * fix(terminal): release oversized recent-output backing strings * fix(terminal): release backing strings held by PTY detectors * fix(memory): own bounded Claude background task labels * fix: detach retained terminal mode scan tails * fix: own retained plugin worker output strings * fix: own incomplete OSC 133 carry strings --------- Co-authored-by: m4air <m4air@Mac.localdomain> Co-authored-by: m4air <m4air@m4airs-MacBook-Air.local>
This commit is contained in:
@@ -1,4 +1,5 @@
|
||||
import type { Readable } from 'node:stream'
|
||||
import { ownRetainedString } from '../../shared/own-retained-string'
|
||||
|
||||
type PluginWorkerOutputSink = (level: 'info' | 'warn' | 'error', line: string) => void
|
||||
|
||||
@@ -21,9 +22,11 @@ export function pipePluginWorkerOutput(
|
||||
if (line.trim().length > 0) {
|
||||
log(
|
||||
level,
|
||||
truncated
|
||||
? `${line.slice(0, PLUGIN_WORKER_OUTPUT_LINE_LIMIT - TRUNCATION_SUFFIX.length)}${TRUNCATION_SUFFIX}`
|
||||
: line
|
||||
ownRetainedString(
|
||||
truncated
|
||||
? `${line.slice(0, PLUGIN_WORKER_OUTPUT_LINE_LIMIT - TRUNCATION_SUFFIX.length)}${TRUNCATION_SUFFIX}`
|
||||
: line
|
||||
)
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -49,7 +52,7 @@ export function pipePluginWorkerOutput(
|
||||
buffered = ''
|
||||
discarding = newline === -1
|
||||
} else {
|
||||
buffered += segment
|
||||
buffered += newline === -1 ? ownRetainedString(segment) : segment
|
||||
if (newline !== -1) {
|
||||
emit(buffered)
|
||||
buffered = ''
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
import { once } from 'node:events'
|
||||
import { PassThrough } from 'node:stream'
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { PluginLogBuffer } from './plugin-log-buffer'
|
||||
import { pipePluginWorkerOutput } from './plugin-worker-output-buffer'
|
||||
|
||||
async function heapAfterGc(): Promise<number> {
|
||||
if (!('gc' in globalThis) || typeof globalThis.gc !== 'function') {
|
||||
throw new Error('The test runner must enable --expose-gc')
|
||||
}
|
||||
for (let round = 0; round < 3; round++) {
|
||||
await new Promise<void>((resolve) => setImmediate(resolve))
|
||||
globalThis.gc()
|
||||
}
|
||||
return process.memoryUsage().heapUsed
|
||||
}
|
||||
|
||||
async function endStream(stream: PassThrough): Promise<void> {
|
||||
const ended = once(stream, 'end')
|
||||
stream.end()
|
||||
await ended
|
||||
}
|
||||
|
||||
function writeTail(stream: PassThrough, index: number): void {
|
||||
stream.write(`${' '.repeat(4 * 1024 * 1024)}\nretained output ${index}`)
|
||||
}
|
||||
|
||||
function writeLine(stream: PassThrough, index: number, truncated: boolean): void {
|
||||
const prefix = String(index).padStart(4, '0')
|
||||
stream.write(
|
||||
truncated
|
||||
? `${prefix}${'x'.repeat(64 * 1024)}\n`
|
||||
: `${' '.repeat(64 * 1024)}\nretained output ${prefix}\n`
|
||||
)
|
||||
}
|
||||
|
||||
describe('plugin worker retained output', () => {
|
||||
it('keeps unfinished output after consuming a large chunk without retaining the parent', async () => {
|
||||
const lines: string[] = []
|
||||
const before = await heapAfterGc()
|
||||
const streams = Array.from({ length: 8 }, (_value, index) => {
|
||||
const stream = new PassThrough()
|
||||
pipePluginWorkerOutput(stream, 'info', (_level, line) => lines.push(line))
|
||||
writeTail(stream, index)
|
||||
return stream
|
||||
})
|
||||
|
||||
expect((await heapAfterGc()) - before).toBeLessThan(2 * 1024 * 1024)
|
||||
expect(lines).toEqual([])
|
||||
for (const stream of streams) {
|
||||
await endStream(stream)
|
||||
}
|
||||
expect(lines).toEqual(Array.from({ length: 8 }, (_value, index) => `retained output ${index}`))
|
||||
})
|
||||
|
||||
it.each([false, true])(
|
||||
'owns emitted log text without retaining consumed chunks (truncated=%s)',
|
||||
async (truncated) => {
|
||||
const logs = new PluginLogBuffer()
|
||||
const stream = new PassThrough()
|
||||
pipePluginWorkerOutput(stream, 'error', (level, line) => logs.append('plugin', level, line))
|
||||
const before = await heapAfterGc()
|
||||
for (let index = 0; index < 205; index++) {
|
||||
writeLine(stream, index, truncated)
|
||||
}
|
||||
await endStream(stream)
|
||||
|
||||
// Compare text after the heap check: comparisons can flatten concatenated strings.
|
||||
expect((await heapAfterGc()) - before).toBeLessThan(5 * 1024 * 1024)
|
||||
expect(logs.get('plugin')).toHaveLength(200)
|
||||
for (const [index, row] of logs.get('plugin').entries()) {
|
||||
const prefix = String(index + 5).padStart(4, '0')
|
||||
expect(row.level).toBe('error')
|
||||
expect(row.line).toBe(
|
||||
truncated
|
||||
? `${prefix}${'x'.repeat(8192 - 4 - '… [truncated]'.length)}… [truncated]`
|
||||
: `retained output ${prefix}`
|
||||
)
|
||||
}
|
||||
}
|
||||
)
|
||||
})
|
||||
Reference in New Issue
Block a user