From f691dfae7fedbb82a53ea783481434c8610dfec9 Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Tue, 1 Sep 2026 01:23:31 -0700 Subject: [PATCH] fix(ssh): preserve identity evidence epochs and admission --- .../ssh-pty-notification-rejection.ts | 43 +++++++++++ .../providers/ssh-pty-notification-routing.ts | 47 +----------- .../dispatcher-notification-publication.ts | 4 +- src/relay/pty-handler-attach-replay.test.ts | 73 +++++++++++++++++++ src/relay/pty-handler.ts | 48 ++++++++++-- .../pty-consumer-session-capabilities.ts | 4 +- src/shared/pty-consumer-session.test.ts | 14 ++++ 7 files changed, 178 insertions(+), 55 deletions(-) create mode 100644 src/main/providers/ssh-pty-notification-rejection.ts diff --git a/src/main/providers/ssh-pty-notification-rejection.ts b/src/main/providers/ssh-pty-notification-rejection.ts new file mode 100644 index 00000000000..5bfe20aaa17 --- /dev/null +++ b/src/main/providers/ssh-pty-notification-rejection.ts @@ -0,0 +1,43 @@ +import type { SshPtyRejectedSourceRecovery } from './ssh-pty-source-delivery-state' + +export function rejectedRecoveryPriority(recovery: SshPtyRejectedSourceRecovery): number { + if (recovery === 'reconnect-channel') { + return 3 + } + return recovery === 'fresh-activation' ? 2 : 1 +} + +export function rejectedSourceIdentity(params: { + deliveryToken?: unknown + clientGeneration?: unknown + ownerGeneration?: unknown + ptyIncarnation?: unknown +}): + | Readonly<{ + deliveryToken: string + clientGeneration: number + ownerGeneration: number + ptyIncarnation: string + }> + | undefined { + if ( + typeof params.deliveryToken !== 'string' || + params.deliveryToken.length === 0 || + !positiveSafeInteger(params.clientGeneration) || + !positiveSafeInteger(params.ownerGeneration) || + typeof params.ptyIncarnation !== 'string' || + params.ptyIncarnation.length === 0 + ) { + return undefined + } + return Object.freeze({ + deliveryToken: params.deliveryToken, + clientGeneration: params.clientGeneration, + ownerGeneration: params.ownerGeneration, + ptyIncarnation: params.ptyIncarnation + }) +} + +function positiveSafeInteger(value: unknown): value is number { + return Number.isSafeInteger(value) && (value as number) > 0 +} diff --git a/src/main/providers/ssh-pty-notification-routing.ts b/src/main/providers/ssh-pty-notification-routing.ts index d749ac375ac..9ac12f2ac4f 100644 --- a/src/main/providers/ssh-pty-notification-routing.ts +++ b/src/main/providers/ssh-pty-notification-routing.ts @@ -4,9 +4,9 @@ import type { PtySourceReceivingActivation } from '../../shared/pty-source-recei import type { SshPtyDataCallback, SshPtyExitCallback, - SshPtyReplayCallback + SshPtyReplayCallback, + SshPtyIdentityEvidenceCallback } from './ssh-pty-provider-contract' -import type { SshPtyIdentityEvidenceCallback } from './ssh-pty-provider-contract' import { parseSshPtySourceFrame } from './ssh-pty-source-frame' import { parseSshIdentityEvidenceNotification } from './ssh-pty-identity-notification' import { SshPtySourceDeliveryLedger } from './ssh-pty-source-delivery-ledger' @@ -14,6 +14,7 @@ import type { PendingSshPtySourceData, SshPtyRejectedSourceRecovery } from './ssh-pty-source-delivery-state' +import { rejectedRecoveryPriority, rejectedSourceIdentity } from './ssh-pty-notification-rejection' export type { SshPtyDataCallback, @@ -264,45 +265,3 @@ export function subscribeSshPtyNotifications(args: { } }) } - -function rejectedRecoveryPriority(recovery: SshPtyRejectedSourceRecovery): number { - if (recovery === 'reconnect-channel') { - return 3 - } - return recovery === 'fresh-activation' ? 2 : 1 -} - -function rejectedSourceIdentity(params: { - deliveryToken?: unknown - clientGeneration?: unknown - ownerGeneration?: unknown - ptyIncarnation?: unknown -}): - | Readonly<{ - deliveryToken: string - clientGeneration: number - ownerGeneration: number - ptyIncarnation: string - }> - | undefined { - if ( - typeof params.deliveryToken !== 'string' || - params.deliveryToken.length === 0 || - !positiveSafeInteger(params.clientGeneration) || - !positiveSafeInteger(params.ownerGeneration) || - typeof params.ptyIncarnation !== 'string' || - params.ptyIncarnation.length === 0 - ) { - return undefined - } - return Object.freeze({ - deliveryToken: params.deliveryToken, - clientGeneration: params.clientGeneration, - ownerGeneration: params.ownerGeneration, - ptyIncarnation: params.ptyIncarnation - }) -} - -function positiveSafeInteger(value: unknown): value is number { - return Number.isSafeInteger(value) && (value as number) > 0 -} diff --git a/src/relay/dispatcher-notification-publication.ts b/src/relay/dispatcher-notification-publication.ts index e6908344fba..8146c55f0aa 100644 --- a/src/relay/dispatcher-notification-publication.ts +++ b/src/relay/dispatcher-notification-publication.ts @@ -3,9 +3,9 @@ import type { JsonRpcNotification } from './protocol' import { DROPPED_NOTIFICATION_LOG_KEY_LIMIT, type PreparedRelayFrame, - type RelayClient + type RelayClient, + type PtyIdentityEvidencePublicationAdmission } from './dispatcher-contract' -import type { PtyIdentityEvidencePublicationAdmission } from './dispatcher-contract' import { RelayDispatcherPtyPublication } from './dispatcher-pty-publication' export abstract class RelayDispatcherNotificationPublication extends RelayDispatcherPtyPublication { diff --git a/src/relay/pty-handler-attach-replay.test.ts b/src/relay/pty-handler-attach-replay.test.ts index 2b87d338025..3c4503c7265 100644 --- a/src/relay/pty-handler-attach-replay.test.ts +++ b/src/relay/pty-handler-attach-replay.test.ts @@ -209,6 +209,79 @@ describe('PtyHandler', () => { expect(handler.getIdentityEvidenceDebugSnapshot().processTableReads).toBe(1) }) + it('publishes mixed held evidence as epoch-consistent batches', async () => { + const dataCallbacks: ((data: string) => void)[] = [] + mockPtySpawn.mockReturnValue({ + ...mockPtyInstance, + onData: vi.fn((cb: (data: string) => void) => { + dataCallbacks.push(cb) + }), + onExit: vi.fn() + }) + vi.spyOn(processTableSnapshot, 'getStrictProcessTableSnapshot').mockResolvedValue([]) + const publications: Record[] = [] + Object.assign(dispatcher, { + activeClientIds: () => [1], + admitsPtyIdentityEvidencePublication: () => true, + publishProducerNotification: vi.fn((_clientId: number, method: string, params: unknown) => { + if (method === 'pty.identityEvidence' && params && typeof params === 'object') { + publications.push(params as Record) + } + return true + }) + }) + + await spawnPty() + await spawnPty() + const boundary = '\x1b]133;C\x07' + dataCallbacks[0]?.(boundary) + await vi.advanceTimersByTimeAsync(0) + dataCallbacks[1]?.(boundary) + await vi.advanceTimersByTimeAsync(0) + + expect(handler.getIdentityEvidenceDebugSnapshot().processTableReads).toBe(2) + expect(handler.getIdentityEvidenceDebugSnapshot().heldRows).toBe(2) + publications.length = 0 + await dispatcher.callRequest( + 'pty.identityEvidence.setVisibility', + { ids: [PTY_1, testPtyId(2)] }, + { clientId: 1, isStale: () => false } as never + ) + + expect(publications).toHaveLength(2) + for (const publication of publications) { + const epoch = publication.observationEpoch + expect(typeof epoch).toBe('number') + expect( + (publication.rows as { foregroundProcessEvidence: { observationEpoch: number } }[]).every( + (row) => row.foregroundProcessEvidence.observationEpoch === epoch + ) + ).toBe(true) + } + }) + + it('does not read the process table for clients without the identity capability', async () => { + let dataCallback: ((data: string) => void) | undefined + mockPtySpawn.mockReturnValue({ + ...mockPtyInstance, + onData: vi.fn((cb: (data: string) => void) => { + dataCallback = cb + }), + onExit: vi.fn() + }) + vi.spyOn(processTableSnapshot, 'getStrictProcessTableSnapshot').mockResolvedValue([]) + Object.assign(dispatcher, { + activeClientIds: () => [1], + admitsPtyIdentityEvidencePublication: () => false + }) + + await spawnPty() + dataCallback?.('\x1b]133;C\x07') + await vi.advanceTimersByTimeAsync(0) + + expect(handler.getIdentityEvidenceDebugSnapshot().processTableReads).toBe(0) + }) + it('suppresses legacy replay after the V1 owner is already active', async () => { let dataCallback: ((data: string) => void) | undefined mockPtySpawn.mockReturnValue({ diff --git a/src/relay/pty-handler.ts b/src/relay/pty-handler.ts index 6c875932afa..c90c8fbede6 100644 --- a/src/relay/pty-handler.ts +++ b/src/relay/pty-handler.ts @@ -2512,12 +2512,28 @@ export class PtyHandler { if (!this.dispatcher.publishProducerNotification) { return } - const epoch = Math.max(...rows.map((row) => row.foregroundProcessEvidence.observationEpoch)) - this.dispatcher.publishProducerNotification(clientId, 'pty.identityEvidence', { - authorityGeneration: this.ptyIdMintEpoch, - observationEpoch: epoch, - rows - }) + // Rows can come from different host reads (a boundary-triggered read refreshes only one PTY), + // so one envelope epoch would make the provider reject the entire reseed as inconsistent. + // Preserve each row's observation epoch and send one internally consistent batch per epoch. + const rowsByEpoch = new Map() + for (const row of rows) { + const epoch = row.foregroundProcessEvidence.observationEpoch + const batch = rowsByEpoch.get(epoch) + if (batch) { + batch.push(row) + } else { + rowsByEpoch.set(epoch, [row]) + } + } + for (const [observationEpoch, epochRows] of [...rowsByEpoch.entries()].sort( + ([left], [right]) => left - right + )) { + this.dispatcher.publishProducerNotification(clientId, 'pty.identityEvidence', { + authorityGeneration: this.ptyIdMintEpoch, + observationEpoch, + rows: epochRows + }) + } } private reconcileVisibleIdentityEvidence(): void { @@ -2590,11 +2606,13 @@ export class PtyHandler { this.identityEvidenceReadQueued = true return } - if (this.dispatcher.hasConnectedClients && !this.dispatcher.hasConnectedClients()) { + if (!this.hasIdentityEvidenceConsumer()) { this.identityEvidencePendingIds.clear() return } - const requestedIds = this.identityEvidencePendingIds + // Snapshot before clearing; retaining the mutable set here would make `clear()` erase the + // filter too, turning every boundary-triggered read into a full-pool scan. + const requestedIds = new Set(this.identityEvidencePendingIds) this.identityEvidencePendingIds.clear() const entries = Array.from(this.ptys.values()).filter( (managed) => @@ -2687,6 +2705,20 @@ export class PtyHandler { } } + /** Skip process-table work when no connected client negotiated identity evidence. */ + private hasIdentityEvidenceConsumer(): boolean { + const activeClientIds = this.dispatcher.activeClientIds + const admits = this.dispatcher.admitsPtyIdentityEvidencePublication + // Test doubles and pre-capability dispatchers omit these methods; retain their historical + // behaviour so a relay can still serve local PTY reads while the optional capability rolls out. + if (typeof activeClientIds !== 'function' || typeof admits !== 'function') { + return true + } + return activeClientIds + .call(this.dispatcher) + .some((clientId) => admits.call(this.dispatcher, clientId)) + } + private async listProcesses(): Promise { const results: PtyProcessSummary[] = [] // Why (SSH-v3 P2 — the host is the authoritative liveness source, so it has to look): this diff --git a/src/shared/pty-consumer-session-capabilities.ts b/src/shared/pty-consumer-session-capabilities.ts index b9ea9909025..2c04ff50efd 100644 --- a/src/shared/pty-consumer-session-capabilities.ts +++ b/src/shared/pty-consumer-session-capabilities.ts @@ -45,7 +45,9 @@ export function intersectPtyConsumerCapabilities( const outputSupported = Boolean( offer && support && offer.versions.includes(1) && support.versions.includes(1) ) - const identitySupported = Boolean(identityOffer && identitySupport?.versions.includes(1)) + const identitySupported = Boolean( + identityOffer && identityOffer.versions.includes(1) && identitySupport?.versions.includes(1) + ) if (!outputSupported && !identitySupported) { return {} } diff --git a/src/shared/pty-consumer-session.test.ts b/src/shared/pty-consumer-session.test.ts index c53991a0f50..615efb55aeb 100644 --- a/src/shared/pty-consumer-session.test.ts +++ b/src/shared/pty-consumer-session.test.ts @@ -531,6 +531,20 @@ describe('PtyConsumerSession', () => { }) }) + it('does not grant identity evidence when the client omits V1', () => { + const session = new PtyConsumerSession({ + serverBuildId: 'relay-build', + identityEvidence: { versions: [1] }, + createLease: () => 'lease' + }) + const admission = session.admit( + ownerHello({ capabilities: { identityEvidence: { versions: [2] } } }), + auth('connection-1') + ) + + expect(admission.grant.capabilities).toBeUndefined() + }) + it('makes token-free bounded legacy an explicit capability omission', () => { const session = createSession() const admission = session.admit(ownerHello(), auth('connection-1'))