diff --git a/apps/web/src/app/routes/sender.tsx b/apps/web/src/app/routes/sender.tsx index 7af63b9..596517a 100644 --- a/apps/web/src/app/routes/sender.tsx +++ b/apps/web/src/app/routes/sender.tsx @@ -35,6 +35,17 @@ import { normalizeRaptorQRepairPercent, type FecCodec, } from '@raptorqr/core/fec/codec'; +import { + DEFAULT_RAPTORQ_PLAYBACK_STRATEGY, + RAPTORQ_PLAYBACK_STRATEGIES, + createRaptorQPlaybackOrders, + formatRaptorQPlaybackStrategy, + getRaptorQPlaybackWindowPacketIndices, + normalizeRaptorQPlaybackStrategy, + type RaptorQPlaybackOrders, + type RaptorQPlaybackPhase, + type RaptorQPlaybackStrategy, +} from '@raptorqr/core/sender/raptorq_playback'; import { QrWorkerPool } from '@/lib/qr_worker_pool'; // ─── Types ─────────────────────────────────────────────────────────────────── @@ -53,6 +64,9 @@ interface GifResult { interface LiveTransfer { packets: Uint8Array[]; + initialOrder: number[]; + loopOrder: number[]; + activeOrder: RaptorQPlaybackPhase; width: number; height: number; tileWidth: number; @@ -270,6 +284,9 @@ export function SenderPage() { const [qrEncoder, setQrEncoder] = useState(DEFAULT_QR_ENCODER); const [fecCodec, setFecCodec] = useState(DEFAULT_FEC_CODEC); const [raptorqRepairPercent, setRaptorqRepairPercent] = useState(DEFAULT_RAPTORQ_REPAIR_PERCENT); + const [raptorqPlaybackStrategy, setRaptorqPlaybackStrategy] = useState( + DEFAULT_RAPTORQ_PLAYBACK_STRATEGY, + ); const [advancedGenerateOpen, setAdvancedGenerateOpen] = useState(false); const [frameRateFps, setFrameRateFps] = useState(DEFAULT_FRAME_RATE_FPS); const [actualLiveFps, setActualLiveFps] = useState(0); @@ -286,6 +303,7 @@ export function SenderPage() { qrEncoder: QREncoder; fecCodec: FecCodec; raptorqRepairPercent: number; + raptorqPlaybackStrategy: RaptorQPlaybackStrategy | null; } | null>(null); const [error, setError] = useState(''); const [fullscreenActive, setFullscreenActive] = useState(false); @@ -382,19 +400,21 @@ export function SenderPage() { const renderWindowFrames = getRenderWindowDisplayFrames(transfer, frameCacheRef.current); const maxOutstandingPackets = getMaxOutstandingRenderPackets(transfer, frameCacheRef.current); - - pruneFrameCacheForPlaybackWindow( - frameCacheRef.current, + const windowPacketIndexes = getPlaybackWindowPacketIndexes( transfer, startFrameIndex, renderWindowFrames, ); + const windowPacketIndexSet = new Set(windowPacketIndexes); + + pruneFrameCacheForPlaybackWindow( + frameCacheRef.current, + windowPacketIndexSet, + ); pruneRequestedPacketsForPlaybackWindow( requestedPacketImagesRef.current, - transfer, - startFrameIndex, - renderWindowFrames, + windowPacketIndexSet, ); let availableSlots = maxOutstandingPackets @@ -404,34 +424,15 @@ export function SenderPage() { if (availableSlots <= 0) return; const missingPacketIndexes: number[] = []; - const endFrameIndex = startFrameIndex + renderWindowFrames; + for (const packetIndex of windowPacketIndexes) { + if (frameCacheRef.current.frames.has(packetIndex)) continue; + if (requestedPacketImagesRef.current.has(packetIndex)) continue; - for (let frameIndex = startFrameIndex; frameIndex < endFrameIndex; frameIndex++) { - const normalizedFrameIndex = frameIndex % transfer.displayFrameCount; + requestedPacketImagesRef.current.add(packetIndex); + missingPacketIndexes.push(packetIndex); + availableSlots--; - for (let tileIndex = 0; tileIndex < transfer.parallelCount; tileIndex++) { - const packetIndex = getPacketIndexForDisplayFrame( - transfer, - normalizedFrameIndex, - tileIndex, - ); - - if (packetIndex === null) continue; - if (frameCacheRef.current.frames.has(packetIndex)) continue; - if (requestedPacketImagesRef.current.has(packetIndex)) continue; - - requestedPacketImagesRef.current.add(packetIndex); - missingPacketIndexes.push(packetIndex); - availableSlots--; - - if (missingPacketIndexes.length >= QR_RENDER_BATCH_LIMIT || availableSlots <= 0) { - break; - } - } - - if (missingPacketIndexes.length >= QR_RENDER_BATCH_LIMIT || availableSlots <= 0) { - break; - } + if (missingPacketIndexes.length >= QR_RENDER_BATCH_LIMIT || availableSlots <= 0) break; } if (missingPacketIndexes.length === 0) return; @@ -469,13 +470,13 @@ export function SenderPage() { activeTransfer, frameCacheRef.current, ); - - if (!isPacketInPlaybackWindow( + const activeWindowPackets = new Set(getPlaybackWindowPacketIndexes( activeTransfer, - packetIndex, currentFrameIndex, activeRenderWindowFrames, - )) { + )); + + if (!activeWindowPackets.has(packetIndex)) { return; } @@ -486,12 +487,12 @@ export function SenderPage() { return; } - if (!isPacketInPlaybackWindow( + const latestWindowPackets = new Set(getPlaybackWindowPacketIndexes( activeTransfer, - packetIndex, liveFrameIndexRef.current % activeTransfer.displayFrameCount, activeRenderWindowFrames, - )) { + )); + if (!latestWindowPackets.has(packetIndex)) { closeRenderedTile(renderedTile); return; } @@ -567,8 +568,20 @@ export function SenderPage() { recordLiveFrameDraw(); - requestRenderWindow(transfer, frameIndex + 1); liveFrameIndexRef.current = (frameIndex + 1) % transfer.displayFrameCount; + + // Balanced uses a source-first first pass for Live, then switches only + // after the final initial display frame was actually drawn. This keeps + // the render cache keyed by canonical packet index while the playback + // order changes underneath it. + if ( + transfer.activeOrder === 'initial' && + frameIndex === transfer.displayFrameCount - 1 + ) { + transfer.activeOrder = 'loop'; + } + + requestRenderWindow(transfer, liveFrameIndexRef.current); advanced++; } @@ -738,6 +751,11 @@ export function SenderPage() { resetOutput(); }, [resetOutput]); + const handleRaptorQPlaybackStrategyChange = useCallback((value: string) => { + setRaptorqPlaybackStrategy(normalizeRaptorQPlaybackStrategy(value)); + resetOutput(); + }, [resetOutput]); + const handleFrameRateChange = useCallback((value: string) => { const nextFrameRate = clampFrameRate(Number(value)); frameRateFpsRef.current = nextFrameRate; @@ -788,6 +806,8 @@ export function SenderPage() { const encoded = await new Promise<{ packets: Uint8Array[]; + sourcePacketIndices?: number[]; + repairPacketIndices?: number[]; totalGenerations: number; stats: { originalSize: number; preprocessedSize: number; frameCount: number }; }>((resolve, reject) => { @@ -827,12 +847,18 @@ export function SenderPage() { qrEncoder, fecCodec, raptorqRepairPercent, + raptorqPlaybackStrategy: fecCodec === 'wasm-raptorq' + ? raptorqPlaybackStrategy + : null, }); const nextLiveTransfer = createLiveTransfer( encoded.packets, selectedQRProfile, qrEncoder, parallelQRCount, + fecCodec === 'wasm-raptorq' ? raptorqPlaybackStrategy : null, + encoded.sourcePacketIndices, + encoded.repairPacketIndices, ); setLiveTransfer(nextLiveTransfer); setStatus( @@ -853,7 +879,19 @@ export function SenderPage() { setEncodingLive(false); } } - }, [mode, text, file, resetOutput, qrVersion, eccLevel, qrEncoder, fecCodec, raptorqRepairPercent, parallelQRCount]); + }, [ + mode, + text, + file, + resetOutput, + qrVersion, + eccLevel, + qrEncoder, + fecCodec, + raptorqRepairPercent, + raptorqPlaybackStrategy, + parallelQRCount, + ]); const handlePrepareGif = useCallback(async () => { const transfer = liveTransferRef.current; @@ -907,6 +945,7 @@ export function SenderPage() { eccLevel: transfer.eccLevel, qrEncoder: transfer.qrEncoder, parallelCount: transfer.parallelCount, + packetOrder: transfer.loopOrder, }, ); }); @@ -1103,26 +1142,46 @@ export function SenderPage() { {fecCodec === 'wasm-raptorq' && ( - + <> + + + )} )} @@ -1260,6 +1319,12 @@ export function SenderPage() { ? `${formatFecCodec(stats.fecCodec)} · ${stats.raptorqRepairPercent}% repair` : formatFecCodec(stats.fecCodec)} + RaptorQ playback + + {stats.raptorqPlaybackStrategy + ? formatRaptorQPlaybackStrategy(stats.raptorqPlaybackStrategy) + : 'Not applicable (JS RLNC)'} + QR packets {stats.frameCount} Parallel QR @@ -1394,6 +1459,9 @@ function createLiveTransfer( profile: QRTransferProfile, qrEncoder: QREncoder, parallelCount: ParallelQRCount, + raptorqStrategy: RaptorQPlaybackStrategy | null, + sourcePacketIndices?: number[], + repairPacketIndices?: number[], ): LiveTransfer { if (packets.length === 0) { throw new Error('No QR packets were generated.'); @@ -1404,9 +1472,28 @@ function createLiveTransfer( const scale = Math.max(2, Math.round(LIVE_TARGET_PX / totalModules)); const tileSize = totalModules * scale; const layout = getParallelLayout(parallelCount); + const packetIndexes = Array.from({ length: packets.length }, (_, index) => index); + let playbackOrders: RaptorQPlaybackOrders; + if (raptorqStrategy) { + if (!sourcePacketIndices || !repairPacketIndices) { + throw new Error('RaptorQ packet classification metadata is missing.'); + } + playbackOrders = createRaptorQPlaybackOrders( + sourcePacketIndices, + repairPacketIndices, + raptorqStrategy, + ); + } else { + playbackOrders = { + initialOrder: packetIndexes, + loopOrder: packetIndexes, + }; + } return { packets, + ...playbackOrders, + activeOrder: 'initial', width: tileSize * layout.columns, height: tileSize * layout.rows, tileWidth: tileSize, @@ -1498,16 +1585,10 @@ function getMaxOutstandingRenderPackets(transfer: LiveTransfer, cache: FrameCach function pruneFrameCacheForPlaybackWindow( cache: FrameCache, - transfer: LiveTransfer, - startFrameIndex: number, - renderWindowFrames = getRenderWindowDisplayFrames(transfer, cache), + windowPacketIndexes: ReadonlySet, ): void { - if (transfer.displayFrameCount <= renderWindowFrames) return; - for (const [packetIndex, tile] of cache.frames) { - if (isPacketInPlaybackWindow(transfer, packetIndex, startFrameIndex, renderWindowFrames)) { - continue; - } + if (windowPacketIndexes.has(packetIndex)) continue; closeRenderedTile(tile); cache.frames.delete(packetIndex); } @@ -1515,41 +1596,26 @@ function pruneFrameCacheForPlaybackWindow( function pruneRequestedPacketsForPlaybackWindow( requestedPackets: Set, - transfer: LiveTransfer, - startFrameIndex: number, - renderWindowFrames: number, + windowPacketIndexes: ReadonlySet, ): void { - if (transfer.displayFrameCount <= renderWindowFrames) return; - for (const packetIndex of requestedPackets) { - if (isPacketInPlaybackWindow(transfer, packetIndex, startFrameIndex, renderWindowFrames)) { - continue; - } + if (windowPacketIndexes.has(packetIndex)) continue; requestedPackets.delete(packetIndex); } } -function isPacketInPlaybackWindow( +function getPlaybackWindowPacketIndexes( transfer: LiveTransfer, - packetIndex: number, startFrameIndex: number, renderWindowFrames: number, -): boolean { - if (packetIndex < 0 || packetIndex >= transfer.packets.length) return false; - if (transfer.displayFrameCount <= renderWindowFrames) return true; - - const packetFrameIndex = Math.floor(packetIndex / transfer.parallelCount); - const normalizedStartFrame = normalizeFrameIndex(startFrameIndex, transfer.displayFrameCount); - const distance = ( - packetFrameIndex - - normalizedStartFrame - + transfer.displayFrameCount - ) % transfer.displayFrameCount; - return distance < renderWindowFrames; -} - -function normalizeFrameIndex(frameIndex: number, frameCount: number): number { - return ((frameIndex % frameCount) + frameCount) % frameCount; +): number[] { + return getRaptorQPlaybackWindowPacketIndices( + transfer, + transfer.activeOrder, + transfer.parallelCount, + startFrameIndex, + Math.min(renderWindowFrames, transfer.displayFrameCount), + ); } function getParallelLayout(parallelCount: ParallelQRCount): { columns: number; rows: number } { @@ -1565,10 +1631,16 @@ function getPacketIndexForDisplayFrame( frameIndex: number, tileIndex: number, ): number | null { - return stripedPacketIndex( - transfer.packets.length, + const orderedPosition = stripedPacketIndex( + transfer.initialOrder.length, transfer.parallelCount, frameIndex, tileIndex, ); + if (orderedPosition === null) return null; + + const order = transfer.activeOrder === 'initial' + ? transfer.initialOrder + : transfer.loopOrder; + return order[orderedPosition] ?? null; } diff --git a/apps/web/src/tests/gif_roundtrip.test.ts b/apps/web/src/tests/gif_roundtrip.test.ts index d7ddafb..7bdf005 100644 --- a/apps/web/src/tests/gif_roundtrip.test.ts +++ b/apps/web/src/tests/gif_roundtrip.test.ts @@ -11,6 +11,7 @@ import { decodeQRFromCanvas } from '@raptorqr/core/qr/qr_decode'; import { parsePacket } from '@raptorqr/core/protocol/packet'; import { MAX_PAYLOAD_SIZE, QR_VERSION, ECC_LEVEL, FRAME_DELAY_MS } from '@raptorqr/core/protocol/constants'; import { packetizeRaptorQ } from '@raptorqr/core/sender/raptorq_packetizer'; +import { createRaptorQPlaybackOrders } from '@raptorqr/core/sender/raptorq_playback'; describe('GIF Roundtrip', () => { it('should encode and decode a GIF', async () => { @@ -26,7 +27,13 @@ describe('GIF Roundtrip', () => { repairPercent: DEFAULT_RAPTORQ_REPAIR_PERCENT, }, ); - const frames = result.packets; + const { loopOrder } = createRaptorQPlaybackOrders( + result.sourcePacketIndices, + result.repairPacketIndices, + 'balanced', + ); + const frames = loopOrder.map((packetIndex) => result.packets[packetIndex]!); + expect(loopOrder).toHaveLength(result.packets.length); // Generate QR image frames const imageFrames: Uint8Array[] = []; diff --git a/apps/web/src/workers/encode.worker.ts b/apps/web/src/workers/encode.worker.ts index 1628b04..419622d 100644 --- a/apps/web/src/workers/encode.worker.ts +++ b/apps/web/src/workers/encode.worker.ts @@ -34,6 +34,8 @@ interface EncodeInput { interface EncodeOutput { type: 'encoded'; packets: Uint8Array[]; + sourcePacketIndices?: number[]; + repairPacketIndices?: number[]; totalGenerations: number; stats: { originalSize: number; @@ -88,6 +90,8 @@ async function handleEncode(input: EncodeInput): Promise { return { type: 'encoded', packets: result.packets, + sourcePacketIndices: result.sourcePacketIndices, + repairPacketIndices: result.repairPacketIndices, totalGenerations: result.totalGenerations, stats: { originalSize: originalBytes.length, diff --git a/apps/web/src/workers/gif.worker.ts b/apps/web/src/workers/gif.worker.ts index 584bd58..8f68e3b 100644 --- a/apps/web/src/workers/gif.worker.ts +++ b/apps/web/src/workers/gif.worker.ts @@ -15,7 +15,7 @@ import { createQRGif } from '@raptorqr/core/gif/gif_render'; import { QR_VERSION, ECC_LEVEL, FRAME_DELAY_MS } from '@raptorqr/core/protocol/constants'; import { stripedFrameCount, - stripedPacketIndex, + stripedOrderedPacketIndex, type ParallelQRCount, } from '@raptorqr/core/sender/parallel_striping'; @@ -24,6 +24,8 @@ import { interface GenerateInput { type: 'generate'; packets: Uint8Array[]; + /** Canonical packet indexes in playback order. */ + packetOrder?: number[]; frameDelayMs?: number; qrVersion?: number; eccLevel?: EccLevel; @@ -62,6 +64,7 @@ self.onmessage = (e: MessageEvent) => { async function handleGenerate(input: GenerateInput): Promise { const { packets } = input; + const packetOrder = normalizePacketOrder(input.packetOrder, packets.length); const frameDelayMs = normalizeFrameDelayMs(input.frameDelayMs); const qrVersion = normalizeQRVersion(input.qrVersion); const eccLevel = normalizeEccLevel(input.eccLevel); @@ -82,14 +85,14 @@ async function handleGenerate(input: GenerateInput): Promise { const frames: Uint8Array[] = []; const width = tileSize * layout.columns; const height = tileSize * layout.rows; - const frameCount = stripedFrameCount(packets.length, parallelCount); + const frameCount = stripedFrameCount(packetOrder.length, parallelCount); for (let frameIndex = 0; frameIndex < frameCount; frameIndex++) { const composite = new Uint8ClampedArray(width * height * 4); composite.fill(255); for (let tileIndex = 0; tileIndex < parallelCount; tileIndex++) { - const packetIndex = stripedPacketIndex(packets.length, parallelCount, frameIndex, tileIndex); + const packetIndex = stripedOrderedPacketIndex(packetOrder, parallelCount, frameIndex, tileIndex); if (packetIndex === null) continue; const imageData = await renderQRCodeImageData( packets[packetIndex]!, @@ -139,6 +142,25 @@ function normalizeParallelQRCount(value: number | undefined): ParallelQRCount { return value === 1 || value === 2 || value === 4 || value === 6 || value === 8 ? value : 4; } +function normalizePacketOrder(order: number[] | undefined, packetCount: number): number[] { + if (!order) return Array.from({ length: packetCount }, (_, index) => index); + if (order.length !== packetCount) { + throw new RangeError(`Invalid GIF packet order length: ${order.length}, expected ${packetCount}`); + } + + const seen = new Set(); + for (const packetIndex of order) { + if (!Number.isInteger(packetIndex) || packetIndex < 0 || packetIndex >= packetCount) { + throw new RangeError(`Invalid GIF packet index: ${packetIndex}`); + } + if (seen.has(packetIndex)) { + throw new RangeError(`Duplicate GIF packet index: ${packetIndex}`); + } + seen.add(packetIndex); + } + return order; +} + function getParallelLayout(parallelCount: ParallelQRCount): { columns: number; rows: number } { if (parallelCount === 1) return { columns: 1, rows: 1 }; if (parallelCount === 2) return { columns: 2, rows: 1 }; diff --git a/benchmark.final.test.ts b/benchmark.final.test.ts index 46764d2..010f6e9 100644 --- a/benchmark.final.test.ts +++ b/benchmark.final.test.ts @@ -35,6 +35,7 @@ import { packetCodec, parsePacket } from '@raptorqr/core/protocol/packet'; import { decodeQRCodesFromCanvas } from '@raptorqr/core/qr/qr_decode'; import { renderQRCodeImageData } from '@raptorqr/core/qr/qr_encoder_browser'; import { packetizeRaptorQ } from '@raptorqr/core/sender/raptorq_packetizer'; +import { createRaptorQPlaybackOrders } from '@raptorqr/core/sender/raptorq_playback'; const QR_VERSION = 30; const ECC_LEVEL = 'L'; @@ -161,6 +162,11 @@ describe('final V30 4-way transfer benchmark', () => { expect(packetized.packets).toHaveLength(TOTAL_QR_SYMBOLS); expect(packetized.symbolSize).toBe(profile.maxPayloadSize); expect(packetized.dataLength).toBe(payload.length); + const { loopOrder } = createRaptorQPlaybackOrders( + packetized.sourcePacketIndices, + packetized.repairPacketIndices, + 'balanced', + ); const decoder = await RaptorQWasmDecoder.create(packetized.dataLength, profile.maxPayloadSize); const seenPayloadIds = new Set(); @@ -182,7 +188,7 @@ describe('final V30 4-way transfer benchmark', () => { const symbolIndex = displayFrame * PARALLEL_QR_COUNT + tileIndex; if (dropSet.has(symbolIndex)) continue; - const packet = packetized.packets[symbolIndex]!; + const packet = packetized.packets[loopOrder[symbolIndex]!]!; const qrImage = await renderQRCodeImageData( packet, QR_VERSION, diff --git a/benchmark.md b/benchmark.md index 3b15475..49b57f0 100644 --- a/benchmark.md +++ b/benchmark.md @@ -22,6 +22,7 @@ The final benchmark is implemented in `benchmark.final.test.ts` and exercises the end-to-end transfer stack: - RaptorQ packetization through `@raptorqr/raptorq-wasm` +- default Balanced loop ordering with repair packets evenly interleaved - QR rendering through `@raptorqr/fast-qr-wasm` - 4 QR symbols per display frame in a 2x2 composite image - QR parsing through `decodeQRCodesFromCanvas` diff --git a/benchmark.test.ts b/benchmark.test.ts index 0a6944c..a55b66b 100644 --- a/benchmark.test.ts +++ b/benchmark.test.ts @@ -50,6 +50,7 @@ import type { Packet as ParsedPacket } from '@raptorqr/core/protocol/packet'; import { decodeQRFromCanvas } from '@raptorqr/core/qr/qr_decode'; import { renderQRCodeImageData } from '@raptorqr/core/qr/qr_encoder_browser'; import { packetizeRaptorQ } from '@raptorqr/core/sender/raptorq_packetizer'; +import { createRaptorQPlaybackOrders } from '@raptorqr/core/sender/raptorq_playback'; interface BenchConfig { payloadBytes: number; @@ -212,7 +213,13 @@ async function runFullChainOnce(iteration: number, payload: Uint8Array, config: repairPercent: DEFAULT_RAPTORQ_REPAIR_PERCENT, }, )); - const parsedOriginalPackets = packetized.packets.map((packet) => parsePacket(packet)); + const { loopOrder } = createRaptorQPlaybackOrders( + packetized.sourcePacketIndices, + packetized.repairPacketIndices, + 'balanced', + ); + const scheduledPackets = loopOrder.map((packetIndex) => packetized.packets[packetIndex]!); + const parsedOriginalPackets = scheduledPackets.map((packet) => parsePacket(packet)); expect(packetized.dataLength).toBe(payload.length); expect(packetized.symbolSize).toBe(MAX_PAYLOAD_SIZE); @@ -220,7 +227,7 @@ async function runFullChainOnce(iteration: number, payload: Uint8Array, config: expect(parsedOriginalPackets.every((packet) => packetCodec(packet.header) === 'wasm-raptorq')).toBe(true); expect(parsedOriginalPackets.every((packet) => packet.header.symbolIndex === RAPTORQ_SYMBOL_INDEX)).toBe(true); - const [images, qrRenderMs] = await timed(() => renderFrames(packetized.packets, config.scale)); + const [images, qrRenderMs] = await timed(() => renderFrames(scheduledPackets, config.scale)); const [decodedPackets, qrDecodeMs] = await timed(() => decodeFrames(images)); expect(decodedPackets).toHaveLength(parsedOriginalPackets.length); diff --git a/packages/raptorqr-core/README.md b/packages/raptorqr-core/README.md index 1a8adff..cb9e95b 100644 --- a/packages/raptorqr-core/README.md +++ b/packages/raptorqr-core/README.md @@ -65,6 +65,8 @@ Important result fields: ```ts result.packets; // Uint8Array[] ready for QR rendering +result.sourcePacketIndices; // canonical indexes of systematic source packets +result.repairPacketIndices; // canonical indexes of RaptorQ repair packets result.totalGenerations; // RaptorQ packet count in this path result.sourceGenerations;// estimated source packet count result.dataLength; // preprocessed payload length @@ -72,6 +74,32 @@ result.isCompressed; // whether deflate-raw was applied result.symbolSize; // transport payload size used by the RaptorQ codec ``` +## Schedule RaptorQ Playback + +Playback scheduling changes packet order without copying or modifying packet bytes. + +```ts +import { + createRaptorQPlaybackOrders, + type RaptorQPlaybackStrategy, +} from '@raptorqr/core/sender/raptorq_playback'; + +const strategy: RaptorQPlaybackStrategy = 'balanced'; +const orders = createRaptorQPlaybackOrders( + result.sourcePacketIndices, + result.repairPacketIndices, + strategy, +); + +const firstLivePass = orders.initialOrder.map((index) => result.packets[index]!); +const repeatedPass = orders.loopOrder.map((index) => result.packets[index]!); +``` + +`fast-start` uses source-first ordering for every pass. `balanced` uses a +source-first initial pass and evenly interleaved repeated passes. `even-spread` +interleaves source and repair packets from the first pass. The interleave ratio +is derived from the packets actually generated by the configured repair rate. + ## Decode RaptorQ Packets Use the packet parser to validate the transport wrapper, then pass each RaptorQ payload to the decoder. diff --git a/packages/raptorqr-core/src/index.ts b/packages/raptorqr-core/src/index.ts index 70270db..0bc800d 100644 --- a/packages/raptorqr-core/src/index.ts +++ b/packages/raptorqr-core/src/index.ts @@ -12,3 +12,4 @@ export * from './sender/legacy_rlnc'; export * from './sender/parallel_striping'; export * from './sender/preprocess_payload'; export * from './sender/raptorq_packetizer'; +export * from './sender/raptorq_playback'; diff --git a/packages/raptorqr-core/src/sender/parallel_striping.ts b/packages/raptorqr-core/src/sender/parallel_striping.ts index 08a91db..15335ca 100644 --- a/packages/raptorqr-core/src/sender/parallel_striping.ts +++ b/packages/raptorqr-core/src/sender/parallel_striping.ts @@ -23,3 +23,22 @@ export function stripedPacketIndex( const packetIndex = frameIndex * parallelCount + tileIndex; return packetIndex < packetCount ? packetIndex : null; } + +/** + * Resolve a display tile through a canonical packet order before striping. + * This keeps parallel QR grouping independent from packet classification. + */ +export function stripedOrderedPacketIndex( + packetOrder: readonly number[], + parallelCount: ParallelQRCount, + frameIndex: number, + tileIndex: number, +): number | null { + const orderedPosition = stripedPacketIndex( + packetOrder.length, + parallelCount, + frameIndex, + tileIndex, + ); + return orderedPosition === null ? null : packetOrder[orderedPosition] ?? null; +} diff --git a/packages/raptorqr-core/src/sender/raptorq_packetizer.ts b/packages/raptorqr-core/src/sender/raptorq_packetizer.ts index 3344003..5de3fdf 100644 --- a/packages/raptorqr-core/src/sender/raptorq_packetizer.ts +++ b/packages/raptorqr-core/src/sender/raptorq_packetizer.ts @@ -6,8 +6,18 @@ import { type PreprocessResult, } from '@raptorqr/core/sender/preprocess_payload'; +const RAPTORQ_PAYLOAD_ID_BYTES = 4; +const RAPTORQ_MAX_SOURCE_SYMBOLS_PER_BLOCK = 56_403; + +export interface RaptorQPacketIndexMetadata { + sourcePacketIndices: number[]; + repairPacketIndices: number[]; +} + export interface RaptorQPacketizerResult { packets: Uint8Array[]; + sourcePacketIndices: number[]; + repairPacketIndices: number[]; totalGenerations: number; sourceGenerations: number; dataLength: number; @@ -51,6 +61,11 @@ export function buildRaptorQTransportPackets( symbolSize: number, ): RaptorQPacketizerResult { const totalPackets = serializedPackets.length; + const packetIndexMetadata = classifyRaptorQPackets( + serializedPackets, + preprocessed.dataLength, + symbolSize, + ); const packets = serializedPackets.map((payload, index) => { const header: PacketHeader = { generationIndex: 0, @@ -71,6 +86,7 @@ export function buildRaptorQTransportPackets( return { packets, + ...packetIndexMetadata, totalGenerations: totalPackets, sourceGenerations, dataLength: preprocessed.dataLength, @@ -79,3 +95,70 @@ export function buildRaptorQTransportPackets( symbolSize, }; } + +/** + * Classify raw RaptorQ codec packets from their Payload IDs. + * + * Payload ID is four bytes: one source block number followed by a 24-bit ESI. + * The source symbol count for each block is derived from the same RFC 6330 + * source-block geometry used by the WASM wrapper. This intentionally does not + * assume that the packet array has one source prefix or copy packet payloads. + */ +export function classifyRaptorQPackets( + serializedPackets: readonly Uint8Array[], + dataLength: number, + maxTransportPayloadSize: number, +): RaptorQPacketIndexMetadata { + const sourceSymbolSize = maxTransportPayloadSize - RAPTORQ_PAYLOAD_ID_BYTES; + if (!Number.isInteger(sourceSymbolSize) || sourceSymbolSize <= 0) { + throw new RangeError(`Invalid RaptorQ transport payload size: ${maxTransportPayloadSize}`); + } + if (!Number.isInteger(dataLength) || dataLength < 0) { + throw new RangeError(`Invalid RaptorQ data length: ${dataLength}`); + } + + const totalSourceSymbols = Math.max(1, Math.ceil(dataLength / sourceSymbolSize)); + const sourceBlockCount = Math.max( + 1, + Math.ceil(totalSourceSymbols / RAPTORQ_MAX_SOURCE_SYMBOLS_PER_BLOCK), + ); + const largestBlockSourceCount = Math.ceil(totalSourceSymbols / sourceBlockCount); + const smallestBlockSourceCount = largestBlockSourceCount - 1; + const largerBlockCount = totalSourceSymbols - smallestBlockSourceCount * sourceBlockCount; + + const sourcePacketIndices: number[] = []; + const repairPacketIndices: number[] = []; + + serializedPackets.forEach((payload, packetIndex) => { + if (payload.length < RAPTORQ_PAYLOAD_ID_BYTES) { + throw new Error( + `RaptorQ packet ${packetIndex} is too short for a ${RAPTORQ_PAYLOAD_ID_BYTES}-byte Payload ID`, + ); + } + + const sourceBlockNumber = payload[0]!; + if (sourceBlockNumber >= sourceBlockCount) { + throw new Error( + `RaptorQ packet ${packetIndex} references source block ${sourceBlockNumber}, ` + + `but the transfer has ${sourceBlockCount} source blocks`, + ); + } + + const encodingSymbolId = ( + (payload[1]! << 16) | + (payload[2]! << 8) | + payload[3]! + ) >>> 0; + const sourceCountForBlock = sourceBlockNumber < largerBlockCount + ? largestBlockSourceCount + : smallestBlockSourceCount; + + if (encodingSymbolId < sourceCountForBlock) { + sourcePacketIndices.push(packetIndex); + } else { + repairPacketIndices.push(packetIndex); + } + }); + + return { sourcePacketIndices, repairPacketIndices }; +} diff --git a/packages/raptorqr-core/src/sender/raptorq_playback.ts b/packages/raptorqr-core/src/sender/raptorq_playback.ts new file mode 100644 index 0000000..02820f2 --- /dev/null +++ b/packages/raptorqr-core/src/sender/raptorq_playback.ts @@ -0,0 +1,198 @@ +/** + * RaptorQ playback strategy and packet-index scheduling. + * + * The scheduler only moves canonical packet indexes. Packet bytes stay in the + * packetizer's canonical array and are never copied or modified here. + */ + +export type RaptorQPlaybackStrategy = 'fast-start' | 'balanced' | 'even-spread'; + +export const DEFAULT_RAPTORQ_PLAYBACK_STRATEGY: RaptorQPlaybackStrategy = 'balanced'; + +export const RAPTORQ_PLAYBACK_STRATEGIES: readonly RaptorQPlaybackStrategy[] = [ + 'fast-start', + 'balanced', + 'even-spread', +]; + +export interface RaptorQPlaybackOrders { + initialOrder: number[]; + loopOrder: number[]; +} + +export type RaptorQPlaybackPhase = 'initial' | 'loop'; + +/** Normalize persisted or untrusted strategy values to the stable default. */ +export function normalizeRaptorQPlaybackStrategy(value: unknown): RaptorQPlaybackStrategy { + if (value === 'fast-start' || value === 'balanced' || value === 'even-spread') { + return value; + } + return DEFAULT_RAPTORQ_PLAYBACK_STRATEGY; +} + +export function formatRaptorQPlaybackStrategy(value: RaptorQPlaybackStrategy): string { + switch (value) { + case 'fast-start': + return 'Fast start'; + case 'even-spread': + return 'Even spread'; + case 'balanced': + default: + return 'Balanced'; + } +} + +/** + * Build a source-first order from the packetizer's classification metadata. + */ +export function createSourceFirstPacketIndexOrder( + sourcePacketIndices: readonly number[], + repairPacketIndices: readonly number[], +): number[] { + return validateAndCopyPacketIndexes(sourcePacketIndices, repairPacketIndices); +} + +/** + * Build an even source/repair order using cumulative proportions. + * + * After each source packet, R is accumulated. Once the accumulated value + * reaches S, a repair packet is emitted and S is subtracted. This gives + * source gaps that differ by at most one packet while preserving every input + * index exactly once. + */ +export function createEvenlyInterleavedPacketIndexOrder( + sourcePacketIndices: readonly number[], + repairPacketIndices: readonly number[], +): number[] { + const source = [...sourcePacketIndices]; + const repair = [...repairPacketIndices]; + validatePacketIndexes(source, repair); + + if (source.length === 0) return repair; + if (repair.length === 0) return source; + + const order: number[] = []; + let sourceIndex = 0; + let repairIndex = 0; + let accumulatedRepair = 0; + + while (sourceIndex < source.length) { + order.push(source[sourceIndex++]!); + accumulatedRepair += repair.length; + + while (repairIndex < repair.length && accumulatedRepair >= source.length) { + order.push(repair[repairIndex++]!); + accumulatedRepair -= source.length; + } + } + + // The normal RaptorQ repair range is 0..100%, but keep this total and + // deterministic for callers using synthetic metadata too. + while (repairIndex < repair.length) { + order.push(repair[repairIndex++]!); + } + + return order; +} + +/** + * Create the initial and loop orders for a Live/GIF transfer. + * + * GIFs always use loopOrder. Live playback starts with initialOrder and may + * switch to loopOrder after its first complete display cycle. + */ +export function createRaptorQPlaybackOrders( + sourcePacketIndices: readonly number[], + repairPacketIndices: readonly number[], + strategy: RaptorQPlaybackStrategy = DEFAULT_RAPTORQ_PLAYBACK_STRATEGY, +): RaptorQPlaybackOrders { + const normalized = normalizeRaptorQPlaybackStrategy(strategy); + const initialSourceFirst = createSourceFirstPacketIndexOrder( + sourcePacketIndices, + repairPacketIndices, + ); + const evenSpread = createEvenlyInterleavedPacketIndexOrder( + sourcePacketIndices, + repairPacketIndices, + ); + + const initialOrder = normalized === 'even-spread' ? evenSpread : initialSourceFirst; + const loopOrder = normalized === 'fast-start' ? initialSourceFirst : evenSpread; + + return { + initialOrder, + loopOrder, + }; +} + +/** + * Return the canonical packets needed by a render window in display order. + * An initial window that crosses the first-cycle boundary uses loopOrder for + * its wrapped portion, so Balanced can pre-render its first repair-spread loop. + */ +export function getRaptorQPlaybackWindowPacketIndices( + orders: Pick, + activePhase: RaptorQPlaybackPhase, + parallelCount: number, + startFrameIndex: number, + windowFrameCount: number, +): number[] { + if (!Number.isInteger(parallelCount) || parallelCount <= 0) { + throw new RangeError(`Invalid parallel QR count: ${parallelCount}`); + } + if (!Number.isInteger(startFrameIndex) || startFrameIndex < 0) { + throw new RangeError(`Invalid playback start frame: ${startFrameIndex}`); + } + if (!Number.isInteger(windowFrameCount) || windowFrameCount < 0) { + throw new RangeError(`Invalid playback window frame count: ${windowFrameCount}`); + } + if (orders.initialOrder.length !== orders.loopOrder.length) { + throw new RangeError('RaptorQ initial and loop orders must have the same length.'); + } + if (orders.initialOrder.length === 0 || windowFrameCount === 0) return []; + + const displayFrameCount = Math.ceil(orders.initialOrder.length / parallelCount); + const packetIndices: number[] = []; + const seen = new Set(); + + for (let offset = 0; offset < windowFrameCount; offset++) { + const absoluteFrameIndex = startFrameIndex + offset; + const useLoopOrder = activePhase === 'loop' || absoluteFrameIndex >= displayFrameCount; + const order = useLoopOrder ? orders.loopOrder : orders.initialOrder; + const normalizedFrameIndex = absoluteFrameIndex % displayFrameCount; + const firstPacketPosition = normalizedFrameIndex * parallelCount; + + for (let tileIndex = 0; tileIndex < parallelCount; tileIndex++) { + const packetIndex = order[firstPacketPosition + tileIndex]; + if (packetIndex === undefined || seen.has(packetIndex)) continue; + seen.add(packetIndex); + packetIndices.push(packetIndex); + } + } + + return packetIndices; +} + +function validateAndCopyPacketIndexes( + sourcePacketIndices: readonly number[], + repairPacketIndices: readonly number[], +): number[] { + validatePacketIndexes(sourcePacketIndices, repairPacketIndices); + return [...sourcePacketIndices, ...repairPacketIndices]; +} + +function validatePacketIndexes( + sourcePacketIndices: readonly number[], + repairPacketIndices: readonly number[], +): void { + const seen = new Set(); + for (const packetIndex of [...sourcePacketIndices, ...repairPacketIndices]) { + if (!Number.isInteger(packetIndex) || packetIndex < 0) { + throw new RangeError(`Invalid RaptorQ packet index: ${packetIndex}`); + } + if (seen.has(packetIndex)) { + throw new RangeError(`Duplicate RaptorQ packet index: ${packetIndex}`); + } + seen.add(packetIndex); + } +} diff --git a/packages/raptorqr-core/src/tests/parallel_striping.test.ts b/packages/raptorqr-core/src/tests/parallel_striping.test.ts new file mode 100644 index 0000000..d3b209b --- /dev/null +++ b/packages/raptorqr-core/src/tests/parallel_striping.test.ts @@ -0,0 +1,36 @@ +import { describe, expect, it } from 'vitest'; +import { + stripedFrameCount, + stripedOrderedPacketIndex, +} from '@raptorqr/core/sender/parallel_striping'; + +describe('ordered parallel QR striping', () => { + it.each([1, 2, 4, 6, 8] as const)('keeps %s-way frames mixed and complete', (parallelCount) => { + const packetOrder = [0, 6, 1, 7, 2, 8, 3, 9, 4, 10, 5, 11]; + const frameCount = stripedFrameCount(packetOrder.length, parallelCount); + const displayed: number[] = []; + + for (let frameIndex = 0; frameIndex < frameCount; frameIndex++) { + for (let tileIndex = 0; tileIndex < parallelCount; tileIndex++) { + const packetIndex = stripedOrderedPacketIndex( + packetOrder, + parallelCount, + frameIndex, + tileIndex, + ); + if (packetIndex !== null) displayed.push(packetIndex); + } + } + + expect(displayed).toEqual(packetOrder); + expect(new Set(displayed).size).toBe(packetOrder.length); + }); + + it('leaves only tail tile slots empty', () => { + const packetOrder = [0, 1, 2, 3, 4]; + expect(stripedFrameCount(packetOrder.length, 4)).toBe(2); + expect(stripedOrderedPacketIndex(packetOrder, 4, 1, 0)).toBe(4); + expect(stripedOrderedPacketIndex(packetOrder, 4, 1, 1)).toBeNull(); + expect(stripedOrderedPacketIndex(packetOrder, 4, 1, 3)).toBeNull(); + }); +}); diff --git a/packages/raptorqr-core/src/tests/raptorq_playback.test.ts b/packages/raptorqr-core/src/tests/raptorq_playback.test.ts new file mode 100644 index 0000000..a94e158 --- /dev/null +++ b/packages/raptorqr-core/src/tests/raptorq_playback.test.ts @@ -0,0 +1,251 @@ +import { describe, expect, it } from 'vitest'; +import { + createEvenlyInterleavedPacketIndexOrder, + createRaptorQPlaybackOrders, + formatRaptorQPlaybackStrategy, + getRaptorQPlaybackWindowPacketIndices, + normalizeRaptorQPlaybackStrategy, +} from '@raptorqr/core/sender/raptorq_playback'; +import { + classifyRaptorQPackets, + packetizeRaptorQ, +} from '@raptorqr/core/sender/raptorq_packetizer'; +import { packetCodec, parsePacket } from '@raptorqr/core/protocol/packet'; +import { RaptorQWasmDecoder } from '@raptorqr/core/fec/raptorq_wasm'; + +function packetIds( + ids: Array<{ sourceBlock: number; esi: number }>, +): Uint8Array[] { + return ids.map(({ sourceBlock, esi }) => new Uint8Array([ + sourceBlock, + (esi >>> 16) & 0xff, + (esi >>> 8) & 0xff, + esi & 0xff, + ])); +} + +function indexRange(start: number, count: number): number[] { + return Array.from({ length: count }, (_, index) => start + index); +} + +function assertCompleteOrder(order: number[], sourceCount: number, repairCount: number): void { + expect(order).toHaveLength(sourceCount + repairCount); + expect(new Set(order).size).toBe(order.length); + expect(order).toEqual(expect.arrayContaining(indexRange(0, sourceCount + repairCount))); +} + +describe('RaptorQ playback strategy', () => { + it('normalizes and formats the three public strategies', () => { + expect(normalizeRaptorQPlaybackStrategy('fast-start')).toBe('fast-start'); + expect(normalizeRaptorQPlaybackStrategy('balanced')).toBe('balanced'); + expect(normalizeRaptorQPlaybackStrategy('even-spread')).toBe('even-spread'); + expect(normalizeRaptorQPlaybackStrategy('unknown')).toBe('balanced'); + expect(formatRaptorQPlaybackStrategy('fast-start')).toBe('Fast start'); + expect(formatRaptorQPlaybackStrategy('balanced')).toBe('Balanced'); + expect(formatRaptorQPlaybackStrategy('even-spread')).toBe('Even spread'); + }); + + it('creates exact 1000/100 cumulative-proportion spacing', () => { + const source = indexRange(0, 1000); + const repair = indexRange(1000, 100); + const order = createEvenlyInterleavedPacketIndexOrder(source, repair); + + assertCompleteOrder(order, 1000, 100); + expect(order[0]).toBe(0); + + const sourceCountsBetweenRepairs: number[] = []; + let sourceCount = 0; + for (const packetIndex of order) { + if (packetIndex < 1000) { + sourceCount++; + } else { + sourceCountsBetweenRepairs.push(sourceCount); + sourceCount = 0; + } + } + expect(sourceCountsBetweenRepairs).toEqual(new Array(100).fill(10)); + expect(sourceCount).toBe(0); + }); + + it('keeps non-divisible source gaps within one packet', () => { + const order = createEvenlyInterleavedPacketIndexOrder(indexRange(0, 7), indexRange(7, 3)); + assertCompleteOrder(order, 7, 3); + + const gaps: number[] = []; + let sourceCount = 0; + for (const packetIndex of order) { + if (packetIndex < 7) { + sourceCount++; + } else { + gaps.push(sourceCount); + sourceCount = 0; + } + } + expect(gaps).toEqual([3, 2, 2]); + expect(Math.max(...gaps) - Math.min(...gaps)).toBeLessThanOrEqual(1); + }); + + it.each([ + [0, 0], + [10, 10], + [20, 20], + [100, 100], + ])('handles %s%% repair without losing packet indexes', (repairPercent, expectedRepairCount) => { + const source = indexRange(0, 100); + const repair = indexRange(100, expectedRepairCount); + const order = createEvenlyInterleavedPacketIndexOrder(source, repair); + assertCompleteOrder(order, source.length, repair.length); + expect(order[0]).toBe(0); + }); + + it('maps strategy semantics to initial and loop orders', () => { + const source = [0, 2, 4, 6]; + const repair = [1, 3]; + const sourceFirst = [0, 2, 4, 6, 1, 3]; + const even = createEvenlyInterleavedPacketIndexOrder(source, repair); + + const fast = createRaptorQPlaybackOrders(source, repair, 'fast-start'); + expect(fast.initialOrder).toEqual(sourceFirst); + expect(fast.loopOrder).toEqual(sourceFirst); + + const balanced = createRaptorQPlaybackOrders(source, repair, 'balanced'); + expect(balanced.initialOrder).toEqual(sourceFirst); + expect(balanced.loopOrder).toEqual(even); + + const spread = createRaptorQPlaybackOrders(source, repair, 'even-spread'); + expect(spread.initialOrder).toEqual(even); + expect(spread.loopOrder).toEqual(even); + }); + + it('prefetches loop-order packets when an initial window crosses the cycle boundary', () => { + const orders = createRaptorQPlaybackOrders( + indexRange(0, 8), + indexRange(8, 4), + 'balanced', + ); + + const initialTail = orders.initialOrder.slice(8, 12); + const loopHead = orders.loopOrder.slice(0, 8); + const window = getRaptorQPlaybackWindowPacketIndices( + orders, + 'initial', + 4, + 2, + 3, + ); + + expect(window).toEqual([...new Set([...initialTail, ...loopHead])]); + expect(window).not.toEqual([...new Set([ + ...initialTail, + ...orders.initialOrder.slice(0, 8), + ])]); + }); + + it('keeps using loop order when a loop window wraps', () => { + const orders = createRaptorQPlaybackOrders( + indexRange(0, 8), + indexRange(8, 4), + 'balanced', + ); + const window = getRaptorQPlaybackWindowPacketIndices(orders, 'loop', 4, 2, 3); + + expect(window).toEqual([ + ...orders.loopOrder.slice(8, 12), + ...orders.loopOrder.slice(0, 8), + ]); + }); +}); + +describe('RaptorQ packet classification', () => { + it.each([ + [0, 0], + [10, 1], + [20, 2], + [100, 10], + ])('exposes ceil-based source/repair metadata for %s%% repair', async (repairPercent, repairCount) => { + const data = new Uint8Array(1_240); + const result = await packetizeRaptorQ( + data, + false, + false, + undefined, + undefined, + { maxTransportPayloadSize: 128, repairPercent }, + ); + + expect(result.sourcePacketIndices).toHaveLength(10); + expect(result.repairPacketIndices).toHaveLength(repairCount); + expect([ + ...result.sourcePacketIndices, + ...result.repairPacketIndices, + ].sort((a, b) => a - b)).toEqual(indexRange(0, result.packets.length)); + }); + + it('classifies interspersed single-block source and repair packets by ESI', () => { + const packets = packetIds([ + { sourceBlock: 0, esi: 0 }, + { sourceBlock: 0, esi: 10 }, + { sourceBlock: 0, esi: 1 }, + { sourceBlock: 0, esi: 11 }, + ]); + const result = classifyRaptorQPackets(packets, 100, 14); + + expect(result.sourcePacketIndices).toEqual([0, 2]); + expect(result.repairPacketIndices).toEqual([1, 3]); + }); + + it('classifies multiple source blocks using each block K and its ESI', () => { + const totalSourceSymbols = 56_405; + const packets = packetIds([ + { sourceBlock: 0, esi: 28_202 }, + { sourceBlock: 0, esi: 28_203 }, + { sourceBlock: 1, esi: 28_201 }, + { sourceBlock: 1, esi: 28_202 }, + ]); + const result = classifyRaptorQPackets( + packets, + totalSourceSymbols * 10, + 14, + ); + + expect(result.sourcePacketIndices).toEqual([0, 2]); + expect(result.repairPacketIndices).toEqual([1, 3]); + }); +}); + +describe('RaptorQ playback roundtrip', () => { + it.each(['fast-start', 'balanced', 'even-spread'] as const)( + 'recovers after deterministic loss with %s order', + async (strategy) => { + const data = new Uint8Array(30_000); + for (let index = 0; index < data.length; index++) data[index] = (index * 17 + 31) & 0xff; + + const packetized = await packetizeRaptorQ( + data, + false, + false, + undefined, + undefined, + { maxTransportPayloadSize: 128, repairPercent: 20 }, + ); + const { loopOrder } = createRaptorQPlaybackOrders( + packetized.sourcePacketIndices, + packetized.repairPacketIndices, + strategy, + ); + const decoder = await RaptorQWasmDecoder.create(packetized.dataLength, packetized.symbolSize); + + let decoded: Uint8Array | null = null; + for (let position = 0; position < loopOrder.length; position++) { + if (position % 11 === 0) continue; + const packet = parsePacket(packetized.packets[loopOrder[position]!]!); + expect(packetCodec(packet.header)).toBe('wasm-raptorq'); + decoded = decoder.push(packet.payload); + if (decoded) break; + } + + expect(decoded).not.toBeNull(); + expect(decoded!.slice(0, data.length)).toEqual(data); + }, + ); +});