From aaa20293bff1fc2e75e11b87481cc553a96f0eb9 Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Mon, 31 Aug 2026 13:44:22 -0700 Subject: [PATCH] feat(ssh): relay live identity producer, capability push, and SSH gate flip (R2) --- src/main/ipc/pty/ipc/resize-visibility.ts | 3 + src/main/providers/pty-provider-contract.ts | 4 + .../ssh-pty-identity-notification.ts | 57 +++ .../providers/ssh-pty-identity-visibility.ts | 38 ++ ...-pty-notification-routing.identity.test.ts | 49 +++ .../providers/ssh-pty-notification-routing.ts | 43 ++- .../providers/ssh-pty-provider-contract.ts | 11 + .../ssh-pty-provider-output-state.ts | 11 +- src/main/providers/ssh-pty-provider.ts | 14 +- src/main/ssh/ssh-pty-consumer-recovery.ts | 6 +- src/main/ssh/ssh-pty-consumer-session.test.ts | 13 + src/main/ssh/ssh-pty-consumer-session.ts | 26 +- src/main/ssh/ssh-relay-session.ts | 17 +- src/preload/api/pty-api.ts | 4 + src/preload/index.ts | 13 + src/relay/dispatcher-capacity-signals.ts | 5 + src/relay/dispatcher-contract.ts | 5 + .../dispatcher-notification-publication.ts | 40 ++ src/relay/dispatcher.ts | 3 +- src/relay/pty-handler.ts | 356 +++++++++++++++++- src/relay/ssh-pty-consumer-session-adapter.ts | 16 +- .../terminal-pane/ipc-pty-connect.ts | 4 +- .../pty-connection/pane-agent-identity.ts | 26 +- .../pane-ssh-identity-evidence.ts | 71 ++++ .../transport-output-callbacks.ts | 7 + .../terminal-pane/pty-dispatcher.ts | 146 ++----- .../terminal-pane/pty-eager-dispatch.ts | 89 +++++ .../pty-identity-evidence-dispatch.ts | 24 ++ .../terminal-pane/pty-transport-types.ts | 1 + .../components/terminal-pane/pty-transport.ts | 8 +- .../lib/pty-identity-evidence-store.test.ts | 66 ++++ .../src/lib/pty-identity-evidence-store.ts | 115 ++++++ .../pty-consumer-session-capabilities.ts | 32 +- src/shared/pty-consumer-session-contract.ts | 9 + src/shared/pty-consumer-session-hello.ts | 26 +- src/shared/pty-consumer-session.ts | 6 +- src/shared/pty-identity-evidence.test.ts | 22 ++ src/shared/pty-identity-evidence.ts | 63 ++++ src/shared/ssh-types.ts | 3 + 39 files changed, 1286 insertions(+), 166 deletions(-) create mode 100644 src/main/providers/ssh-pty-identity-notification.ts create mode 100644 src/main/providers/ssh-pty-identity-visibility.ts create mode 100644 src/main/providers/ssh-pty-notification-routing.identity.test.ts create mode 100644 src/renderer/src/components/terminal-pane/pty-connection/pane-ssh-identity-evidence.ts create mode 100644 src/renderer/src/components/terminal-pane/pty-eager-dispatch.ts create mode 100644 src/renderer/src/components/terminal-pane/pty-identity-evidence-dispatch.ts create mode 100644 src/renderer/src/lib/pty-identity-evidence-store.test.ts create mode 100644 src/renderer/src/lib/pty-identity-evidence-store.ts create mode 100644 src/shared/pty-identity-evidence.test.ts create mode 100644 src/shared/pty-identity-evidence.ts diff --git a/src/main/ipc/pty/ipc/resize-visibility.ts b/src/main/ipc/pty/ipc/resize-visibility.ts index d0c376fabc8..8f72f419a9b 100644 --- a/src/main/ipc/pty/ipc/resize-visibility.ts +++ b/src/main/ipc/pty/ipc/resize-visibility.ts @@ -239,6 +239,9 @@ export function installPtyResizeVisibilityIpc(session: PtyIpcSession): void { visibleRendererPtys.delete(args.id) } session.syncPtyBackgroundedDelivery(args.id, 'visibility-report') + // Keep the SSH relay's subscription in sync with the renderer's current visible set. The + // provider debounces this fire-and-forget control request and filters ids to its host. + tryGetProviderForPty(args.id)?.setIdentityEvidenceVisibility?.(Array.from(visibleRendererPtys)) }) ipcMain.removeAllListeners('pty:setHiddenRendererPty') diff --git a/src/main/providers/pty-provider-contract.ts b/src/main/providers/pty-provider-contract.ts index 35fca5b0b34..b8f9424a1dd 100644 --- a/src/main/providers/pty-provider-contract.ts +++ b/src/main/providers/pty-provider-contract.ts @@ -12,6 +12,7 @@ import type { import type { PtyProcessInfo } from './pty-process-info' import type { TerminalExitCause } from '../../shared/terminal-exit-cause' import type { TerminalOwner } from '../../shared/terminal-owner' +import type { SshPtyIdentityEvidenceCallback } from './ssh-pty-provider-contract' export type { PtyBackgroundStreamEvent, @@ -227,6 +228,9 @@ export type IPtyProvider = { getProfiles(): Promise<{ name: string; path: string }[]> onData(callback: (payload: PtyDataEvent) => void): () => void onReplay(callback: (payload: { id: string; data: string }) => void): () => void + onIdentityEvidence?: (callback: SshPtyIdentityEvidenceCallback) => () => void + /** Update the host-side visible-pane projection used by identity reconciliation. */ + setIdentityEvidenceVisibility?: (ids: string[]) => void onExit( callback: (payload: { id: string diff --git a/src/main/providers/ssh-pty-identity-notification.ts b/src/main/providers/ssh-pty-identity-notification.ts new file mode 100644 index 00000000000..e13b5064bd4 --- /dev/null +++ b/src/main/providers/ssh-pty-identity-notification.ts @@ -0,0 +1,57 @@ +import { isPtyIncarnationId } from '../../shared/pty-incarnation' +import { isForegroundProcessEvidence } from '../../shared/foreground-process-evidence' +import type { ForegroundProcessEvidence } from '../../shared/foreground-process-evidence' + +export type AdmittedSshIdentityEvidence = { + authorityGeneration: string + observationEpoch: number + rows: { + id: string + incarnationId: string + foregroundProcessEvidence: ForegroundProcessEvidence + }[] +} + +export function parseSshIdentityEvidenceNotification( + params: Readonly>, + toAppPtyId: (id: string) => string +): AdmittedSshIdentityEvidence | null { + const authorityGeneration = params.authorityGeneration + const observationEpoch = params.observationEpoch + const rows = params.rows + if ( + typeof authorityGeneration !== 'string' || + authorityGeneration.length === 0 || + !Number.isSafeInteger(observationEpoch) || + (observationEpoch as number) < 0 || + !Array.isArray(rows) + ) { + return null + } + const admitted: AdmittedSshIdentityEvidence['rows'] = [] + for (const row of rows) { + if (typeof row !== 'object' || row === null) { + return null + } + const input = row as Record + const evidence = input.foregroundProcessEvidence + if ( + typeof input.id !== 'string' || + input.id.length === 0 || + !isPtyIncarnationId(input.incarnationId) || + !isForegroundProcessEvidence(evidence) || + evidence.authorityGeneration !== authorityGeneration || + evidence.observationEpoch !== observationEpoch + ) { + return null + } + let id: string + try { + id = toAppPtyId(input.id) + } catch { + return null + } + admitted.push({ id, incarnationId: input.incarnationId, foregroundProcessEvidence: evidence }) + } + return { authorityGeneration, observationEpoch: observationEpoch as number, rows: admitted } +} diff --git a/src/main/providers/ssh-pty-identity-visibility.ts b/src/main/providers/ssh-pty-identity-visibility.ts new file mode 100644 index 00000000000..0c92f647706 --- /dev/null +++ b/src/main/providers/ssh-pty-identity-visibility.ts @@ -0,0 +1,38 @@ +import type { SshChannelMultiplexer } from '../ssh/ssh-channel-multiplexer' +import { parseAppSshPtyId } from '../../shared/ssh-pty-id' + +export function createSshIdentityVisibilityPublisher( + mux: SshChannelMultiplexer, + connectionId: string +): { set: (ids: string[]) => void; dispose: () => void } { + let timer: ReturnType | null = null + let pending: string[] | null = null + return { + set(ids) { + const relayIds = ids.flatMap((id) => { + const parsed = parseAppSshPtyId(id) + return parsed?.connectionId === connectionId ? [parsed.relayPtyId] : [] + }) + pending = Array.from(new Set(relayIds)).slice(0, 512) + if (timer !== null) { + return + } + timer = setTimeout(() => { + timer = null + const next = pending + pending = null + if (next) { + void mux.request('pty.identityEvidence.setVisibility', { ids: next }).catch(() => {}) + } + }, 50) + timer.unref?.() + }, + dispose() { + if (timer !== null) { + clearTimeout(timer) + } + timer = null + pending = null + } + } +} diff --git a/src/main/providers/ssh-pty-notification-routing.identity.test.ts b/src/main/providers/ssh-pty-notification-routing.identity.test.ts new file mode 100644 index 00000000000..b92cec45198 --- /dev/null +++ b/src/main/providers/ssh-pty-notification-routing.identity.test.ts @@ -0,0 +1,49 @@ +import { describe, expect, it, vi } from 'vitest' +import { SshPtyProvider } from './ssh-pty-provider' + +describe('SSH identity-evidence notification routing', () => { + it('namespaces admitted rows and rejects malformed batches', () => { + const onNotification = vi.fn().mockReturnValue(vi.fn()) + const mux = { + request: vi.fn().mockResolvedValue(undefined), + notify: vi.fn(), + onNotification, + dispose: vi.fn(), + isDisposed: vi.fn().mockReturnValue(false) + } + const provider = new SshPtyProvider('conn-1', mux as never) + const listener = vi.fn() + provider.onIdentityEvidence?.(listener) + const notify = onNotification.mock.calls[0][0] as (method: string, params: unknown) => void + const evidence = { + verdict: 'live', + processName: 'codex', + authorityGeneration: 'relay-generation', + observationEpoch: 4, + capturedAgeMs: 0 + } + notify('pty.identityEvidence', { + authorityGeneration: 'relay-generation', + observationEpoch: 4, + rows: [{ id: 'pty-1', incarnationId: 'incarnation-1', foregroundProcessEvidence: evidence }] + }) + expect(listener).toHaveBeenCalledWith({ + authorityGeneration: 'relay-generation', + observationEpoch: 4, + providerGeneration: 1, + rows: [ + { + id: 'ssh:conn-1@@pty-1', + incarnationId: 'incarnation-1', + foregroundProcessEvidence: evidence + } + ] + }) + notify('pty.identityEvidence', { + authorityGeneration: 'relay-generation', + observationEpoch: 5, + rows: [{ id: 'pty-1', incarnationId: '', foregroundProcessEvidence: evidence }] + }) + expect(listener).toHaveBeenCalledTimes(1) + }) +}) diff --git a/src/main/providers/ssh-pty-notification-routing.ts b/src/main/providers/ssh-pty-notification-routing.ts index b6063799f58..d749ac375ac 100644 --- a/src/main/providers/ssh-pty-notification-routing.ts +++ b/src/main/providers/ssh-pty-notification-routing.ts @@ -6,14 +6,21 @@ import type { SshPtyExitCallback, SshPtyReplayCallback } 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' import type { PendingSshPtySourceData, SshPtyRejectedSourceRecovery } from './ssh-pty-source-delivery-state' -export type { SshPtyDataCallback, SshPtyExitCallback, SshPtyReplayCallback } +export type { + SshPtyDataCallback, + SshPtyExitCallback, + SshPtyReplayCallback, + SshPtyIdentityEvidenceCallback +} export type SshPtyRecoveryActivationLease = Readonly<{ commit: () => void retire: () => void @@ -38,6 +45,7 @@ export function subscribeSshPtyNotifications(args: { dataListeners: Set rejectedDataListeners?: Set replayListeners: Set + identityEvidenceListeners?: Set exitListeners: Set livePtyIds: Set recordExit: (relayPtyId: string, incarnationId: unknown) => void @@ -106,6 +114,8 @@ export function subscribeSshPtyNotifications(args: { } } const sourceDeliveries = new SshPtySourceDeliveryLedger(args.mux, publishData) + let identityAuthorityGeneration: string | undefined + let identityObservationEpoch = -1 const rejectedPublications = new Map< string, { @@ -153,7 +163,36 @@ export function subscribeSshPtyNotifications(args: { // Why: mux delivers every method to generic handlers; non-PTY payloads // (workspace.changed, fs.changed, …) have no `id` and must not reach // toAppPtyId → startsWith. - if (method !== 'pty.exit' && method !== 'pty.data' && method !== 'pty.replay') { + if ( + method !== 'pty.exit' && + method !== 'pty.data' && + method !== 'pty.replay' && + method !== 'pty.identityEvidence' + ) { + return + } + if (method === 'pty.identityEvidence') { + const parsed = parseSshIdentityEvidenceNotification(params, args.toAppPtyId) + if (!parsed) { + return + } + const { authorityGeneration, observationEpoch, rows } = parsed + if ( + identityAuthorityGeneration === authorityGeneration && + (observationEpoch as number) <= identityObservationEpoch + ) { + return + } + identityAuthorityGeneration = authorityGeneration + identityObservationEpoch = observationEpoch as number + for (const listener of args.identityEvidenceListeners ?? []) { + listener({ + authorityGeneration, + observationEpoch: observationEpoch as number, + rows, + providerGeneration: args.providerGeneration + }) + } return } if (typeof params.id !== 'string' || params.id.length === 0) { diff --git a/src/main/providers/ssh-pty-provider-contract.ts b/src/main/providers/ssh-pty-provider-contract.ts index 4ac6f0fb4b2..85afb5cd7b0 100644 --- a/src/main/providers/ssh-pty-provider-contract.ts +++ b/src/main/providers/ssh-pty-provider-contract.ts @@ -1,4 +1,5 @@ import type { PtyIncarnationId } from '../../shared/pty-incarnation' +import type { ForegroundProcessEvidence } from '../../shared/foreground-process-evidence' export type RemoteCliBridgeEnv = { binDir: string @@ -31,6 +32,16 @@ export type SshPtyDataCallback = (payload: { rejectedSourceRecovery?: 'confirm-existing' | 'fresh-activation' | 'reconnect-channel' }) => void export type SshPtyReplayCallback = (payload: { id: string; data: string }) => void +export type SshPtyIdentityEvidenceCallback = (payload: { + authorityGeneration: string + observationEpoch: number + rows: readonly { + id: string + incarnationId: string + foregroundProcessEvidence: ForegroundProcessEvidence + }[] + providerGeneration: number +}) => void export type SshPtyExitCallback = (payload: { id: string code: number diff --git a/src/main/providers/ssh-pty-provider-output-state.ts b/src/main/providers/ssh-pty-provider-output-state.ts index 4d98097b9c6..acba33dfb19 100644 --- a/src/main/providers/ssh-pty-provider-output-state.ts +++ b/src/main/providers/ssh-pty-provider-output-state.ts @@ -3,7 +3,8 @@ import type { SshPtyDataCallback, SshPtyDeliveryPauseAdapter, SshPtyExitCallback, - SshPtyReplayCallback + SshPtyReplayCallback, + SshPtyIdentityEvidenceCallback } from './ssh-pty-provider-contract' import { subscribeSshPtyNotifications, @@ -16,6 +17,7 @@ export class SshPtyProviderOutputState { private readonly dataListeners = new Set() private readonly rejectedDataListeners = new Set() private readonly replayListeners = new Set() + private readonly identityEvidenceListeners = new Set() private readonly exitListeners = new Set() private readonly incarnationByRelayPtyId = new Map() private readonly pausedRelayPtyIds = new Set() @@ -37,6 +39,7 @@ export class SshPtyProviderOutputState { dataListeners: this.dataListeners, rejectedDataListeners: this.rejectedDataListeners, replayListeners: this.replayListeners, + identityEvidenceListeners: this.identityEvidenceListeners, exitListeners: this.exitListeners, providerGeneration, resolvePtyIncarnation: (relayPtyId, incarnationId) => @@ -57,6 +60,7 @@ export class SshPtyProviderOutputState { this.dataListeners.clear() this.rejectedDataListeners.clear() this.replayListeners.clear() + this.identityEvidenceListeners.clear() this.exitListeners.clear() this.incarnationByRelayPtyId.clear() this.deliveryPauseAdapter = null @@ -77,6 +81,11 @@ export class SshPtyProviderOutputState { return () => this.replayListeners.delete(callback) } + onIdentityEvidence(callback: SshPtyIdentityEvidenceCallback): () => void { + this.identityEvidenceListeners.add(callback) + return () => this.identityEvidenceListeners.delete(callback) + } + onExit(callback: SshPtyExitCallback): () => void { this.exitListeners.add(callback) return () => this.exitListeners.delete(callback) diff --git a/src/main/providers/ssh-pty-provider.ts b/src/main/providers/ssh-pty-provider.ts index a8e4218ce63..0f43c95038d 100644 --- a/src/main/providers/ssh-pty-provider.ts +++ b/src/main/providers/ssh-pty-provider.ts @@ -7,7 +7,8 @@ import type { SshPtyDataCallback, SshPtyDeliveryPauseAdapter, SshPtyExitCallback, - SshPtyReplayCallback + SshPtyReplayCallback, + SshPtyIdentityEvidenceCallback } from './ssh-pty-provider-contract' import { SshPtyProviderOutputState } from './ssh-pty-provider-output-state' import { spawnFreshSshPty } from './ssh-agent-session-create-operation' @@ -23,8 +24,8 @@ import { SshPtySpawnExitRaceTracker } from './ssh-pty-spawn-exit-race' import { SshAgentSessionCapabilities } from './ssh-agent-session-capabilities' import type { PtyProcessInspection } from './pty-process-inspection' import { writeToSshPty, writeToSshPtyWithSettlement } from './ssh-pty-write' +import { createSshIdentityVisibilityPublisher } from './ssh-pty-identity-visibility' -// 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 { return deadlineMs === undefined ? undefined : { timeoutMs: Math.max(1, deadlineMs - Date.now()) } } @@ -38,6 +39,7 @@ export class SshPtyProvider implements IPtyProvider { private readonly agentSessionCapabilities: SshAgentSessionCapabilities private spawnExitRaces = new SshPtySpawnExitRaceTracker() private readonly outputState: SshPtyProviderOutputState + private readonly identityVisibility: ReturnType requestHostRpc: NonNullable = (method, params, options) => this.mux.request(method, params as Record, options) @@ -52,6 +54,7 @@ export class SshPtyProvider implements IPtyProvider { this.mux = mux this.agentSessionCapabilities = new SshAgentSessionCapabilities(mux) this.getAppliedSize = createSshPtyAppliedSizeReader(mux, connectionId) + this.identityVisibility = createSshIdentityVisibilityPublisher(mux, connectionId) this.outputState = new SshPtyProviderOutputState(providerGeneration, { mux, @@ -62,14 +65,14 @@ export class SshPtyProvider implements IPtyProvider { } }) } - dispose(): void { + this.identityVisibility.dispose() this.outputState.dispose() this.livePtyIds.clear() } getConnectionId = (): string => this.connectionId - + setIdentityEvidenceVisibility = (ids: string[]): void => this.identityVisibility.set(ids) canProvideAuthoritativeBufferSnapshot = (_id: string): boolean => false private toRelayPtyId = (id: string): string => toRelaySshPtyId(this.connectionId, id) @@ -293,7 +296,6 @@ export class SshPtyProvider implements IPtyProvider { } hasPty = (id: string): boolean => this.livePtyIds.has(id) - async getDefaultShell(): Promise { const result = await this.mux.request('pty.getDefaultShell') return result as string @@ -308,6 +310,8 @@ export class SshPtyProvider implements IPtyProvider { onRejectedData = (callback: SshPtyDataCallback): (() => void) => this.outputState.onRejectedData(callback) onReplay = (callback: SshPtyReplayCallback): (() => void) => this.outputState.onReplay(callback) + onIdentityEvidence = (callback: SshPtyIdentityEvidenceCallback): (() => void) => + this.outputState.onIdentityEvidence(callback) onExit = (callback: SshPtyExitCallback): (() => void) => this.outputState.onExit(callback) setPtyDeliveryPauseAdapter(adapter: SshPtyDeliveryPauseAdapter | null): void { diff --git a/src/main/ssh/ssh-pty-consumer-recovery.ts b/src/main/ssh/ssh-pty-consumer-recovery.ts index 1c357d32a12..7b3c32328a2 100644 --- a/src/main/ssh/ssh-pty-consumer-recovery.ts +++ b/src/main/ssh/ssh-pty-consumer-recovery.ts @@ -23,7 +23,8 @@ function ownerFromPersisted(record: PersistedRecovery): SshPtyConsumerOwnerState clientGeneration: record.clientGeneration, ownerGeneration: record.ownerGeneration, ownerLease: record.ownerLease, - ...(record.outputFlowControl ? { outputFlowControl: record.outputFlowControl } : {}) + ...(record.outputFlowControl ? { outputFlowControl: record.outputFlowControl } : {}), + ...(record.identityEvidence ? { identityEvidence: record.identityEvidence } : {}) } } @@ -78,7 +79,8 @@ export async function rememberSshPtyConsumerRecovery(args: { clientGeneration: args.owner.clientGeneration, ownerGeneration: args.owner.ownerGeneration, ownerLease: args.owner.ownerLease, - ...(args.owner.outputFlowControl ? { outputFlowControl: args.owner.outputFlowControl } : {}) + ...(args.owner.outputFlowControl ? { outputFlowControl: args.owner.outputFlowControl } : {}), + ...(args.owner.identityEvidence ? { identityEvidence: args.owner.identityEvidence } : {}) }) } diff --git a/src/main/ssh/ssh-pty-consumer-session.test.ts b/src/main/ssh/ssh-pty-consumer-session.test.ts index a1fc612f1dd..93d368aae98 100644 --- a/src/main/ssh/ssh-pty-consumer-session.test.ts +++ b/src/main/ssh/ssh-pty-consumer-session.test.ts @@ -124,6 +124,19 @@ describe('openSshPtyConsumerSession', () => { ).rejects.toThrow('did not grant') }) + it('keeps identity evidence optional when an older relay omits the grant', async () => { + const { mux, request } = muxReturning(legacyOwnerGrant()) + const admission = await openSshPtyConsumerSession(mux, { + clientInstanceId: 'client-a', + expectedServerBuildId: 'build-a', + identityEvidence: true + }) + expect(admission.state).not.toHaveProperty('identityEvidence') + expect(request.mock.calls[0][1]).toMatchObject({ + capabilities: { identityEvidence: { versions: [1] } } + }) + }) + it('rejects an unoffered V1 capability in a legacy session', async () => { const { mux } = muxReturning( legacyOwnerGrant({ diff --git a/src/main/ssh/ssh-pty-consumer-session.ts b/src/main/ssh/ssh-pty-consumer-session.ts index aa58534909b..7dbadae6173 100644 --- a/src/main/ssh/ssh-pty-consumer-session.ts +++ b/src/main/ssh/ssh-pty-consumer-session.ts @@ -17,6 +17,9 @@ export type SshPtyConsumerOwnerState = { version: 1 windowSu: number } + identityEvidence?: { + version: 1 + } } export type SshPtyLegacyFallbackState = { @@ -41,6 +44,7 @@ export type OpenSshPtyConsumerSessionOptions = { outputFlowControl?: { requestedWindowSu: number } + identityEvidence?: boolean allowSameBuildLegacyFallback?: boolean } @@ -94,6 +98,10 @@ function validateGrant( } else if (grantedFlow) { throw new Error('Remote relay granted an unoffered PTY output-flow-control capability') } + const grantedIdentity = grant.capabilities?.identityEvidence + if (!options.identityEvidence && grantedIdentity) { + throw new Error('Remote relay granted an unoffered PTY identity-evidence capability') + } return grant as PtyConsumerSessionGrant } @@ -110,13 +118,18 @@ export async function openSshPtyConsumerSession( clientInstanceId: options.clientInstanceId, requestedRole: 'session-owner', ...(options.resume ? { resume: options.resume } : {}), - ...(options.outputFlowControl + ...(options.outputFlowControl || options.identityEvidence ? { capabilities: { - outputFlowControl: { - versions: [1], - requestedWindowSu: options.outputFlowControl.requestedWindowSu - } + ...(options.outputFlowControl + ? { + outputFlowControl: { + versions: [1], + requestedWindowSu: options.outputFlowControl.requestedWindowSu + } + } + : {}), + ...(options.identityEvidence ? { identityEvidence: { versions: [1] } } : {}) } } : {}) @@ -152,6 +165,9 @@ export async function openSshPtyConsumerSession( ownerLease: grant.ownerLease!, ...(grant.capabilities?.outputFlowControl ? { outputFlowControl: grant.capabilities.outputFlowControl } + : {}), + ...(grant.capabilities?.identityEvidence + ? { identityEvidence: grant.capabilities.identityEvidence } : {}) }, resumed: grant.resumed! diff --git a/src/main/ssh/ssh-relay-session.ts b/src/main/ssh/ssh-relay-session.ts index bdc7c4370f2..30ddc8b1f7b 100644 --- a/src/main/ssh/ssh-relay-session.ts +++ b/src/main/ssh/ssh-relay-session.ts @@ -1156,7 +1156,8 @@ export class SshRelaySession { clientInstanceId: this.ptyConsumerClientInstanceId, expectedServerBuildId: serverBuildId, allowSameBuildLegacyFallback: true, - outputFlowControl: { requestedWindowSu: DEFAULT_PTY_SOURCE_WINDOW_SU } + outputFlowControl: { requestedWindowSu: DEFAULT_PTY_SOURCE_WINDOW_SU }, + identityEvidence: true } let admission: SshPtyConsumerAdmission try { @@ -1771,6 +1772,20 @@ export class SshRelaySession { win.webContents.send('pty:replay', payload) } }) + ptyProvider.onIdentityEvidence?.((payload) => { + if ( + this.mux !== mux || + this.activePtyProviderGeneration !== providerGeneration || + payload.providerGeneration !== providerGeneration || + !this.activePtyConsumerOwner()?.identityEvidence + ) { + return + } + const win = this.getMainWindow() + if (win && !win.isDestroyed()) { + win.webContents.send('pty:identityEvidence', payload) + } + }) ptyProvider.onExit((payload) => { if ( this.mux !== mux || diff --git a/src/preload/api/pty-api.ts b/src/preload/api/pty-api.ts index 7eee16ea522..00b4db6d644 100644 --- a/src/preload/api/pty-api.ts +++ b/src/preload/api/pty-api.ts @@ -16,6 +16,7 @@ import type { TerminalSideEffectBatch } from '../../shared/terminal-side-effect- import type { TerminalViewAttributes } from '../../shared/terminal-view-attributes' import type { TuiAgent } from '../../shared/tui-agent' import type { PtyManagementApi } from './pty-management-api' +import type { PtyIdentityEvidenceNotification } from '../../shared/pty-identity-evidence' export type PtyApi = { spawn: (opts: { @@ -104,6 +105,8 @@ export type PtyApi = { /** Ref-counted-on-the-renderer delivery-interest signal that suppresses * the hidden-delivery gate while any raw-byte consumer is registered. */ setPtyDeliveryInterest: (id: string, interested: boolean) => void + /** Updates the SSH relay's visible-pane identity subscription for this connection. */ + setIdentityEvidenceVisibility?: (ids: string[]) => void /** View-attribute bridge (Phase 5 slice 2): app-global composed terminal * appearance push backing main's hidden-PTY OSC/DSR color replies. */ publishTerminalViewAttributes: (attributes: TerminalViewAttributes) => void @@ -188,6 +191,7 @@ export type PtyApi = { }) => void ) => () => void onReplay: (callback: (data: { id: string; data: string }) => void) => () => void + onIdentityEvidence?: (callback: (data: PtyIdentityEvidenceNotification) => void) => () => void /** Out-of-band main→renderer signal that renderer-bound bytes were * dropped (hidden-delivery gate / pending cap); the pane restores from * the model snapshot. Never delivered in-band on pty:data. */ diff --git a/src/preload/index.ts b/src/preload/index.ts index 16f7b82e712..272ea86da3c 100644 --- a/src/preload/index.ts +++ b/src/preload/index.ts @@ -20,6 +20,7 @@ import { } from '../shared/doc-preview-scheme' import type { DocPreviewGrantRequest } from './api/doc-preview-api' import type { AppIdentity } from '../shared/app-identity' +import type { PtyIdentityEvidenceNotification } from '../shared/pty-identity-evidence' import type { MacCapturedDigitRowChord } from '../shared/macos-symbolic-hotkeys' import type { ComputerAwakeStatus } from '../shared/computer-awake-mode' import type { @@ -1174,6 +1175,9 @@ const api = { setPtyDeliveryInterest: (id: string, interested: boolean): void => { ipcRenderer.send('pty:setPtyDeliveryInterest', { id, interested }) }, + setIdentityEvidenceVisibility: (ids: string[]): void => { + ipcRenderer.send('pty:setIdentityEvidenceVisibility', { ids }) + }, /** Push composed terminal appearance so main's model responder can answer OSC 4/10/11/12 and DSR ?996n for hidden-gated PTYs with renderer-true values. */ publishTerminalViewAttributes: (attributes: TerminalViewAttributes): void => { ipcRenderer.send('pty:terminalViewAttributes', attributes) @@ -1297,6 +1301,15 @@ const api = { return () => ipcRenderer.removeListener('pty:replay', listener) }, + onIdentityEvidence: ( + callback: (data: PtyIdentityEvidenceNotification) => void + ): (() => void) => { + const listener = (_event: Electron.IpcRendererEvent, data: PtyIdentityEvidenceNotification) => + callback(data) + ipcRenderer.on('pty:identityEvidence', listener) + return () => ipcRenderer.removeListener('pty:identityEvidence', listener) + }, + /** Out-of-band signal that main dropped renderer-bound bytes (hidden-gate / pending cap); pane restores from the model snapshot. * NOT on pty:data — an in-band marker is ambiguous with chunks fully stripped by OSC-9999 cleaning. */ onModelRestoreNeeded: (callback: (event: PtyModelRestoreNeededEvent) => void): (() => void) => { diff --git a/src/relay/dispatcher-capacity-signals.ts b/src/relay/dispatcher-capacity-signals.ts index 5f945f5c84f..1a835c05924 100644 --- a/src/relay/dispatcher-capacity-signals.ts +++ b/src/relay/dispatcher-capacity-signals.ts @@ -82,6 +82,11 @@ export abstract class RelayDispatcherCapacitySignals extends RelayDispatcherClie return Array.from(this.clients.values()).filter((client) => !client.closed) } + /** Active transport ids for producers that need per-client projections. */ + activeClientIds(): number[] { + return this.activeClients().map((client) => client.id) + } + protected admitsPtyDataPublication( clientId: number, params: Readonly> diff --git a/src/relay/dispatcher-contract.ts b/src/relay/dispatcher-contract.ts index 15d4578c9bf..38ac2f1263a 100644 --- a/src/relay/dispatcher-contract.ts +++ b/src/relay/dispatcher-contract.ts @@ -32,6 +32,11 @@ export type PtyDataPublicationAdmission = ( params: Readonly> ) => boolean +export type PtyIdentityEvidencePublicationAdmission = ( + clientId: number, + params: Readonly> +) => boolean + export type MethodHandler = ( params: Record, context: RequestContext diff --git a/src/relay/dispatcher-notification-publication.ts b/src/relay/dispatcher-notification-publication.ts index ab0851b4203..e6908344fba 100644 --- a/src/relay/dispatcher-notification-publication.ts +++ b/src/relay/dispatcher-notification-publication.ts @@ -5,9 +5,37 @@ import { type PreparedRelayFrame, type RelayClient } from './dispatcher-contract' +import type { PtyIdentityEvidencePublicationAdmission } from './dispatcher-contract' import { RelayDispatcherPtyPublication } from './dispatcher-pty-publication' export abstract class RelayDispatcherNotificationPublication extends RelayDispatcherPtyPublication { + private ptyIdentityEvidenceAdmission: PtyIdentityEvidencePublicationAdmission | null = null + + registerPtyIdentityEvidencePublicationAdmission( + admission: PtyIdentityEvidencePublicationAdmission | null + ): () => void { + if (admission && this.ptyIdentityEvidenceAdmission) { + throw new Error('PTY identity evidence publication admission is already registered') + } + this.ptyIdentityEvidenceAdmission = admission + return () => { + if (this.ptyIdentityEvidenceAdmission === admission) { + this.ptyIdentityEvidenceAdmission = null + } + } + } + + admitsPtyIdentityEvidencePublication( + clientId: number, + params: Readonly> = {} + ): boolean { + return this.ptyIdentityEvidenceAdmission?.(clientId, params) ?? false + } + + hasConnectedClients(): boolean { + return this.activeClients().length > 0 + } + notify(method: string, params?: Record): void { if (this.disposed) { return @@ -26,6 +54,12 @@ export abstract class RelayDispatcherNotificationPublication extends RelayDispat if (method === 'pty.data' && !this.admitsPtyDataPublication(client.id, params ?? {})) { continue } + if ( + method === 'pty.identityEvidence' && + !this.admitsPtyIdentityEvidencePublication(client.id, params ?? {}) + ) { + continue + } frame ??= this.prepareFrame(msg) if (method === 'pty.replay') { // Why: replay is never re-sent, so it takes the control lane where overflow is fatal — the @@ -67,6 +101,12 @@ export abstract class RelayDispatcherNotificationPublication extends RelayDispat if (method === 'pty.data' && !this.admitsPtyDataPublication(client.id, params ?? {})) { return false } + if ( + method === 'pty.identityEvidence' && + !this.admitsPtyIdentityEvidencePublication(client.id, params ?? {}) + ) { + return false + } const frame = this.prepareFrame(msg) if (this.publishPreparedToClient(client, frame, 'ordinary')) { return true diff --git a/src/relay/dispatcher.ts b/src/relay/dispatcher.ts index dbf70abd906..3b2e227b4a0 100644 --- a/src/relay/dispatcher.ts +++ b/src/relay/dispatcher.ts @@ -11,7 +11,8 @@ export type { PtyDataPublicationAdmission, RelayClientSessionIdentity, RelayClientSourceOptions, - RequestContext + RequestContext, + PtyIdentityEvidencePublicationAdmission } from './dispatcher-contract' export class RelayDispatcher extends RelayDispatcherNotificationPublication {} diff --git a/src/relay/pty-handler.ts b/src/relay/pty-handler.ts index 82f0b19b9ff..d2f54289bcb 100644 --- a/src/relay/pty-handler.ts +++ b/src/relay/pty-handler.ts @@ -73,6 +73,11 @@ import { type ProcessTableRow } from '../shared/process-table-snapshot' import type { ForegroundProcessEvidence } from '../shared/foreground-process-evidence' +import { + createPtyIdentityBoundaryScanner, + type PtyIdentityEvidenceNotification, + type PtyIdentityEvidenceRow +} from '../shared/pty-identity-evidence' import { expandWindowsPathEnvironmentVariables } from '../shared/windows-environment-expansion' import { agentSessionOwnerBindingsEqual, @@ -211,6 +216,14 @@ const AGENT_SESSION_CREATE_OPERATION_ID_PATTERN = /^[A-Za-z0-9_-]{43}$/ const AGENT_SESSION_CREATE_OPERATION_RETENTION_MS = 24 * 60 * 60 * 1000 const AGENT_SESSION_CREATE_OPERATION_LIMIT = 4_096 +// Held evidence is only a reconnect seed; bounded retention keeps a quiet relay from +// accumulating one row per recycled PTY forever. +const IDENTITY_EVIDENCE_HELD_MAX_ROWS = 512 +const IDENTITY_EVIDENCE_HELD_MAX_BYTES = 512 * 1024 +const IDENTITY_EVIDENCE_HELD_TTL_MS = 30_000 +const IDENTITY_EVIDENCE_FRESHNESS_MS = 5_000 +const IDENTITY_EVIDENCE_BACKSTOP_MS = 30_000 + type PendingPtyOutput = RelayPtySourceOutput & { data: string interactive?: boolean @@ -451,6 +464,24 @@ export class PtyHandler { private ptys = new Map() private readonly ptyIdMintEpoch: string private foregroundEvidenceEpoch = 0 + private identityEvidenceReadTimer: ReturnType | null = null + private identityEvidenceReadInFlight = false + private identityEvidenceReadQueued = false + private identityEvidenceReadCount = 0 + private readonly identityEvidencePendingIds = new Set() + private identityEvidenceBackstopTimer: ReturnType | null = null + private readonly identityEvidenceHeld = new Map() + private readonly identityEvidenceHeldAt = new Map() + private identityEvidenceHeldBytes = 0 + private readonly identityEvidenceScanners = new Map< + string, + ReturnType + >() + private readonly identityEvidenceVisibleByClient = new Map>() + private readonly identityEvidenceBackoffByPty = new Map< + string, + { attempts: number; nextAt: number } + >() private nextId = 1 private dispatcher: RelayDispatcher private graceTimeMs: number @@ -668,6 +699,10 @@ export class PtyHandler { if (this.ptys.size > 0) { return } + if (this.identityEvidenceBackstopTimer !== null) { + clearInterval(this.identityEvidenceBackstopTimer) + this.identityEvidenceBackstopTimer = null + } this.notifyPoolListener(this.ptyPoolEmptyListener, 'pty-pool-empty') } @@ -839,6 +874,16 @@ export class PtyHandler { private wireAndStore(managed: ManagedPty): void { managed.physicalExit = new PhysicalExitTracker() this.ptys.set(managed.id, managed) + if (this.identityEvidenceBackstopTimer === null) { + this.identityEvidenceBackstopTimer = setInterval( + () => this.reconcileVisibleIdentityEvidence(), + IDENTITY_EVIDENCE_BACKSTOP_MS + ) + this.identityEvidenceBackstopTimer.unref?.() + } + if (this.dispatcher.hasConnectedClients?.()) { + this.scheduleIdentityEvidenceRead() + } // Why: a PTY joining the pool under this paneKey means the surface exists again (reopened pane // or revive), so a prior retirement no longer describes anything and must not mute its hooks. const boundPaneKey = managed.paneKey ?? managed.attachIdentity?.paneKey @@ -864,6 +909,12 @@ export class PtyHandler { write: (data) => managed.pty.write(data), onEmission: emitIngressData }) + // The scanner observes only raw bytes from the live child PTY. Replay is emitted from the + // retained buffer through a separate path and never calls this feed. + this.identityEvidenceScanners.set( + managed.id, + createPtyIdentityBoundaryScanner(() => this.scheduleIdentityEvidenceRead(managed.id)) + ) const startup = managed.startupCommand if (startup?.waitForShellReady) { startup.promptProbe = createShellPromptReadinessProbe({ @@ -882,6 +933,7 @@ export class PtyHandler { }) } managed.pty.onData((data: string) => { + this.identityEvidenceScanners.get(managed.id)?.feed(data) const startup = managed.startupCommand if (startup?.waitForShellReady && startup.outputScanState && !startup.delivered) { const scanned = scanShellStartupOutput(startup.outputScanState, data) @@ -918,6 +970,8 @@ export class PtyHandler { } this.clearStartupCommandTimer(managed) this.releaseRelayIngress(managed) + this.identityEvidenceScanners.delete(managed.id) + this.evictIdentityEvidence(managed.id, managed.incarnationId) this.pausedOutputPtys.delete(managed.id) this.consumerPausedOutputPtys.delete(managed.id) this.flushPtyOutput(managed.id) @@ -979,9 +1033,26 @@ export class PtyHandler { this.dispatcher.onRequest('pty.getCapabilities', async () => ({ startupIngressVersion: PTY_STARTUP_INGRESS_VERSION, agentSessionClaimVersion: AGENT_SESSION_EXECUTION_OWNER_PROTOCOL_VERSION, - agentSessionCreateOperationVersion: AGENT_SESSION_CREATE_OPERATION_PROTOCOL_VERSION + agentSessionCreateOperationVersion: AGENT_SESSION_CREATE_OPERATION_PROTOCOL_VERSION, + identityEvidence: { versions: [1] } })) this.dispatcher.onRequest('pty.listProcesses', () => this.listProcesses()) + this.dispatcher.onRequest('pty.identityEvidence.setVisibility', async (params, context) => { + if (!this.dispatcher.admitsPtyIdentityEvidencePublication(context.clientId)) { + throw new Error('pty_identity_evidence_capability_required') + } + const ids = params.ids + if (!Array.isArray(ids) || ids.length > 512 || ids.some((id) => typeof id !== 'string')) { + throw new Error('invalid_identity_evidence_visibility') + } + const visible = new Set(ids as string[]) + this.identityEvidenceVisibleByClient.set(context.clientId, visible) + this.publishHeldIdentityEvidenceToClient(context.clientId, visible) + return { ok: true } + }) + this.dispatcher.onClientDetached?.((clientId) => { + this.identityEvidenceVisibleByClient.delete(clientId) + }) this.dispatcher.onRequest('pty.getDefaultShell', async () => resolveDefaultShell()) this.dispatcher.onRequest('pty.serialize', (p) => this.serialize(p)) this.dispatcher.onRequest('pty.revive', (p) => this.revive(p)) @@ -1981,6 +2052,12 @@ export class PtyHandler { ) const sourceActivation = context && this.sourcePublication?.receivingActivation?.(id, context.clientId) + if (context) { + this.publishHeldIdentityEvidenceToClient( + context.clientId, + this.identityEvidenceVisibleByClient.get(context.clientId) + ) + } if (typeof activation === 'object') { return { incarnationId: managed.incarnationId, @@ -2370,6 +2447,232 @@ export class PtyHandler { } } + private evictIdentityEvidence(id: string, incarnationId?: string): void { + for (const [key, row] of this.identityEvidenceHeld) { + if (row.id !== id || (incarnationId && row.incarnationId !== incarnationId)) { + continue + } + this.identityEvidenceHeldBytes -= this.identityEvidenceRowBytes(row) + this.identityEvidenceHeld.delete(key) + this.identityEvidenceHeldAt.delete(key) + } + } + + private identityEvidenceRowBytes(row: PtyIdentityEvidenceRow): number { + return chargedPtyRetainedStringBytes(JSON.stringify(row)) + } + + private pruneIdentityEvidenceHeld(now = performance.now()): void { + for (const [key, row] of this.identityEvidenceHeld) { + const age = + row.foregroundProcessEvidence.capturedAgeMs + + Math.max(0, now - (this.identityEvidenceHeldAt.get(key) ?? now)) + if (age > IDENTITY_EVIDENCE_HELD_TTL_MS) { + this.identityEvidenceHeld.delete(key) + this.identityEvidenceHeldAt.delete(key) + this.identityEvidenceHeldBytes -= this.identityEvidenceRowBytes(row) + } + } + // Map insertion order is observation order; evict the oldest rows first when bounded. + while ( + this.identityEvidenceHeld.size > IDENTITY_EVIDENCE_HELD_MAX_ROWS || + this.identityEvidenceHeldBytes > IDENTITY_EVIDENCE_HELD_MAX_BYTES + ) { + const oldest = this.identityEvidenceHeld.keys().next().value as string | undefined + if (!oldest) { + break + } + const row = this.identityEvidenceHeld.get(oldest) + this.identityEvidenceHeld.delete(oldest) + if (row) { + this.identityEvidenceHeldBytes -= this.identityEvidenceRowBytes(row) + } + } + } + + private storeIdentityEvidenceRow(row: PtyIdentityEvidenceRow): void { + const key = `${row.id}\0${row.incarnationId}` + const previous = this.identityEvidenceHeld.get(key) + if (previous) { + this.identityEvidenceHeldBytes -= this.identityEvidenceRowBytes(previous) + } + this.identityEvidenceHeld.set(key, row) + this.identityEvidenceHeldAt.set(key, performance.now()) + this.identityEvidenceHeldBytes += this.identityEvidenceRowBytes(row) + this.pruneIdentityEvidenceHeld() + } + + private publishHeldIdentityEvidenceToClient(clientId: number, ids?: ReadonlySet): void { + this.pruneIdentityEvidenceHeld() + const rows = Array.from(this.identityEvidenceHeld.values()).filter( + (row) => ids === undefined || ids.has(row.id) + ) + if (rows.length === 0) { + return + } + 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 + }) + } + + private reconcileVisibleIdentityEvidence(): void { + if (this.identityEvidenceVisibleByClient.size === 0) { + return + } + const now = performance.now() + const stale = new Set() + for (const visible of this.identityEvidenceVisibleByClient.values()) { + for (const id of visible) { + const managed = this.ptys.get(id) + if (!managed || managed.disposed) { + continue + } + const row = this.identityEvidenceHeld.get(`${managed.id}\0${managed.incarnationId}`) + const heldAt = + this.identityEvidenceHeldAt.get(`${managed.id}\0${managed.incarnationId}`) ?? now + const age = + (row?.foregroundProcessEvidence.capturedAgeMs ?? Number.POSITIVE_INFINITY) + + Math.max(0, now - heldAt) + const backoff = this.identityEvidenceBackoffByPty.get(id) + if ( + (!row || age >= IDENTITY_EVIDENCE_FRESHNESS_MS) && + (!backoff || now >= backoff.nextAt) + ) { + stale.add(id) + } + } + } + for (const id of stale) { + this.scheduleIdentityEvidenceRead(id) + } + } + + private scheduleIdentityEvidenceRead(id?: string): void { + if (id) { + this.identityEvidencePendingIds.add(id) + } + if (this.identityEvidenceReadTimer !== null) { + return + } + this.identityEvidenceReadQueued = true + this.identityEvidenceReadTimer = setTimeout(() => { + this.identityEvidenceReadTimer = null + if (!this.identityEvidenceReadInFlight) { + this.identityEvidenceReadQueued = false + void this.publishIdentityEvidence() + } + }, 0) + this.identityEvidenceReadTimer.unref?.() + } + + private async publishIdentityEvidence(): Promise { + if (this.identityEvidenceReadInFlight) { + this.identityEvidenceReadQueued = true + return + } + if (this.dispatcher.hasConnectedClients && !this.dispatcher.hasConnectedClients()) { + this.identityEvidencePendingIds.clear() + return + } + const requestedIds = this.identityEvidencePendingIds + this.identityEvidencePendingIds.clear() + const entries = Array.from(this.ptys.values()).filter( + (managed) => + !managed.disposed && + managed.pty.pid > 0 && + (requestedIds.size === 0 || requestedIds.has(managed.id)) + ) + if (entries.length === 0) { + return + } + this.identityEvidenceReadInFlight = true + this.identityEvidenceReadCount++ + try { + const epoch = ++this.foregroundEvidenceEpoch + let results: BatchedForegroundProcessResult[] + if (process.platform === 'win32') { + results = entries.map(() => ({ + available: false, + processName: null, + reason: 'unsupported' + })) + } else { + try { + const rows = await getStrictProcessTableSnapshot() + results = await resolveAgentForegroundProcessesBatch( + entries.map((managed) => ({ + rootPid: managed.pty.pid, + fallbackProcess: managed.pty.process || null + })), + { rows } + ) + } catch { + results = entries.map(() => ({ + available: false, + processName: null, + reason: 'table_unreadable' + })) + } + } + const rows: PtyIdentityEvidenceRow[] = entries.map((managed, index) => ({ + id: managed.id, + incarnationId: managed.incarnationId, + foregroundProcessEvidence: toForegroundProcessEvidence( + results[index] ?? { available: false, processName: null, reason: 'table_unreadable' }, + { + authorityGeneration: this.ptyIdMintEpoch, + observationEpoch: epoch, + capturedAgeMs: 0 + } + ) + })) + for (const row of rows) { + this.storeIdentityEvidenceRow(row) + if (row.foregroundProcessEvidence.verdict === 'live') { + this.identityEvidenceBackoffByPty.delete(row.id) + } else { + const attempts = (this.identityEvidenceBackoffByPty.get(row.id)?.attempts ?? 0) + 1 + const delaySeconds = [2, 4, 8, 30][Math.min(attempts - 1, 3)] + this.identityEvidenceBackoffByPty.set(row.id, { + attempts, + nextAt: performance.now() + delaySeconds * 1_000 + }) + } + } + const notification: PtyIdentityEvidenceNotification = { + authorityGeneration: this.ptyIdMintEpoch, + observationEpoch: epoch, + rows + } + this.pruneIdentityEvidenceHeld() + for (const clientId of this.dispatcher.activeClientIds?.() ?? []) { + const visible = this.identityEvidenceVisibleByClient.get(clientId) + // A client that has not sent a visibility report is intentionally uncovered; it still + // receives a held push so reattach can seed its renderer, but no host read is keyed to it. + const projectedRows = visible ? rows.filter((row) => visible.has(row.id)) : rows + if (projectedRows.length === 0) { + continue + } + this.dispatcher.publishProducerNotification(clientId, 'pty.identityEvidence', { + ...notification, + rows: projectedRows + } as unknown as Record) + } + } finally { + this.identityEvidenceReadInFlight = false + if (this.identityEvidenceReadQueued) { + this.identityEvidenceReadQueued = false + this.scheduleIdentityEvidenceRead() + } + } + } + private async listProcesses(): Promise { const results: PtyProcessSummary[] = [] // Why (SSH-v3 P2 — the host is the authoritative liveness source, so it has to look): this @@ -2381,8 +2684,22 @@ export class PtyHandler { // the existing title/liveness path until the measured relay adapter lands. let evidenceRows: readonly ProcessTableRow[] | null = null let evidenceResults: BatchedForegroundProcessResult[] = [] - const evidenceEpoch = ++this.foregroundEvidenceEpoch - if (process.platform !== 'win32' && managedEntries.length > 0) { + const heldEvidence = managedEntries.map(([, managed]) => + this.identityEvidenceHeld.get(`${managed.id}\0${managed.incarnationId}`) + ) + const hasCompleteHeldEvidence = heldEvidence.every((row) => row !== undefined) + const evidenceEpoch = hasCompleteHeldEvidence + ? (heldEvidence[0]?.foregroundProcessEvidence.observationEpoch ?? + this.foregroundEvidenceEpoch) + : ++this.foregroundEvidenceEpoch + if (hasCompleteHeldEvidence) { + evidenceResults = heldEvidence.map((row) => { + const evidence = row!.foregroundProcessEvidence + return evidence.verdict === 'live' + ? { available: true, processName: evidence.processName } + : { available: false, processName: null, reason: evidence.reason } + }) + } else if (process.platform !== 'win32' && managedEntries.length > 0) { try { evidenceRows = await getStrictProcessTableSnapshot() evidenceResults = await resolveAgentForegroundProcessesBatch( @@ -2409,7 +2726,8 @@ export class PtyHandler { : await getForegroundProcessName(managed.pty.pid, managed.pty.process || null)) || 'shell' const foregroundProcessEvidence = process.platform !== 'win32' - ? toForegroundProcessEvidence( + ? (heldEvidence[entryIndex]?.foregroundProcessEvidence ?? + toForegroundProcessEvidence( evidenceResults[entryIndex] ?? { available: false, processName: managed.pty.process || null, @@ -2420,7 +2738,7 @@ export class PtyHandler { observationEpoch: evidenceEpoch, capturedAgeMs: 0 } - ) + )) : undefined results.push({ id, @@ -2714,6 +3032,22 @@ export class PtyHandler { this.pendingOutputByPty.clear() this.pendingProducerBytesByPty.clear() this.pendingExitByPty.clear() + if (this.identityEvidenceReadTimer !== null) { + clearTimeout(this.identityEvidenceReadTimer) + this.identityEvidenceReadTimer = null + } + if (this.identityEvidenceBackstopTimer !== null) { + clearInterval(this.identityEvidenceBackstopTimer) + this.identityEvidenceBackstopTimer = null + } + this.identityEvidencePendingIds.clear() + this.identityEvidenceReadQueued = false + this.identityEvidenceHeld.clear() + this.identityEvidenceHeldAt.clear() + this.identityEvidenceHeldBytes = 0 + this.identityEvidenceScanners.clear() + this.identityEvidenceVisibleByClient.clear() + this.identityEvidenceBackoffByPty.clear() this.pausedOutputPtys.clear() this.consumerPausedOutputPtys.clear() this.lastInputAtByPty.clear() @@ -2796,6 +3130,18 @@ export class PtyHandler { return this.ptys.size } + getIdentityEvidenceDebugSnapshot(): Readonly<{ + heldRows: number + heldBytes: number + processTableReads: number + }> { + return { + heldRows: this.identityEvidenceHeld.size, + heldBytes: this.identityEvidenceHeldBytes, + processTableReads: this.identityEvidenceReadCount + } + } + /** Spawns admitted but not yet in the pool — each already owns a shell the relay must not treat as idle. */ get pendingPtyCreationCount(): number { return this.pendingSpawnCount diff --git a/src/relay/ssh-pty-consumer-session-adapter.ts b/src/relay/ssh-pty-consumer-session-adapter.ts index 83e61f17fa3..ea23148b745 100644 --- a/src/relay/ssh-pty-consumer-session-adapter.ts +++ b/src/relay/ssh-pty-consumer-session-adapter.ts @@ -27,6 +27,7 @@ export class SshPtyConsumerSessionAdapter { private readonly session: PtyConsumerSession private readonly sourceCredit: SshPtySourceCreditAdapter private readonly pausedDeliveryByPty = new Map() + private readonly removeIdentityAdmission: (() => void) | null constructor( private readonly dispatcher: RelayDispatcher, @@ -44,8 +45,20 @@ export class SshPtyConsumerSessionAdapter { ) this.session = new PtyConsumerSession({ serverBuildId, - outputFlowControl: { versions: [1], maxWindowSu: DEFAULT_PTY_SOURCE_WINDOW_SU } + outputFlowControl: { versions: [1], maxWindowSu: DEFAULT_PTY_SOURCE_WINDOW_SU }, + identityEvidence: { versions: [1] } }) + this.removeIdentityAdmission = + dispatcher.registerPtyIdentityEvidencePublicationAdmission?.((clientId, params) => { + const grant = this.session.activeGrant(String(clientId)) + if (grant?.capabilities?.identityEvidence?.version !== 1) { + return false + } + if (params.clientGeneration !== undefined) { + return params.clientGeneration === grant.clientGeneration + } + return true + }) ?? null // Why the admission is consulted again at drain time (see isStillAdmitted): a frame can sit // queued behind a saturated socket long enough for the grant or the delivery to be retired, and // publishing it then hands the client output from an owner it no longer is. @@ -104,6 +117,7 @@ export class SshPtyConsumerSessionAdapter { this.sourceCredit.cancel(params, this.session.activeGrant(String(context.clientId))) ) dispatcher.onDisposed(() => { + this.removeIdentityAdmission?.() for (const id of this.pausedDeliveryByPty.keys()) { this.setDeliveryPaused?.(id, false) } diff --git a/src/renderer/src/components/terminal-pane/ipc-pty-connect.ts b/src/renderer/src/components/terminal-pane/ipc-pty-connect.ts index 70d3a0fd047..0e671e2b4f4 100644 --- a/src/renderer/src/components/terminal-pane/ipc-pty-connect.ts +++ b/src/renderer/src/components/terminal-pane/ipc-pty-connect.ts @@ -24,7 +24,7 @@ type IpcPtyConnectContext = { transportOptions: IpcPtyTransportOptions handlers: IpcPtySessionHandlers isDestroyed: () => boolean - bind: (id: string) => void + bind: (id: string, incarnationId?: string) => void isCurrent: (id: string) => boolean setCallbacks: (callbacks: PtyConnectOptions['callbacks']) => void getCallbacks: () => PtyConnectOptions['callbacks'] @@ -103,7 +103,7 @@ export async function connectIpcPty( // buffered exit is the real thing. A fresh spawn's PTY did not exist yet. discardPreHandlerPtyStateFromPriorIncarnation(spawnResult.id, priorIncarnationFence) } - context.bind(spawnResult.id) + context.bind(spawnResult.id, spawnResult.incarnationId) if (!spawnResult.isReattach && !spawnResult.coldRestore) { onPtySpawn?.(spawnResult.id) } diff --git a/src/renderer/src/components/terminal-pane/pty-connection/pane-agent-identity.ts b/src/renderer/src/components/terminal-pane/pty-connection/pane-agent-identity.ts index f2c4290af76..cb5697e7225 100644 --- a/src/renderer/src/components/terminal-pane/pty-connection/pane-agent-identity.ts +++ b/src/renderer/src/components/terminal-pane/pty-connection/pane-agent-identity.ts @@ -9,17 +9,28 @@ import { } from '@/lib/pane-manager/windows-pty-compatibility' import { createTerminalCommandLifecycle } from '../terminal-command-lifecycle' import { createPaneForegroundAgentTracker } from '../pane-foreground-agent-tracker' -import { isRemoteExecutionHostPtyId } from '../remote-execution-host-pty' import { dispatchTerminalCommandFinishedEvent } from '@/hooks/terminal-command-finished-event' import { getExecutionHostIdForWorktree } from '@/lib/worktree-runtime-owner' import { resolveCommittedTitleAgentType } from '@/lib/pane-agent-evidence' import type { TuiAgent } from '../../../../../shared/tui-agent' import { isTuiAgent, TUI_AGENT_CONFIG } from '../../../../../shared/tui-agent-config' +import { parseAppSshPtyId } from '../../../../../shared/ssh-pty-id' +import { bindPaneSshIdentityEvidence } from './pane-ssh-identity-evidence' import type { ConnectPanePtySession } from './connect-pane-pty-session' /** Pane agent identity, foreground-agent sampling, and command lifecycle handling. */ export function installPaneAgentIdentity(session: ConnectPanePtySession): void { + const bindSshIdentity = bindPaneSshIdentityEvidence(session) + const isDirectSshPtyId = (id: string): boolean => parseAppSshPtyId(id) !== null + session.disposeIdentityEvidence = session.disposeIdentityEvidence ?? null + session.bindIdentityEvidence = (ptyId: string, incarnationId?: string): void => { + session.disposeIdentityEvidence?.() + if (!isDirectSshPtyId(ptyId)) { + return + } + session.disposeIdentityEvidence = bindSshIdentity(ptyId, incarnationId) + } // Why: the 133;D confirmation guard and the visible-pane resampler both key off // "does this pane expect an agent"; derive each signal once so the two callers // can't drift and silently reintroduce the icon bug this fix closes. @@ -106,10 +117,6 @@ export function installPaneAgentIdentity(session: ConnectPanePtySession): void { if (current && current.acceptedStatusSeq !== armedAcceptedStatusSeq) { return } - // Why: main-side only. The renderer row and launch config are already owned by the deferred - // drop above; what that path cannot reach is the hook server's per-pane Claude latches, which - // `agentStatus:drop` deliberately preserves for a still-live pane. Main echoes its own clear - // back through the pane-status-cleared channel, so both sides stay consistent. window.api?.agentStatus?.reconcileEndedProcess?.(session.cacheKey) } session.visibleForegroundSamplePending = false @@ -127,7 +134,8 @@ export function installPaneAgentIdentity(session: ConnectPanePtySession): void { } } session.isForegroundTrackingAllowed = (id: string): boolean => { - if (isRemoteExecutionHostPtyId(id)) { + if (isDirectSshPtyId(id)) { + session.bindIdentityEvidence?.(id) return false } if (!navigator.userAgent.includes('Windows')) { @@ -153,8 +161,10 @@ export function installPaneAgentIdentity(session: ConnectPanePtySession): void { session.paneForegroundAgentTracker = createPaneForegroundAgentTracker({ getPtyId: () => session.transport.getPtyId(), isTrackablePtyId: session.isForegroundTrackingAllowed, - readForegroundProcess: (id) => window.api.pty.getForegroundProcess(id), - confirmForegroundProcess: (id) => window.api.pty.confirmForegroundProcess(id), + readForegroundProcess: (id) => + isDirectSshPtyId(id) ? Promise.resolve(null) : window.api.pty.getForegroundProcess(id), + confirmForegroundProcess: (id) => + isDirectSshPtyId(id) ? Promise.resolve(null) : window.api.pty.confirmForegroundProcess(id), publish: (entry) => useAppStore.getState().setPaneForegroundAgent(session.cacheKey, entry), hasKnownAgentIdentity: session.paneHasKnownAgentIdentity, onConfirmedShellForeground: (reason) => { diff --git a/src/renderer/src/components/terminal-pane/pty-connection/pane-ssh-identity-evidence.ts b/src/renderer/src/components/terminal-pane/pty-connection/pane-ssh-identity-evidence.ts new file mode 100644 index 00000000000..8f79b0962de --- /dev/null +++ b/src/renderer/src/components/terminal-pane/pty-connection/pane-ssh-identity-evidence.ts @@ -0,0 +1,71 @@ +import { useAppStore } from '@/store' +import { registerPtyIdentityEvidenceHandler } from '../pty-dispatcher' +import { recognizeAgentProcess } from '../../../../../shared/agent-process-recognition' +import { parseAppSshPtyId } from '../../../../../shared/ssh-pty-id' +import { ptyIdentityEvidenceStore } from '@/lib/pty-identity-evidence-store' +import type { ConnectPanePtySession } from './connect-pane-pty-session' + +export function bindPaneSshIdentityEvidence( + session: ConnectPanePtySession +): (ptyId: string, incarnationId?: string) => (() => void) | null { + let dispose: (() => void) | null = null + let epoch = -1 + let providerGeneration = -1 + return (ptyId, incarnationId) => { + dispose?.() + const parsed = parseAppSshPtyId(ptyId) + if (!parsed) { + return null + } + dispose = registerPtyIdentityEvidenceHandler(ptyId, (notification) => { + const incomingProviderGeneration = notification.providerGeneration ?? 0 + if (incomingProviderGeneration < providerGeneration) { + return + } + if (incomingProviderGeneration > providerGeneration) { + providerGeneration = incomingProviderGeneration + epoch = -1 + ptyIdentityEvidenceStore.activateGeneration( + parsed.connectionId, + notification.authorityGeneration + ) + } + if (notification.observationEpoch <= epoch) { + return + } + const row = notification.rows.find((candidate) => candidate.id === ptyId) + if (!row) { + return + } + if (incarnationId && row.incarnationId !== incarnationId) { + return + } + epoch = notification.observationEpoch + ptyIdentityEvidenceStore.applyPush({ + hostId: parsed.connectionId, + ptyId, + incarnationId: row.incarnationId, + authorityGeneration: notification.authorityGeneration, + observationEpoch: notification.observationEpoch, + evidence: row.foregroundProcessEvidence + }) + const evidence = row.foregroundProcessEvidence + if (evidence.verdict === 'live') { + const recognized = evidence.processName ? recognizeAgentProcess(evidence.processName) : null + useAppStore.getState().setPaneForegroundAgent(session.cacheKey, { + agent: recognized?.agent ?? null, + shellForeground: false, + routingTrusted: recognized !== null + }) + } else { + const current = useAppStore.getState().paneForegroundAgentByPaneKey[session.cacheKey] + useAppStore.getState().setPaneForegroundAgent(session.cacheKey, { + agent: current?.agent ?? null, + shellForeground: current?.shellForeground ?? false, + routingTrusted: false + }) + } + }) + return dispose + } +} diff --git a/src/renderer/src/components/terminal-pane/pty-connection/transport-output-callbacks.ts b/src/renderer/src/components/terminal-pane/pty-connection/transport-output-callbacks.ts index a1e2e97e070..419aec0dccb 100644 --- a/src/renderer/src/components/terminal-pane/pty-connection/transport-output-callbacks.ts +++ b/src/renderer/src/components/terminal-pane/pty-connection/transport-output-callbacks.ts @@ -37,6 +37,13 @@ export function bindCaptureTransportOutputCallbacks(session: ConnectPanePtySessi }, onConnect: (): void => { if (isCurrent()) { + const ptyId = session.transport.getPtyId() + if (ptyId) { + session.bindIdentityEvidence?.( + ptyId, + session.transport.getPtyIncarnationId?.() ?? undefined + ) + } session.reportRemoteRendererSerializerReady() // Re-derive the pause bit after a rebind; visibility can change while no PTY is bound. session.syncHiddenRendererPtyDelivery() diff --git a/src/renderer/src/components/terminal-pane/pty-dispatcher.ts b/src/renderer/src/components/terminal-pane/pty-dispatcher.ts index ccb38ae2f55..11a4087eee6 100644 --- a/src/renderer/src/components/terminal-pane/pty-dispatcher.ts +++ b/src/renderer/src/components/terminal-pane/pty-dispatcher.ts @@ -1,18 +1,10 @@ -/** Singleton PTY event dispatcher and eager buffer helpers, split out from pty-transport.ts. */ -import { TERMINAL_SCROLLBACK_SESSION_BUFFER_BYTE_LIMIT } from '../../../../shared/terminal-scrollback-limits' import { clearProcessedPtyCharTotal, deliverPtyDataWithDeferredAck, exposeE2eTerminalPtyAckGate, getProcessedPtyCharTotals } from './terminal-pty-ack-gate' -import { clampUtf8Tail, type EagerBufferChunk } from './pty-eager-buffer-clamp' -import { - bufferPreHandlerPtyData, - clearPreHandlerPtyState, - drainPreHandlerPtyData, - drainPreHandlerPtyExit -} from './pty-pre-handler-buffer' +import { bufferPreHandlerPtyData } from './pty-pre-handler-buffer' import { deliverPtyExitToHandlers } from './pty-exit-delivery' import { clearReceivedPtyCharTotal, @@ -32,6 +24,20 @@ import { ptyReplayHandlers } from './pty-shutdown-data-suspension' import { markCommittedPtyShutdowns } from './pty-shutdown-exit-deferral' +import { + dispatchPtyIdentityEvidence, + ptyIdentityEvidenceHandlers, + registerPtyIdentityEvidenceHandler as registerPtyIdentityEvidenceHandlerInternal +} from './pty-identity-evidence-dispatch' +import { + getEagerPtyBufferHandle, + hasEagerPtyHandles, + registerEagerPtyBuffer +} from './pty-eager-dispatch' + +export { ptyIdentityEvidenceHandlers } +export type { EagerPtyHandle } from './pty-eager-dispatch' +export { getEagerPtyBufferHandle, registerEagerPtyBuffer } export { ptyDataHandlers, @@ -46,20 +52,15 @@ export { unregisterPtyDataHandlers } from './pty-shutdown-data-suspension' -// ── Singleton PTY event dispatcher ─────────────────────────────────── -// One global IPC listener per channel (routed by PTY ID) avoids the N-listener MaxListenersExceededWarning with many panes. - export type PtyDataMeta = { seq?: number rawLength?: number transformed?: boolean background?: boolean - /** Main dropped this PTY's buffered output at the pending cap; repaint from the main-owned snapshot, not the live stream. */ droppedOutput?: boolean } -/** Sidecar PTY-data observers, invoked AFTER the primary handler so a side-effect-only watcher can't delay xterm rendering. */ -/** Per-PTY replay handlers on a dedicated pty:replay channel so the renderer can engage the replay guard and suppress xterm auto-replies. */ +/** Sidecar PTY-data observers. */ const ptyExitSidecars = new Map< string, Set<(code: number, context: { hadPrimary: boolean }) => void> @@ -69,7 +70,14 @@ let ptyDispatcherAttached = false let pushListenerUnsubscribes: (() => void)[] = [] -/** Detach and re-subscribe every push-channel listener; called by the delivery watchdog on a confirmed wedge. */ +export function registerPtyIdentityEvidenceHandler( + ptyId: string, + handler: Parameters[1] +): () => void { + ensurePtyDispatcher() + return registerPtyIdentityEvidenceHandlerInternal(ptyId, handler) +} + export function reattachPtyDispatcherPushListeners(): void { recordTerminalFreezeBreadcrumb('push-listeners-reattach', { staleListenerCount: pushListenerUnsubscribes.length @@ -92,7 +100,7 @@ export function ensurePtyDispatcher(): void { attachPtyPushListeners() startTerminalDeliveryWatchdog({ reattachPushListeners: reattachPtyDispatcherPushListeners, - hasAttachedPtys: () => ptyDataHandlers.size > 0 || eagerPtyHandles.size > 0 + hasAttachedPtys: () => ptyDataHandlers.size > 0 || hasEagerPtyHandles() }) } @@ -143,7 +151,6 @@ function handleDispatchedPtyData(payload: { const chars = payload.rawLength ?? payload.data.length const dispatch = (): void => { if (isPtyDataHandlerShutdownPending(payload.id)) { - // Why: teardown output is speculative until the owner verifies sleep; retain it so a failed attempt resumes without losing terminal data. bufferPtyShutdownData(payload.id, payload.data, meta) return } @@ -155,7 +162,6 @@ function handleDispatchedPtyData(payload: { } const sidecars = ptyDataSidecars.get(payload.id) if (sidecars && sidecars.size > 0) { - // Why: snapshot before iterating — watchers often unsubscribe (or subscribe siblings) mid-iteration, and mutating the live Set would skip or double-fire. const snapshot = Array.from(sidecars) for (const watcher of snapshot) { watcher(payload.data) @@ -163,7 +169,6 @@ function handleDispatchedPtyData(payload: { } } recordPtyDataReceived(payload.id, chars) - // Why deferred: main budgets by bytes PARSED not received; ACK fires when xterm consumes, and undelivered chunks settle at return so no PTY stays backpressured. deliverPtyDataWithDeferredAck(payload.id, chars, dispatch) } @@ -182,6 +187,10 @@ function attachPtySecondaryPushListeners(unsubscribes: (() => void)[]): void { ptyReplayHandlers.get(payload.id)?.(payload.data) }) ) + const unsubscribeIdentity = window.api.pty.onIdentityEvidence?.(dispatchPtyIdentityEvidence) + if (unsubscribeIdentity) { + unsubscribes.push(unsubscribeIdentity) + } unsubscribes.push( window.api.pty.onExit((payload) => { if (payload.preserveRendererBinding === true) { @@ -195,6 +204,7 @@ function attachPtySecondaryPushListeners(unsubscribes: (() => void)[]): void { if (sidecars) { ptyExitSidecars.delete(payload.id) } + ptyIdentityEvidenceHandlers.delete(payload.id) const primary = ptyExitHandlers.get(payload.id) if (primary) { // Why: one-shot owner — remove before invoking so a throwing callback can't stay registered for a duplicate exit. @@ -221,7 +231,6 @@ function attachPtySecondaryPushListeners(unsubscribes: (() => void)[]): void { if (unsubscribeResync) { unsubscribes.push(unsubscribeResync) } - // Why: tell main the pty:data listener is live; until it fires, bytes to a listener-less page are dropped-but-counted and pin the delivery gate. window.api.pty.rendererDispatcherReady?.() } @@ -247,98 +256,3 @@ export function subscribeToPtyExit( } } } - -// ─── Eager PTY buffer for reconnection on restart ──────────────────── -// Why: PTYs spawn before TerminalPane mounts; buffer the early shell output (prompt/MOTD) so attach() can replay it. - -export type EagerPtyHandle = { flush: () => string; dispose: () => void } -const eagerPtyHandles = new Map() - -export function getEagerPtyBufferHandle(ptyId: string): EagerPtyHandle | undefined { - return eagerPtyHandles.get(ptyId) -} - -// Why: cap matches TerminalPane's scrollback serialization limit so a restored shell (e.g. tail -f) can't grow unbounded. -const EAGER_BUFFER_MAX_BYTES = TERMINAL_SCROLLBACK_SESSION_BUFFER_BYTE_LIMIT - -/** `incarnationId` names the lifetime the caller just spawned. Without it a background launch that - * is handed a relay-recycled id drains whatever the id's PREVIOUS owner left here and tears its own - * freshly started agent session down seconds after launch. */ -export function registerEagerPtyBuffer( - ptyId: string, - onExit: (ptyId: string, code: number) => void, - incarnationId?: string -): EagerPtyHandle { - ensurePtyDispatcher() - // Why: head index instead of Array.shift() (O(n)) so pre-attach buffering isn't quadratic under many small chunks. - const chunks: EagerBufferChunk[] = [] - let head = 0 - let bufferBytes = 0 - - const dataHandler = (data: string): void => { - // Why: a single over-cap chunk would bypass the trim loop below; keep only its most-recent tail. - const chunk = clampUtf8Tail(data, EAGER_BUFFER_MAX_BYTES) - chunks.push(chunk) - bufferBytes += chunk.bytes - // Drop whole leading chunks (keeping the prompt-bearing tail) until within cap. - while (bufferBytes > EAGER_BUFFER_MAX_BYTES && head < chunks.length - 1) { - bufferBytes -= chunks[head].bytes - chunks[head] = { data: '', bytes: 0 } - head += 1 - } - // Compact when dead slots reach half the array so it can't grow unbounded. - if (head > 0 && head * 2 >= chunks.length) { - chunks.splice(0, head) - head = 0 - } - } - const exitHandler = (code: number): void => { - // Shell died before attach; identity-guard so we never evict a handler a transport re-registered for this id (#7894 detach/attach race). - if (ptyDataHandlers.get(ptyId) === dataHandler) { - ptyDataHandlers.delete(ptyId) - ptyReplayHandlers.delete(ptyId) - } - ptyExitHandlers.delete(ptyId) - eagerPtyHandles.delete(ptyId) - onExit(ptyId, code) - } - - ptyDataHandlers.set(ptyId, dataHandler) - ptyExitHandlers.set(ptyId, exitHandler) - - const handle: EagerPtyHandle = { - flush() { - const data = chunks - .slice(head) - .map((chunk) => chunk.data) - .join('') - chunks.length = 0 - head = 0 - bufferBytes = 0 - return data - }, - dispose() { - // Why: identity-guard removal — after attach() swaps in its own handler this must no-op, not evict it. - if (ptyDataHandlers.get(ptyId) === dataHandler) { - ptyDataHandlers.delete(ptyId) - ptyReplayHandlers.delete(ptyId) - } - if (ptyExitHandlers.get(ptyId) === exitHandler) { - ptyExitHandlers.delete(ptyId) - } - eagerPtyHandles.delete(ptyId) - } - } - - eagerPtyHandles.set(ptyId, handle) - drainPreHandlerPtyData(ptyId, dataHandler) - // Why: defer the pre-handler exit one microtask so the caller receives the returned handle before onExit fires. - queueMicrotask(() => { - if (ptyExitHandlers.get(ptyId) === exitHandler) { - drainPreHandlerPtyExit(ptyId, exitHandler, incarnationId) - } else { - clearPreHandlerPtyState(ptyId) - } - }) - return handle -} diff --git a/src/renderer/src/components/terminal-pane/pty-eager-dispatch.ts b/src/renderer/src/components/terminal-pane/pty-eager-dispatch.ts new file mode 100644 index 00000000000..43039cf1ef7 --- /dev/null +++ b/src/renderer/src/components/terminal-pane/pty-eager-dispatch.ts @@ -0,0 +1,89 @@ +import { TERMINAL_SCROLLBACK_SESSION_BUFFER_BYTE_LIMIT } from '../../../../shared/terminal-scrollback-limits' +import { clampUtf8Tail, type EagerBufferChunk } from './pty-eager-buffer-clamp' +import { ptyDataHandlers, ptyExitHandlers, ptyReplayHandlers } from './pty-shutdown-data-suspension' +import { + clearPreHandlerPtyState, + drainPreHandlerPtyData, + drainPreHandlerPtyExit +} from './pty-pre-handler-buffer' +import { ensurePtyDispatcher } from './pty-dispatcher' + +export type EagerPtyHandle = { flush: () => string; dispose: () => void } +const eagerPtyHandles = new Map() +const EAGER_BUFFER_MAX_BYTES = TERMINAL_SCROLLBACK_SESSION_BUFFER_BYTE_LIMIT + +export function getEagerPtyBufferHandle(ptyId: string): EagerPtyHandle | undefined { + return eagerPtyHandles.get(ptyId) +} + +export function hasEagerPtyHandles(): boolean { + return eagerPtyHandles.size > 0 +} + +export function registerEagerPtyBuffer( + ptyId: string, + onExit: (ptyId: string, code: number) => void, + incarnationId?: string +): EagerPtyHandle { + ensurePtyDispatcher() + const chunks: EagerBufferChunk[] = [] + let head = 0 + let bufferBytes = 0 + const dataHandler = (data: string): void => { + const chunk = clampUtf8Tail(data, EAGER_BUFFER_MAX_BYTES) + chunks.push(chunk) + bufferBytes += chunk.bytes + while (bufferBytes > EAGER_BUFFER_MAX_BYTES && head < chunks.length - 1) { + bufferBytes -= chunks[head].bytes + chunks[head] = { data: '', bytes: 0 } + head += 1 + } + if (head > 0 && head * 2 >= chunks.length) { + chunks.splice(0, head) + head = 0 + } + } + const exitHandler = (code: number): void => { + if (ptyDataHandlers.get(ptyId) === dataHandler) { + ptyDataHandlers.delete(ptyId) + ptyReplayHandlers.delete(ptyId) + } + ptyExitHandlers.delete(ptyId) + eagerPtyHandles.delete(ptyId) + onExit(ptyId, code) + } + ptyDataHandlers.set(ptyId, dataHandler) + ptyExitHandlers.set(ptyId, exitHandler) + const handle: EagerPtyHandle = { + flush() { + const data = chunks + .slice(head) + .map((chunk) => chunk.data) + .join('') + chunks.length = 0 + head = 0 + bufferBytes = 0 + return data + }, + dispose() { + if (ptyDataHandlers.get(ptyId) === dataHandler) { + ptyDataHandlers.delete(ptyId) + ptyReplayHandlers.delete(ptyId) + } + if (ptyExitHandlers.get(ptyId) === exitHandler) { + ptyExitHandlers.delete(ptyId) + } + eagerPtyHandles.delete(ptyId) + } + } + eagerPtyHandles.set(ptyId, handle) + drainPreHandlerPtyData(ptyId, dataHandler) + queueMicrotask(() => { + if (ptyExitHandlers.get(ptyId) === exitHandler) { + drainPreHandlerPtyExit(ptyId, exitHandler, incarnationId) + } else { + clearPreHandlerPtyState(ptyId) + } + }) + return handle +} diff --git a/src/renderer/src/components/terminal-pane/pty-identity-evidence-dispatch.ts b/src/renderer/src/components/terminal-pane/pty-identity-evidence-dispatch.ts new file mode 100644 index 00000000000..710fe1cc5df --- /dev/null +++ b/src/renderer/src/components/terminal-pane/pty-identity-evidence-dispatch.ts @@ -0,0 +1,24 @@ +import type { PtyIdentityEvidenceNotification } from '../../../../shared/pty-identity-evidence' + +export const ptyIdentityEvidenceHandlers = new Map< + string, + (notification: PtyIdentityEvidenceNotification) => void +>() + +export function registerPtyIdentityEvidenceHandler( + ptyId: string, + handler: (notification: PtyIdentityEvidenceNotification) => void +): () => void { + ptyIdentityEvidenceHandlers.set(ptyId, handler) + return () => { + if (ptyIdentityEvidenceHandlers.get(ptyId) === handler) { + ptyIdentityEvidenceHandlers.delete(ptyId) + } + } +} + +export function dispatchPtyIdentityEvidence(notification: PtyIdentityEvidenceNotification): void { + for (const row of notification.rows) { + ptyIdentityEvidenceHandlers.get(row.id)?.(notification) + } +} diff --git a/src/renderer/src/components/terminal-pane/pty-transport-types.ts b/src/renderer/src/components/terminal-pane/pty-transport-types.ts index ae88b4bf698..4f466fbfee7 100644 --- a/src/renderer/src/components/terminal-pane/pty-transport-types.ts +++ b/src/renderer/src/components/terminal-pane/pty-transport-types.ts @@ -201,6 +201,7 @@ export type PtyTransport = { /** The user dismissed the error surface; the next occurrence of the same message must surface again. */ notifyErrorSurfaceDismissed?: () => void getPtyId: () => string | null + getPtyIncarnationId?: () => string | null getConnectionId?: () => string | null | undefined /** The runtime captured by this transport; legacy remote PTY ids do not * encode their owner, and current worktree settings may have changed. */ diff --git a/src/renderer/src/components/terminal-pane/pty-transport.ts b/src/renderer/src/components/terminal-pane/pty-transport.ts index 79f8e0d1e6d..7d35fa5b154 100644 --- a/src/renderer/src/components/terminal-pane/pty-transport.ts +++ b/src/renderer/src/components/terminal-pane/pty-transport.ts @@ -46,6 +46,7 @@ export function createIpcPtyTransport(opts: IpcPtyTransportOptions = {}): PtyTra let connected = false let destroyed = false let ptyId: string | null = null + let ptyIncarnationId: string | null = null let suppressAttentionEvents = false let storedCallbacks: Parameters[0]['callbacks'] = {} @@ -78,11 +79,13 @@ export function createIpcPtyTransport(opts: IpcPtyTransportOptions = {}): PtyTra markExited: () => { connected = false ptyId = null + ptyIncarnationId = null }, onPtyExit }) - const bind = (id: string): void => { + const bind = (id: string, incarnationId?: string): void => { ptyId = id + ptyIncarnationId = incarnationId ?? null connected = true } const setCallbacks = (callbacks: typeof storedCallbacks): void => { @@ -122,6 +125,7 @@ export function createIpcPtyTransport(opts: IpcPtyTransportOptions = {}): PtyTra window.api.pty.kill(id) connected = false ptyId = null + ptyIncarnationId = null handlers.unregisterAll(id) storedCallbacks.onDisconnect?.() } @@ -140,6 +144,7 @@ export function createIpcPtyTransport(opts: IpcPtyTransportOptions = {}): PtyTra } connected = false ptyId = null + ptyIncarnationId = null storedCallbacks = {} }, @@ -188,6 +193,7 @@ export function createIpcPtyTransport(opts: IpcPtyTransportOptions = {}): PtyTra isConnected: () => connected, getPtyId: () => ptyId, + getPtyIncarnationId: () => ptyIncarnationId, getConnectionId: () => connectionId ?? null, getLocalSessionMetadata: () => connectionId diff --git a/src/renderer/src/lib/pty-identity-evidence-store.test.ts b/src/renderer/src/lib/pty-identity-evidence-store.test.ts new file mode 100644 index 00000000000..d2a6ea34555 --- /dev/null +++ b/src/renderer/src/lib/pty-identity-evidence-store.test.ts @@ -0,0 +1,66 @@ +import { describe, expect, it } from 'vitest' +import { createPtyIdentityEvidenceStore } from './pty-identity-evidence-store' + +const evidence = (overrides: Record = {}) => ({ + verdict: 'live' as const, + processName: 'codex', + authorityGeneration: 'g1', + observationEpoch: 1, + capturedAgeMs: 0, + ...overrides +}) + +describe('PTY identity evidence store', () => { + it('rejects non-increasing pushes and rebases serialized age', () => { + let clock = 100 + const store = createPtyIdentityEvidenceStore({ now: () => clock }) + const row = { + hostId: 'ssh:box', + ptyId: 'pty-1', + incarnationId: 'inc-1', + authorityGeneration: 'g1', + observationEpoch: 1, + evidence: evidence({ capturedAgeMs: 4_000 }) + } + expect(store.applyPush(row)).toBe(true) + expect(store.applyPush(row)).toBe(false) + expect(store.get('ssh:box', 'pty-1', 'inc-1')?.receivedAtMs).toBe(-3_900) + clock = 6_000 + expect(store.get('ssh:box', 'pty-1', 'inc-1')?.evidence.verdict).toBe('unverifiable') + }) + + it('fences a late row from a retired authority generation', () => { + const store = createPtyIdentityEvidenceStore({ now: () => 0 }) + const base = { + hostId: 'ssh:box', + ptyId: 'pty-1', + incarnationId: 'inc-1', + observationEpoch: 1, + evidence: evidence() + } + expect(store.applyPush({ ...base, authorityGeneration: 'g1' })).toBe(true) + expect( + store.applyPush({ + ...base, + authorityGeneration: 'g2', + evidence: evidence({ authorityGeneration: 'g2' }) + }) + ).toBe(false) + expect( + store.applyPush({ + ...base, + authorityGeneration: 'g1', + observationEpoch: 2, + evidence: evidence({ observationEpoch: 2 }) + }) + ).toBe(true) + store.activateGeneration('ssh:box', 'g2') + expect( + store.applyPush({ + ...base, + authorityGeneration: 'g2', + evidence: evidence({ authorityGeneration: 'g2' }) + }) + ).toBe(true) + }) +}) diff --git a/src/renderer/src/lib/pty-identity-evidence-store.ts b/src/renderer/src/lib/pty-identity-evidence-store.ts new file mode 100644 index 00000000000..246140f1dcb --- /dev/null +++ b/src/renderer/src/lib/pty-identity-evidence-store.ts @@ -0,0 +1,115 @@ +import type { ForegroundProcessEvidence } from '../../../shared/foreground-process-evidence' + +export type PtyIdentityEvidenceStoreRow = { + hostId: string + ptyId: string + incarnationId: string + authorityGeneration: string + observationEpoch: number + evidence: ForegroundProcessEvidence + receivedAtMs: number + presentationAgent?: string | null +} + +export type PtyIdentityEvidenceStore = ReturnType + +/** Renderer-wide evidence projection; keyed by execution host to fence reconnects. */ +export const ptyIdentityEvidenceStore = createPtyIdentityEvidenceStore() + +export function createPtyIdentityEvidenceStore( + options: { + now?: () => number + freshnessMs?: number + } = {} +) { + const now = options.now ?? (() => performance.now()) + const freshnessMs = options.freshnessMs ?? 5_000 + const rows = new Map() + const generationByHost = new Map() + const epochByHost = new Map() + const key = (hostId: string, ptyId: string, incarnationId: string): string => + `${hostId}\0${ptyId}\0${incarnationId}` + + const apply = ( + row: Omit, + restore = false + ): boolean => { + const currentGeneration = generationByHost.get(row.hostId) + const currentEpoch = epochByHost.get(row.hostId) ?? -1 + if (currentGeneration !== undefined && currentGeneration !== row.authorityGeneration) { + return false + } + if (!restore && row.observationEpoch <= currentEpoch) { + return false + } + generationByHost.set(row.hostId, row.authorityGeneration) + epochByHost.set(row.hostId, Math.max(currentEpoch, row.observationEpoch)) + rows.set(key(row.hostId, row.ptyId, row.incarnationId), { + ...row, + receivedAtMs: now() - row.evidence.capturedAgeMs + }) + return true + } + + return { + activateGeneration: (hostId: string, authorityGeneration: string): void => { + for (const [rowKey, row] of rows) { + if (row.hostId === hostId) { + rows.delete(rowKey) + } + } + generationByHost.set(hostId, authorityGeneration) + epochByHost.set(hostId, -1) + }, + applyPush: (row: Omit): boolean => apply(row), + applySeed: (row: Omit): boolean => + apply(row, true), + get: ( + hostId: string, + ptyId: string, + incarnationId: string + ): PtyIdentityEvidenceStoreRow | null => { + const row = rows.get(key(hostId, ptyId, incarnationId)) + if (!row) { + return null + } + if (now() - row.receivedAtMs > freshnessMs) { + return { + ...row, + evidence: { + authorityGeneration: row.evidence.authorityGeneration, + observationEpoch: row.evidence.observationEpoch, + capturedAgeMs: row.evidence.capturedAgeMs, + verdict: 'unverifiable', + reason: 'stale' + } + } + } + return row + }, + markHostUnverifiable: (hostId: string): void => { + for (const [rowKey, row] of rows) { + if (row.hostId !== hostId) { + continue + } + rows.set(rowKey, { + ...row, + evidence: { ...row.evidence, verdict: 'unverifiable', reason: 'disconnected' } + }) + } + }, + evict: (hostId: string, ptyId: string, incarnationId?: string): void => { + for (const [rowKey, row] of rows) { + if ( + row.hostId === hostId && + row.ptyId === ptyId && + (!incarnationId || row.incarnationId === incarnationId) + ) { + rows.delete(rowKey) + } + } + }, + size: (): number => rows.size, + snapshot: (): PtyIdentityEvidenceStoreRow[] => Array.from(rows.values()) + } +} diff --git a/src/shared/pty-consumer-session-capabilities.ts b/src/shared/pty-consumer-session-capabilities.ts index 89962f6b3b3..b9ea9909025 100644 --- a/src/shared/pty-consumer-session-capabilities.ts +++ b/src/shared/pty-consumer-session-capabilities.ts @@ -18,6 +18,15 @@ export function assertPtyConsumerSessionOptions(options: PtyConsumerSessionOptio ) { throw new Error('outputFlowControl support is invalid') } + if ( + options.identityEvidence && + (options.identityEvidence.versions.length > MAX_CAPABILITY_VERSIONS || + options.identityEvidence.versions.some( + (version) => !Number.isSafeInteger(version) || version <= 0 + )) + ) { + throw new Error('identityEvidence support is invalid') + } if ( options.ownerGraceMs !== undefined && (!Number.isSafeInteger(options.ownerGraceMs) || options.ownerGraceMs < 0) @@ -28,18 +37,29 @@ export function assertPtyConsumerSessionOptions(options: PtyConsumerSessionOptio export function intersectPtyConsumerCapabilities( hello: PtyConsumerSessionHello, - support: PtyConsumerSessionOptions['outputFlowControl'] + support: PtyConsumerSessionOptions['outputFlowControl'], + identitySupport: PtyConsumerSessionOptions['identityEvidence'] = undefined ): Pick { const offer = hello.capabilities?.outputFlowControl - if (!offer || !support || !offer.versions.includes(1) || !support.versions.includes(1)) { + const identityOffer = hello.capabilities?.identityEvidence + const outputSupported = Boolean( + offer && support && offer.versions.includes(1) && support.versions.includes(1) + ) + const identitySupported = Boolean(identityOffer && identitySupport?.versions.includes(1)) + if (!outputSupported && !identitySupported) { return {} } return { capabilities: { - outputFlowControl: { - version: 1, - windowSu: Math.min(offer.requestedWindowSu, support.maxWindowSu) - } + ...(outputSupported + ? { + outputFlowControl: { + version: 1, + windowSu: Math.min(offer!.requestedWindowSu, support!.maxWindowSu) + } + } + : {}), + ...(identitySupported ? { identityEvidence: { version: 1 as const } } : {}) } } } diff --git a/src/shared/pty-consumer-session-contract.ts b/src/shared/pty-consumer-session-contract.ts index 49527f83959..8bd393b81b5 100644 --- a/src/shared/pty-consumer-session-contract.ts +++ b/src/shared/pty-consumer-session-contract.ts @@ -37,6 +37,9 @@ export type PtyConsumerSessionHello = { versions: number[] requestedWindowSu: number } + identityEvidence?: { + versions: number[] + } } } @@ -55,6 +58,9 @@ export type PtyConsumerSessionGrant = { version: 1 windowSu: number } + identityEvidence?: { + version: 1 + } } } @@ -85,6 +91,9 @@ export type PtyConsumerSessionOptions = { versions: readonly number[] maxWindowSu: number } + identityEvidence?: { + versions: readonly number[] + } ownerGraceMs?: number now?: () => number createLease?: () => string diff --git a/src/shared/pty-consumer-session-hello.ts b/src/shared/pty-consumer-session-hello.ts index 9dc54adf30c..40102211a9b 100644 --- a/src/shared/pty-consumer-session-hello.ts +++ b/src/shared/pty-consumer-session-hello.ts @@ -20,17 +20,25 @@ export function validateHello(hello: PtyConsumerSessionHello): void { assertNonEmptyString(hello.resume.ownerLease, 'resume.ownerLease') } const flow = hello.capabilities?.outputFlowControl - if (!flow) { - return + if (flow) { + if ( + !Array.isArray(flow.versions) || + flow.versions.length > MAX_CAPABILITY_VERSIONS || + flow.versions.some((version) => !Number.isSafeInteger(version) || version <= 0) + ) { + throw new Error('outputFlowControl.versions must contain positive safe integers') + } + if (!Number.isSafeInteger(flow.requestedWindowSu) || flow.requestedWindowSu <= 0) { + throw new Error('outputFlowControl.requestedWindowSu must be a positive safe integer') + } } + const identity = hello.capabilities?.identityEvidence if ( - !Array.isArray(flow.versions) || - flow.versions.length > MAX_CAPABILITY_VERSIONS || - flow.versions.some((version) => !Number.isSafeInteger(version) || version <= 0) + identity && + (!Array.isArray(identity.versions) || + identity.versions.length > MAX_CAPABILITY_VERSIONS || + identity.versions.some((version) => !Number.isSafeInteger(version) || version <= 0)) ) { - throw new Error('outputFlowControl.versions must contain positive safe integers') - } - if (!Number.isSafeInteger(flow.requestedWindowSu) || flow.requestedWindowSu <= 0) { - throw new Error('outputFlowControl.requestedWindowSu must be a positive safe integer') + throw new Error('identityEvidence.versions must contain positive safe integers') } } diff --git a/src/shared/pty-consumer-session.ts b/src/shared/pty-consumer-session.ts index 6b040be47cf..f01c645e359 100644 --- a/src/shared/pty-consumer-session.ts +++ b/src/shared/pty-consumer-session.ts @@ -89,7 +89,11 @@ export class PtyConsumerSession { ...(owner ? { ownerGeneration: owner.generation, ownerLease: owner.lease, resumed: owner.resumed } : {}), - ...intersectPtyConsumerCapabilities(hello, this.options.outputFlowControl) + ...intersectPtyConsumerCapabilities( + hello, + this.options.outputFlowControl, + this.options.identityEvidence + ) }) const client: ClientRecord = { principal: authentication.principal, diff --git a/src/shared/pty-identity-evidence.test.ts b/src/shared/pty-identity-evidence.test.ts new file mode 100644 index 00000000000..e2789207186 --- /dev/null +++ b/src/shared/pty-identity-evidence.test.ts @@ -0,0 +1,22 @@ +import { describe, expect, it, vi } from 'vitest' +import { createPtyIdentityBoundaryScanner } from './pty-identity-evidence' + +describe('PTY identity boundary scanner', () => { + it('handles split OSC 133/777 markers and ignores ordinary bytes', () => { + const seen: string[] = [] + const scanner = createPtyIdentityBoundaryScanner((boundary) => seen.push(boundary)) + scanner.feed('prompt\x1b]133;') + scanner.feed('C\x07output\x1b]133;D;0\x1b\\') + scanner.feed('\x1b]777;agent;started\x07') + expect(seen).toEqual(['C', 'D', 'C']) + }) + + it('does not retain a marker after reset', () => { + const onBoundary = vi.fn() + const scanner = createPtyIdentityBoundaryScanner(onBoundary) + scanner.feed('\x1b]133;') + scanner.reset() + scanner.feed('C\x07') + expect(onBoundary).not.toHaveBeenCalled() + }) +}) diff --git a/src/shared/pty-identity-evidence.ts b/src/shared/pty-identity-evidence.ts new file mode 100644 index 00000000000..82c4039c770 --- /dev/null +++ b/src/shared/pty-identity-evidence.ts @@ -0,0 +1,63 @@ +import type { ForegroundProcessEvidence } from './foreground-process-evidence' + +export const PTY_IDENTITY_EVIDENCE_CAPABILITY = 'pty.identityEvidence' as const +export const PTY_IDENTITY_EVIDENCE_VERSION = 1 as const + +export type PtyIdentityEvidenceRow = { + id: string + incarnationId: string + foregroundProcessEvidence: ForegroundProcessEvidence +} + +export type PtyIdentityEvidenceNotification = { + authorityGeneration: string + observationEpoch: number + rows: PtyIdentityEvidenceRow[] + /** Client-side SSH provider generation; absent on the relay wire. */ + providerGeneration?: number +} + +export type PtyIdentityBoundary = 'A' | 'C' | 'D' + +/** Chunk-safe scanner for live shell markers. Replay callers intentionally never feed this. */ +export function createPtyIdentityBoundaryScanner( + onBoundary: (boundary: PtyIdentityBoundary) => void +): { feed: (data: string) => void; reset: () => void } { + let carry = '' + const maxCarry = 256 + const feed = (data: string): void => { + if (!data) { + return + } + const input = carry + data + let cursor = 0 + while (cursor < input.length) { + const osc = input.indexOf('\x1b]', cursor) + if (osc === -1) { + break + } + const endBel = input.indexOf('\x07', osc + 2) + const endSt = input.indexOf('\x1b\\', osc + 2) + const end = endBel === -1 ? endSt : endSt === -1 ? endBel : Math.min(endBel, endSt) + if (end < 0) { + carry = input.slice(osc).slice(-maxCarry) + return + } + const payload = input.slice(osc + 2, end) + const marker = payload.match(/^133;([ACD])(?:;|$)/) + if (marker) { + onBoundary(marker[1] as PtyIdentityBoundary) + } else if (/^777(?:;|$)/.test(payload)) { + onBoundary('C') + } + cursor = end + (input[end] === '\x07' ? 1 : 2) + } + carry = input.slice(Math.max(cursor, input.length - maxCarry)) + } + return { + feed, + reset: () => { + carry = '' + } + } +} diff --git a/src/shared/ssh-types.ts b/src/shared/ssh-types.ts index cbcce254922..ce5e138c752 100644 --- a/src/shared/ssh-types.ts +++ b/src/shared/ssh-types.ts @@ -230,6 +230,9 @@ export type SshPtyConsumerRecovery = { version: 1 windowSu: number } + identityEvidence?: { + version: 1 + } } // ─── Port Forwarding Types ─────────────────────────────────────────