mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 08:03:12 +00:00
fix(ssh): stop fencing an agent-session create the mux never dispatched
This commit is contained in:
@@ -1,4 +1,7 @@
|
||||
import type { SshChannelMultiplexer } from '../ssh/ssh-channel-multiplexer'
|
||||
import {
|
||||
isUndispatchedSshRequestError,
|
||||
type SshChannelMultiplexer
|
||||
} from '../ssh/ssh-channel-multiplexer'
|
||||
import { AGENT_SESSION_CREATE_OPERATION_PROTOCOL_VERSION } from '../../shared/agent-session-host-authority'
|
||||
import { isPtyIncarnationId } from '../../shared/pty-incarnation'
|
||||
import type { PtySpawnResult } from './pty-spawn-result'
|
||||
@@ -66,7 +69,10 @@ export async function requestSshAgentSessionCreate(args: {
|
||||
: undefined
|
||||
return await args.mux.request('pty.spawn', args.params, options)
|
||||
} catch (error) {
|
||||
if (!args.operationId) {
|
||||
// Why: a spawn the mux never framed cannot have left a PTY on the host, so fencing it would
|
||||
// retain the 24h replay tombstone over an operation id whose create was never issued — every
|
||||
// later replay then rethrows this failure instead of starting the agent the user asked for.
|
||||
if (!args.operationId || isUndispatchedSshRequestError(error)) {
|
||||
throw error
|
||||
}
|
||||
const spawnError = error instanceof Error ? error : new Error(String(error))
|
||||
|
||||
@@ -235,6 +235,79 @@ describe('SSH fresh agent-session create operations', () => {
|
||||
expect(request).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('leaves the replay fence off a create the disposed mux never dispatched', async () => {
|
||||
const transport = createTransport()
|
||||
const mux = new SshChannelMultiplexer(transport)
|
||||
const liveProvider = new SshPtyProvider('conn-1', mux)
|
||||
const operationId = 'e'.repeat(43)
|
||||
|
||||
// Why: prime the memoized capability probe, so the second spawn reaches the dispatch seam
|
||||
// without a round trip — the ordering that lets a mid-preflight dispose land before it.
|
||||
const primed = liveProvider.spawn({
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
command: 'codex',
|
||||
agentSessionCreateOperationId: operationId
|
||||
})
|
||||
const capabilityRequest = await waitForRequest(transport, 'pty.getCapabilities')
|
||||
transport.deliver(
|
||||
responseFrame(
|
||||
capabilityRequest.id as number,
|
||||
{ agentSessionCreateOperationVersion: AGENT_SESSION_CREATE_OPERATION_PROTOCOL_VERSION },
|
||||
1
|
||||
)
|
||||
)
|
||||
const spawnRequest = await waitForRequest(transport, 'pty.spawn')
|
||||
transport.deliver(
|
||||
responseFrame(spawnRequest.id as number, { id: 'pty-1', incarnationId: 'incarnation-1' }, 2)
|
||||
)
|
||||
await expect(primed).resolves.toMatchObject({ incarnationId: 'incarnation-1' })
|
||||
|
||||
mux.dispose('connection_lost')
|
||||
const failure = await liveProvider
|
||||
.spawn({
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
command: 'codex',
|
||||
agentSessionCreateOperationId: 'f'.repeat(43)
|
||||
})
|
||||
.catch((error: unknown) => error)
|
||||
|
||||
expect((failure as { code?: unknown }).code).toBe('CONNECTION_LOST')
|
||||
expect(failure).not.toHaveProperty('agentSessionOperationOutcome')
|
||||
expect(
|
||||
requestPayloads(transport).filter((payload) => payload.method === 'pty.spawn')
|
||||
).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('keeps the replay fence on a create the link lost after dispatch', async () => {
|
||||
const transport = createTransport()
|
||||
const mux = new SshChannelMultiplexer(transport)
|
||||
const liveProvider = new SshPtyProvider('conn-1', mux)
|
||||
|
||||
const spawn = liveProvider.spawn({
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
command: 'codex',
|
||||
agentSessionCreateOperationId: 'g'.repeat(43)
|
||||
})
|
||||
const capabilityRequest = await waitForRequest(transport, 'pty.getCapabilities')
|
||||
transport.deliver(
|
||||
responseFrame(
|
||||
capabilityRequest.id as number,
|
||||
{ agentSessionCreateOperationVersion: AGENT_SESSION_CREATE_OPERATION_PROTOCOL_VERSION },
|
||||
1
|
||||
)
|
||||
)
|
||||
await waitForRequest(transport, 'pty.spawn')
|
||||
|
||||
mux.dispose('connection_lost')
|
||||
const failure = await spawn.catch((error: unknown) => error)
|
||||
|
||||
expect((failure as { code?: unknown }).code).toBe('CONNECTION_LOST')
|
||||
expect(failure).toMatchObject({ agentSessionOperationOutcome: 'unknown' })
|
||||
})
|
||||
|
||||
it('withholds same-turn source data until claim validation and isolates rollback', async () => {
|
||||
const transport = createTransport()
|
||||
const mux = new SshChannelMultiplexer(transport)
|
||||
|
||||
@@ -74,6 +74,20 @@ function sshMuxRequestTimeoutError(method: string, timeoutMs: number): Error {
|
||||
})
|
||||
}
|
||||
|
||||
// Why: `request` rejects a disposed mux and an already-aborted signal before it frames anything,
|
||||
// so the peer provably never saw those calls. A caller fencing a possibly-applied mutation needs
|
||||
// that apart from a mid-flight loss, which carries the identical code and message.
|
||||
export function isUndispatchedSshRequestError(error: unknown): boolean {
|
||||
return (
|
||||
(error as { sshRequestUndispatched?: unknown } | null | undefined)?.sshRequestUndispatched ===
|
||||
true
|
||||
)
|
||||
}
|
||||
|
||||
function markUndispatched<T extends Error>(error: T): T {
|
||||
return Object.assign(error, { sshRequestUndispatched: true as const })
|
||||
}
|
||||
|
||||
/**
|
||||
* True when a request may have run on the host despite failing here.
|
||||
*
|
||||
@@ -82,9 +96,13 @@ function sshMuxRequestTimeoutError(method: string, timeoutMs: number): Error {
|
||||
* wedged link lost at TIMEOUT_MS turned what used to surface as SSH_MUX_REQUEST_TIMEOUT into
|
||||
* CONNECTION_LOST, so callers that phrase the verdict to a user must branch on this rather than on
|
||||
* the timeout alone or they silently start reporting absence
|
||||
* (docs/reference/ssh-execution-boundary.md).
|
||||
* (docs/reference/ssh-execution-boundary.md). Only a request the multiplexer refused before framing
|
||||
* anything is provably un-run — hence the undispatched carve-out.
|
||||
*/
|
||||
export function isSshRequestOutcomeUnverifiable(error: unknown): boolean {
|
||||
if (isUndispatchedSshRequestError(error)) {
|
||||
return false
|
||||
}
|
||||
const code = error instanceof Error ? (error as Error & { code?: unknown }).code : undefined
|
||||
return code === SSH_MUX_REQUEST_TIMEOUT_CODE || code === 'CONNECTION_LOST'
|
||||
}
|
||||
@@ -230,12 +248,12 @@ export class SshChannelMultiplexer {
|
||||
options?: SshMultiplexerRequestOptions
|
||||
): Promise<unknown> {
|
||||
if (this.disposed) {
|
||||
throw this.disposedError()
|
||||
throw markUndispatched(this.disposedError())
|
||||
}
|
||||
if (options?.signal?.aborted) {
|
||||
const error = new Error(`Request "${method}" was cancelled`) as Error & { name: string }
|
||||
error.name = 'AbortError'
|
||||
throw error
|
||||
throw markUndispatched(error)
|
||||
}
|
||||
|
||||
const id = this.nextRequestId++
|
||||
|
||||
@@ -23,6 +23,14 @@ describe('SSH request outcome verdict', () => {
|
||||
expect(isSshRequestOutcomeUnverifiable(createSshDisposalError('connection_lost'))).toBe(true)
|
||||
})
|
||||
|
||||
it('does not claim unverifiable for a request that never reached the wire', () => {
|
||||
// A disposed mux rejects before framing anything, so the peer provably never saw it.
|
||||
const undispatched = Object.assign(createSshDisposalError('connection_lost'), {
|
||||
sshRequestUndispatched: true
|
||||
})
|
||||
expect(isSshRequestOutcomeUnverifiable(undispatched)).toBe(false)
|
||||
})
|
||||
|
||||
it('does not claim unverifiable for a deliberate shutdown', () => {
|
||||
expect(isSshRequestOutcomeUnverifiable(createSshDisposalError('shutdown'))).toBe(false)
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user