mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 16:02:03 +00:00
fix(ssh): preserve identity evidence epochs and admission
This commit is contained in:
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<string, unknown>[] = []
|
||||
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<string, unknown>)
|
||||
}
|
||||
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({
|
||||
|
||||
@@ -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<number, PtyIdentityEvidenceRow[]>()
|
||||
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<PtyProcessSummary[]> {
|
||||
const results: PtyProcessSummary[] = []
|
||||
// Why (SSH-v3 P2 — the host is the authoritative liveness source, so it has to look): this
|
||||
|
||||
@@ -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 {}
|
||||
}
|
||||
|
||||
@@ -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'))
|
||||
|
||||
Reference in New Issue
Block a user