Skip to content
Merged
2 changes: 1 addition & 1 deletion eslint.config.js
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ const featureBoundaries = (features) => {
{
paths: ['pinia', 'vue'].map((name) => ({
name,
message: `The ${feature.dir} pure layer must stay framework-free — no ${name}.`,
message: `The ${feature.dir} pure layer must stay framework-free: no ${name}.`,
})),
patterns: [
{
Expand Down
110 changes: 78 additions & 32 deletions src/core/streaming/__tests__/cachedStreamFetcher.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,51 +2,97 @@ import { RequestPool } from '@/src/core/streaming/requestPool';
import {
CachedStreamFetcher,
sliceChunks,
StopSignal,
} from '@/src/core/streaming/cachedStreamFetcher';
import { describe, expect, it } from 'vitest';

const readStream = async (
stream: ReadableStream<Uint8Array>,
stopAfter = Infinity
) => {
const reader = stream.getReader();
const chunks: Uint8Array[] = [];
let size = 0;
let completed = false;
try {
while (size <= stopAfter) {
const result = await reader.read();
if (result.done) {
completed = true;
break;
}
chunks.push(result.value);
size += result.value.length;
}
} finally {
if (!completed) await reader.cancel();
reader.releaseLock();
}

const bytes = new Uint8Array(size);
let offset = 0;
chunks.forEach((chunk) => {
bytes.set(chunk, offset);
offset += chunk.length;
});
return bytes;
};

describe('CachedStreamFetcher', () => {
it('should support stopping and resuming', async () => {
const pool = new RequestPool();
const fetcher = new CachedStreamFetcher(
'https://data.kitware.com/api/v1/file/57b5d4648d777f10f2693e7e/download',
{
fetch: pool.fetch,
}
const source = Uint8Array.from(
{ length: 32 * 1024 + 123 },
(_, index) => (index * 31) % 251
);
const requestedRanges: Array<string | null> = [];
const fetchRange: typeof fetch = async (_input, init) => {
const range = new Headers(init?.headers).get('Range');
requestedRanges.push(range);
const start = range ? Number(range.match(/^bytes=(\d+)-$/)?.[1]) : 0;
let offset = start;
const body = new ReadableStream<Uint8Array>({
pull(controller) {
if (offset === source.length) {
controller.close();
return;
}
const end = Math.min(offset + 4096, source.length);
controller.enqueue(source.slice(offset, end));
offset = end;
},
});
return new Response(body, {
status: start === 0 ? 200 : 206,
headers: {
'content-length': String(source.length - start),
...(start === 0
? {}
: {
'content-range': `bytes ${start}-${source.length - 1}/${source.length}`,
}),
},
});
};
const pool = new RequestPool(1, fetchRange);
const fetcher = new CachedStreamFetcher('https://example.test/data', {
fetch: pool.fetch,
});

await fetcher.connect();
let stream = fetcher.getStream();
let size = 0;
try {
// @ts-ignore
for await (const chunk of stream) {
size += chunk.length;
if (size > 4096 * 3) {
break;
}
}
} catch (err) {
if (err !== StopSignal) throw err;
} finally {
fetcher.close();
}
const partial = await readStream(fetcher.getStream(), 4096 * 3);
expect(partial).toEqual(source.slice(0, partial.length));
const resumeAt = fetcher.size;
expect(resumeAt).toBeGreaterThanOrEqual(partial.length);
expect(resumeAt).toBeLessThan(source.length);
fetcher.close();

await fetcher.connect();
expect(requestedRanges).toEqual([null, `bytes=${resumeAt}-`]);

// ensure we can read the stream multiple times
for (let i = 0; i < 2; i++) {
stream = fetcher.getStream();
size = 0;
// @ts-ignore

for await (const chunk of stream) {
size += chunk.length;
}

expect(size).to.equal(fetcher.size);
expect(await readStream(fetcher.getStream())).toEqual(source);
}
expect(fetcher.size).toBe(source.length);
expect(requestedRanges).toHaveLength(2);

fetcher.close();
});
Expand Down
Loading
Loading