diff --git a/src/main/ssh/ssh-file-stream-inactivity-deadline.test.ts b/src/main/ssh/ssh-file-stream-inactivity-deadline.test.ts new file mode 100644 index 00000000000..f87d1e6f778 --- /dev/null +++ b/src/main/ssh/ssh-file-stream-inactivity-deadline.test.ts @@ -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) + }) +}) diff --git a/src/main/ssh/ssh-file-stream-inactivity-deadline.ts b/src/main/ssh/ssh-file-stream-inactivity-deadline.ts index 5470e8121fe..cbb9839f9b3 100644 --- a/src/main/ssh/ssh-file-stream-inactivity-deadline.ts +++ b/src/main/ssh/ssh-file-stream-inactivity-deadline.ts @@ -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?.() } diff --git a/src/main/ssh/ssh-git-response-stream-reader.ts b/src/main/ssh/ssh-git-response-stream-reader.ts index 6a50dc28a72..0a8b26aa779 100644 --- a/src/main/ssh/ssh-git-response-stream-reader.ts +++ b/src/main/ssh/ssh-git-response-stream-reader.ts @@ -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( diff --git a/src/main/ssh/ssh-git-stream-idle-timer.test.ts b/src/main/ssh/ssh-git-stream-idle-timer.test.ts new file mode 100644 index 00000000000..7cd5207c63a --- /dev/null +++ b/src/main/ssh/ssh-git-stream-idle-timer.test.ts @@ -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) => 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) => 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() + } + } +)