fix(ssh): scope every activate() release path to the record its caller owns

Two of the three release/cancel sites in RelayPtySourcePublication.activate()
acted on `current` unconditionally. A superseded transport re-entering activate()
therefore released — or cancelled and deleted — the delivery its own replacement
had just opened: releasing the fence resumes a send the replacement is still
rotating, and retiring it blanks the pane that owns it. The live path is the
unadmitted/subscriber branch, so guarding only the first site leaves the defect
exactly as it was; all three now act only on a record the caller still owns.

Also give the restore-required token its own toast copy. It must not join
UNREATTACHABLE_SESSION_SOURCES: that copy says "Open a new terminal", which here
abandons a running agent on a PTY the relay has just proven alive
(docs/reference/ssh-execution-boundary.md).
This commit is contained in:
Neil
2026-09-01 22:51:50 -07:00
parent 7f4a17d8eb
commit 1fc250a9e5
5 changed files with 210 additions and 13 deletions
+12 -12
View File
@@ -65,20 +65,24 @@ export class RelayPtySourcePublication {
context: RequestContext | undefined,
recovery?: PtySourceRecoveryRequest
): false | 'opened' | 'rotated' | 'existing' | PtySourceRecoveryResult {
let current = this.deliveries.get(id)
// A superseded request can find the delivery its own replacement opened: releasing that fence
// resumes a send the replacement is still rotating, and cancelling it blanks the pane that owns
// it. So every bail-out below acts only on a record this caller still owns.
const owned = current?.clientId === context?.clientId ? current : undefined
if (!context?.onResponseSettled) {
this.sender.releaseRotationFence(this.deliveries.get(id))
this.sender.releaseRotationFence(owned)
return false
}
const mode = this.session.deliveryMode(context.clientId)
let current = this.deliveries.get(id)
if (mode === 'unadmitted' || mode === 'subscriber') {
this.sender.releaseRotationFence(current)
this.sender.releaseRotationFence(owned)
return false
}
if (mode === 'legacy-owner') {
if (current) {
this.session.cancelDelivery(current.identity, 'source-credit-disabled')
this.sender.wakeSendWaiters(current)
if (owned) {
this.session.cancelDelivery(owned.identity, 'source-credit-disabled')
this.sender.wakeSendWaiters(owned)
this.deliveries.delete(id)
this.onCapacity(id)
}
@@ -88,7 +92,7 @@ export class RelayPtySourcePublication {
current?.clientId === context.clientId &&
!current.restoreRequired &&
current.sourceExitState !== 'pending' &&
this.deliveryClosedUnderRecord(current)
ptySourceDeliveryClosed(this.session, current.identity)
) {
// Why: a canceled delivery can never resume as 'existing'; retire it so re-attach opens fresh.
this.sender.wakeSendWaiters(current)
@@ -219,7 +223,7 @@ export class RelayPtySourcePublication {
}
if (!output.sourceAccepted && !appendPtySourceOutput(this.session, record, output)) {
this.counters.appendDenied++
if (this.deliveryClosedUnderRecord(record)) {
if (ptySourceDeliveryClosed(this.session, record.identity)) {
this.sender.wakeSendWaiters(record)
this.deliveries.delete(id)
// Why: deferred — publish() can run inside flushPendingOutput's captured-queue drain,
@@ -275,10 +279,6 @@ export class RelayPtySourcePublication {
this.sender.dispose()
}
private deliveryClosedUnderRecord(record: RelayPtySourceDeliveryRecord): boolean {
return ptySourceDeliveryClosed(this.session, record.identity)
}
private registerActivationSettlement(
id: string,
record: RelayPtySourceDeliveryRecord,
@@ -0,0 +1,161 @@
import { afterEach, describe, expect, it, vi } from 'vitest'
import {
RelayDispatcher,
type RelayClientSessionIdentity,
type RequestContext,
type SinkWriteSettlement
} from './dispatcher'
import { encodeJsonRpcFrame, MessageType } from './protocol'
import { RelayPtySourcePublication } from './relay-pty-source-publication'
import { SshPtyConsumerSessionAdapter } from './ssh-pty-consumer-session-adapter'
const endpointIdentity: RelayClientSessionIdentity = {
principal: 'endpoint-principal',
authenticated: true,
allowSessionOwner: true,
authenticationKind: 'endpoint-credential'
}
function requestFrame(id: number, method: string, params: Record<string, unknown>): Buffer {
return encodeJsonRpcFrame({ jsonrpc: '2.0', id, method, params }, id, 0)
}
function responseResult(buffer: Buffer): Record<string, unknown> | null {
if (buffer[0] !== MessageType.Regular) {
return null
}
const length = buffer.readUInt32BE(9)
const message = JSON.parse(buffer.subarray(13, 13 + length).toString('utf8'))
return message.id === undefined ? null : (message.result ?? null)
}
async function flushRequests(): Promise<void> {
await new Promise((resolve) => setImmediate(resolve))
}
describe('PTY source activation from a superseded owner', () => {
let dispatcher: RelayDispatcher | null = null
afterEach(() => {
dispatcher?.dispose()
dispatcher = null
})
async function createHarness() {
const writes: Buffer[] = []
dispatcher = new RelayDispatcher(
(data, onSettled) => {
writes.push(Buffer.from(data))
onSettled({ ok: true })
return true
},
{ supportsWriteCallback: true },
endpointIdentity
)
let publication: RelayPtySourcePublication
const adapter = new SshPtyConsumerSessionAdapter(dispatcher, 'build-a', undefined, (id) =>
publication.onCreditAvailable(id)
)
publication = new RelayPtySourcePublication(dispatcher, adapter, () => {})
dispatcher.feed(
requestFrame(1, 'pty.openClient', {
protocolVersion: 1,
clientInstanceId: 'client-1',
requestedRole: 'session-owner',
capabilities: { outputFlowControl: { versions: [1], requestedWindowSu: 4 } }
})
)
await flushRequests()
return { adapter, publication, writes }
}
function contextFor(
clientId: number,
settlements: ((result: SinkWriteSettlement) => void)[]
): RequestContext {
return {
clientId,
isStale: () => false,
sessionIdentity: endpointIdentity,
onResponseSettled: (callback) => settlements.push(callback)
}
}
/**
* The superseded transport must not release, cancel or retire the delivery its own replacement
* opened: releasing the fence resumes a send the replacement is still rotating, and retiring it
* blanks the pane that owns it.
*/
it('leaves the replacement delivery intact when the superseded owner re-activates', async () => {
const { publication, adapter, writes } = await createHarness()
const settlements: ((result: SinkWriteSettlement) => void)[] = []
expect(publication.activate('pty-1', 'incarnation-1', contextFor(1, settlements))).toBe(
'opened'
)
settlements[0]({ ok: true })
const activation = publication.receivingActivation('pty-1', 1)!
const ownerGrant = writes.map(responseResult).find((result) => result?.ownerLease)!
// The original transport arms its rotation fence while waiting for a checkpoint-safe send.
await expect(publication.waitForPendingSend('pty-1')).resolves.toBe(true)
const replacementWrites: Buffer[] = []
const replacementClientId = dispatcher!.attachClient(
(data, onSettled) => {
replacementWrites.push(Buffer.from(data))
onSettled({ ok: true })
return true
},
{ supportsWriteCallback: true },
endpointIdentity
)
dispatcher!.feedClient(
replacementClientId,
requestFrame(2, 'pty.openClient', {
protocolVersion: 1,
clientInstanceId: 'client-1',
requestedRole: 'session-owner',
resume: {
ownerGeneration: ownerGrant.ownerGeneration,
ownerLease: ownerGrant.ownerLease
},
capabilities: { outputFlowControl: { versions: [1], requestedWindowSu: 4 } }
})
)
await flushRequests()
const replacementSettlements: ((result: SinkWriteSettlement) => void)[] = []
const recovery = {
status: 'checkpoint' as const,
clientGeneration: activation.clientGeneration,
ownerGeneration: activation.ownerGeneration,
ptyIncarnation: activation.ptyIncarnation,
deliveryToken: activation.deliveryToken,
acceptedSourceEndSu: 0
}
expect(
publication.activate(
'pty-1',
'incarnation-1',
contextFor(replacementClientId, replacementSettlements),
recovery
)
).toMatchObject({ status: 'pending' })
replacementSettlements[0]({ ok: true })
// Arm the replacement fence, then let the superseded transport's activation resume. isStale()
// stays false on purpose: the superseded client is still attached, so production reaches the
// delivery-mode bail-outs rather than any stale early-out.
await expect(publication.waitForPendingSend('pty-1')).resolves.toBe(true)
const cancelDelivery = vi.spyOn(adapter, 'cancelDelivery')
expect(publication.activate('pty-1', 'incarnation-1', contextFor(1, []), recovery)).toBe(false)
expect(cancelDelivery).not.toHaveBeenCalled()
expect(publication.publish('pty-1', { data: 'replacement-output' }, false)).toBe(false)
// Release the replacement fence through its owning transport so no parked work is left behind.
expect(
publication.activate('pty-1', 'incarnation-1', contextFor(replacementClientId, []))
).toBe('existing')
expect(replacementWrites.length).toBeGreaterThan(0)
})
})
@@ -185,6 +185,19 @@ describe('humanizeTerminalError', () => {
expect(humanized).not.toContain('exited')
})
// A live PTY whose delivery was retired must never get the "open a new terminal" copy: acting on
// that abandons a running agent on the host.
it('describes a retired output source as reconnecting, not as a lost session', () => {
const humanized = humanizeTerminalError(
'SSH_PTY_SOURCE_RESTORE_REQUIRED: remote:2f1c:pty-7 checkpointUnavailable'
)
expect(humanized).not.toContain('SSH_PTY_SOURCE_RESTORE_REQUIRED')
expect(humanized).not.toContain('remote:2f1c:pty-7')
expect(humanized).not.toContain('checkpointUnavailable')
expect(humanized).not.toContain('Open a new terminal')
expect(humanized).toContain('still running')
})
it('replaces only the unreattachable line in an aggregated error', () => {
const humanized = humanizeTerminalError('Paste failed.\nSSH_SESSION_EXPIRED: orca:2f1c@@pty-7')
expect(humanized.startsWith('Paste failed.\n')).toBe(true)
@@ -216,6 +229,14 @@ describe('isExplainedTerminalError', () => {
).toBe(true)
})
it('suppresses the issue link while a live session restores its output source', () => {
expect(
isExplainedTerminalError(
'SSH_PTY_SOURCE_RESTORE_REQUIRED: remote:2f1c:pty-7 checkpointUnavailable'
)
).toBe(true)
})
it('suppresses the issue link for a session the host cannot reattach', () => {
expect(isExplainedTerminalError('SSH_SESSION_EXPIRED: orca:2f1c@@pty-7')).toBe(true)
expect(
@@ -35,6 +35,13 @@ const UNREATTACHABLE_SESSION_SOURCES = [
'PTY "[^"\\r\\n]*" not found(?: \\(identity mismatch\\))?',
'(?:SessionNotFoundError: )?Session not found: \\S+'
]
// The relay answered and proved the shell is still running — only its output delivery was retired.
// Deliberately NOT one of the sources above: that copy says to open a new terminal, which here
// abandons a live agent. Same lastIndex hazard, so keep the test and replace forms separate.
const SOURCE_RESTORE_REQUIRED_SOURCE =
'SSH_PTY_SOURCE_RESTORE_REQUIRED(?::[ \\t]*\\S*(?:[ \\t]+\\S+)?)?'
const SOURCE_RESTORE_REQUIRED_PATTERN = new RegExp(SOURCE_RESTORE_REQUIRED_SOURCE)
const SOURCE_RESTORE_REQUIRED_REPLACE_PATTERN = new RegExp(SOURCE_RESTORE_REQUIRED_SOURCE, 'g')
const UNREATTACHABLE_SESSION_PATTERNS = UNREATTACHABLE_SESSION_SOURCES.map(
(source) => new RegExp(source)
)
@@ -77,6 +84,7 @@ export function isExplainedTerminalError(error: string): boolean {
(line) =>
TERMINAL_HOST_GONE_PATTERN.test(line) ||
LEGACY_TERMINAL_HOST_GONE_PATTERN.test(line) ||
SOURCE_RESTORE_REQUIRED_PATTERN.test(line) ||
UNREATTACHABLE_SESSION_PATTERNS.some((pattern) => pattern.test(line))
)
}
@@ -113,6 +121,12 @@ export function humanizeTerminalError(error: string): string {
)
humanized = humanized.replaceAll(PANE_OWNER_UNVERIFIED_MARKER, () => explanation)
}
humanized = humanized.replace(SOURCE_RESTORE_REQUIRED_REPLACE_PATTERN, () =>
translate(
'auto.components.terminal.pane.TerminalErrorToast.sourceRestoring',
'Reconnecting this terminal — its output is being re-established. The session is still running.'
)
)
humanized = humanizeUnreattachableSession(humanized)
if (!isExplainedTerminalError(humanized)) {
return humanized
+2 -1
View File
@@ -3087,7 +3087,8 @@
"ownerUnknown": "Orca couldn't verify this terminal's owner.",
"42b283ecfc": "Orca couldn't safely reconnect this terminal because the host couldn't verify its saved session. Orca left the saved session unchanged. Click Retry to try reconnecting now. If it still cannot reconnect, open a new terminal.",
"e16012e31e": "The terminal daemon that owned this session exited, so the session and its scrollback could not be recovered. Open a new terminal to continue.",
"sessionUnavailable": "Orca couldn't reattach to this pane's terminal session on the host. Open a new terminal to continue."
"sessionUnavailable": "Orca couldn't reattach to this pane's terminal session on the host. Open a new terminal to continue.",
"sourceRestoring": "Reconnecting this terminal — its output is being re-established. The session is still running."
},
"TerminalProcessExitOverlay": {
"capacityTitle": "Git Bash console limit reached",