From 7b0e4f2b99e915b60df7c03b51696cef3e5f7afc Mon Sep 17 00:00:00 2001 From: Neil <4138956+nwparker@users.noreply.github.com> Date: Sun, 30 Aug 2026 16:35:36 -0700 Subject: [PATCH] perf(relay): index PTY source-credit send spans --- src/relay/pty-source-credit-ledger.test.ts | 148 ++++++++++++++++++++- src/relay/pty-source-credit-ledger.ts | 7 +- src/relay/pty-source-credit-record.ts | 12 ++ src/relay/pty-source-credit-settlement.ts | 5 + 4 files changed, 167 insertions(+), 5 deletions(-) diff --git a/src/relay/pty-source-credit-ledger.test.ts b/src/relay/pty-source-credit-ledger.test.ts index 37ad35a8f1e..0bb6a52334f 100644 --- a/src/relay/pty-source-credit-ledger.test.ts +++ b/src/relay/pty-source-credit-ledger.test.ts @@ -1,5 +1,9 @@ import { describe, expect, it } from 'vitest' -import type { PtySourceDeliveryIdentity } from '../shared/pty-source-credit-contract' +import { + ptySourceDeliveryKey, + type PtySourceDeliveryIdentity, + type PtySourceSpan +} from '../shared/pty-source-credit-contract' import { RelayPtySourceCreditLedger } from './pty-source-credit-ledger' import { CLOSED_DELIVERY_TOMBSTONE_LIMIT } from './pty-source-credit-record' @@ -51,7 +55,149 @@ function drainOne( return reservation } +type CursorRecord = { + spans: PtySourceSpan[] + sendSpanIndex: number +} + +function getCursorRecord( + ledger: RelayPtySourceCreditLedger, + owner: PtySourceDeliveryIdentity +): CursorRecord { + const internals = ledger as unknown as { deliveries: Map } + const record = internals.deliveries.get(ptySourceDeliveryKey(owner)) + if (!record) { + throw new Error('test delivery record missing') + } + return record +} + describe('RelayPtySourceCreditLedger', () => { + it('keeps send-span lookup near-linear across a retained-frame burst', () => { + const spanCount = 1_024 + const ledger = new RelayPtySourceCreditLedger({ + maxRetainedSourceSu: spanCount * 2, + maxAggregateRetainedSourceSu: spanCount * 2, + maxRetainedDataBytes: spanCount * 1_024, + maxAggregateRetainedDataBytes: spanCount * 1_024, + maxRetainedSpans: spanCount, + maxAggregateRetainedSpans: spanCount + }) + const owner = identity() + ledger.open(owner, spanCount * 2) + for (let index = 0; index < spanCount; index += 1) { + append(ledger, owner, 'x', `span-${index}`) + } + + const record = getCursorRecord(ledger, owner) + let indexedReads = 0 + const spans = record.spans + record.spans = new Proxy(spans, { + get(target, property, receiver) { + if (typeof property === 'string' && /^\d+$/.test(property)) { + indexedReads += 1 + } + return Reflect.get(target, property, receiver) + } + }) + + let sends = 0 + while (true) { + const reservation = ledger.reserveNextSend(owner, 1) + if (!reservation) { + break + } + ledger.commitSend(reservation) + sends += 1 + } + + const naivePredicateVisits = (spanCount * (spanCount + 1)) / 2 + expect(sends).toBe(spanCount) + // The final span remains the cursor until another source append arrives. + expect(record.sendSpanIndex).toBe(spanCount - 1) + expect(indexedReads).toBe(spanCount * 2 - 1) + // The old Array.find path evaluated one predicate per retained prefix. + expect(naivePredicateVisits).toBe(524_800) + console.log( + `[bench] send-span cursor: indexed span reads ${indexedReads}; naive find predicate visits ${naivePredicateVisits}` + ) + }) + + it('keeps the uncovered-cursor diagnostic when a retained span has a gap', () => { + const ledger = new RelayPtySourceCreditLedger() + const owner = identity('token-gap') + ledger.open(owner, 8) + append(ledger, owner, 'ab', 'span-gap') + const record = getCursorRecord(ledger, owner) + record.spans = [ + Object.freeze({ + ...record.spans[0], + sourceStartSu: 2, + sourceEndSu: 4, + transform: Object.freeze({ ...record.spans[0].transform, rawLengthSu: 2 }) + }) + ] + + expect(() => ledger.reserveNextSend(owner, 2)).toThrow( + 'PTY source delivery cursor is not covered by the retained ledger' + ) + }) + + it('keeps the cursor correct across ACK reclaim, rollback, rotation, and close', () => { + const ledger = new RelayPtySourceCreditLedger() + const oldOwner = identity() + const replacement = identity('token-replacement', { + clientGeneration: 4, + ownerGeneration: 5 + }) + ledger.open(oldOwner, 16) + append(ledger, oldOwner, 'ab', 'span-a') + append(ledger, oldOwner, 'cd', 'span-b') + append(ledger, oldOwner, 'ef', 'span-c') + + const first = ledger.reserveNextSend(oldOwner, 2)! + ledger.commitSend(first) + expect(getCursorRecord(ledger, oldOwner).sendSpanIndex).toBe(0) + ledger.acknowledge(oldOwner, { + id: oldOwner.id, + clientGeneration: oldOwner.clientGeneration, + ownerGeneration: oldOwner.ownerGeneration, + deliveryToken: oldOwner.deliveryToken, + creditedEndSu: 2 + }) + expect(getCursorRecord(ledger, oldOwner).sendSpanIndex).toBe(0) + + const attempted = ledger.reserveNextSend(oldOwner, 1)! + expect(attempted.span.data).toBe('c') + ledger.rollbackSend(attempted) + const retried = ledger.reserveNextSend(oldOwner, 1)! + expect(retried.span.data).toBe('c') + ledger.commitSend(retried) + + ledger.commitSend(ledger.reserveNextSend(oldOwner, 2)!) + const rotation = ledger.rotate(oldOwner, replacement, 2, 16) + expect(rotation.recovery.map((span) => span.data).join('')).toBe('cdef') + expect(getCursorRecord(ledger, replacement).sendSpanIndex).toBe(0) + ledger.commitSend(ledger.reserveNextSend(replacement, 16)!) + ledger.commitSend(ledger.reserveNextSend(replacement, 16)!) + ledger.seal(replacement) + ledger.settleExitPublication(replacement, { ok: true }) + ledger.acknowledge(replacement, { + id: replacement.id, + clientGeneration: replacement.clientGeneration, + ownerGeneration: replacement.ownerGeneration, + deliveryToken: replacement.deliveryToken, + creditedEndSu: 6 + }) + expect(ledger.snapshotIfKnown(replacement)?.state).toBe('closed') + + const canceled = identity('token-canceled') + ledger.open(canceled, 8) + append(ledger, canceled, 'x', 'span-cancel') + ledger.cancel(canceled, 'test-close') + expect(ledger.snapshotIfKnown(canceled)?.state).toBe('closed') + }) + it('never exceeds a token source window across generated send/ACK sequences', () => { for (let seed = 1; seed <= 40; seed++) { const ledger = new RelayPtySourceCreditLedger() diff --git a/src/relay/pty-source-credit-ledger.ts b/src/relay/pty-source-credit-ledger.ts index be628672c91..8999b84bfa7 100644 --- a/src/relay/pty-source-credit-ledger.ts +++ b/src/relay/pty-source-credit-ledger.ts @@ -25,6 +25,7 @@ import { createDeliveryCancellation, createDeliveryRecord, createReplacementDeliveryRecord, + findPtySourceSpanForSend, matchingDeliverySnapshot, MAX_SOURCE_SPAN_DATA_BYTES, ptyOwnerKey, @@ -134,10 +135,8 @@ export class RelayPtySourceCreditLedger { if (remainingWindowSu <= 0 || record.sentEndSu >= record.receivedEndSu) { return null } - const containing = record.spans.find( - (span) => span.sourceStartSu <= record.sentEndSu && span.sourceEndSu > record.sentEndSu - ) - if (!containing) { + const containing = findPtySourceSpanForSend(record) + if (!containing || containing.sourceStartSu > record.sentEndSu) { throw new Error('PTY source delivery cursor is not covered by the retained ledger') } const reservedLengthSu = reservedPtySourceSendLength(record, remainingWindowSu, maxSourceSu) diff --git a/src/relay/pty-source-credit-record.ts b/src/relay/pty-source-credit-record.ts index d2730d25dc7..d07cd2f9354 100644 --- a/src/relay/pty-source-credit-record.ts +++ b/src/relay/pty-source-credit-record.ts @@ -42,6 +42,8 @@ export type DeliveryRecord = { creditedEndSu: number retainedDataBytes: number spans: PtySourceSpan[] + /** First retained span that may contain the next unsent source unit. */ + sendSpanIndex: number sentBoundaries: Set pendingSend: PtySourceSendReservation | null reservedAckEndSu: number | null @@ -68,6 +70,7 @@ export function createDeliveryRecord( creditedEndSu: checkpointSourceEndSu, retainedDataBytes: 0, spans: [], + sendSpanIndex: 0, sentBoundaries: new Set([checkpointSourceEndSu]), pendingSend: null, reservedAckEndSu: null, @@ -160,6 +163,15 @@ export function sliceForSend( }) } +export function findPtySourceSpanForSend(record: DeliveryRecord): PtySourceSpan | undefined { + let span = record.spans[record.sendSpanIndex] + while (span && span.sourceEndSu <= record.sentEndSu) { + record.sendSpanIndex += 1 + span = record.spans[record.sendSpanIndex] + } + return span +} + export function snapshotDeliveryRecord(record: DeliveryRecord): PtySourceDeliverySnapshot { return Object.freeze({ ...record.identity, diff --git a/src/relay/pty-source-credit-settlement.ts b/src/relay/pty-source-credit-settlement.ts index 464ea9cefe1..e12e73521e5 100644 --- a/src/relay/pty-source-credit-settlement.ts +++ b/src/relay/pty-source-credit-settlement.ts @@ -18,8 +18,13 @@ export function reservedPtySourceSendLength( } function reclaimCreditedSpans(record: DeliveryRecord): void { + let removed = 0 while (record.spans[0]?.sourceEndSu <= record.creditedEndSu) { record.retainedDataBytes -= chargedPtyRetainedStringBytes(record.spans.shift()!.data) + removed += 1 + } + if (removed > 0) { + record.sendSpanIndex = Math.max(0, record.sendSpanIndex - removed) } }