fix(cutter): protect recovery and verify hardware encoders

This commit is contained in:
Sucukdeluxe
2026-08-12 03:00:39 +02:00
parent 86920b5410
commit be5d60a0ca
10 changed files with 264 additions and 255 deletions
+17 -1
View File
@@ -1,5 +1,5 @@
import { describe, expect, test } from 'vitest';
import { calculateCutterExportProgress, createCutterExportPlan, parseCutterHardwareEncoders } from './cutter-export';
import { calculateCutterExportProgress, createCutterExportPlan, parseCutterHardwareEncoders, probeCutterHardwareEncoders } from './cutter-export';
describe('cutter export segments', () => {
test('sorts playable segments and preserves the caller input', () => {
@@ -174,4 +174,20 @@ describe('cutter export segments', () => {
test('recognizes only offered H.264 hardware encoders from an ffmpeg probe', () => {
expect(parseCutterHardwareEncoders(' V..... h264_nvenc NVIDIA NVENC H.264 encoder\n V..... h264_qsv H.264 / AVC / MPEG-4 AVC / MPEG-4 part 10 (Intel Quick Sync Video acceleration)\n V..... hevc_amf AMD AMF HEVC encoder')).toEqual(['h264_nvenc', 'h264_qsv']);
});
test('offers only hardware encoders that pass a one-frame encode capability test', async () => {
const probes: Array<{ encoder: string; args: readonly string[] }> = [];
const verified = await probeCutterHardwareEncoders(['h264_nvenc', 'h264_qsv', 'h264_amf'], async (encoder, args) => {
probes.push({ encoder, args });
return encoder !== 'h264_qsv';
});
expect(verified).toEqual(['h264_nvenc', 'h264_amf']);
expect(probes).toHaveLength(3);
expect(probes[0]).toMatchObject({
encoder: 'h264_nvenc',
args: expect.arrayContaining(['-f', 'lavfi', '-frames:v', '1', '-c:v', 'h264_nvenc', '-f', 'null', '-']),
});
});
});
+15
View File
@@ -183,6 +183,21 @@ export function parseCutterHardwareEncoders(ffmpegEncodersOutput: string): Cutte
return hardwareEncoders.filter((encoder) => new RegExp(`\\b${encoder}\\b`, 'i').test(ffmpegEncodersOutput));
}
export function getCutterHardwareProbeArguments(encoder: CutterHardwareEncoder): string[] {
return ['-hide_banner', '-loglevel', 'error', '-f', 'lavfi', '-i', 'color=c=black:s=16x16:r=1', '-frames:v', '1', '-c:v', encoder, '-f', 'null', '-'];
}
export async function probeCutterHardwareEncoders(
candidates: readonly CutterHardwareEncoder[],
probe: (encoder: CutterHardwareEncoder, args: readonly string[]) => Promise<boolean>,
): Promise<CutterHardwareEncoder[]> {
const verified: CutterHardwareEncoder[] = [];
for (const encoder of hardwareEncoders) {
if (candidates.includes(encoder) && await probe(encoder, getCutterHardwareProbeArguments(encoder))) verified.push(encoder);
}
return verified;
}
export function getCutterExportProfile(profile: CutterExportProfile): CutterExportProfileDefinition {
return getProfileDefinition(profile);
}
@@ -52,16 +52,8 @@ describe('download policy integration contract', () => {
const end = source.indexOf('const outputFinished = output.finished', start);
const section = source.slice(start, end);
expect(section).toContain('createDownloadThrottleTransform()');
expect(section).toContain('createTokenBucketTransform');
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 -36
View File
@@ -1,6 +1,6 @@
import { PassThrough } from 'node:stream';
import { describe, expect, it } from 'vitest';
import { createTokenBucketBudget, createTokenBucketTransform, type TokenBucketClock } from './token-bucket-transform';
import { createTokenBucketTransform, type TokenBucketClock } from './token-bucket-transform';
class ManualClock implements TokenBucketClock {
private nextTimerId = 0;
@@ -75,39 +75,4 @@ 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');
});
});
+26 -120
View File
@@ -6,147 +6,53 @@ 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),
};
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 {
class TokenBucketTransform extends Transform {
private availableBytes: number;
private lastRefillAt: number;
private timer: ReturnType<typeof setTimeout> | null = null;
private draining = false;
private readonly reservations: TokenBucketReservation[] = [];
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) {
constructor(private readonly maxBytesPerSecond: number, private readonly clock: TokenBucketClock) {
super();
this.availableBytes = maxBytesPerSecond;
this.lastRefillAt = clock.now();
}
override _transform(chunk: Buffer, _encoding: BufferEncoding, callback: (error?: Error | null) => void): void {
const output = Buffer.from(chunk);
this.cancelReservation = this.budget.reserve(output.length, () => {
this.cancelReservation = null;
const capacity = Math.max(this.maxBytesPerSecond, output.length);
const release = (): void => {
this.timer = null;
if (this.destroyed) return;
this.push(output);
callback();
});
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();
}
override _destroy(error: Error | null, callback: (error: Error | null) => void): void {
this.cancelReservation?.();
this.cancelReservation = null;
if (this.timer) this.clock.clearTimeout(this.timer);
this.timer = null;
callback(error);
}
}
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);
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);
}