mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 16:02:56 +00:00
* 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>
106 lines
3.5 KiB
TypeScript
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
|
|
)
|
|
}
|