mirror of
https://github.com/stablyai/orca.git
synced 2026-09-22 00:02:31 +00:00
perf(ssh): reuse streamed response idle timers (#18956)
This commit is contained in:
@@ -0,0 +1,87 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import { createSshFileStreamInactivityDeadline } from './ssh-file-stream-inactivity-deadline'
|
||||
import type { SystemPowerLifecycleListener } from '../system-power-lifecycle'
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks()
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
describe('SSH file stream inactivity timer', () => {
|
||||
it('reuses one timer while retaining the deadline of the latest chunk', () => {
|
||||
vi.useFakeTimers()
|
||||
const allocate = vi.spyOn(globalThis, 'setTimeout')
|
||||
const onTimeout = vi.fn()
|
||||
const unsubscribe = vi.fn()
|
||||
const deadline = createSshFileStreamInactivityDeadline(onTimeout, (listener) => {
|
||||
listener.onResume()
|
||||
return unsubscribe
|
||||
})
|
||||
deadline.reset()
|
||||
vi.advanceTimersByTime(30_000)
|
||||
for (let chunk = 0; chunk < 1000; chunk += 1) {
|
||||
deadline.reset()
|
||||
}
|
||||
expect(allocate).toHaveBeenCalledTimes(1)
|
||||
vi.advanceTimersByTime(59_999)
|
||||
expect(onTimeout).not.toHaveBeenCalled()
|
||||
vi.advanceTimersByTime(1)
|
||||
expect(onTimeout).toHaveBeenCalledTimes(1)
|
||||
deadline.clear()
|
||||
expect(unsubscribe).toHaveBeenCalledTimes(1)
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
|
||||
it('releases on suspend and creates a fresh timer on resume', () => {
|
||||
vi.useFakeTimers()
|
||||
const allocate = vi.spyOn(globalThis, 'setTimeout')
|
||||
const onTimeout = vi.fn()
|
||||
let power!: SystemPowerLifecycleListener
|
||||
const deadline = createSshFileStreamInactivityDeadline(onTimeout, (listener) => {
|
||||
power = listener
|
||||
listener.onResume()
|
||||
return vi.fn()
|
||||
})
|
||||
deadline.reset()
|
||||
vi.advanceTimersByTime(30_000)
|
||||
power.onSuspend()
|
||||
for (let chunk = 0; chunk < 1000; chunk += 1) {
|
||||
deadline.reset()
|
||||
}
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
vi.advanceTimersByTime(120_000)
|
||||
expect(onTimeout).not.toHaveBeenCalled()
|
||||
power.onResume()
|
||||
expect(allocate).toHaveBeenCalledTimes(2)
|
||||
vi.advanceTimersByTime(59_999)
|
||||
expect(onTimeout).not.toHaveBeenCalled()
|
||||
vi.advanceTimersByTime(1)
|
||||
expect(onTimeout).toHaveBeenCalledTimes(1)
|
||||
deadline.clear()
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
|
||||
it('clears the timer and subscription and supports a later reset', () => {
|
||||
vi.useFakeTimers()
|
||||
const onTimeout = vi.fn()
|
||||
const unsubscribe = vi.fn()
|
||||
const subscribe = vi.fn((listener: SystemPowerLifecycleListener) => {
|
||||
listener.onResume()
|
||||
return unsubscribe
|
||||
})
|
||||
const deadline = createSshFileStreamInactivityDeadline(onTimeout, subscribe)
|
||||
deadline.reset()
|
||||
deadline.clear()
|
||||
deadline.clear()
|
||||
expect(unsubscribe).toHaveBeenCalledTimes(1)
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
vi.advanceTimersByTime(120_000)
|
||||
expect(onTimeout).not.toHaveBeenCalled()
|
||||
deadline.reset()
|
||||
expect(subscribe).toHaveBeenCalledTimes(2)
|
||||
expect(vi.getTimerCount()).toBe(1)
|
||||
deadline.clear()
|
||||
expect(unsubscribe).toHaveBeenCalledTimes(2)
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
})
|
||||
@@ -24,10 +24,13 @@ export function createSshFileStreamInactivityDeadline(
|
||||
}
|
||||
}
|
||||
const arm = (): void => {
|
||||
clearTimer()
|
||||
if (suspended) {
|
||||
return
|
||||
}
|
||||
if (timer) {
|
||||
timer.refresh()
|
||||
return
|
||||
}
|
||||
timer = setTimeout(onTimeout, SSH_FILE_STREAM_INACTIVITY_TIMEOUT_MS)
|
||||
timer.unref?.()
|
||||
}
|
||||
|
||||
@@ -90,7 +90,10 @@ export function requestGitStreamable(
|
||||
// killed, but a wedged stream (no frames arriving) rejects instead of
|
||||
// hanging the caller forever.
|
||||
const armInactivity = (): void => {
|
||||
clearInactivity()
|
||||
if (inactivityTimer) {
|
||||
inactivityTimer.refresh()
|
||||
return
|
||||
}
|
||||
inactivityTimer = setTimeout(() => {
|
||||
fail(
|
||||
new GitResponseStreamError(
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
import { expect, it, vi } from 'vitest'
|
||||
import type { SshChannelMultiplexer } from './ssh-channel-multiplexer'
|
||||
import { requestGitStreamable } from './ssh-git-response-stream-reader'
|
||||
|
||||
it.each(['end', 'abort', 'timeout'] as const)(
|
||||
'reuses the idle deadline across 1000 chunks and cleans up on %s',
|
||||
async (finish) => {
|
||||
vi.useFakeTimers()
|
||||
const setTimer = vi.spyOn(globalThis, 'setTimeout')
|
||||
try {
|
||||
const listeners = new Map<string, (params: Record<string, unknown>) => void>()
|
||||
const controller = new AbortController()
|
||||
const content = 'x'.repeat(998)
|
||||
const encoded = Buffer.from(JSON.stringify(content))
|
||||
const notify = vi.fn()
|
||||
const mux = {
|
||||
request: vi.fn(async () => ({
|
||||
__orcaGitResponseStream: { streamId: 7, totalBytes: encoded.length, chunkCount: 1000 }
|
||||
})),
|
||||
isDisposed: () => false,
|
||||
notify,
|
||||
onDispose: () => () => {},
|
||||
onNotificationByMethod: (
|
||||
method: string,
|
||||
callback: (params: Record<string, unknown>) => void
|
||||
) => {
|
||||
listeners.set(method, callback)
|
||||
return () => listeners.delete(method)
|
||||
}
|
||||
}
|
||||
const promise = requestGitStreamable(
|
||||
mux as unknown as SshChannelMultiplexer,
|
||||
'git.diff',
|
||||
{},
|
||||
{
|
||||
signal: controller.signal
|
||||
}
|
||||
)
|
||||
const outcome = promise.then(
|
||||
(value) => ({ value }),
|
||||
(error: Error) => ({ error: error.message })
|
||||
)
|
||||
await vi.advanceTimersByTimeAsync(15_000)
|
||||
for (let seq = 0; seq < encoded.length; seq++) {
|
||||
listeners.get('git.responseChunk')!({
|
||||
streamId: 7,
|
||||
seq,
|
||||
data: encoded.subarray(seq, seq + 1).toString('base64')
|
||||
})
|
||||
}
|
||||
await vi.advanceTimersByTimeAsync(29_999)
|
||||
expect(listeners.size).toBe(3)
|
||||
expect(notify.mock.calls.filter(([method]) => method === 'git.responseAck')).toHaveLength(
|
||||
1000
|
||||
)
|
||||
const allocations = setTimer.mock.calls.filter(([, delay]) => delay === 30_000).length
|
||||
if (finish === 'end') {
|
||||
listeners.get('git.responseEnd')!({ streamId: 7 })
|
||||
expect(await outcome).toEqual({ value: content })
|
||||
} else if (finish === 'abort') {
|
||||
controller.abort()
|
||||
expect(await outcome).toEqual({ error: 'Request was cancelled' })
|
||||
} else {
|
||||
await vi.advanceTimersByTimeAsync(1)
|
||||
expect(await outcome).toEqual({
|
||||
error: 'Git response stream stalled (>30000ms without data)'
|
||||
})
|
||||
}
|
||||
expect(allocations).toBe(1)
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
expect(listeners.size).toBe(0)
|
||||
expect(
|
||||
notify.mock.calls.filter(([method]) => method === 'git.cancelResponseStream')
|
||||
).toHaveLength(finish === 'end' ? 0 : 1)
|
||||
} finally {
|
||||
setTimer.mockRestore()
|
||||
vi.useRealTimers()
|
||||
}
|
||||
}
|
||||
)
|
||||
Reference in New Issue
Block a user