feat(raptorq): add configurable playback scheduling strategies

This commit is contained in:
infrost
2026-07-13 01:49:53 +01:00
parent 4d09dd0a3a
commit bb726f800d
14 changed files with 839 additions and 104 deletions
+169 -97
View File
@@ -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<QREncoder>(DEFAULT_QR_ENCODER);
const [fecCodec, setFecCodec] = useState<FecCodec>(DEFAULT_FEC_CODEC);
const [raptorqRepairPercent, setRaptorqRepairPercent] = useState(DEFAULT_RAPTORQ_REPAIR_PERCENT);
const [raptorqPlaybackStrategy, setRaptorqPlaybackStrategy] = useState<RaptorQPlaybackStrategy>(
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() {
</select>
</label>
{fecCodec === 'wasm-raptorq' && (
<label>
<div style={{ ...S.row, justifyContent: 'space-between', alignItems: 'baseline' }}>
<span style={S.label}>RaptorQ repair</span>
<span style={S.infoValue}>{raptorqRepairPercent}%</span>
</div>
<input
type="range"
min={MIN_RAPTORQ_REPAIR_PERCENT}
max={MAX_RAPTORQ_REPAIR_PERCENT}
step={1}
value={raptorqRepairPercent}
style={S.slider}
disabled={encodingLive}
onInput={(e) => handleRaptorQRepairPercentChange((e.target as HTMLInputElement).value)}
/>
<div style={S.sliderLabels}>
<span>Less QR</span>
<span>More repair</span>
</div>
</label>
<>
<label>
<div style={{ ...S.row, justifyContent: 'space-between', alignItems: 'baseline' }}>
<span style={S.label}>RaptorQ repair</span>
<span style={S.infoValue}>{raptorqRepairPercent}%</span>
</div>
<input
type="range"
min={MIN_RAPTORQ_REPAIR_PERCENT}
max={MAX_RAPTORQ_REPAIR_PERCENT}
step={1}
value={raptorqRepairPercent}
style={S.slider}
disabled={encodingLive || liveTransfer !== null}
onInput={(e) => handleRaptorQRepairPercentChange((e.target as HTMLInputElement).value)}
/>
<div style={S.sliderLabels}>
<span>Less QR</span>
<span>More repair</span>
</div>
</label>
<label>
<div style={{ ...S.row, justifyContent: 'space-between', alignItems: 'baseline' }}>
<span style={S.label}>RaptorQ playback</span>
<span style={S.infoValue}>{formatRaptorQPlaybackStrategy(raptorqPlaybackStrategy)}</span>
</div>
<select
value={raptorqPlaybackStrategy}
style={S.select}
disabled={encodingLive || liveTransfer !== null}
onChange={(e) => handleRaptorQPlaybackStrategyChange((e.target as HTMLSelectElement).value)}
>
{RAPTORQ_PLAYBACK_STRATEGIES.map((strategy) => (
<option key={strategy} value={strategy}>
{formatRaptorQPlaybackStrategy(strategy)}
</option>
))}
</select>
</label>
</>
)}
</div>
)}
@@ -1260,6 +1319,12 @@ export function SenderPage() {
? `${formatFecCodec(stats.fecCodec)} · ${stats.raptorqRepairPercent}% repair`
: formatFecCodec(stats.fecCodec)}
</span>
<span style={S.infoLabel}>RaptorQ playback</span>
<span style={S.infoValue}>
{stats.raptorqPlaybackStrategy
? formatRaptorQPlaybackStrategy(stats.raptorqPlaybackStrategy)
: 'Not applicable (JS RLNC)'}
</span>
<span style={S.infoLabel}>QR packets</span>
<span style={S.infoValue}>{stats.frameCount}</span>
<span style={S.infoLabel}>Parallel QR</span>
@@ -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<number>,
): 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<number>,
transfer: LiveTransfer,
startFrameIndex: number,
renderWindowFrames: number,
windowPacketIndexes: ReadonlySet<number>,
): 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;
}
+8 -1
View File
@@ -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[] = [];
+4
View File
@@ -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<EncodeOutput> {
return {
type: 'encoded',
packets: result.packets,
sourcePacketIndices: result.sourcePacketIndices,
repairPacketIndices: result.repairPacketIndices,
totalGenerations: result.totalGenerations,
stats: {
originalSize: originalBytes.length,
+25 -3
View File
@@ -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<GenerateInput>) => {
async function handleGenerate(input: GenerateInput): Promise<GifOutput> {
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<GifOutput> {
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<number>();
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 };
+7 -1
View File
@@ -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<string>();
@@ -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,
+1
View File
@@ -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`
+9 -2
View File
@@ -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);
+28
View File
@@ -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.
+1
View File
@@ -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';
@@ -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;
}
@@ -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 };
}
@@ -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<RaptorQPlaybackOrders, 'initialOrder' | 'loopOrder'>,
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<number>();
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<number>();
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);
}
}
@@ -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();
});
});
@@ -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);
},
);
});