mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
* fix(ssh): recover instead of wedging when the relay channel dies mid-connect Three coupled defects made a dropped SSH relay look like a permanent bug: 1. SshRelaySession.establish()/reconnect() ran their last liveness gate before configureRelayGraceTime(), whose mux.notify() can dispose the mux synchronously (writer control-lane admission cap, or a throwing transport). The session then latched _state='ready' + _onReady (status bar "connected") while watchMuxForRelayLoss() silently no-op'd on the dead mux, so the bounded relay backoff in ipc/ssh.ts never ran and the fs/pty/git providers stayed registered against a dead multiplexer. Both sites now re-check mux.isDisposed() after the notify and take the existing failure path. 2. SshChannelMultiplexer.request()/notifyWithSettlement() always reported the permanent-shutdown string 'Multiplexer disposed' with no code, even when the recorded dispose reason was connection_lost. The reason is now recorded and a shared disposedError() factory serves dispose(), request(), notifyWithSettlement(), so a transient drop reports 'SSH connection lost, reconnecting...' / CONNECTION_LOST. onDispose() on an already-disposed mux now fires the handler synchronously with that reason instead of returning a silent no-op (without retaining it). 3. TerminalErrorToast no longer renders a transient relay drop in the red "please file an issue" style. The marker is matched with includes() because the message reaches the toast IPC-wrapped. ssh-git-response-stream-reader registers its onDispose subscriber after the abort wiring, since an already-dead mux now fails synchronously there and the cleanup must be able to drop the caller's abort listener. Closes #11953 * fix(ssh): treat a mux killed during PTY reattach as relay loss reconnect()'s post-reattach gate bare-returned when ownsAttempt() went false, and reattachKnownPtys swallows every per-PTY error, so a control-lane failure during a large reattach burst disposed the mux without ever reaching the catch: providers stayed bound to the dead mux, no relay-loss watcher was installed, and the session wedged in 'reconnecting' until restart. Take the failure path when our own mux is the one that died so ssh.ts's bounded backoff retries. Co-authored-by: Orca <help@stably.ai> * fix(ssh): recover when relay dies during setup instead of wedging Introduce verifyRelayAttempt() to detect mux disposal at each setup phase (consumer session, home resolution, provider registration, PTY reattach). Routes mid-setup connection loss into relay-loss recovery instead of hanging in reconnecting state. * Extract SSH disposal error factory Multiple sites were duplicating the disposal error creation logic with specific message and code values. The renderer uses these to distinguish temporary disconnects (show reconnection overlay) from permanent shutdown (show error toast), so all producers must use the same factory to avoid silent UI degradation. --------- Co-authored-by: Orca <help@stably.ai> Co-authored-by: Jinjing <6427696+AmethystLiang@users.noreply.github.com>
305 lines
10 KiB
TypeScript
305 lines
10 KiB
TypeScript
import type { SshChannelMultiplexer } from './ssh-channel-multiplexer'
|
|
import { createSshDisposalError } from './ssh-channel-multiplexer'
|
|
import { RelayErrorCode, isGitResponseStreamMarker } from './relay-protocol'
|
|
|
|
const SENTINEL_STREAM_ID = -1
|
|
|
|
/** Reject if no stream frame (chunk/end/error) arrives within this window,
|
|
* reset on each frame. mux.request's own timeout only bounds the fast sentinel
|
|
* response; without this, a relay pump that breaks on staleness (which sends no
|
|
* responseEnd) while the SSH channel stays up would hang the client forever. */
|
|
const STREAM_INACTIVITY_TIMEOUT_MS = 30_000
|
|
|
|
/** Bound transient buffering of other concurrent streams' chunks while this
|
|
* reader awaits its sentinel: every reader sees all git.responseChunk frames
|
|
* and can't filter by streamId until its own sentinel resolves. Foreign frames
|
|
* are dropped on drain anyway; this just caps the pre-sentinel backlog. */
|
|
const MAX_PENDING_FRAMES = 64
|
|
|
|
export class GitResponseStreamError extends Error {
|
|
readonly code = RelayErrorCode.StreamProtocolError
|
|
constructor(message: string) {
|
|
super(message)
|
|
}
|
|
}
|
|
|
|
type PendingFrame =
|
|
| { kind: 'chunk'; params: Record<string, unknown> }
|
|
| { kind: 'end'; params: Record<string, unknown> }
|
|
| { kind: 'error'; params: Record<string, unknown> }
|
|
|
|
/**
|
|
* Request a git method that may return a large payload, opting into response
|
|
* streaming so a big diff/exec response is chunked onto the relay's bulk lane
|
|
* instead of one JSON-RPC frame (which would head-of-line-block pty.data echo
|
|
* on the shared SSH channel).
|
|
*
|
|
* Cross-version behavior:
|
|
* - New relay + big result → returns the stream sentinel; we reassemble chunks.
|
|
* - New relay + small result, or old client → plain single-frame result.
|
|
* - Old relay (ignores `__streamResponse`) → returns the plain result; the
|
|
* marker check fails and we return it directly, i.e. today's behavior.
|
|
*/
|
|
export function requestGitStreamable(
|
|
mux: SshChannelMultiplexer,
|
|
method: string,
|
|
params: Record<string, unknown>,
|
|
options?: {
|
|
signal?: AbortSignal
|
|
/** Bounds only the sentinel request (forwarded to mux.request), like today. */
|
|
timeoutMs?: number
|
|
/** Bounds the post-sentinel reassembly stall; resets on each chunk. */
|
|
inactivityTimeoutMs?: number
|
|
}
|
|
): Promise<unknown> {
|
|
// Why: subscribe to chunk/end/error BEFORE awaiting the sentinel response so a
|
|
// chunk that lands in the same dispatch tick as the response is not dropped
|
|
// (mirrors readFileViaStream). streamIdRef stays SENTINEL until the sentinel
|
|
// resolves; frames are queued until then and drained.
|
|
const streamIdRef = { current: SENTINEL_STREAM_ID }
|
|
const unsubscribers: (() => void)[] = []
|
|
const cleanup = (): void => {
|
|
while (unsubscribers.length > 0) {
|
|
try {
|
|
unsubscribers.pop()?.()
|
|
} catch {
|
|
// best-effort
|
|
}
|
|
}
|
|
}
|
|
|
|
return new Promise<unknown>((resolve, reject) => {
|
|
const parts: Buffer[] = []
|
|
let expectedSeq = 0
|
|
let receivedBytes = 0
|
|
let totalBytes = 0
|
|
let chunkCount = 0
|
|
let settled = false
|
|
let metadataReady = false
|
|
const pending: PendingFrame[] = []
|
|
|
|
const inactivityMs = options?.inactivityTimeoutMs ?? STREAM_INACTIVITY_TIMEOUT_MS
|
|
let inactivityTimer: ReturnType<typeof setTimeout> | null = null
|
|
const clearInactivity = (): void => {
|
|
if (inactivityTimer) {
|
|
clearTimeout(inactivityTimer)
|
|
inactivityTimer = null
|
|
}
|
|
}
|
|
// Why: reset on every stream frame so a legitimately long stream is not
|
|
// killed, but a wedged stream (no frames arriving) rejects instead of
|
|
// hanging the caller forever.
|
|
const armInactivity = (): void => {
|
|
clearInactivity()
|
|
inactivityTimer = setTimeout(() => {
|
|
fail(
|
|
new GitResponseStreamError(
|
|
`Git response stream stalled (>${inactivityMs}ms without data)`
|
|
)
|
|
)
|
|
}, inactivityMs)
|
|
inactivityTimer.unref?.()
|
|
}
|
|
|
|
const cancel = (): void => {
|
|
if (streamIdRef.current !== SENTINEL_STREAM_ID && !mux.isDisposed()) {
|
|
try {
|
|
mux.notify('git.cancelResponseStream', { streamId: streamIdRef.current })
|
|
} catch {
|
|
// best-effort
|
|
}
|
|
}
|
|
}
|
|
const fail = (err: Error): void => {
|
|
if (settled) {
|
|
return
|
|
}
|
|
settled = true
|
|
clearInactivity()
|
|
cancel()
|
|
cleanup()
|
|
reject(err)
|
|
}
|
|
const succeed = (value: unknown): void => {
|
|
if (settled) {
|
|
return
|
|
}
|
|
settled = true
|
|
clearInactivity()
|
|
cleanup()
|
|
resolve(value)
|
|
}
|
|
|
|
const handleChunk = (p: Record<string, unknown>): void => {
|
|
if (settled || p.streamId !== streamIdRef.current) {
|
|
return
|
|
}
|
|
const seq = p.seq as number
|
|
const data = p.data as string
|
|
if (typeof seq !== 'number' || typeof data !== 'string') {
|
|
fail(new GitResponseStreamError(`Malformed chunk for git stream ${streamIdRef.current}`))
|
|
return
|
|
}
|
|
if (seq !== expectedSeq) {
|
|
fail(
|
|
new GitResponseStreamError(
|
|
`Out-of-order chunk for git stream ${streamIdRef.current}: expected ${expectedSeq}, got ${seq}`
|
|
)
|
|
)
|
|
return
|
|
}
|
|
const decoded = Buffer.from(data, 'base64')
|
|
parts.push(decoded)
|
|
receivedBytes += decoded.length
|
|
expectedSeq += 1
|
|
armInactivity()
|
|
// Why: credit-based flow control — the relay caps unacked chunks so a big
|
|
// response cannot queue unbounded ahead of interactive pty.data frames.
|
|
if (!mux.isDisposed()) {
|
|
try {
|
|
mux.notify('git.responseAck', { streamId: streamIdRef.current, seq })
|
|
} catch {
|
|
// Disposal can race the check; the ACK is best-effort during teardown.
|
|
}
|
|
}
|
|
}
|
|
|
|
const handleEnd = (p: Record<string, unknown>): void => {
|
|
if (settled || p.streamId !== streamIdRef.current) {
|
|
return
|
|
}
|
|
if (expectedSeq !== chunkCount || receivedBytes !== totalBytes) {
|
|
fail(
|
|
new GitResponseStreamError(
|
|
`Git stream ${streamIdRef.current} incomplete: chunks ${expectedSeq}/${chunkCount}, bytes ${receivedBytes}/${totalBytes}`
|
|
)
|
|
)
|
|
return
|
|
}
|
|
try {
|
|
succeed(JSON.parse(Buffer.concat(parts).toString('utf-8')))
|
|
} catch (err) {
|
|
fail(
|
|
new GitResponseStreamError(
|
|
`Git stream ${streamIdRef.current} JSON parse failed: ${String(err)}`
|
|
)
|
|
)
|
|
}
|
|
}
|
|
|
|
const handleStreamError = (p: Record<string, unknown>): void => {
|
|
if (settled || p.streamId !== streamIdRef.current) {
|
|
return
|
|
}
|
|
fail(new Error((p.message as string | undefined) ?? 'git response stream error'))
|
|
}
|
|
|
|
const drainPending = (): void => {
|
|
while (!settled && pending.length > 0) {
|
|
const frame = pending.shift()!
|
|
if (frame.kind === 'chunk') {
|
|
handleChunk(frame.params)
|
|
} else if (frame.kind === 'end') {
|
|
handleEnd(frame.params)
|
|
} else {
|
|
handleStreamError(frame.params)
|
|
}
|
|
}
|
|
}
|
|
|
|
// Why: pre-sentinel we cannot filter by streamId (our id is unknown yet), so
|
|
// every concurrent reader transiently buffers all readers' chunks. Cap the
|
|
// backlog by dropping the oldest; foreign frames are dropped on drain anyway,
|
|
// and if our own seq-0 were ever dropped the seq check fails loudly rather
|
|
// than corrupting. The sentinel normally resolves long before this cap.
|
|
const pushPending = (frame: PendingFrame): void => {
|
|
pending.push(frame)
|
|
if (pending.length > MAX_PENDING_FRAMES) {
|
|
pending.shift()
|
|
}
|
|
}
|
|
|
|
unsubscribers.push(
|
|
mux.onNotificationByMethod('git.responseChunk', (p) => {
|
|
if (!metadataReady) {
|
|
pushPending({ kind: 'chunk', params: p })
|
|
return
|
|
}
|
|
handleChunk(p)
|
|
})
|
|
)
|
|
unsubscribers.push(
|
|
mux.onNotificationByMethod('git.responseEnd', (p) => {
|
|
if (!metadataReady) {
|
|
pushPending({ kind: 'end', params: p })
|
|
return
|
|
}
|
|
handleEnd(p)
|
|
})
|
|
)
|
|
unsubscribers.push(
|
|
mux.onNotificationByMethod('git.responseError', (p) => {
|
|
if (!metadataReady) {
|
|
pushPending({ kind: 'error', params: p })
|
|
return
|
|
}
|
|
handleStreamError(p)
|
|
})
|
|
)
|
|
if (options?.signal) {
|
|
const signal = options.signal
|
|
if (signal.aborted) {
|
|
const err = new Error('Request was cancelled') as Error & { name: string }
|
|
err.name = 'AbortError'
|
|
fail(err)
|
|
return
|
|
}
|
|
const onAbort = (): void => {
|
|
const err = new Error('Request was cancelled') as Error & { name: string }
|
|
err.name = 'AbortError'
|
|
fail(err)
|
|
}
|
|
signal.addEventListener('abort', onAbort, { once: true })
|
|
unsubscribers.push(() => signal.removeEventListener('abort', onAbort))
|
|
}
|
|
|
|
// Why: registered last because an already-disposed mux fails synchronously here,
|
|
// and that cleanup must be able to drop the abort listener above (#11953).
|
|
unsubscribers.push(mux.onDispose((reason) => fail(createSshDisposalError(reason))))
|
|
|
|
// Why: forward only the mux-request options (signal/timeoutMs) and omit them
|
|
// entirely when absent, so callers that previously issued a 2-arg
|
|
// mux.request keep the same call shape (and their tests). inactivityTimeoutMs
|
|
// governs reassembly here, not the sentinel request.
|
|
const streamParams = { ...params, __streamResponse: true }
|
|
const requestOptions =
|
|
options?.signal !== undefined || options?.timeoutMs !== undefined
|
|
? { signal: options.signal, timeoutMs: options.timeoutMs }
|
|
: undefined
|
|
const requestPromise = requestOptions
|
|
? mux.request(method, streamParams, requestOptions)
|
|
: mux.request(method, streamParams)
|
|
void requestPromise
|
|
.then((result) => {
|
|
if (settled) {
|
|
return
|
|
}
|
|
// Old relay / small result: plain single-frame value, no stream follows.
|
|
if (!isGitResponseStreamMarker(result)) {
|
|
succeed(result)
|
|
return
|
|
}
|
|
const marker = result.__orcaGitResponseStream
|
|
totalBytes = marker.totalBytes
|
|
chunkCount = marker.chunkCount
|
|
streamIdRef.current = marker.streamId
|
|
metadataReady = true
|
|
// Why: start the inactivity deadline now — mux.request's timeout only
|
|
// covered the sentinel; the reassembly phase needs its own guard.
|
|
armInactivity()
|
|
drainPending()
|
|
})
|
|
.catch((err) => fail(err as Error))
|
|
})
|
|
}
|