fix(downloads): share throttle budget globally
This commit is contained in:
+11
-5
@@ -4,6 +4,7 @@ import * as fs from 'fs';
|
||||
import { spawn, ChildProcess, execSync, spawnSync } from 'child_process';
|
||||
import { connect as tlsConnect, TLSSocket } from 'node:tls';
|
||||
import { pathToFileURL } from 'node:url';
|
||||
import type { Transform } from 'node:stream';
|
||||
import axios from 'axios';
|
||||
import { autoUpdater } from 'electron-updater';
|
||||
import { compareUpdateVersions, isNewerUpdateVersion, normalizeUpdateVersion } from './main/domain/update-version-utils';
|
||||
@@ -19,7 +20,7 @@ import {
|
||||
import { tBackend as tBackendCore, type BackendMessageKey } from './main/domain/i18n-backend';
|
||||
import { watchRendererChanges } from './main/dev-reload';
|
||||
import { createPausableOutput, type PausableOutput } from './main/domain/pausable-output';
|
||||
import { createTokenBucketTransform } from './main/domain/token-bucket-transform';
|
||||
import { createTokenBucketBudget, createTokenBucketTransform } from './main/domain/token-bucket-transform';
|
||||
import { decideDownloadStart, normalizeDownloadPolicy, type DownloadPolicy } from './main/domain/download-policy';
|
||||
import { PartialDownloadRegistry } from './main/domain/partial-download';
|
||||
import { QueueProcessRegistry, QueueRunLifecycle, waitForChildProcessExit } from './main/queue/process-registry';
|
||||
@@ -418,6 +419,11 @@ function getStreamlinkStreamArg(): string {
|
||||
return `${choice},best`;
|
||||
}
|
||||
|
||||
function createDownloadThrottleTransform(): Transform | undefined {
|
||||
const maxBytesPerSecond = config.download_policy.throttle?.maxBytesPerSecond ?? null;
|
||||
downloadThrottleBudget.setMaxBytesPerSecond(maxBytesPerSecond);
|
||||
return maxBytesPerSecond ? createTokenBucketTransform(maxBytesPerSecond, undefined, downloadThrottleBudget) : undefined;
|
||||
}
|
||||
function normalizeConfigTemplates(input: Config): Config {
|
||||
// downloaded_vod_ids is bounded so a long-running app doesn't accumulate
|
||||
// an unbounded list across years of downloads. Latest entries kept.
|
||||
@@ -809,6 +815,7 @@ const activeDownloads = new Map<string, ActiveDownloadTracking>();
|
||||
const cancelledItemIds = new Set<string>();
|
||||
const queueProcessRegistry = new QueueProcessRegistry();
|
||||
const queueRunLifecycle = new QueueRunLifecycle(queueProcessRegistry);
|
||||
const downloadThrottleBudget = createTokenBucketBudget(null);
|
||||
let downloadPolicyWakeTimer: NodeJS.Timeout | null = null;
|
||||
let lastDownloadPolicyStatusFingerprint = '';
|
||||
|
||||
@@ -4039,11 +4046,10 @@ function downloadVODPart(
|
||||
resolve({ success: false, error: tBackend('unknownDownloadError') });
|
||||
return;
|
||||
}
|
||||
const maxBytesPerSecond = config.download_policy.throttle?.maxBytesPerSecond;
|
||||
const output = createPausableOutput(
|
||||
proc.stdout,
|
||||
outputStream,
|
||||
maxBytesPerSecond ? createTokenBucketTransform(maxBytesPerSecond) : undefined,
|
||||
createDownloadThrottleTransform(),
|
||||
);
|
||||
const outputFinished = output.finished.then(() => null, (error) => error);
|
||||
const processRegistration = queueProcessRegistry.register(itemId, 'streamlink', {
|
||||
@@ -7515,6 +7521,7 @@ ipcMain.handle('save-config', (event, newConfig: Partial<Config>, fileCapability
|
||||
}
|
||||
const nextConfig = normalizeConfigTemplates({ ...config, ...acceptedConfig });
|
||||
config = persistStateChange(config, () => nextConfig, saveConfig);
|
||||
downloadThrottleBudget.setMaxBytesPerSecond(config.download_policy.throttle?.maxBytesPerSecond ?? null);
|
||||
if (JSON.stringify(config.download_policy) !== previousDownloadPolicy && !isDownloading && downloadQueue.some((item) => item.status === 'pending')) {
|
||||
scheduleQueueProcessing();
|
||||
} else {
|
||||
@@ -8165,11 +8172,10 @@ registerTrustedIpcHandler(ipcMain, 'download-clip', isTrustedRendererEvent, () =
|
||||
resolve({ success: false, error: tBackend('unknownDownloadError') });
|
||||
return;
|
||||
}
|
||||
const maxBytesPerSecond = config.download_policy.throttle?.maxBytesPerSecond;
|
||||
const output = createPausableOutput(
|
||||
proc.stdout,
|
||||
fs.createWriteStream(partialFilename, { flags: 'w' }),
|
||||
maxBytesPerSecond ? createTokenBucketTransform(maxBytesPerSecond) : undefined,
|
||||
createDownloadThrottleTransform(),
|
||||
);
|
||||
const outputFinished = output.finished.then(() => null, (error) => error);
|
||||
|
||||
|
||||
@@ -52,8 +52,16 @@ describe('download policy integration contract', () => {
|
||||
const end = source.indexOf('const outputFinished = output.finished', start);
|
||||
const section = source.slice(start, end);
|
||||
|
||||
expect(section).toContain('createTokenBucketTransform');
|
||||
expect(section).toContain('createDownloadThrottleTransform()');
|
||||
expect(section).toContain("const args = [...streamlinkCmd.prefixArgs, url, getStreamlinkStreamArg(), '--stdout'];");
|
||||
expect(section).not.toMatch(/args\.push\([^\n]*(?:bandwidth|rate-limit|max-rate|throttle)/i);
|
||||
});
|
||||
|
||||
it('routes queue and clip stdout through one app-wide token bucket budget', () => {
|
||||
const source = readFileSync(join(process.cwd(), 'src', 'main.ts'), 'utf8');
|
||||
|
||||
expect(source).toContain('const downloadThrottleBudget = createTokenBucketBudget(null);');
|
||||
expect(source).toContain('function createDownloadThrottleTransform(): Transform | undefined');
|
||||
expect(source.match(/createDownloadThrottleTransform\(\)/g)).toHaveLength(3);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { PassThrough } from 'node:stream';
|
||||
import { describe, expect, it } from 'vitest';
|
||||
import { createTokenBucketTransform, type TokenBucketClock } from './token-bucket-transform';
|
||||
import { createTokenBucketBudget, createTokenBucketTransform, type TokenBucketClock } from './token-bucket-transform';
|
||||
|
||||
class ManualClock implements TokenBucketClock {
|
||||
private nextTimerId = 0;
|
||||
@@ -75,4 +75,39 @@ describe('app-side token bucket transform', () => {
|
||||
expect(clock.timerCount).toBe(0);
|
||||
expect(Buffer.concat(output).toString()).toBe('a');
|
||||
});
|
||||
|
||||
it('shares one byte budget across concurrent output transforms', () => {
|
||||
const clock = new ManualClock();
|
||||
const budget = createTokenBucketBudget(2, clock);
|
||||
const first = createTokenBucketTransform(2, clock, budget);
|
||||
const second = createTokenBucketTransform(2, clock, budget);
|
||||
const firstOutput: Buffer[] = [];
|
||||
const secondOutput: Buffer[] = [];
|
||||
first.on('data', (chunk: Buffer) => firstOutput.push(Buffer.from(chunk)));
|
||||
second.on('data', (chunk: Buffer) => secondOutput.push(Buffer.from(chunk)));
|
||||
|
||||
first.write(Buffer.from('ab'));
|
||||
second.write(Buffer.from('cd'));
|
||||
|
||||
expect(Buffer.concat(firstOutput).toString()).toBe('ab');
|
||||
expect(Buffer.concat(secondOutput).toString()).toBe('');
|
||||
expect(clock.timerCount).toBe(1);
|
||||
|
||||
clock.advance(1_000);
|
||||
|
||||
expect(Buffer.concat(secondOutput).toString()).toBe('cd');
|
||||
});
|
||||
|
||||
it('seeds an app-wide budget when throttling is enabled after startup', () => {
|
||||
const clock = new ManualClock();
|
||||
const budget = createTokenBucketBudget(null, clock);
|
||||
budget.setMaxBytesPerSecond(2);
|
||||
const transform = createTokenBucketTransform(2, clock, budget);
|
||||
const output: Buffer[] = [];
|
||||
transform.on('data', (chunk: Buffer) => output.push(Buffer.from(chunk)));
|
||||
|
||||
transform.write(Buffer.from('ab'));
|
||||
|
||||
expect(Buffer.concat(output).toString()).toBe('ab');
|
||||
});
|
||||
});
|
||||
|
||||
@@ -6,53 +6,147 @@ export interface TokenBucketClock {
|
||||
clearTimeout(handle: ReturnType<typeof setTimeout>): void;
|
||||
}
|
||||
|
||||
export interface TokenBucketBudget {
|
||||
reserve(bytes: number, release: () => void): () => void;
|
||||
setMaxBytesPerSecond(maxBytesPerSecond: number | null): void;
|
||||
}
|
||||
|
||||
interface TokenBucketReservation {
|
||||
bytes: number;
|
||||
release: () => void;
|
||||
cancelled: boolean;
|
||||
}
|
||||
|
||||
const systemClock: TokenBucketClock = {
|
||||
now: () => Date.now(),
|
||||
setTimeout: (callback, delayMs) => setTimeout(callback, delayMs),
|
||||
clearTimeout: (handle) => clearTimeout(handle),
|
||||
};
|
||||
|
||||
class TokenBucketTransform extends Transform {
|
||||
function assertRate(maxBytesPerSecond: number): void {
|
||||
if (!Number.isSafeInteger(maxBytesPerSecond) || maxBytesPerSecond <= 0) throw new RangeError('maxBytesPerSecond must be a positive safe integer');
|
||||
}
|
||||
|
||||
class SharedTokenBucketBudget implements TokenBucketBudget {
|
||||
private availableBytes: number;
|
||||
private lastRefillAt: number;
|
||||
private timer: ReturnType<typeof setTimeout> | null = null;
|
||||
private draining = false;
|
||||
private readonly reservations: TokenBucketReservation[] = [];
|
||||
|
||||
constructor(private readonly maxBytesPerSecond: number, private readonly clock: TokenBucketClock) {
|
||||
super();
|
||||
this.availableBytes = maxBytesPerSecond;
|
||||
constructor(private maxBytesPerSecond: number | null, private readonly clock: TokenBucketClock) {
|
||||
if (maxBytesPerSecond !== null) assertRate(maxBytesPerSecond);
|
||||
this.availableBytes = maxBytesPerSecond ?? 0;
|
||||
this.lastRefillAt = clock.now();
|
||||
}
|
||||
|
||||
reserve(bytes: number, release: () => void): () => void {
|
||||
const reservation: TokenBucketReservation = { bytes, release, cancelled: false };
|
||||
this.reservations.push(reservation);
|
||||
this.drain();
|
||||
return () => {
|
||||
if (reservation.cancelled) return;
|
||||
reservation.cancelled = true;
|
||||
this.drain();
|
||||
};
|
||||
}
|
||||
|
||||
setMaxBytesPerSecond(maxBytesPerSecond: number | null): void {
|
||||
if (maxBytesPerSecond !== null) assertRate(maxBytesPerSecond);
|
||||
if (this.maxBytesPerSecond === maxBytesPerSecond) return;
|
||||
const wasUnlimited = this.maxBytesPerSecond === null;
|
||||
this.maxBytesPerSecond = maxBytesPerSecond;
|
||||
this.availableBytes = maxBytesPerSecond === null ? 0 : wasUnlimited ? maxBytesPerSecond : Math.min(this.availableBytes, maxBytesPerSecond);
|
||||
this.lastRefillAt = this.clock.now();
|
||||
this.drain();
|
||||
}
|
||||
|
||||
private refill(capacity: number): void {
|
||||
if (this.maxBytesPerSecond === null) return;
|
||||
const now = this.clock.now();
|
||||
const elapsed = Math.max(0, now - this.lastRefillAt);
|
||||
this.availableBytes = Math.min(capacity, this.availableBytes + (elapsed * this.maxBytesPerSecond) / 1000);
|
||||
this.lastRefillAt = now;
|
||||
}
|
||||
|
||||
private clearTimer(): void {
|
||||
if (!this.timer) return;
|
||||
this.clock.clearTimeout(this.timer);
|
||||
this.timer = null;
|
||||
}
|
||||
|
||||
private removeCancelledReservations(): void {
|
||||
while (this.reservations[0]?.cancelled) this.reservations.shift();
|
||||
}
|
||||
|
||||
private drain(): void {
|
||||
if (this.draining) return;
|
||||
this.draining = true;
|
||||
try {
|
||||
this.clearTimer();
|
||||
while (true) {
|
||||
this.removeCancelledReservations();
|
||||
const reservation = this.reservations[0];
|
||||
if (!reservation) return;
|
||||
if (this.maxBytesPerSecond === null) {
|
||||
this.reservations.shift();
|
||||
reservation.release();
|
||||
continue;
|
||||
}
|
||||
const capacity = Math.max(this.maxBytesPerSecond, reservation.bytes);
|
||||
this.refill(capacity);
|
||||
if (this.availableBytes >= reservation.bytes) {
|
||||
this.availableBytes -= reservation.bytes;
|
||||
this.reservations.shift();
|
||||
reservation.release();
|
||||
continue;
|
||||
}
|
||||
const delayMs = Math.max(1, Math.ceil(((reservation.bytes - this.availableBytes) * 1000) / this.maxBytesPerSecond));
|
||||
this.timer = this.clock.setTimeout(() => {
|
||||
this.timer = null;
|
||||
this.drain();
|
||||
}, delayMs);
|
||||
return;
|
||||
}
|
||||
} finally {
|
||||
this.draining = false;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
class TokenBucketTransform extends Transform {
|
||||
private cancelReservation: (() => void) | null = null;
|
||||
|
||||
constructor(private readonly budget: TokenBucketBudget) {
|
||||
super();
|
||||
}
|
||||
|
||||
override _transform(chunk: Buffer, _encoding: BufferEncoding, callback: (error?: Error | null) => void): void {
|
||||
const output = Buffer.from(chunk);
|
||||
const capacity = Math.max(this.maxBytesPerSecond, output.length);
|
||||
const release = (): void => {
|
||||
this.timer = null;
|
||||
this.cancelReservation = this.budget.reserve(output.length, () => {
|
||||
this.cancelReservation = null;
|
||||
if (this.destroyed) return;
|
||||
const now = this.clock.now();
|
||||
const elapsed = Math.max(0, now - this.lastRefillAt);
|
||||
this.availableBytes = Math.min(capacity, this.availableBytes + (elapsed * this.maxBytesPerSecond) / 1000);
|
||||
this.lastRefillAt = now;
|
||||
if (this.availableBytes >= output.length) {
|
||||
this.availableBytes -= output.length;
|
||||
this.push(output);
|
||||
callback();
|
||||
return;
|
||||
}
|
||||
const delayMs = Math.max(1, Math.ceil(((output.length - this.availableBytes) * 1000) / this.maxBytesPerSecond));
|
||||
this.timer = this.clock.setTimeout(release, delayMs);
|
||||
};
|
||||
release();
|
||||
this.push(output);
|
||||
callback();
|
||||
});
|
||||
}
|
||||
|
||||
override _destroy(error: Error | null, callback: (error: Error | null) => void): void {
|
||||
if (this.timer) this.clock.clearTimeout(this.timer);
|
||||
this.timer = null;
|
||||
this.cancelReservation?.();
|
||||
this.cancelReservation = null;
|
||||
callback(error);
|
||||
}
|
||||
}
|
||||
|
||||
export function createTokenBucketTransform(maxBytesPerSecond: number, clock: TokenBucketClock = systemClock): Transform {
|
||||
if (!Number.isSafeInteger(maxBytesPerSecond) || maxBytesPerSecond <= 0) throw new RangeError('maxBytesPerSecond must be a positive safe integer');
|
||||
return new TokenBucketTransform(maxBytesPerSecond, clock);
|
||||
export function createTokenBucketBudget(maxBytesPerSecond: number | null, clock: TokenBucketClock = systemClock): TokenBucketBudget {
|
||||
return new SharedTokenBucketBudget(maxBytesPerSecond, clock);
|
||||
}
|
||||
|
||||
export function createTokenBucketTransform(
|
||||
maxBytesPerSecond: number,
|
||||
clock: TokenBucketClock = systemClock,
|
||||
budget: TokenBucketBudget = createTokenBucketBudget(maxBytesPerSecond, clock),
|
||||
): Transform {
|
||||
assertRate(maxBytesPerSecond);
|
||||
return new TokenBucketTransform(budget);
|
||||
}
|
||||
|
||||
@@ -17,7 +17,8 @@ const inputIds = [
|
||||
'deletePartsAfterMergeToggle', 'discordWebhookUrl', 'discordNotifyLiveStartToggle', 'discordNotifyLiveEndToggle',
|
||||
'discordNotifyVodCompleteToggle', 'discordNotifyVodAutoQueuedToggle', 'autoVodPollMinutes', 'autoVodMaxAgeHours',
|
||||
'autoCleanupEnabledToggle', 'autoCleanupDays', 'autoCleanupTarget', 'autoCleanupAction', 'streamlinkQuality',
|
||||
'metadataCacheMinutes', 'vodFilenameTemplate', 'partsFilenameTemplate', 'defaultClipFilenameTemplate'
|
||||
'metadataCacheMinutes', 'vodFilenameTemplate', 'partsFilenameTemplate', 'defaultClipFilenameTemplate',
|
||||
'downloadThrottleMiBps', 'downloadWindows', 'downloadPolicyValidation'
|
||||
];
|
||||
|
||||
function createInput(value = '', checked = false): Input {
|
||||
@@ -25,6 +26,52 @@ function createInput(value = '', checked = false): Input {
|
||||
}
|
||||
|
||||
describe('renderer settings autosave orchestration', () => {
|
||||
it('persists a pure download policy change through the real autosave fingerprint', async () => {
|
||||
const inputs = new Map(inputIds.map((id) => [id, createInput()]));
|
||||
inputs.get('downloadThrottleMiBps')!.value = '1';
|
||||
inputs.get('downloadWindows')!.value = '22:00-06:00';
|
||||
const saveConfigCalls: Array<Record<string, unknown>> = [];
|
||||
const window = {
|
||||
api: {
|
||||
setClientSecret: () => Promise.resolve({ encryptionAvailable: true, clientSecretConfigured: false, discordWebhookConfigured: false }),
|
||||
clearClientSecret: () => Promise.resolve({ encryptionAvailable: true, clientSecretConfigured: false, discordWebhookConfigured: false }),
|
||||
setDiscordWebhook: () => Promise.resolve({ encryptionAvailable: true, clientSecretConfigured: false, discordWebhookConfigured: false }),
|
||||
clearDiscordWebhook: () => Promise.resolve({ encryptionAvailable: true, clientSecretConfigured: false, discordWebhookConfigured: false }),
|
||||
saveConfig(payload: Record<string, unknown>) {
|
||||
saveConfigCalls.push(payload);
|
||||
return Promise.resolve(payload);
|
||||
},
|
||||
},
|
||||
};
|
||||
const sandbox = {
|
||||
window,
|
||||
config: { download_policy: { throttle: { maxBytesPerSecond: 1_048_576 }, windows: [{ start: '22:00', end: '06:00' }] } },
|
||||
UI_TEXT: { status: {}, static: {}, streamers: {} },
|
||||
byId: (id: string) => inputs.get(id) ?? createInput(),
|
||||
collectUnknownTemplatePlaceholders: () => [],
|
||||
document: { hidden: false, querySelector: () => null, getElementById: () => null },
|
||||
setTimeout,
|
||||
clearTimeout,
|
||||
console,
|
||||
};
|
||||
const context = vm.createContext(sandbox);
|
||||
const source = fs.readFileSync(path.join(process.cwd(), 'src', 'renderer-settings.ts'), 'utf8');
|
||||
const compiled = ts.transpileModule(source, {
|
||||
compilerOptions: { target: ts.ScriptTarget.ES2022, module: ts.ModuleKind.None },
|
||||
}).outputText;
|
||||
vm.runInContext(compiled, context);
|
||||
|
||||
vm.runInContext('lastPersistedSettingsFingerprint = getSettingsFingerprint(collectAutoSavePayload())', context);
|
||||
inputs.get('downloadThrottleMiBps')!.value = '1.5';
|
||||
await (vm.runInContext('flushSettingsAutoSave(false)', context) as Promise<void>);
|
||||
|
||||
expect(saveConfigCalls).toHaveLength(1);
|
||||
expect(saveConfigCalls[0].download_policy).toEqual({
|
||||
throttle: { maxBytesPerSecond: 1_572_864 },
|
||||
windows: [{ start: '22:00', end: '06:00' }]
|
||||
});
|
||||
});
|
||||
|
||||
it('persists a newer secret after an earlier asynchronous save settles', async () => {
|
||||
const inputs = new Map(inputIds.map((id) => [id, createInput()]));
|
||||
inputs.get('clientSecret')!.value = 'A';
|
||||
|
||||
@@ -859,6 +859,8 @@ function getSettingsFingerprint(payload: Partial<AppConfig>): string {
|
||||
effective.auto_cleanup_action ?? 'archive',
|
||||
effective.streamlink_quality ?? 'best',
|
||||
effective.metadata_cache_minutes ?? 10,
|
||||
effective.download_policy?.throttle?.maxBytesPerSecond ?? null,
|
||||
effective.download_policy?.windows ?? [],
|
||||
effective.filename_template_vod ?? '{title}.mp4',
|
||||
effective.filename_template_parts ?? '{date}_Part{part_padded}.mp4',
|
||||
effective.filename_template_clip ?? '{date}_{part}.mp4'
|
||||
|
||||
Reference in New Issue
Block a user