mirror of
https://github.com/stablyai/orca.git
synced 2026-10-03 00:02:19 +00:00
* feat(ssh): add pty.resumeClient and split SSH PTY process listing (#16741 T2 P5+P6)
P6: a relay now answers pty.resumeClient, which admits only an exact resume of
the existing session owner (it never mints a fresh claim when that owner is
gone). resumeSshPtyConsumerSession calls it with cancellation and authority
checks; an old relay's method-not-found becomes a pty_consumer_resume_unsupported
refusal that leaves the channel usable for pty.openClient. Owner grant
publication now rolls back if the response-settlement hook cannot be armed, and
the adapter exposes read-only owner and publication-settled queries.
P5: SshPtyProvider.listProcesses moves to ssh-pty-process-list unchanged, and
the notification-routing tests split into a shared fixture plus recovery and
recovery-activation files.
Nothing calls pty.resumeClient yet (T6). Ported by hunk from #16741
(a68b6f3531) without ownership-transfer (T7) or Bun runtime hunks.
* refactor(ssh): type the consumer-session transport as the members it uses
---------
Co-authored-by: m4air <m4air@m4airs-Air.localdomain>
This commit is contained in:
@@ -0,0 +1,326 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { createSubscription, sourceActivation } from './ssh-pty-notification-routing-test-fixture'
|
||||
|
||||
describe('subscribeSshPtyNotifications', () => {
|
||||
it('accepts non-empty recovery from the activation checkpoint', () => {
|
||||
const { handler, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ checkpointSourceEndSu: 4, recoveryEndSu: 8 })
|
||||
)
|
||||
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'next',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 8,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
lease.commit()
|
||||
expect(onData).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
data: 'next',
|
||||
source: expect.objectContaining({ sourceStartSu: 4, sourceEndSu: 8 })
|
||||
})
|
||||
)
|
||||
})
|
||||
|
||||
it('routes held and later recovery frames only to the private sink until commit', () => {
|
||||
const { handler, dataListeners, livePtyIds, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
const onRecoveryData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ checkpointSourceEndSu: 4, recoveryEndSu: 12 })
|
||||
)
|
||||
const publishSource = (data: string, sourceEndSu: number): void => {
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data,
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
}
|
||||
|
||||
publishSource('held', 8)
|
||||
const recoveryLease = lease.transferToRecovery(onRecoveryData)
|
||||
publishSource('next', 12)
|
||||
|
||||
expect(onRecoveryData.mock.calls.map(([payload]) => payload.data)).toEqual(['held', 'next'])
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
expect(livePtyIds).not.toContain('ssh:conn@@pty-1')
|
||||
|
||||
recoveryLease.commit()
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
publishSource('live', 16)
|
||||
|
||||
expect(onRecoveryData).toHaveBeenCalledTimes(2)
|
||||
expect(onData).toHaveBeenCalledWith(expect.objectContaining({ data: 'live' }))
|
||||
expect(livePtyIds).toContain('ssh:conn@@pty-1')
|
||||
})
|
||||
|
||||
it('retires an exited private recovery when its activation commits', () => {
|
||||
const { handler, mux, dataListeners, livePtyIds, installReceivingActivation } =
|
||||
createSubscription()
|
||||
const onData = vi.fn()
|
||||
const onRecoveryData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation('pty-1', sourceActivation({ recoveryEndSu: 4 }))
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'held',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 4,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
const recoveryLease = lease.transferToRecovery(onRecoveryData)
|
||||
|
||||
handler('pty.exit', { id: 'pty-1', code: 0, incarnationId: 'incarnation-1' })
|
||||
recoveryLease.commit()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'late',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 8,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
|
||||
expect(onRecoveryData).toHaveBeenCalledOnce()
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
expect(livePtyIds).not.toContain('ssh:conn@@pty-1')
|
||||
expect(mux.request).toHaveBeenCalledWith(
|
||||
'pty.cancelDelivery',
|
||||
expect.objectContaining({ id: 'pty-1', deliveryToken: 'token-1' })
|
||||
)
|
||||
})
|
||||
|
||||
it('retires private recovery locally and restores the exact predecessor', () => {
|
||||
const { handler, mux, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
const onRecoveryData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ deliveryToken: 'token-old', recoveryEndSu: 3 })
|
||||
).commit()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'pre',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
const replacement = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-new',
|
||||
checkpointSourceEndSu: 3,
|
||||
recoveryEndSu: 6
|
||||
})
|
||||
)
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'new',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-new',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
replacement.transferToRecovery(onRecoveryData).retire()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'old',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
expect(onRecoveryData).toHaveBeenCalledWith(expect.objectContaining({ data: 'new' }))
|
||||
expect(onData.mock.calls.map(([payload]) => payload.data)).toEqual(['pre', 'old'])
|
||||
expect(mux.request).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('rejects a stale activation without disturbing current continuity', () => {
|
||||
const { handler, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
installReceivingActivation('pty-1', sourceActivation({ recoveryEndSu: 3 })).commit()
|
||||
|
||||
expect(() =>
|
||||
installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 1,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-stale',
|
||||
checkpointSourceEndSu: 3,
|
||||
recoveryEndSu: 3
|
||||
})
|
||||
)
|
||||
).toThrow('ssh_source_receiving_activation_stale')
|
||||
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'one',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
expect(onData).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('drops provisional frames and settles cancellation before rollback completes', async () => {
|
||||
const { handler, mux, dataListeners, livePtyIds, installReceivingActivation } =
|
||||
createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ checkpointSourceEndSu: 4, recoveryEndSu: 8 })
|
||||
)
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'next',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 8,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
|
||||
await expect(lease.rollback()).resolves.toBe(true)
|
||||
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
expect(livePtyIds).not.toContain('ssh:conn@@pty-1')
|
||||
expect(mux.request).toHaveBeenCalledWith('pty.cancelDelivery', {
|
||||
id: 'pty-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
deliveryToken: 'token-1'
|
||||
})
|
||||
})
|
||||
|
||||
it('restores the exact prior cursor when a replacement rolls back after frames', async () => {
|
||||
const { handler, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
installReceivingActivation('pty-1', sourceActivation({ deliveryToken: 'token-old' })).commit()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'pre',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
const replacement = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-new',
|
||||
checkpointSourceEndSu: 3,
|
||||
recoveryEndSu: 3
|
||||
})
|
||||
)
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'new',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-new',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
await replacement.rollback()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'old',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
expect(onData.mock.calls.map(([payload]) => payload.data)).toEqual(['pre', 'old'])
|
||||
})
|
||||
|
||||
it('does not let an older lease rollback replace a newer activation', async () => {
|
||||
const { handler, mux, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const older = installReceivingActivation('pty-1', sourceActivation())
|
||||
const newer = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-new'
|
||||
})
|
||||
)
|
||||
|
||||
await older.rollback()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'new',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-new',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
newer.commit()
|
||||
|
||||
expect(onData).toHaveBeenCalledWith(expect.objectContaining({ data: 'new' }))
|
||||
expect(mux.request).toHaveBeenCalledWith('pty.cancelDelivery', {
|
||||
id: 'pty-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
deliveryToken: 'token-1'
|
||||
})
|
||||
expect(mux.request).not.toHaveBeenCalledWith(
|
||||
'pty.cancelDelivery',
|
||||
expect.objectContaining({ deliveryToken: 'token-new' })
|
||||
)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,281 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { subscribeSshPtyNotifications } from './ssh-pty-notification-routing'
|
||||
import type { PtySourceReceivingActivation } from '../../shared/pty-source-receiving-activation'
|
||||
|
||||
type MockMux = {
|
||||
onNotification: ReturnType<typeof vi.fn>
|
||||
request: ReturnType<typeof vi.fn>
|
||||
}
|
||||
|
||||
function createSubscription() {
|
||||
const mux: MockMux = {
|
||||
onNotification: vi.fn(),
|
||||
request: vi.fn(async () => ({ canceled: true, sentEndSu: 0, creditedEndSu: 0 }))
|
||||
}
|
||||
const dataListeners = new Set<(payload: { id: string; data: string }) => void>()
|
||||
const replayListeners = new Set<(payload: { id: string; data: string }) => void>()
|
||||
const exitListeners = new Set<(payload: { id: string; code: number }) => void>()
|
||||
const livePtyIds = new Set<string>()
|
||||
const recordExit = vi.fn()
|
||||
const toAppPtyId = vi.fn((id: string) => `ssh:conn@@${id}`)
|
||||
const resolvePtyIncarnation = vi.fn((id: string) => `incarnation:${id}`)
|
||||
|
||||
const subscription = subscribeSshPtyNotifications({
|
||||
mux: mux as never,
|
||||
toAppPtyId,
|
||||
dataListeners: dataListeners as never,
|
||||
replayListeners: replayListeners as never,
|
||||
exitListeners: exitListeners as never,
|
||||
livePtyIds,
|
||||
recordExit,
|
||||
providerGeneration: 7,
|
||||
resolvePtyIncarnation,
|
||||
peekPtyIncarnation: () => undefined
|
||||
})
|
||||
const handler = mux.onNotification.mock.calls[0]?.[0] as (
|
||||
method: string,
|
||||
params: Record<string, unknown>
|
||||
) => void
|
||||
if (!handler) {
|
||||
throw new Error('notification handler was not registered')
|
||||
}
|
||||
return {
|
||||
handler,
|
||||
mux,
|
||||
toAppPtyId,
|
||||
dataListeners,
|
||||
replayListeners,
|
||||
exitListeners,
|
||||
livePtyIds,
|
||||
recordExit,
|
||||
resolvePtyIncarnation,
|
||||
installReceivingActivation: subscription.installReceivingActivation
|
||||
}
|
||||
}
|
||||
|
||||
function sourceActivation(
|
||||
overrides: Partial<PtySourceReceivingActivation> = {}
|
||||
): PtySourceReceivingActivation {
|
||||
return Object.freeze({
|
||||
status: 'pending',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
checkpointSourceEndSu: 0,
|
||||
recoveryEndSu: 0,
|
||||
...overrides
|
||||
})
|
||||
}
|
||||
|
||||
describe('SSH PTY notification recovery routing', () => {
|
||||
it('rejects a stale activation without disturbing current continuity', () => {
|
||||
const { handler, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
installReceivingActivation('pty-1', sourceActivation({ recoveryEndSu: 3 })).commit()
|
||||
|
||||
expect(() =>
|
||||
installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 1,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-stale',
|
||||
checkpointSourceEndSu: 3,
|
||||
recoveryEndSu: 3
|
||||
})
|
||||
)
|
||||
).toThrow('ssh_source_receiving_activation_stale')
|
||||
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'one',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
expect(onData).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('drops provisional frames and settles cancellation before rollback completes', async () => {
|
||||
const { handler, mux, dataListeners, livePtyIds, installReceivingActivation } =
|
||||
createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ checkpointSourceEndSu: 4, recoveryEndSu: 8 })
|
||||
)
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'next',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 8,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
|
||||
await expect(lease.rollback()).resolves.toBe(true)
|
||||
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
expect(livePtyIds).not.toContain('ssh:conn@@pty-1')
|
||||
expect(mux.request).toHaveBeenCalledWith('pty.cancelDelivery', {
|
||||
id: 'pty-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
deliveryToken: 'token-1'
|
||||
})
|
||||
})
|
||||
|
||||
it('restores the exact prior cursor when a replacement rolls back after frames', async () => {
|
||||
const { handler, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
installReceivingActivation('pty-1', sourceActivation({ deliveryToken: 'token-old' })).commit()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'pre',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
const replacement = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-new',
|
||||
checkpointSourceEndSu: 3,
|
||||
recoveryEndSu: 3
|
||||
})
|
||||
)
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'new',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-new',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
await replacement.rollback()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'old',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
expect(onData.mock.calls.map(([payload]) => payload.data)).toEqual(['pre', 'old'])
|
||||
})
|
||||
|
||||
it('does not let an older lease rollback replace a newer activation', async () => {
|
||||
const { handler, mux, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const older = installReceivingActivation('pty-1', sourceActivation())
|
||||
const newer = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-new'
|
||||
})
|
||||
)
|
||||
|
||||
await older.rollback()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'new',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-new',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
newer.commit()
|
||||
|
||||
expect(onData).toHaveBeenCalledWith(expect.objectContaining({ data: 'new' }))
|
||||
expect(mux.request).toHaveBeenCalledWith('pty.cancelDelivery', {
|
||||
id: 'pty-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
deliveryToken: 'token-1'
|
||||
})
|
||||
expect(mux.request).not.toHaveBeenCalledWith(
|
||||
'pty.cancelDelivery',
|
||||
expect.objectContaining({ deliveryToken: 'token-new' })
|
||||
)
|
||||
})
|
||||
|
||||
it('ignores PTY methods with missing ids', () => {
|
||||
const { handler, toAppPtyId, dataListeners } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
|
||||
expect(() => handler('pty.data', { data: 'orphan' })).not.toThrow()
|
||||
expect(toAppPtyId).not.toHaveBeenCalled()
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('leaves recovery and cancellation control methods to their dedicated handlers', () => {
|
||||
const {
|
||||
handler,
|
||||
mux,
|
||||
toAppPtyId,
|
||||
dataListeners,
|
||||
replayListeners,
|
||||
exitListeners,
|
||||
livePtyIds,
|
||||
recordExit,
|
||||
resolvePtyIncarnation
|
||||
} = createSubscription()
|
||||
const onData = vi.fn()
|
||||
const onReplay = vi.fn()
|
||||
const onExit = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
replayListeners.add(onReplay)
|
||||
exitListeners.add(onExit)
|
||||
livePtyIds.add('ssh:conn@@unrelated')
|
||||
|
||||
for (const method of [
|
||||
'pty.recoveryData',
|
||||
'pty.recoveryComplete',
|
||||
'pty.restoreRequired',
|
||||
'pty.deliveryCanceled'
|
||||
]) {
|
||||
handler(method, {
|
||||
id: 'pty-1',
|
||||
data: 'control',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3
|
||||
})
|
||||
}
|
||||
|
||||
expect(toAppPtyId).not.toHaveBeenCalled()
|
||||
expect(resolvePtyIncarnation).not.toHaveBeenCalled()
|
||||
expect(recordExit).not.toHaveBeenCalled()
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
expect(onReplay).not.toHaveBeenCalled()
|
||||
expect(onExit).not.toHaveBeenCalled()
|
||||
expect(livePtyIds).toEqual(new Set(['ssh:conn@@unrelated']))
|
||||
expect(mux.request).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,94 @@
|
||||
import { vi, type Mock } from 'vitest'
|
||||
import {
|
||||
subscribeSshPtyNotifications,
|
||||
type SshPtyNotificationSubscription
|
||||
} from './ssh-pty-notification-routing'
|
||||
import type { PtySourceReceivingActivation } from '../../shared/pty-source-receiving-activation'
|
||||
|
||||
type NotificationHandler = (method: string, params: Record<string, unknown>) => void
|
||||
type CancellationRequest = (
|
||||
method: string,
|
||||
params: Record<string, unknown>
|
||||
) => Promise<{ canceled: boolean; sentEndSu: number; creditedEndSu: number }>
|
||||
|
||||
type MockMux = {
|
||||
onNotification: Mock<(handler: NotificationHandler) => void>
|
||||
request: Mock<CancellationRequest>
|
||||
}
|
||||
|
||||
export type SshPtyNotificationTestSubscription = {
|
||||
handler: NotificationHandler
|
||||
mux: MockMux
|
||||
toAppPtyId: Mock<(id: string) => string>
|
||||
dataListeners: Set<(payload: { id: string; data: string }) => void>
|
||||
replayListeners: Set<(payload: { id: string; data: string }) => void>
|
||||
exitListeners: Set<(payload: { id: string; code: number }) => void>
|
||||
livePtyIds: Set<string>
|
||||
recordExit: Mock<(relayPtyId: string, incarnationId: unknown) => void>
|
||||
resolvePtyIncarnation: Mock<(id: string) => string>
|
||||
installReceivingActivation: SshPtyNotificationSubscription['installReceivingActivation']
|
||||
}
|
||||
|
||||
export function createSubscription(): SshPtyNotificationTestSubscription {
|
||||
const mux: MockMux = {
|
||||
onNotification: vi.fn<(handler: NotificationHandler) => void>(),
|
||||
request: vi.fn<CancellationRequest>(async () => ({
|
||||
canceled: true,
|
||||
sentEndSu: 0,
|
||||
creditedEndSu: 0
|
||||
}))
|
||||
}
|
||||
const dataListeners = new Set<(payload: { id: string; data: string }) => void>()
|
||||
const replayListeners = new Set<(payload: { id: string; data: string }) => void>()
|
||||
const exitListeners = new Set<(payload: { id: string; code: number }) => void>()
|
||||
const livePtyIds = new Set<string>()
|
||||
const recordExit = vi.fn<(relayPtyId: string, incarnationId: unknown) => void>()
|
||||
const toAppPtyId = vi.fn((id: string) => `ssh:conn@@${id}`)
|
||||
const resolvePtyIncarnation = vi.fn((id: string) => `incarnation:${id}`)
|
||||
|
||||
const subscription = subscribeSshPtyNotifications({
|
||||
mux: mux as never,
|
||||
toAppPtyId,
|
||||
dataListeners: dataListeners as never,
|
||||
replayListeners: replayListeners as never,
|
||||
exitListeners: exitListeners as never,
|
||||
livePtyIds,
|
||||
recordExit,
|
||||
providerGeneration: 7,
|
||||
resolvePtyIncarnation,
|
||||
peekPtyIncarnation: () => undefined
|
||||
})
|
||||
|
||||
const handler = mux.onNotification.mock.calls[0]?.[0]
|
||||
if (!handler) {
|
||||
throw new Error('notification handler was not registered')
|
||||
}
|
||||
|
||||
return {
|
||||
handler,
|
||||
mux,
|
||||
toAppPtyId,
|
||||
dataListeners,
|
||||
replayListeners,
|
||||
exitListeners,
|
||||
livePtyIds,
|
||||
recordExit,
|
||||
resolvePtyIncarnation,
|
||||
installReceivingActivation: subscription.installReceivingActivation
|
||||
}
|
||||
}
|
||||
|
||||
export function sourceActivation(
|
||||
overrides: Partial<PtySourceReceivingActivation> = {}
|
||||
): PtySourceReceivingActivation {
|
||||
return Object.freeze({
|
||||
status: 'pending',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
checkpointSourceEndSu: 0,
|
||||
recoveryEndSu: 0,
|
||||
...overrides
|
||||
})
|
||||
}
|
||||
@@ -1,74 +1,5 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { subscribeSshPtyNotifications } from './ssh-pty-notification-routing'
|
||||
import type { PtySourceReceivingActivation } from '../../shared/pty-source-receiving-activation'
|
||||
|
||||
type MockMux = {
|
||||
onNotification: ReturnType<typeof vi.fn>
|
||||
request: ReturnType<typeof vi.fn>
|
||||
}
|
||||
|
||||
function createSubscription() {
|
||||
const mux: MockMux = {
|
||||
onNotification: vi.fn(),
|
||||
request: vi.fn(async () => ({ canceled: true, sentEndSu: 0, creditedEndSu: 0 }))
|
||||
}
|
||||
const dataListeners = new Set<(payload: { id: string; data: string }) => void>()
|
||||
const replayListeners = new Set<(payload: { id: string; data: string }) => void>()
|
||||
const exitListeners = new Set<(payload: { id: string; code: number }) => void>()
|
||||
const livePtyIds = new Set<string>()
|
||||
const recordExit = vi.fn()
|
||||
const toAppPtyId = vi.fn((id: string) => `ssh:conn@@${id}`)
|
||||
const resolvePtyIncarnation = vi.fn((id: string) => `incarnation:${id}`)
|
||||
|
||||
const subscription = subscribeSshPtyNotifications({
|
||||
mux: mux as never,
|
||||
toAppPtyId,
|
||||
dataListeners: dataListeners as never,
|
||||
replayListeners: replayListeners as never,
|
||||
exitListeners: exitListeners as never,
|
||||
livePtyIds,
|
||||
recordExit,
|
||||
providerGeneration: 7,
|
||||
resolvePtyIncarnation,
|
||||
peekPtyIncarnation: () => undefined
|
||||
})
|
||||
|
||||
const handler = mux.onNotification.mock.calls[0]?.[0] as (
|
||||
method: string,
|
||||
params: Record<string, unknown>
|
||||
) => void
|
||||
if (!handler) {
|
||||
throw new Error('notification handler was not registered')
|
||||
}
|
||||
|
||||
return {
|
||||
handler,
|
||||
mux,
|
||||
toAppPtyId,
|
||||
dataListeners,
|
||||
replayListeners,
|
||||
exitListeners,
|
||||
livePtyIds,
|
||||
recordExit,
|
||||
resolvePtyIncarnation,
|
||||
installReceivingActivation: subscription.installReceivingActivation
|
||||
}
|
||||
}
|
||||
|
||||
function sourceActivation(
|
||||
overrides: Partial<PtySourceReceivingActivation> = {}
|
||||
): PtySourceReceivingActivation {
|
||||
return Object.freeze({
|
||||
status: 'pending',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
checkpointSourceEndSu: 0,
|
||||
recoveryEndSu: 0,
|
||||
...overrides
|
||||
})
|
||||
}
|
||||
import { createSubscription, sourceActivation } from './ssh-pty-notification-routing-test-fixture'
|
||||
|
||||
describe('subscribeSshPtyNotifications', () => {
|
||||
it('ignores non-PTY notifications without mapping params.id', () => {
|
||||
@@ -480,328 +411,6 @@ describe('subscribeSshPtyNotifications', () => {
|
||||
expect(mux.request).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('accepts non-empty recovery from the activation checkpoint', () => {
|
||||
const { handler, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ checkpointSourceEndSu: 4, recoveryEndSu: 8 })
|
||||
)
|
||||
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'next',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 8,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
lease.commit()
|
||||
expect(onData).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
data: 'next',
|
||||
source: expect.objectContaining({ sourceStartSu: 4, sourceEndSu: 8 })
|
||||
})
|
||||
)
|
||||
})
|
||||
|
||||
it('routes held and later recovery frames only to the private sink until commit', () => {
|
||||
const { handler, dataListeners, livePtyIds, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
const onRecoveryData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ checkpointSourceEndSu: 4, recoveryEndSu: 12 })
|
||||
)
|
||||
const publishSource = (data: string, sourceEndSu: number): void => {
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data,
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
}
|
||||
|
||||
publishSource('held', 8)
|
||||
const recoveryLease = lease.transferToRecovery(onRecoveryData)
|
||||
publishSource('next', 12)
|
||||
|
||||
expect(onRecoveryData.mock.calls.map(([payload]) => payload.data)).toEqual(['held', 'next'])
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
expect(livePtyIds).not.toContain('ssh:conn@@pty-1')
|
||||
|
||||
recoveryLease.commit()
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
publishSource('live', 16)
|
||||
|
||||
expect(onRecoveryData).toHaveBeenCalledTimes(2)
|
||||
expect(onData).toHaveBeenCalledWith(expect.objectContaining({ data: 'live' }))
|
||||
expect(livePtyIds).toContain('ssh:conn@@pty-1')
|
||||
})
|
||||
|
||||
it('retires an exited private recovery when its activation commits', () => {
|
||||
const { handler, mux, dataListeners, livePtyIds, installReceivingActivation } =
|
||||
createSubscription()
|
||||
const onData = vi.fn()
|
||||
const onRecoveryData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation('pty-1', sourceActivation({ recoveryEndSu: 4 }))
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'held',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 4,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
const recoveryLease = lease.transferToRecovery(onRecoveryData)
|
||||
|
||||
handler('pty.exit', { id: 'pty-1', code: 0, incarnationId: 'incarnation-1' })
|
||||
recoveryLease.commit()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'late',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 8,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
|
||||
expect(onRecoveryData).toHaveBeenCalledOnce()
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
expect(livePtyIds).not.toContain('ssh:conn@@pty-1')
|
||||
expect(mux.request).toHaveBeenCalledWith(
|
||||
'pty.cancelDelivery',
|
||||
expect.objectContaining({ id: 'pty-1', deliveryToken: 'token-1' })
|
||||
)
|
||||
})
|
||||
|
||||
it('retires private recovery locally and restores the exact predecessor', () => {
|
||||
const { handler, mux, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
const onRecoveryData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ deliveryToken: 'token-old', recoveryEndSu: 3 })
|
||||
).commit()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'pre',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
const replacement = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-new',
|
||||
checkpointSourceEndSu: 3,
|
||||
recoveryEndSu: 6
|
||||
})
|
||||
)
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'new',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-new',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
replacement.transferToRecovery(onRecoveryData).retire()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'old',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
expect(onRecoveryData).toHaveBeenCalledWith(expect.objectContaining({ data: 'new' }))
|
||||
expect(onData.mock.calls.map(([payload]) => payload.data)).toEqual(['pre', 'old'])
|
||||
expect(mux.request).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('rejects a stale activation without disturbing current continuity', () => {
|
||||
const { handler, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
installReceivingActivation('pty-1', sourceActivation({ recoveryEndSu: 3 })).commit()
|
||||
|
||||
expect(() =>
|
||||
installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 1,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-stale',
|
||||
checkpointSourceEndSu: 3,
|
||||
recoveryEndSu: 3
|
||||
})
|
||||
)
|
||||
).toThrow('ssh_source_receiving_activation_stale')
|
||||
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'one',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
expect(onData).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('drops provisional frames and settles cancellation before rollback completes', async () => {
|
||||
const { handler, mux, dataListeners, livePtyIds, installReceivingActivation } =
|
||||
createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const lease = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({ checkpointSourceEndSu: 4, recoveryEndSu: 8 })
|
||||
)
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'next',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 8,
|
||||
sourceLengthSu: 4
|
||||
})
|
||||
|
||||
await expect(lease.rollback()).resolves.toBe(true)
|
||||
|
||||
expect(onData).not.toHaveBeenCalled()
|
||||
expect(livePtyIds).not.toContain('ssh:conn@@pty-1')
|
||||
expect(mux.request).toHaveBeenCalledWith('pty.cancelDelivery', {
|
||||
id: 'pty-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
deliveryToken: 'token-1'
|
||||
})
|
||||
})
|
||||
|
||||
it('restores the exact prior cursor when a replacement rolls back after frames', async () => {
|
||||
const { handler, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
installReceivingActivation('pty-1', sourceActivation({ deliveryToken: 'token-old' })).commit()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'pre',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
const replacement = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-new',
|
||||
checkpointSourceEndSu: 3,
|
||||
recoveryEndSu: 3
|
||||
})
|
||||
)
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'new',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-new',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
await replacement.rollback()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'old',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-old',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
sourceEndSu: 6,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
|
||||
expect(onData.mock.calls.map(([payload]) => payload.data)).toEqual(['pre', 'old'])
|
||||
})
|
||||
|
||||
it('does not let an older lease rollback replace a newer activation', async () => {
|
||||
const { handler, mux, dataListeners, installReceivingActivation } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
dataListeners.add(onData)
|
||||
const older = installReceivingActivation('pty-1', sourceActivation())
|
||||
const newer = installReceivingActivation(
|
||||
'pty-1',
|
||||
sourceActivation({
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
deliveryToken: 'token-new'
|
||||
})
|
||||
)
|
||||
|
||||
await older.rollback()
|
||||
handler('pty.data', {
|
||||
id: 'pty-1',
|
||||
data: 'new',
|
||||
ptyIncarnation: 'incarnation-1',
|
||||
deliveryToken: 'token-new',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 4,
|
||||
sourceEndSu: 3,
|
||||
sourceLengthSu: 3
|
||||
})
|
||||
newer.commit()
|
||||
|
||||
expect(onData).toHaveBeenCalledWith(expect.objectContaining({ data: 'new' }))
|
||||
expect(mux.request).toHaveBeenCalledWith('pty.cancelDelivery', {
|
||||
id: 'pty-1',
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 3,
|
||||
deliveryToken: 'token-1'
|
||||
})
|
||||
expect(mux.request).not.toHaveBeenCalledWith(
|
||||
'pty.cancelDelivery',
|
||||
expect.objectContaining({ deliveryToken: 'token-new' })
|
||||
)
|
||||
})
|
||||
|
||||
it('ignores PTY methods with missing ids', () => {
|
||||
const { handler, toAppPtyId, dataListeners } = createSubscription()
|
||||
const onData = vi.fn()
|
||||
|
||||
@@ -0,0 +1,52 @@
|
||||
import type { SshChannelMultiplexer } from '../ssh/ssh-channel-multiplexer'
|
||||
import type { IPtyProvider, PtyProcessInfo } from './types'
|
||||
import { toAppSshPtyId, toRelaySshPtyId } from './ssh-pty-id'
|
||||
import { mapSshPtyProcessList } from './ssh-agent-session-process-list'
|
||||
import type { SshPtyProviderOutputState } from './ssh-pty-provider-output-state'
|
||||
|
||||
export function createSshPtyProcessLister(
|
||||
args: Pick<
|
||||
Parameters<typeof listSshPtyProcesses>[0],
|
||||
'mux' | 'connectionId' | 'livePtyIds' | 'outputState'
|
||||
>
|
||||
): IPtyProvider['listProcesses'] {
|
||||
return (options) =>
|
||||
listSshPtyProcesses({
|
||||
...args,
|
||||
includeForegroundProcessEvidence: options?.includeForegroundProcessEvidence,
|
||||
deadlineMs: options?.deadlineMs
|
||||
})
|
||||
}
|
||||
|
||||
export async function listSshPtyProcesses(
|
||||
args: Readonly<{
|
||||
mux: SshChannelMultiplexer
|
||||
connectionId: string
|
||||
livePtyIds: Set<string>
|
||||
outputState: SshPtyProviderOutputState
|
||||
includeForegroundProcessEvidence?: boolean
|
||||
deadlineMs?: number
|
||||
}>
|
||||
): Promise<PtyProcessInfo[]> {
|
||||
const result = await args.mux.request(
|
||||
'pty.listProcesses',
|
||||
args.includeForegroundProcessEvidence === undefined
|
||||
? undefined
|
||||
: { includeForegroundProcessEvidence: args.includeForegroundProcessEvidence },
|
||||
args.deadlineMs === undefined
|
||||
? undefined
|
||||
: { timeoutMs: Math.max(1, args.deadlineMs - Date.now()) }
|
||||
)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the relay answers pty.listProcesses with PtyProcessInfo rows; the mapper rejects unproven ownership.
|
||||
const processes = mapSshPtyProcessList(result as PtyProcessInfo[], (id) =>
|
||||
toAppSshPtyId(args.connectionId, id)
|
||||
)
|
||||
for (const process of processes) {
|
||||
args.livePtyIds.add(process.id)
|
||||
args.outputState.rememberPtyIncarnation(
|
||||
toRelaySshPtyId(args.connectionId, process.id),
|
||||
process.incarnationId
|
||||
)
|
||||
}
|
||||
return processes
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
import type { SshChannelMultiplexer } from '../ssh/ssh-channel-multiplexer'
|
||||
import type { IPtyProvider, PtyProcessInfo, PtySpawnOptions, PtySpawnResult } from './types'
|
||||
import type { IPtyProvider, PtySpawnOptions, PtySpawnResult } from './types'
|
||||
import type { WriteSettlement } from '../../shared/pty-write-settlement'
|
||||
import type { TerminalOscColorQueryReplyColors } from '../../shared/terminal-osc-color-reply'
|
||||
import { toAppSshPtyId, toRelaySshPtyId } from './ssh-pty-id'
|
||||
@@ -13,7 +13,6 @@ import type {
|
||||
} from './ssh-pty-provider-contract'
|
||||
import { SshPtyProviderOutputState } from './ssh-pty-provider-output-state'
|
||||
import { spawnFreshSshPty } from './ssh-agent-session-create-operation'
|
||||
import { mapSshPtyProcessList } from './ssh-agent-session-process-list'
|
||||
import {
|
||||
requestSshPtyAttach,
|
||||
reattachSshPtySessionForSpawn,
|
||||
@@ -26,6 +25,7 @@ import { SshAgentSessionCapabilities } from './ssh-agent-session-capabilities'
|
||||
import type { PtyProcessInspection } from './pty-process-inspection'
|
||||
import { spawnWithTerminalRuntimeRepair, type TerminalRepairHook } from './ssh-pty-spawn-repair'
|
||||
import { createSshPtyProviderRpcOperations } from './ssh-pty-provider-rpc-operations'
|
||||
import { createSshPtyProcessLister } from './ssh-pty-process-list'
|
||||
|
||||
// Why: sequential relay teardown calls share one absolute budget; convert to the mux-relative timeout only at dispatch.
|
||||
function relayTimeoutOptions(deadlineMs: number | undefined): { timeoutMs: number } | undefined {
|
||||
@@ -38,6 +38,7 @@ export class SshPtyProvider implements IPtyProvider {
|
||||
private connectionId: string
|
||||
private livePtyIds = new Set<string>()
|
||||
readonly getAppliedSize: NonNullable<IPtyProvider['getAppliedSize']>
|
||||
readonly listProcesses: IPtyProvider['listProcesses']
|
||||
private readonly agentSessionCapabilities: SshAgentSessionCapabilities
|
||||
private spawnExitRaces = new SshPtySpawnExitRaceTracker()
|
||||
private readonly outputState: SshPtyProviderOutputState
|
||||
@@ -101,6 +102,12 @@ export class SshPtyProvider implements IPtyProvider {
|
||||
this.spawnExitRaces.recordExit(relayPtyId, incarnationId)
|
||||
}
|
||||
})
|
||||
this.listProcesses = createSshPtyProcessLister({
|
||||
mux,
|
||||
connectionId,
|
||||
livePtyIds: this.livePtyIds,
|
||||
outputState: this.outputState
|
||||
})
|
||||
}
|
||||
|
||||
dispose(): void {
|
||||
@@ -270,26 +277,6 @@ export class SshPtyProvider implements IPtyProvider {
|
||||
this.livePtyIds.delete(id)
|
||||
}
|
||||
|
||||
async listProcesses(opts?: {
|
||||
deadlineMs?: number
|
||||
includeForegroundProcessEvidence?: boolean
|
||||
}): Promise<PtyProcessInfo[]> {
|
||||
const result = await this.mux.request(
|
||||
'pty.listProcesses',
|
||||
opts?.includeForegroundProcessEvidence === undefined
|
||||
? undefined
|
||||
: { includeForegroundProcessEvidence: opts.includeForegroundProcessEvidence },
|
||||
relayTimeoutOptions(opts?.deadlineMs)
|
||||
)
|
||||
const processes = mapSshPtyProcessList(result as PtyProcessInfo[], (id) => this.toAppPtyId(id))
|
||||
for (const process of processes) {
|
||||
this.livePtyIds.add(process.id)
|
||||
const relayPtyId = this.toRelayPtyId(process.id)
|
||||
this.outputState.rememberPtyIncarnation(relayPtyId, process.incarnationId)
|
||||
}
|
||||
return processes
|
||||
}
|
||||
|
||||
hasPty = (id: string): boolean => this.livePtyIds.has(id)
|
||||
|
||||
onData = (callback: SshPtyDataCallback): (() => void) => this.outputState.onData(callback)
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { buildSshPtySpawnEnv } from './ssh-pty-spawn-env'
|
||||
|
||||
describe('buildSshPtySpawnEnv relay bridge', () => {
|
||||
it('prepends the CLI bin dir once and publishes the Node relay bridge', () => {
|
||||
const env = buildSshPtySpawnEnv({
|
||||
env: { PATH: '/home/me/.orca-relay/bin:/usr/bin' },
|
||||
remoteCliBridgeEnv: {
|
||||
binDir: '/home/me/.orca-relay/bin',
|
||||
relayDir: '/home/me/.orca-relay/relay-v1',
|
||||
nodePath: '/usr/bin/node',
|
||||
sockPath: '/home/me/.orca-relay/relay.sock'
|
||||
}
|
||||
})
|
||||
|
||||
expect(env).toMatchObject({
|
||||
PATH: '/home/me/.orca-relay/bin:/usr/bin',
|
||||
ORCA_REMOTE_CLI_BIN_DIR: '/home/me/.orca-relay/bin',
|
||||
ORCA_RELAY_DIR: '/home/me/.orca-relay/relay-v1',
|
||||
ORCA_RELAY_NODE_PATH: '/usr/bin/node',
|
||||
ORCA_RELAY_SOCKET_PATH: '/home/me/.orca-relay/relay.sock'
|
||||
})
|
||||
expect(env).not.toHaveProperty('ORCA_RELAY_CREDENTIAL_FILE')
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,213 @@
|
||||
import { expect, it, vi } from 'vitest'
|
||||
import { PTY_CONSUMER_SESSION_PROTOCOL_VERSION } from '../../shared/pty-consumer-session'
|
||||
import type { SshChannelMultiplexer } from './ssh-channel-multiplexer'
|
||||
import { openSshPtyConsumerSession, resumeSshPtyConsumerSession } from './ssh-pty-consumer-session'
|
||||
|
||||
function fixture() {
|
||||
const controller = new AbortController()
|
||||
const grant = {
|
||||
protocolVersion: PTY_CONSUMER_SESSION_PROTOCOL_VERSION,
|
||||
serverBuildId: 'build',
|
||||
role: 'session-owner',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 3,
|
||||
ownerLease: 'owner',
|
||||
resumed: true,
|
||||
capabilities: { outputFlowControl: { version: 1, windowSu: 256 } }
|
||||
}
|
||||
const request = vi.fn().mockResolvedValue(grant)
|
||||
const isDisposed = vi.fn(() => false)
|
||||
const mux: Pick<SshChannelMultiplexer, 'request' | 'isDisposed'> = { request, isDisposed }
|
||||
const options = {
|
||||
clientInstanceId: 'client',
|
||||
expectedServerBuildId: 'build',
|
||||
resume: { ownerGeneration: 2, ownerLease: 'owner' },
|
||||
outputFlowControl: { requestedWindowSu: 512 },
|
||||
signal: controller.signal,
|
||||
assertAuthority: vi.fn()
|
||||
}
|
||||
const run = () => resumeSshPtyConsumerSession(mux, options)
|
||||
return { grant, request, isDisposed, mux, options, run, controller }
|
||||
}
|
||||
|
||||
it('resumes through only the dedicated method and passes cancellation to transport', async () => {
|
||||
const f = fixture()
|
||||
await expect(f.run()).resolves.toEqual({
|
||||
resumed: true,
|
||||
state: {
|
||||
mode: 'negotiated',
|
||||
clientInstanceId: 'client',
|
||||
clientGeneration: 3,
|
||||
ownerGeneration: 3,
|
||||
ownerLease: 'owner',
|
||||
outputFlowControl: { version: 1, windowSu: 256 }
|
||||
}
|
||||
})
|
||||
expect(f.request).toHaveBeenCalledExactlyOnceWith(
|
||||
'pty.resumeClient',
|
||||
{
|
||||
protocolVersion: PTY_CONSUMER_SESSION_PROTOCOL_VERSION,
|
||||
clientInstanceId: 'client',
|
||||
requestedRole: 'session-owner',
|
||||
resume: { ownerGeneration: 2, ownerLease: 'owner' },
|
||||
capabilities: { outputFlowControl: { versions: [1], requestedWindowSu: 512 } }
|
||||
},
|
||||
{ timeoutMs: 10_000, signal: f.controller.signal }
|
||||
)
|
||||
expect(f.options.assertAuthority).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it.each([-32601, -32000])(
|
||||
'never falls back on host error %s, even with legacy runtime option',
|
||||
async (code) => {
|
||||
const f = fixture()
|
||||
const error = Object.assign(new Error('host-refusal'), { code })
|
||||
f.request.mockRejectedValue(error)
|
||||
const options = { ...f.options, allowSameBuildLegacyFallback: true }
|
||||
const refusal = resumeSshPtyConsumerSession(f.mux, options)
|
||||
// An old relay's method-not-found becomes a typed refusal; other host errors pass through.
|
||||
const expected =
|
||||
code === -32601 ? { message: 'pty_consumer_resume_unsupported', cause: error } : error
|
||||
await expect(refusal).rejects.toMatchObject(expected)
|
||||
expect(f.request).toHaveBeenCalledOnce()
|
||||
expect(f.request.mock.calls[0][0]).toBe('pty.resumeClient')
|
||||
}
|
||||
)
|
||||
|
||||
it.each([
|
||||
{ ownerGeneration: 0 },
|
||||
{ ownerGeneration: -1 },
|
||||
{ ownerGeneration: 1.5 },
|
||||
{ ownerGeneration: Number.MAX_SAFE_INTEGER + 1 },
|
||||
{ ownerLease: '' },
|
||||
{ ownerLease: 'x'.repeat(513) }
|
||||
])('refuses invalid resume proof before RPC %j', async (proof) => {
|
||||
const f = fixture()
|
||||
Object.assign(f.options.resume, proof)
|
||||
await expect(f.run()).rejects.toThrow('resume_required')
|
||||
expect(f.request).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('refuses absent resume proof before RPC', async () => {
|
||||
const f = fixture()
|
||||
Reflect.deleteProperty(f.options, 'resume')
|
||||
await expect(f.run()).rejects.toThrow('resume_required')
|
||||
expect(f.request).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it.each(['clientInstanceId', 'expectedServerBuildId'] as const)(
|
||||
'refuses missing %s before RPC',
|
||||
async (field) => {
|
||||
const f = fixture()
|
||||
f.options[field] = ''
|
||||
await expect(f.run()).rejects.toThrow('resume_required')
|
||||
expect(f.request).not.toHaveBeenCalled()
|
||||
}
|
||||
)
|
||||
|
||||
it.each([
|
||||
{ ownerLease: 'different' },
|
||||
{ ownerGeneration: 2 },
|
||||
{ ownerGeneration: 1 },
|
||||
{ resumed: false },
|
||||
{ resumed: undefined },
|
||||
{ ownerGeneration: 0 },
|
||||
{ clientGeneration: 0 },
|
||||
{ serverBuildId: 'other' },
|
||||
{ role: 'observer' },
|
||||
{ capabilities: undefined }
|
||||
])('refuses invalid or non-resumed grant %j', async (change) => {
|
||||
const f = fixture()
|
||||
Object.assign(f.grant, change)
|
||||
await expect(f.run()).rejects.toThrow()
|
||||
expect(f.request).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it.each(['before', 'after'] as const)('refuses lost native authority %s RPC', async (when) => {
|
||||
const f = fixture()
|
||||
const revoke = () =>
|
||||
f.options.assertAuthority.mockImplementation(() => {
|
||||
throw new Error('native-revoked')
|
||||
})
|
||||
if (when === 'before') {
|
||||
revoke()
|
||||
} else {
|
||||
f.request.mockImplementation(async () => {
|
||||
await Promise.resolve()
|
||||
revoke()
|
||||
return f.grant
|
||||
})
|
||||
}
|
||||
await expect(f.run()).rejects.toThrow('native-revoked')
|
||||
expect(f.request).toHaveBeenCalledTimes(when === 'before' ? 0 : 1)
|
||||
})
|
||||
|
||||
it.each(['before', 'after'] as const)('refuses abort %s RPC', async (when) => {
|
||||
const f = fixture()
|
||||
if (when === 'before') {
|
||||
f.controller.abort(new Error('aborted'))
|
||||
} else {
|
||||
f.request.mockImplementation(async () => {
|
||||
await Promise.resolve()
|
||||
f.controller.abort(new Error('aborted'))
|
||||
return f.grant
|
||||
})
|
||||
}
|
||||
await expect(f.run()).rejects.toThrow('aborted')
|
||||
expect(f.request).toHaveBeenCalledTimes(when === 'before' ? 0 : 1)
|
||||
})
|
||||
|
||||
it.each(['before', 'after'] as const)('refuses disposed mux %s RPC', async (when) => {
|
||||
const f = fixture()
|
||||
if (when === 'before') {
|
||||
f.isDisposed.mockReturnValue(true)
|
||||
} else {
|
||||
f.request.mockImplementation(async () => {
|
||||
await Promise.resolve()
|
||||
f.isDisposed.mockReturnValue(true)
|
||||
return f.grant
|
||||
})
|
||||
}
|
||||
await expect(f.run()).rejects.toThrow('transport_closed')
|
||||
expect(f.request).toHaveBeenCalledTimes(when === 'before' ? 0 : 1)
|
||||
})
|
||||
|
||||
it('pins the original proof across await rather than accepting caller mutation', async () => {
|
||||
const f = fixture()
|
||||
f.request.mockImplementation(async () => {
|
||||
await Promise.resolve()
|
||||
f.options.resume.ownerLease = 'replacement'
|
||||
f.options.resume.ownerGeneration = 100
|
||||
return f.grant
|
||||
})
|
||||
await expect(f.run()).resolves.toMatchObject({
|
||||
resumed: true,
|
||||
state: { ownerLease: 'owner', ownerGeneration: 3 }
|
||||
})
|
||||
expect(f.request.mock.calls[0][1].resume).toEqual({ ownerLease: 'owner', ownerGeneration: 2 })
|
||||
})
|
||||
|
||||
it('leaves an old relay channel usable for openClient after it refuses pty.resumeClient', async () => {
|
||||
const f = fixture()
|
||||
f.request.mockRejectedValueOnce(Object.assign(new Error('Method not found'), { code: -32601 }))
|
||||
await expect(f.run()).rejects.toThrow('pty_consumer_resume_unsupported')
|
||||
expect(f.isDisposed()).toBe(false)
|
||||
const { capabilities: _flow, ...plainGrant } = f.grant
|
||||
f.request.mockResolvedValueOnce({
|
||||
...plainGrant,
|
||||
clientGeneration: 1,
|
||||
ownerGeneration: 1,
|
||||
resumed: false
|
||||
})
|
||||
|
||||
const admission = await openSshPtyConsumerSession(f.mux, {
|
||||
clientInstanceId: 'client',
|
||||
expectedServerBuildId: 'build'
|
||||
})
|
||||
|
||||
expect(admission).toMatchObject({ resumed: false, state: { mode: 'negotiated' } })
|
||||
expect(f.request.mock.calls.map(([method]) => method)).toEqual([
|
||||
'pty.resumeClient',
|
||||
'pty.openClient'
|
||||
])
|
||||
})
|
||||
@@ -1,5 +1,6 @@
|
||||
import {
|
||||
PTY_CONSUMER_SESSION_PROTOCOL_VERSION,
|
||||
PTY_CONSUMER_RESUME_CLIENT_METHOD,
|
||||
type PtyConsumerSessionGrant
|
||||
} from '../../shared/pty-consumer-session'
|
||||
import type { SshChannelMultiplexer } from './ssh-channel-multiplexer'
|
||||
@@ -98,13 +99,90 @@ function validateGrant(
|
||||
}
|
||||
|
||||
export async function openSshPtyConsumerSession(
|
||||
mux: SshChannelMultiplexer,
|
||||
mux: Pick<SshChannelMultiplexer, 'request'>,
|
||||
options: OpenSshPtyConsumerSessionOptions
|
||||
): Promise<SshPtyConsumerAdmission> {
|
||||
return requestPtyConsumerSession(mux, options, SSH_PTY_OPEN_CLIENT_METHOD)
|
||||
}
|
||||
|
||||
/** Dedicated RPC: old hosts refuse without minting a replacement claim. */
|
||||
export async function resumeSshPtyConsumerSession(
|
||||
mux: Pick<SshChannelMultiplexer, 'request' | 'isDisposed'>,
|
||||
options: Omit<OpenSshPtyConsumerSessionOptions, 'resume' | 'allowSameBuildLegacyFallback'> & {
|
||||
resume: NonNullable<OpenSshPtyConsumerSessionOptions['resume']>
|
||||
signal: AbortSignal
|
||||
assertAuthority: () => void
|
||||
}
|
||||
): Promise<SshPtyConsumerAdmission> {
|
||||
const resume = { ...options.resume }
|
||||
const signal = options.signal
|
||||
const assertAuthority = options.assertAuthority
|
||||
const isDisposed = mux.isDisposed.bind(mux)
|
||||
const assertCurrent = () => {
|
||||
signal.throwIfAborted()
|
||||
assertAuthority()
|
||||
if (isDisposed()) {
|
||||
throw new Error('pty_consumer_resume_transport_closed')
|
||||
}
|
||||
}
|
||||
assertCurrent()
|
||||
if (
|
||||
!Number.isSafeInteger(resume.ownerGeneration) ||
|
||||
resume.ownerGeneration <= 0 ||
|
||||
typeof resume.ownerLease !== 'string' ||
|
||||
!resume.ownerLease ||
|
||||
resume.ownerLease.length > 512 ||
|
||||
!options.clientInstanceId ||
|
||||
!options.expectedServerBuildId
|
||||
) {
|
||||
throw new Error('pty_consumer_resume_required')
|
||||
}
|
||||
let admission: SshPtyConsumerAdmission
|
||||
try {
|
||||
admission = await requestPtyConsumerSession(
|
||||
mux,
|
||||
{
|
||||
clientInstanceId: options.clientInstanceId,
|
||||
expectedServerBuildId: options.expectedServerBuildId,
|
||||
resume,
|
||||
...(options.outputFlowControl
|
||||
? { outputFlowControl: { ...options.outputFlowControl } }
|
||||
: {}),
|
||||
allowSameBuildLegacyFallback: false
|
||||
},
|
||||
PTY_CONSUMER_RESUME_CLIENT_METHOD,
|
||||
signal
|
||||
)
|
||||
} catch (error) {
|
||||
// Why: a relay that predates pty.resumeClient answers method-not-found; that is a refusal of
|
||||
// resume on a still-open channel, not a transport failure, so callers can fall back to openClient.
|
||||
if (typeof error === 'object' && error !== null && 'code' in error && error.code === -32601) {
|
||||
throw new Error('pty_consumer_resume_unsupported', { cause: error })
|
||||
}
|
||||
throw error
|
||||
}
|
||||
assertCurrent()
|
||||
if (
|
||||
!admission.resumed ||
|
||||
admission.state.mode !== 'negotiated' ||
|
||||
admission.state.ownerLease !== resume.ownerLease ||
|
||||
admission.state.ownerGeneration <= resume.ownerGeneration
|
||||
) {
|
||||
throw new Error('pty_consumer_resume_grant_mismatch')
|
||||
}
|
||||
return admission
|
||||
}
|
||||
|
||||
async function requestPtyConsumerSession(
|
||||
mux: Pick<SshChannelMultiplexer, 'request'>,
|
||||
options: OpenSshPtyConsumerSessionOptions,
|
||||
method: string,
|
||||
signal?: AbortSignal
|
||||
): Promise<SshPtyConsumerAdmission> {
|
||||
let result: unknown
|
||||
try {
|
||||
result = await mux.request(
|
||||
SSH_PTY_OPEN_CLIENT_METHOD,
|
||||
method,
|
||||
{
|
||||
protocolVersion: PTY_CONSUMER_SESSION_PROTOCOL_VERSION,
|
||||
clientInstanceId: options.clientInstanceId,
|
||||
@@ -121,7 +199,7 @@ export async function openSshPtyConsumerSession(
|
||||
}
|
||||
: {})
|
||||
},
|
||||
{ timeoutMs: SSH_PTY_OPEN_CLIENT_TIMEOUT_MS }
|
||||
{ timeoutMs: SSH_PTY_OPEN_CLIENT_TIMEOUT_MS, ...(signal ? { signal } : {}) }
|
||||
)
|
||||
} catch (error) {
|
||||
const code = (error as { code?: unknown })?.code
|
||||
|
||||
@@ -0,0 +1,131 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import { RelayDispatcher, type RelayClientSessionIdentity } from './dispatcher'
|
||||
import { encodeJsonRpcFrame } from './protocol'
|
||||
import { SshPtyConsumerSessionAdapter } from './ssh-pty-consumer-session-adapter'
|
||||
|
||||
const endpointIdentity: RelayClientSessionIdentity = {
|
||||
principal: 'endpoint-principal',
|
||||
authenticated: true,
|
||||
allowSessionOwner: true,
|
||||
authenticationKind: 'endpoint-credential'
|
||||
}
|
||||
|
||||
function frame(id: number, method: string, overrides: Record<string, unknown> = {}): Buffer {
|
||||
return encodeJsonRpcFrame(
|
||||
{
|
||||
jsonrpc: '2.0',
|
||||
id,
|
||||
method,
|
||||
params: {
|
||||
protocolVersion: 1,
|
||||
clientInstanceId: 'stable-client',
|
||||
requestedRole: 'session-owner',
|
||||
...overrides
|
||||
}
|
||||
},
|
||||
id,
|
||||
0
|
||||
)
|
||||
}
|
||||
|
||||
function response(buffer: Buffer) {
|
||||
return JSON.parse(buffer.subarray(13, 13 + buffer.readUInt32BE(9)).toString('utf8'))
|
||||
}
|
||||
|
||||
async function flushRequests(): Promise<void> {
|
||||
await new Promise((resolve) => setImmediate(resolve))
|
||||
}
|
||||
|
||||
describe('pty.resumeClient through RelayDispatcher', () => {
|
||||
let dispatcher: RelayDispatcher | undefined
|
||||
afterEach(() => dispatcher?.dispose())
|
||||
|
||||
function fixture() {
|
||||
const writes: Buffer[] = []
|
||||
dispatcher = new RelayDispatcher(
|
||||
(data, settled) => {
|
||||
writes.push(Buffer.from(data))
|
||||
settled({ ok: true })
|
||||
return true
|
||||
},
|
||||
{ supportsWriteCallback: true },
|
||||
endpointIdentity
|
||||
)
|
||||
const adapter = new SshPtyConsumerSessionAdapter(dispatcher, 'build-a')
|
||||
return { dispatcher, adapter, writes }
|
||||
}
|
||||
|
||||
it('refuses absent historical ownership without minting a fresh owner or reserving the connection', async () => {
|
||||
const { dispatcher, adapter, writes } = fixture()
|
||||
const resume = { ownerGeneration: 1, ownerLease: 'forgotten-lease' }
|
||||
dispatcher.feed(frame(1, 'pty.resumeClient', { resume }))
|
||||
await flushRequests()
|
||||
expect(response(writes[0]).error.message).toContain('pty_consumer_resume_owner_missing')
|
||||
expect(adapter.activeSessionOwner(1)).toBeNull()
|
||||
expect(() => adapter.assertOwnerPublicationSettled()).not.toThrow()
|
||||
dispatcher.feed(frame(2, 'pty.openClient', { resume }))
|
||||
await flushRequests()
|
||||
expect(response(writes[1]).result).toMatchObject({
|
||||
role: 'session-owner',
|
||||
clientGeneration: 1,
|
||||
ownerGeneration: 1,
|
||||
resumed: false
|
||||
})
|
||||
expect(adapter.activeSessionOwner(1)).not.toBeNull()
|
||||
})
|
||||
|
||||
it.each([true, false])(
|
||||
'settles exact-owner resume authority only on response publication (success=%s)',
|
||||
async (ok) => {
|
||||
const { dispatcher, adapter, writes } = fixture()
|
||||
dispatcher.feed(frame(1, 'pty.openClient'))
|
||||
await flushRequests()
|
||||
const grant = response(writes[0]).result
|
||||
const incumbent = adapter.activeSessionOwner(1)
|
||||
expect(incumbent).not.toBeNull()
|
||||
const resumedWrites: Buffer[] = []
|
||||
const settlements: ((result: { ok: true } | { ok: false; error: Error }) => void)[] = []
|
||||
const successor = dispatcher.attachClient(
|
||||
(data, settled) => {
|
||||
resumedWrites.push(Buffer.from(data))
|
||||
settlements.push(settled)
|
||||
return true
|
||||
},
|
||||
{ supportsWriteCallback: true },
|
||||
endpointIdentity
|
||||
)
|
||||
const release = vi.spyOn(dispatcher, 'releaseDisplacedClient')
|
||||
dispatcher.feedClient(
|
||||
successor,
|
||||
frame(2, 'pty.resumeClient', {
|
||||
resume: { ownerGeneration: grant.ownerGeneration, ownerLease: grant.ownerLease }
|
||||
})
|
||||
)
|
||||
await flushRequests()
|
||||
expect(response(resumedWrites[0]).result).toMatchObject({
|
||||
role: 'session-owner',
|
||||
resumed: true,
|
||||
ownerGeneration: 2,
|
||||
ownerLease: grant.ownerLease
|
||||
})
|
||||
expect(adapter.activeSessionOwner(1)).toEqual(incumbent)
|
||||
expect(adapter.activeSessionOwner(successor)).toBeNull()
|
||||
expect(release).not.toHaveBeenCalled()
|
||||
expect(() => adapter.assertOwnerPublicationSettled()).toThrow(
|
||||
'pty_consumer_owner_publication_pending'
|
||||
)
|
||||
settlements[0](ok ? { ok: true } : { ok: false, error: new Error('response failed') })
|
||||
expect(() => adapter.assertOwnerPublicationSettled()).not.toThrow()
|
||||
if (ok) {
|
||||
expect(adapter.activeSessionOwner(1)).toBeNull()
|
||||
expect(adapter.activeSessionOwner(successor)).not.toBeNull()
|
||||
expect(release).toHaveBeenCalledTimes(1)
|
||||
} else {
|
||||
expect(adapter.activeSessionOwner(1)).toEqual(incumbent)
|
||||
expect(adapter.activeSessionOwner(successor)).toBeNull()
|
||||
expect(release).not.toHaveBeenCalled()
|
||||
}
|
||||
release.mockRestore()
|
||||
}
|
||||
)
|
||||
})
|
||||
@@ -61,6 +61,60 @@ describe('SshPtyConsumerSessionAdapter', () => {
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('rolls back replacement ownership when response callback registration throws', async () => {
|
||||
const writes: Buffer[] = []
|
||||
dispatcher = new RelayDispatcher(
|
||||
(data, onSettled) => {
|
||||
writes.push(Buffer.from(data))
|
||||
onSettled({ ok: true })
|
||||
return true
|
||||
},
|
||||
{ supportsWriteCallback: true },
|
||||
endpointIdentity
|
||||
)
|
||||
const registration = vi.spyOn(dispatcher, 'onRequest')
|
||||
const adapter = new SshPtyConsumerSessionAdapter(dispatcher, 'build-a')
|
||||
const openClient = registration.mock.calls.find(([method]) => method === 'pty.openClient')![1]
|
||||
registration.mockRestore()
|
||||
dispatcher.feed(openFrame(1))
|
||||
await flushRequests()
|
||||
const grant = responseResult(writes[0])
|
||||
const incumbent = adapter.activeSessionOwner(1)
|
||||
expect(incumbent).not.toBeNull()
|
||||
const replacement = dispatcher.attachClient(() => true, {}, endpointIdentity)
|
||||
const closeIncumbent = vi.spyOn(dispatcher, 'releaseDisplacedClient')
|
||||
const params = {
|
||||
protocolVersion: 1,
|
||||
clientInstanceId: 'client-1',
|
||||
requestedRole: 'session-owner',
|
||||
resume: { ownerGeneration: grant.ownerGeneration, ownerLease: grant.ownerLease }
|
||||
}
|
||||
await expect(
|
||||
openClient(params, {
|
||||
clientId: replacement,
|
||||
isStale: () => false,
|
||||
sessionIdentity: endpointIdentity,
|
||||
onResponseSettled: () => {
|
||||
throw new Error('registration failed')
|
||||
}
|
||||
})
|
||||
).rejects.toThrow('registration failed')
|
||||
expect(() => adapter.assertOwnerPublicationSettled()).not.toThrow()
|
||||
expect(adapter.activeSessionOwner(1)).toEqual(incumbent)
|
||||
expect(adapter.activeSessionOwner(replacement)).toBeNull()
|
||||
expect(closeIncumbent).not.toHaveBeenCalled()
|
||||
await openClient(params, {
|
||||
clientId: replacement,
|
||||
isStale: () => false,
|
||||
sessionIdentity: endpointIdentity,
|
||||
onResponseSettled: (settle) => settle({ ok: false, error: new Error('cancel retry') })
|
||||
})
|
||||
expect(() => adapter.assertOwnerPublicationSettled()).not.toThrow()
|
||||
expect(adapter.activeSessionOwner(1)).toEqual(incumbent)
|
||||
expect(closeIncumbent).not.toHaveBeenCalled()
|
||||
closeIncumbent.mockRestore()
|
||||
})
|
||||
|
||||
it('does not activate owner authority until the grant write settles', async () => {
|
||||
const firstWrites: Buffer[] = []
|
||||
const firstSettlements: ((result: { ok: true } | { ok: false; error: Error }) => void)[] = []
|
||||
@@ -73,10 +127,13 @@ describe('SshPtyConsumerSessionAdapter', () => {
|
||||
{ supportsWriteCallback: true },
|
||||
endpointIdentity
|
||||
)
|
||||
new SshPtyConsumerSessionAdapter(dispatcher, 'build-a')
|
||||
const adapter = new SshPtyConsumerSessionAdapter(dispatcher, 'build-a')
|
||||
|
||||
dispatcher.feed(openFrame(1))
|
||||
await flushRequests()
|
||||
expect(() => adapter.assertOwnerPublicationSettled()).toThrow(
|
||||
'pty_consumer_owner_publication_pending'
|
||||
)
|
||||
|
||||
const secondWrites: Buffer[] = []
|
||||
const secondId = dispatcher.attachClient(
|
||||
@@ -102,6 +159,7 @@ describe('SshPtyConsumerSessionAdapter', () => {
|
||||
code: PTY_CONSUMER_OWNER_RECOVERY_PENDING_ERROR
|
||||
})
|
||||
firstSettlements[0]({ ok: true })
|
||||
expect(() => adapter.assertOwnerPublicationSettled()).not.toThrow()
|
||||
})
|
||||
|
||||
it('rolls back owner election when the grant write fails', async () => {
|
||||
@@ -113,9 +171,10 @@ describe('SshPtyConsumerSessionAdapter', () => {
|
||||
{ supportsWriteCallback: true },
|
||||
endpointIdentity
|
||||
)
|
||||
new SshPtyConsumerSessionAdapter(dispatcher, 'build-a')
|
||||
const adapter = new SshPtyConsumerSessionAdapter(dispatcher, 'build-a')
|
||||
dispatcher.feed(openFrame(1))
|
||||
await flushRequests()
|
||||
expect(() => adapter.assertOwnerPublicationSettled()).not.toThrow()
|
||||
|
||||
const retryWrites: Buffer[] = []
|
||||
const retryId = dispatcher.attachClient(
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import {
|
||||
PTY_CONSUMER_SESSION_PROTOCOL_VERSION,
|
||||
PTY_CONSUMER_RESUME_CLIENT_METHOD,
|
||||
PtyConsumerSession,
|
||||
type PtyConsumerSessionAdmission,
|
||||
type PtyConsumerSessionGrant
|
||||
@@ -27,6 +28,7 @@ export class SshPtyConsumerSessionAdapter {
|
||||
private readonly session: PtyConsumerSession
|
||||
private readonly sourceCredit: SshPtySourceCreditAdapter
|
||||
private readonly pausedDeliveryByPty = new Map<string, PtySourceDeliveryIdentity>()
|
||||
private readonly pendingPublications = new Set<PtyConsumerSessionAdmission>()
|
||||
|
||||
constructor(
|
||||
private readonly dispatcher: RelayDispatcher,
|
||||
@@ -59,6 +61,9 @@ export class SshPtyConsumerSessionAdapter {
|
||||
dispatcher.onRequest(SSH_PTY_OPEN_CLIENT_METHOD, (params, context) =>
|
||||
this.openClient(params, context)
|
||||
)
|
||||
dispatcher.onRequest(PTY_CONSUMER_RESUME_CLIENT_METHOD, (params, context) =>
|
||||
this.openClient(params, context, true)
|
||||
)
|
||||
dispatcher.onClientDetached((clientId, cause) => {
|
||||
const connectionKey = String(clientId)
|
||||
const grant = this.session.activeGrant(connectionKey)
|
||||
@@ -214,9 +219,35 @@ export class SshPtyConsumerSessionAdapter {
|
||||
return sshPtyDeliveryMode(this.session.activeGrant(String(clientId)))
|
||||
}
|
||||
|
||||
activeSessionOwner(
|
||||
clientId: number
|
||||
): Readonly<{ ownerGeneration: number; ownerLease: string }> | null {
|
||||
const grant = this.session.activeGrant(String(clientId))
|
||||
const ownerGeneration = grant?.ownerGeneration
|
||||
if (
|
||||
grant?.role !== 'session-owner' ||
|
||||
typeof ownerGeneration !== 'number' ||
|
||||
!Number.isSafeInteger(ownerGeneration) ||
|
||||
!grant.ownerLease
|
||||
) {
|
||||
return null
|
||||
}
|
||||
return Object.freeze({
|
||||
ownerGeneration,
|
||||
ownerLease: grant.ownerLease
|
||||
})
|
||||
}
|
||||
|
||||
assertOwnerPublicationSettled(): void {
|
||||
if (this.pendingPublications.size > 0) {
|
||||
throw new Error('pty_consumer_owner_publication_pending')
|
||||
}
|
||||
}
|
||||
|
||||
private async openClient(
|
||||
rawParams: Record<string, unknown>,
|
||||
context: RequestContext
|
||||
context: RequestContext,
|
||||
resumeOnly = false
|
||||
): Promise<PtyConsumerSessionGrant> {
|
||||
const params = parseOpenClientParams(rawParams)
|
||||
if (params.protocolVersion !== PTY_CONSUMER_SESSION_PROTOCOL_VERSION) {
|
||||
@@ -225,24 +256,38 @@ export class SshPtyConsumerSessionAdapter {
|
||||
)
|
||||
}
|
||||
const identity = requireIdentity(context)
|
||||
const admission = this.session.admit(params, {
|
||||
const authenticate = {
|
||||
connectionId: String(context.clientId),
|
||||
principal: identity.principal,
|
||||
authenticated: identity.authenticated,
|
||||
allowSessionOwner: identity.allowSessionOwner
|
||||
})
|
||||
}
|
||||
const admission = resumeOnly
|
||||
? this.session.admitResumed(params, authenticate)
|
||||
: this.session.admit(params, authenticate)
|
||||
if (!context.onResponseSettled) {
|
||||
admission.rollbackPublication()
|
||||
throw new Error('SSH PTY consumer response publication fence is unavailable')
|
||||
}
|
||||
context.onResponseSettled((result) => {
|
||||
if (!result.ok) {
|
||||
admission.rollbackPublication()
|
||||
return
|
||||
}
|
||||
admission.commitPublication()
|
||||
this.closeDisplacedOwner(admission.displacedOwner)
|
||||
})
|
||||
this.pendingPublications.add(admission)
|
||||
try {
|
||||
context.onResponseSettled((result) => {
|
||||
try {
|
||||
if (!result.ok) {
|
||||
admission.rollbackPublication()
|
||||
return
|
||||
}
|
||||
admission.commitPublication()
|
||||
this.closeDisplacedOwner(admission.displacedOwner)
|
||||
} finally {
|
||||
this.pendingPublications.delete(admission)
|
||||
}
|
||||
})
|
||||
} catch (error) {
|
||||
this.pendingPublications.delete(admission)
|
||||
admission.rollbackPublication()
|
||||
throw error
|
||||
}
|
||||
return admission.grant
|
||||
}
|
||||
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
export const PTY_CONSUMER_SESSION_PROTOCOL_VERSION = 1
|
||||
export const PTY_CONSUMER_RESUME_CLIENT_METHOD = 'pty.resumeClient'
|
||||
export const PTY_CONSUMER_OWNER_GRACE_MS = 30_000
|
||||
export const PTY_CONSUMER_STALE_OWNER_RECOVERY_ERROR = -32041
|
||||
// Why: recovery is blocked only while the incumbent owner's grant publication is still settling — a
|
||||
|
||||
@@ -0,0 +1,239 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import {
|
||||
PTY_CONSUMER_OWNER_RECOVERY_PENDING_ERROR,
|
||||
PtyConsumerSession,
|
||||
type PtyConsumerAuthentication,
|
||||
type PtyConsumerSessionHello
|
||||
} from './pty-consumer-session'
|
||||
|
||||
function authentication(overrides: Partial<PtyConsumerAuthentication> = {}) {
|
||||
return {
|
||||
connectionId: 'successor',
|
||||
principal: 'desktop',
|
||||
authenticated: true,
|
||||
allowSessionOwner: true,
|
||||
...overrides
|
||||
}
|
||||
}
|
||||
|
||||
function hello(overrides: Partial<PtyConsumerSessionHello> = {}): PtyConsumerSessionHello {
|
||||
return {
|
||||
clientInstanceId: 'client-a',
|
||||
requestedRole: 'session-owner',
|
||||
resume: { ownerGeneration: 1, ownerLease: 'lease-1' },
|
||||
...overrides
|
||||
}
|
||||
}
|
||||
|
||||
function fixture() {
|
||||
let now = 0
|
||||
let lease = 0
|
||||
const createLease = vi.fn(() => `lease-${++lease}`)
|
||||
const session = new PtyConsumerSession({
|
||||
serverBuildId: 'relay-build',
|
||||
createLease,
|
||||
ownerGraceMs: 30_000,
|
||||
now: () => now
|
||||
})
|
||||
return { session, createLease, advance: () => (now += 30_001) }
|
||||
}
|
||||
|
||||
function activate(session: PtyConsumerSession) {
|
||||
const admission = session.admit(
|
||||
hello({ resume: undefined }),
|
||||
authentication({ connectionId: 'source' })
|
||||
)
|
||||
admission.commitPublication()
|
||||
return admission
|
||||
}
|
||||
|
||||
describe('PtyConsumerSession resume-only admission', () => {
|
||||
it('does not shorten disconnected-owner grace after a refused strict claim', () => {
|
||||
let now = 0
|
||||
const session = new PtyConsumerSession({
|
||||
serverBuildId: 'relay-build',
|
||||
ownerGraceMs: 30_000,
|
||||
createLease: () => 'lease-1',
|
||||
now: () => now
|
||||
})
|
||||
activate(session)
|
||||
session.close('source', 'peer-closed')
|
||||
expect(() =>
|
||||
session.admitResumed(
|
||||
hello({ resume: { ownerGeneration: 1, ownerLease: 'wrong' } }),
|
||||
authentication()
|
||||
)
|
||||
).toThrow()
|
||||
now = 1_000
|
||||
expect(session.admitResumed(hello(), authentication()).grant.resumed).toBe(true)
|
||||
})
|
||||
it('refuses a missing owner without minting a lease or consuming generations or the connection', () => {
|
||||
const { session, createLease } = fixture()
|
||||
expect(() => session.admitResumed(hello(), authentication())).toThrow(
|
||||
'pty_consumer_resume_owner_missing'
|
||||
)
|
||||
expect(createLease).not.toHaveBeenCalled()
|
||||
expect(session.activeGrant('successor')).toBeNull()
|
||||
const ordinary = session.admit(hello(), authentication())
|
||||
expect(ordinary.grant).toMatchObject({
|
||||
resumed: false,
|
||||
clientGeneration: 1,
|
||||
ownerGeneration: 1,
|
||||
ownerLease: 'lease-1'
|
||||
})
|
||||
expect(createLease).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('publishes a valid resumed owner atomically without minting a new lease', () => {
|
||||
const { session, createLease } = fixture()
|
||||
const source = activate(session)
|
||||
const successor = session.admitResumed(hello(), authentication())
|
||||
expect(successor.grant).toMatchObject({
|
||||
resumed: true,
|
||||
ownerGeneration: 2,
|
||||
clientGeneration: 2,
|
||||
ownerLease: 'lease-1'
|
||||
})
|
||||
expect(successor.displacedOwner).toEqual({ connectionId: 'source', grant: source.grant })
|
||||
expect(session.activeGrant('source')).toBe(source.grant)
|
||||
expect(session.activeGrant('successor')).toBeNull()
|
||||
successor.commitPublication()
|
||||
expect(session.activeGrant('source')).toBeNull()
|
||||
expect(session.activeGrant('successor')).toBe(successor.grant)
|
||||
expect(createLease).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('restores the incumbent on rollback and permits retry on the same connection', () => {
|
||||
const { session } = fixture()
|
||||
const source = activate(session)
|
||||
const successor = session.admitResumed(hello(), authentication())
|
||||
successor.rollbackPublication()
|
||||
successor.commitPublication()
|
||||
expect(session.activeGrant('source')).toBe(source.grant)
|
||||
expect(session.activeGrant('successor')).toBeNull()
|
||||
const retry = session.admitResumed(hello(), authentication())
|
||||
expect(retry.grant).toMatchObject({ resumed: true, ownerGeneration: 3, ownerLease: 'lease-1' })
|
||||
retry.commitPublication()
|
||||
expect(session.activeGrant('successor')).toBe(retry.grant)
|
||||
})
|
||||
|
||||
it.each([
|
||||
['lease', hello({ resume: { ownerGeneration: 1, ownerLease: 'wrong' } }), authentication()],
|
||||
['client', hello({ clientInstanceId: 'other' }), authentication()],
|
||||
['principal', hello(), authentication({ principal: 'other' })],
|
||||
[
|
||||
'generation',
|
||||
hello({ resume: { ownerGeneration: 99, ownerLease: 'lease-1' } }),
|
||||
authentication()
|
||||
]
|
||||
])(
|
||||
'refuses a wrong %s without disturbing the owner or consuming the connection',
|
||||
(_name, request, auth) => {
|
||||
const { session, createLease } = fixture()
|
||||
const source = activate(session)
|
||||
expect(() => session.admitResumed(request, auth)).toThrow()
|
||||
expect(session.activeGrant('source')).toBe(source.grant)
|
||||
const valid = session.admitResumed(hello(), authentication())
|
||||
expect(valid.grant).toMatchObject({ resumed: true, clientGeneration: 2, ownerGeneration: 2 })
|
||||
expect(createLease).toHaveBeenCalledTimes(1)
|
||||
}
|
||||
)
|
||||
|
||||
it('refuses expired ownership while leaving ordinary fresh-claim fallback intact', () => {
|
||||
const { session, createLease, advance } = fixture()
|
||||
activate(session)
|
||||
session.close('source')
|
||||
advance()
|
||||
expect(() => session.admitResumed(hello(), authentication())).toThrow(
|
||||
'pty_consumer_resume_owner_missing'
|
||||
)
|
||||
expect(createLease).toHaveBeenCalledTimes(1)
|
||||
expect(session.admit(hello(), authentication()).grant).toMatchObject({
|
||||
resumed: false,
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 2,
|
||||
ownerLease: 'lease-2'
|
||||
})
|
||||
})
|
||||
|
||||
it.each([
|
||||
[
|
||||
'missing proof',
|
||||
hello({ resume: undefined }),
|
||||
authentication(),
|
||||
'pty_consumer_resume_required'
|
||||
],
|
||||
[
|
||||
'subscriber',
|
||||
hello({ requestedRole: 'subscriber' }),
|
||||
authentication(),
|
||||
'pty_consumer_resume_required'
|
||||
],
|
||||
[
|
||||
'owner-ineligible',
|
||||
hello(),
|
||||
authentication({ allowSessionOwner: false }),
|
||||
'pty_consumer_resume_required'
|
||||
],
|
||||
[
|
||||
'unauthenticated',
|
||||
hello(),
|
||||
authentication({ authenticated: false }),
|
||||
'authentication required'
|
||||
]
|
||||
])('refuses %s without reserving the connection', (_name, request, auth, error) => {
|
||||
const { session } = fixture()
|
||||
activate(session)
|
||||
expect(() => session.admitResumed(request, auth)).toThrow(error)
|
||||
expect(session.admitResumed(hello(), authentication()).grant).toMatchObject({
|
||||
resumed: true,
|
||||
clientGeneration: 2,
|
||||
ownerGeneration: 2
|
||||
})
|
||||
})
|
||||
|
||||
it('preserves pending-recovery refusal and permits retry after rollback', () => {
|
||||
const { session } = fixture()
|
||||
const source = activate(session)
|
||||
const pending = session.admitResumed(hello(), authentication({ connectionId: 'pending' }))
|
||||
expect(() => session.admitResumed(hello(), authentication())).toThrow(
|
||||
expect.objectContaining({ code: PTY_CONSUMER_OWNER_RECOVERY_PENDING_ERROR })
|
||||
)
|
||||
expect(session.activeGrant('source')).toBe(source.grant)
|
||||
pending.rollbackPublication()
|
||||
const retry = session.admitResumed(hello(), authentication())
|
||||
expect(retry.grant).toMatchObject({ resumed: true, clientGeneration: 3, ownerGeneration: 3 })
|
||||
retry.commitPublication()
|
||||
expect(session.activeGrant('source')).toBeNull()
|
||||
})
|
||||
|
||||
it('recovers a disconnected owner within grace', () => {
|
||||
const { session } = fixture()
|
||||
activate(session)
|
||||
session.close('source')
|
||||
const successor = session.admitResumed(hello(), authentication())
|
||||
expect(successor.displacedOwner).toBeUndefined()
|
||||
expect(successor.grant).toMatchObject({
|
||||
resumed: true,
|
||||
ownerGeneration: 2,
|
||||
ownerLease: 'lease-1'
|
||||
})
|
||||
successor.commitPublication()
|
||||
expect(session.activeGrant('successor')).toBe(successor.grant)
|
||||
})
|
||||
|
||||
it('recovers the last durable generation only after the unpersisted successor disconnects', () => {
|
||||
const { session, createLease } = fixture()
|
||||
activate(session)
|
||||
const first = session.admitResumed(hello(), authentication())
|
||||
first.commitPublication()
|
||||
const retryAuthentication = authentication({ connectionId: 'retry' })
|
||||
expect(() => session.admitResumed(hello(), retryAuthentication)).toThrow('superseded')
|
||||
session.close('successor', 'peer-closed')
|
||||
const retry = session.admitResumed(hello(), retryAuthentication)
|
||||
expect(retry.grant).toMatchObject({ resumed: true, ownerGeneration: 3, ownerLease: 'lease-1' })
|
||||
retry.commitPublication()
|
||||
expect(createLease).toHaveBeenCalledOnce()
|
||||
expect(session.activeGrant('retry')).toBe(retry.grant)
|
||||
})
|
||||
})
|
||||
@@ -64,6 +64,22 @@ export class PtyConsumerSession {
|
||||
admit(
|
||||
hello: PtyConsumerSessionHello,
|
||||
authentication: PtyConsumerAuthentication
|
||||
): PtyConsumerSessionAdmission {
|
||||
return this.admitInternal(hello, authentication, false)
|
||||
}
|
||||
|
||||
/** Migration recovery must not turn an absent historical owner into a fresh claim. */
|
||||
admitResumed(
|
||||
hello: PtyConsumerSessionHello,
|
||||
authentication: PtyConsumerAuthentication
|
||||
): PtyConsumerSessionAdmission {
|
||||
return this.admitInternal(hello, authentication, true)
|
||||
}
|
||||
|
||||
private admitInternal(
|
||||
hello: PtyConsumerSessionHello,
|
||||
authentication: PtyConsumerAuthentication,
|
||||
requireResume: boolean
|
||||
): PtyConsumerSessionAdmission {
|
||||
validateHello(hello)
|
||||
assertNonEmptyString(authentication.connectionId, 'connectionId')
|
||||
@@ -71,6 +87,14 @@ export class PtyConsumerSession {
|
||||
if (!authentication.authenticated) {
|
||||
throw new Error('PTY consumer authentication required')
|
||||
}
|
||||
if (
|
||||
requireResume &&
|
||||
(!hello.resume ||
|
||||
hello.requestedRole !== 'session-owner' ||
|
||||
!authentication.allowSessionOwner)
|
||||
) {
|
||||
throw new Error('pty_consumer_resume_required')
|
||||
}
|
||||
this.expireOwner()
|
||||
|
||||
// Why even an identical repeat is rejected: the two responses settle their publications
|
||||
@@ -80,6 +104,13 @@ export class PtyConsumerSession {
|
||||
throw new Error('pty.openClient may be used only once per transport connection')
|
||||
}
|
||||
|
||||
if (requireResume) {
|
||||
if (!this.owner) {
|
||||
throw new Error('pty_consumer_resume_owner_missing')
|
||||
}
|
||||
// Refused recovery must not shorten an incumbent's disconnected-owner grace.
|
||||
assertPtyConsumerOwnerRecovery(hello, authentication, this.owner)
|
||||
}
|
||||
const owner = this.selectOwner(hello, authentication)
|
||||
const grant = Object.freeze({
|
||||
protocolVersion: PTY_CONSUMER_SESSION_PROTOCOL_VERSION,
|
||||
|
||||
Reference in New Issue
Block a user