feat(downloads): queue season episode downloads (#1357)

* docs(downloads): specify season queueing

* docs(downloads): plan season queue implementation

* feat(downloads): define episode queue identity

* fix(downloads): align episode identity contract

* feat(downloads): coordinate season queue submissions

* fix(downloads): keep queue coordination provider neutral

* fix(downloads): reconcile legacy episode identities

* fix(downloads): fail closed on invalid stored coordinates

* refactor(downloads): adapt Xtream episode requests

* fix(downloads): use canonical Stalker episode ids

* test(downloads): cover Stalker adapter reactivity

* feat(downloads): add selected season queue action

* refactor(downloads): extract season download presenter

* feat(downloads): localize season queue feedback

* test(downloads): cover series batch queue flow

* test(downloads): harden series queue fixtures

* docs(downloads): describe season queueing

* docs(downloads): clarify season queue IPC contract

* fix(downloads): isolate season header build warnings

* fix(downloads): label season view toggles

* fix(downloads): preserve Xtream episode headers

* fix(downloads): fail closed on stale episode state

* fix(downloads): align renderer queue safeguards

* fix(downloads): block ambiguous episode actions

* fix(downloads): accept nullable legacy coordinates

* fix(downloads): preserve scoped episode ownership

* fix(downloads): probe restored files asynchronously

* fix(downloads): bound restored file probes

* fix(downloads): release timed out file probes

* fix(downloads): bound file probe callers

* fix(downloads): refresh stable season skips

* fix(downloads): fail closed before provider prep

* fix(downloads): preserve retained partial ownership

* fix(downloads): reconcile partial cleanup completion

* fix(downloads): await authoritative list refresh

* fix(downloads): coalesce list refreshes

* fix(downloads): preserve specials season identity

* fix(stalker): preserve specials season mapping

* fix(downloads): distinguish missing Xtream seasons
This commit is contained in:
4gray authored and GitHub committed 2026-08-03 08:53:44 +02:00
1 parent d3cc18dc72
commit 96facd6f49
79 files changed
+8352 -1255

No files matched your search

@@ -21,6 +21,8 @@ export interface ScenarioConfig {
performanceFixture?: 'catalog-100k';
/** Build series details on demand instead of during portal initialization. */
deferSeriesDetails?: true;
/** Optional local stream fixture for deterministic download queue tests. */
downloadStreamFixture?: 'slow-series';
}
/**
@@ -99,6 +101,19 @@ export const SCENARIOS: Record<string, ScenarioConfig> = {
accountStatus: 'Active',
expiryDate: '2099-12-31',
},
'downloadqueue:downloadqueue': {
name: 'download-queue',
description:
'Download queue — 1 series category, 4 series with 4 episodes each',
seed: 8080,
categoryCount: { live: 0, vod: 0, series: 1 },
itemsPerCategory: 4,
seasonsPerSeries: 1,
episodesPerSeason: 4,
accountStatus: 'Active',
expiryDate: '2099-12-31',
downloadStreamFixture: 'slow-series',
},
'epg:epg': {
name: 'epg-fixture',
description:
@@ -1,9 +1,14 @@
import { connect } from 'node:net';
import express from 'express';
import { resetAll } from './data-store.js';
import {
createXtreamMockApp,
parseXtreamMockServerEnvironment,
} from './server.js';
import {
type SlowSeriesDownloadOptions,
streamSlowSeriesDownload,
} from './slow-series-download.js';
import { startLoopbackServer } from './testing/http-server.fixture.js';
jest.mock('@faker-js/faker', () => {
@@ -117,6 +122,38 @@ describe('Xtream mock server factory', () => {
}
});
it('serves the download queue series fixture locally without changing ordinary series redirects', async () => {
const running = await startLoopbackServer(
createXtreamMockApp({ host: '127.0.0.1', port: 0 })
);
try {
const localSeries = await fetch(
`${running.origin}/series/downloadqueue/downloadqueue/80000.mkv`,
{ redirect: 'manual' }
);
const ordinarySeries = await fetch(
`${running.origin}/series/user1/pass1/80000.mkv`,
{ redirect: 'manual' }
);
expect(localSeries.status).toBe(200);
expect(localSeries.headers.get('content-type')).toContain(
'video/mp4'
);
expect(
Number(localSeries.headers.get('content-length'))
).toBeGreaterThan(1024 * 1024);
await localSeries.body?.cancel();
expect(ordinarySeries.status).toBe(302);
expect(ordinarySeries.headers.get('location')).toBe(
'https://test-streams.mux.dev/x36xhzz/x36xhzz.m3u8'
);
} finally {
await running.close();
}
});
it('keeps a non-EPG timezone stream empty across repeated short-EPG requests', async () => {
const running = await startLoopbackServer(
createXtreamMockApp({ host: '127.0.0.1', port: 0 })
@@ -192,6 +229,64 @@ describe('Xtream mock server factory', () => {
});
});
describe('Slow series download stream', () => {
it('completes the configured byte count over a loopback response', async () => {
const options = {
chunkSize: 1_024,
intervalMs: 1,
totalBytes: 10 * 1_024 + 7,
} satisfies SlowSeriesDownloadOptions;
const running = await startLoopbackServer(
createSlowSeriesDownloadApp(options)
);
let closed = false;
try {
const response = await fetch(`${running.origin}/slow-series`);
const body = await response.arrayBuffer();
expect(response.status).toBe(200);
expect(response.headers.get('content-type')).toContain('video/mp4');
expect(Number(response.headers.get('content-length'))).toBe(
options.totalBytes
);
expect(body.byteLength).toBe(options.totalBytes);
await within(running.close(), 1_000);
closed = true;
} finally {
if (!closed) await running.close().catch(() => undefined);
}
});
it('stops a longer recursive stream after the response body is cancelled', async () => {
const options = {
chunkSize: 1_024,
intervalMs: 25,
totalBytes: 64 * 1_024,
} satisfies SlowSeriesDownloadOptions;
const running = await startLoopbackServer(
createSlowSeriesDownloadApp(options)
);
let closed = false;
try {
const response = await fetch(`${running.origin}/slow-series`);
const reader = response.body?.getReader();
if (!reader) throw new Error('Expected a streaming response body.');
const firstChunk = await within(reader.read(), 1_000);
expect(firstChunk.done).toBe(false);
expect(firstChunk.value?.byteLength).toBeGreaterThan(0);
await within(reader.cancel(), 1_000);
await within(running.close(), 1_000);
closed = true;
} finally {
if (!closed) await running.close().catch(() => undefined);
}
});
});
describe('Xtream mock environment parsing', () => {
it('uses safe defaults and enables control only for the exact flag', () => {
expect(parseXtreamMockServerEnvironment({})).toEqual({
@@ -271,3 +366,28 @@ async function rawLoopbackGet(
});
});
}
function createSlowSeriesDownloadApp(
options: SlowSeriesDownloadOptions
): express.Express {
const app = express();
app.get('/slow-series', (request, response) => {
streamSlowSeriesDownload(request, response, options);
});
return app;
}
async function within<T>(promise: Promise<T>, timeoutMs: number): Promise<T> {
let timer: NodeJS.Timeout | undefined;
const timeout = new Promise<never>((_resolve, reject) => {
timer = setTimeout(
() => reject(new Error(`Operation exceeded ${timeoutMs}ms.`)),
timeoutMs
);
});
try {
return await Promise.race([promise, timeout]);
} finally {
if (timer) clearTimeout(timer);
}
}
+22 -1
View File
@@ -14,6 +14,8 @@ import {
XtreamPerformanceController,
} from './performance-control.js';
import { dispatchAction } from './routes/dispatch.js';
import { getScenario } from './scenarios.js';
import { streamSlowSeriesDownload } from './slow-series-download.js';
export { createXtreamMockServerShutdown } from './server-lifecycle.js';
@@ -251,10 +253,21 @@ function installStreamRoutes(
}
response.redirect(HLS_STUB);
};
const seriesResponse = (request: Request, response: Response) => {
if (isPerformanceMediaRequest(request, controlEnabled)) {
response.status(410).json({ error: 'performance-media-disabled' });
return;
}
if (isSlowSeriesDownloadRequest(request)) {
streamSlowSeriesDownload(request, response);
return;
}
response.redirect(HLS_STUB);
};
app.get('/live/:username/:password/:streamId.m3u8', streamResponse);
app.get('/live/:username/:password/:streamId.ts', streamResponse);
app.get('/movie/:username/:password/:streamId.:ext', streamResponse);
app.get('/series/:username/:password/:streamId.:ext', streamResponse);
app.get('/series/:username/:password/:streamId.:ext', seriesResponse);
app.all(
'/timeshift/:username/:password/:duration/:start/:streamId.ts',
streamResponse
@@ -262,6 +275,14 @@ function installStreamRoutes(
app.all('/streaming/timeshift.php', streamResponse);
}
function isSlowSeriesDownloadRequest(request: Request): boolean {
const username = String(request.params['username'] ?? '');
const password = String(request.params['password'] ?? '');
return (
getScenario(username, password).downloadStreamFixture === 'slow-series'
);
}
function isPerformanceMediaRequest(
request: Request,
controlEnabled: boolean
@@ -0,0 +1,93 @@
import type { Request, Response } from 'express';
export interface SlowSeriesDownloadOptions {
readonly chunkSize: number;
readonly intervalMs: number;
readonly totalBytes: number;
}
export const DEFAULT_SLOW_SERIES_DOWNLOAD_OPTIONS: SlowSeriesDownloadOptions = {
chunkSize: 32 * 1024,
intervalMs: 100,
totalBytes: 8 * 1024 * 1024,
};
export function streamSlowSeriesDownload(
request: Request,
response: Response,
options: SlowSeriesDownloadOptions = DEFAULT_SLOW_SERIES_DOWNLOAD_OPTIONS
): void {
validateOptions(options);
const chunkBuffer = Buffer.alloc(options.chunkSize);
let sentBytes = 0;
let stopped = false;
let timer: NodeJS.Timeout | undefined;
function stop(): void {
if (stopped) return;
stopped = true;
if (timer) {
clearTimeout(timer);
timer = undefined;
}
request.off('close', stop);
response.off('close', stop);
response.off('finish', stop);
response.off('drain', scheduleChunk);
}
function scheduleChunk(): void {
if (!stopped && !timer) {
timer = setTimeout(writeChunk, options.intervalMs);
}
}
function writeChunk(): void {
timer = undefined;
if (stopped || response.destroyed || response.writableEnded) {
stop();
return;
}
const remainingBytes = options.totalBytes - sentBytes;
const chunk =
remainingBytes >= chunkBuffer.length
? chunkBuffer
: chunkBuffer.subarray(0, remainingBytes);
sentBytes += chunk.length;
const canContinue = response.write(chunk);
if (sentBytes >= options.totalBytes) {
response.end();
} else if (canContinue) {
scheduleChunk();
} else {
response.once('drain', scheduleChunk);
}
}
request.once('close', stop);
response.once('close', stop);
response.once('finish', stop);
response
.status(200)
.type('video/mp4')
.set('Content-Length', String(options.totalBytes))
.set('Cache-Control', 'no-store')
.flushHeaders();
scheduleChunk();
}
function validateOptions(options: SlowSeriesDownloadOptions): void {
if (!Number.isSafeInteger(options.totalBytes) || options.totalBytes <= 0) {
throw new Error('Slow download totalBytes must be a positive integer.');
}
if (!Number.isSafeInteger(options.chunkSize) || options.chunkSize <= 0) {
throw new Error('Slow download chunkSize must be a positive integer.');
}
if (!Number.isSafeInteger(options.intervalMs) || options.intervalMs < 0) {
throw new Error(
'Slow download intervalMs must be a non-negative integer.'
);
}
}