Files
orca/src/relay/relay-pty-source-activation.ts
JinjingandOrcaWin 5f7807497e feat(ssh): bound relay PTY output end to end (#11005)
* docs: design SSH relay PTY backpressure

* fix(ssh): bound relay frame decoding

* fix(relay): bound PTY output publication

* fix(ssh): bound PTY model admission

* fix(ssh): settle closed model admissions

* feat(ssh): negotiate bounded PTY consumer sessions

* fix(ssh): fence exit on renderer settlement

* feat(ssh): track PTY source credit end to end

* fix(ssh): recover bounded PTY output across reconnect

* feat(ssh): complete relay PTY output backpressure

* fix(ssh): close final PTY source credit races

* docs(ssh): record final backpressure validation

* feat(ssh): complete relay PTY source-credit lifecycle

* test(ssh): complete provider notification fixture

* fix(ssh): preserve terminal source credit across rotation

* fix(ssh): fail closed on recovery cancellation

* fix(ssh): prioritize mux control writes after drain

* fix(ssh): retire canceled relay restore deliveries

* fix(ssh): order exit cancellation cleanup

* fix(ssh): gate provisional source activation

* test(ssh): register mux drain-priority coverage

* fix(ssh): type stale owner recovery mismatches

* fix(ssh): close projection replacement races

* fix(relay): contain streaming edge failures

* fix(ssh): secure relay endpoint credentials

* docs(ssh): reconcile final backpressure lifecycle

* fix(ssh): bound main IPC output lifecycle

* fix(ssh): close recovery ownership gaps

* docs(ssh): record exact artifact validation

* fix(ssh): reject reclaimed snapshot replacements

* fix(ssh): fence model admission across reconnect

* fix(ssh): contain migration failure per PTY

* docs(ssh): record final exact-head validation

* test(ssh): align deploy fixtures with credential publication

* feat(ssh): add per-target bounded output setting

* fix(ssh): close source recovery review gaps

* fix(ssh): latch source credit environment override

* feat(ssh): make PTY source credit the default

* docs(ssh): record always-on relay validation

* docs(ssh): bind validation to current main

* test(ssh): grant source credit in IPC fixture

* test(ssh): grant source credit in fake relay

---------

Co-authored-by: OrcaWin <293788423+OrcaWin@users.noreply.github.com>
2026-07-29 17:03:15 -07:00

106 lines
3.5 KiB
TypeScript

import type {
PtySourceRecoveryCheckpoint,
PtySourceRecoveryRequest,
PtySourceRecoveryResult
} from '../shared/pty-source-recovery-contract'
import type { PtySourceReceivingActivation } from '../shared/pty-source-receiving-activation'
import type { PtySourceDeliveryIdentity } from '../shared/pty-source-credit-contract'
import type { RequestContext } from './dispatcher'
import type {
RelayPtySourceDeliveryRecord,
RelayPtySourceSendScheduler
} from './relay-pty-source-send-scheduler'
import type { SshPtyConsumerSessionAdapter } from './ssh-pty-consumer-session-adapter'
export function createPtySourceReceivingActivation(
identity: PtySourceDeliveryIdentity,
checkpointSourceEndSu: number,
recoveryEndSu: number
): PtySourceReceivingActivation {
return Object.freeze({
status: 'pending',
clientGeneration: identity.clientGeneration,
ownerGeneration: identity.ownerGeneration,
ptyIncarnation: identity.ptyIncarnation,
deliveryToken: identity.deliveryToken,
checkpointSourceEndSu,
recoveryEndSu
})
}
export function registerPtySourceActivationSettlement(options: {
id: string
record: RelayPtySourceDeliveryRecord
context: RequestContext
deliveries: Map<string, RelayPtySourceDeliveryRecord>
session: SshPtyConsumerSessionAdapter
sender: RelayPtySourceSendScheduler
onCapacity: (id: string) => void
}): void {
const { id, record, context, deliveries, session, sender, onCapacity } = options
context.onResponseSettled!((result) => {
if (deliveries.get(id) !== record || !record.activating) {
return
}
if (!result.ok) {
if (!record.activationRecoveryRequest) {
session.cancelDelivery(record.identity, 'activation-publication-failed')
deliveries.delete(id)
}
return
}
record.activating = false
sender.completeRecoveryIfReady(record)
sender.pump(record)
onCapacity(id)
})
}
export function registerCanceledPtySourceRetirement(
record: RelayPtySourceDeliveryRecord,
context: RequestContext,
deliveries: Map<string, RelayPtySourceDeliveryRecord>,
onCapacity: (id: string) => void
): void {
record.recoveryCheckpointSourceEndSu = null
record.recoveryEndSu = null
context.onResponseSettled!(() => {
if (deliveries.get(record.identity.id) !== record) {
return
}
deliveries.delete(record.identity.id)
onCapacity(record.identity.id)
})
}
export function pendingPtySourceRecoveryResult(
record: RelayPtySourceDeliveryRecord
): PtySourceRecoveryResult {
if (record.recoveryCheckpointSourceEndSu === null || record.recoveryEndSu === null) {
return Object.freeze({ status: 'restoreRequired', reason: 'checkpointUnavailable' })
}
return Object.freeze({
status: 'pending',
clientGeneration: record.identity.clientGeneration,
ownerGeneration: record.identity.ownerGeneration,
ptyIncarnation: record.identity.ptyIncarnation,
deliveryToken: record.identity.deliveryToken,
checkpointSourceEndSu: record.recoveryCheckpointSourceEndSu,
recoveryEndSu: record.recoveryEndSu
})
}
export function samePtySourceRecoveryRequest(
expected: PtySourceRecoveryCheckpoint,
received: PtySourceRecoveryRequest | undefined
): boolean {
return (
received?.status === 'checkpoint' &&
received.deliveryToken === expected.deliveryToken &&
received.clientGeneration === expected.clientGeneration &&
received.ownerGeneration === expected.ownerGeneration &&
received.ptyIncarnation === expected.ptyIncarnation &&
received.acceptedSourceEndSu === expected.acceptedSourceEndSu
)
}