← Files ModRetro Chromatic PluginARCHIVED FILE
dist/device-capture/linux-helper.js
24.1 KB · Oct 2, 2026 · 00:37 UTC
import { createHash, randomUUID } from "node:crypto";
import { constants } from "node:fs";
import { lstat, mkdtemp, open, realpath, unlink } from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import { NativeHelperError, verifyExecutable } from "./native-helper.js";
import { NativeCaptureFiles } from "./native-files.js";
import { scanChromaticVideoNodes, openBoundLinuxVideo } from "./linux-device.js";
import { LinuxFrameStream } from "./linux-frame-stream.js";
import { LinuxPacketStats } from "./linux-packet-stats.js";
import { LinuxCaptureProcess } from "./linux-process.js";
import { NATIVE_LIMITS } from "./native-protocol.js";
const permissions = { video: "not_required", audio: "not_supported" };
const prefix = ["-hide_banner", "-loglevel", "info", "-nostats", "-xerror", "-max_alloc", "67108864"];
const deviceAccess = { scan: scanChromaticVideoNodes, open: openBoundLinuxVideo, launch: LinuxCaptureProcess.launch };
/** Linux implementation of the same capture protocol, using a pinned V4L2/VP8 executable. */
export class LinuxNativeHelper {
executable;
events;
io;
ready;
sequence = 0;
nodes = new Map();
inventory;
node;
producer;
recording;
recordings = 0;
screenshots = 0;
frameId = 0;
sourceGeneration = 0;
latest;
latestImage;
previewFiles = [];
connection = "disconnected";
lastError = null;
closing;
cancelled = false;
unresolved;
operation;
processes = [];
owned = [];
diagnostics = { bytes: 0 };
constructor(executable, directory, events, io = deviceAccess) {
this.executable = executable;
this.events = events;
this.io = io;
this.ready = { event: "ready", protocolVersion: 1, sessionId: randomUUID(), stateSequence: 0, sessionDir: directory, platform: "Linux", limits: NATIVE_LIMITS };
}
static async launch(options) {
if (process.platform !== "linux")
throw new Error("The Linux capture backend requires a Linux host.");
await verifyExecutable(options.executable, options.sha256);
const directory = await mkdtemp(path.join(await realpath(os.tmpdir()), "chromatic-native-capture-"));
await NativeCaptureFiles.create(directory);
return new LinuxNativeHelper(options.executable, directory, options.onEvent);
}
/** Inert test seam. The production launch path always uses the fixed device implementation above. */
static async attach(directory, events, io) {
await NativeCaptureFiles.create(directory);
return new LinuxNativeHelper("/test-only/ffmpeg", directory, events, io);
}
next() { return { sessionId: this.ready.sessionId, stateSequence: ++this.sequence }; }
assertOpen() { if (this.cancelled)
throw new Error("Linux capture is closing."); }
remember(process) {
this.processes.push(process.record);
this.owned.push(process);
return process;
}
capacity() { if (this.processes.length >= 64)
throw new Error("This capture session reached its process-evidence bound. Close it before starting another session."); }
status() {
const r = this.recording;
return { ...this.next(), connection: this.connection,
selection: this.node ? { videoDeviceId: this.node.id, videoLabel: this.node.label, player: this.node.player } : null,
permissions, source: this.latest ? { width: this.latest.width, height: this.latest.height, firstPTS: this.producer?.firstPTS ?? null, lastPTS: this.latest.sourcePTS, videoSamples: this.producer?.frameCount ?? 0, audioSamples: 0 } : null,
recording: r?.final ?? (r ? { captureId: r.id, state: "recording", durationMs: Math.min(performance.now() - r.started, r.requested + 15000), requestedDurationMs: r.requested, deadlineUnixMs: r.deadline,
videoSamples: r.stats.count, audioSamples: 0, droppedVideoSamples: null, droppedAudioSamples: 0 } : null),
latestFrame: this.latestImage ? { frameId: this.latestImage.frameId, basename: this.latestImage.basename, bytes: this.latestImage.bytes, sha256: this.latestImage.sha256,
width: this.latestImage.width, height: this.latestImage.height, sourcePTS: this.latestImage.sourcePTS, format: "jpeg" } : null, lastError: this.lastError,
processes: this.processes.map(record => ({ ...record })), sampleEvidence: "encoded-output", sourceGeneration: this.sourceGeneration, sourceClock: "ffmpeg-relative-input" };
}
async isCapture(node) {
this.capacity();
this.assertOpen();
const fd = await this.io.open(node);
let process;
try {
process = this.remember(this.io.launch(this.executable, [...prefix, "-nostdin", "-f", "v4l2", "-list_formats", "all", "-i", "/proc/self/fd/6"], "video-capabilities", { 1: bytes => { if (bytes.length)
throw new Error("Unexpected video capability output."); } }, fd.fd, this.diagnostics));
}
finally {
try {
await fd.close();
}
catch {
throw new NativeHelperError("Linux video descriptor closure is unconfirmed.", true, "LINUX_CLOSE_UNKNOWN");
}
}
{
const result = await process.wait(5000);
const text = result.record.stderr;
if (result.error || !result.record.spawnObserved || !Number.isSafeInteger(result.record.pid) || result.record.pid <= 0 || !result.record.streamsSettled || !result.record.closedAt || result.record.signal !== null)
throw result.error ?? new Error("Linux video query did not close completely.");
if (/Not a video capture device|The device does not support the streaming I\/O method/.test(text))
return false;
if (result.record.exitCode !== 0)
throw new Error(`Cannot query the selected Chromatic video node. Check Linux USB/video access. ${text.slice(-1900)}`);
// FFmpeg can suppress the terminal ENUM_FMT errno. Empty output is not
// evidence of a metadata node; only the explicit capability errors above
// establish that classification.
if (!/(?:Raw |Compressed):\s+(?!Unsupported\b)\S+\s+:.*:\s+\d+x\d+/.test(text))
throw new Error("The Linux video query returned no supported formats. Device availability is unconfirmed; reconnect or check the device before another attempt.");
return true;
}
}
async list() {
if (this.producer)
throw new Error("Disconnect before refreshing Linux video devices.");
const { total, candidates } = await this.io.scan();
const accepted = [];
for (const node of candidates.slice(0, 8))
if (await this.isCapture(node))
accepted.push(node);
this.assertOpen();
this.nodes = new Map(accepted.slice(0, 4).map(node => [node.id, node]));
this.inventory = { ...this.next(), inventoryId: randomUUID(), permissions,
videoDevices: [...this.nodes.values()].map(node => ({ deviceId: node.id, label: node.label, player: node.player, formats: [] })), audioDevices: [],
counts: { video: { total, unlabeled: 0, unsupported: total - candidates.length + Math.min(candidates.length, 8) - accepted.length, omitted: Math.max(candidates.length - 8, 0) + Math.max(accepted.length - 4, 0) }, audio: { total: 0, unlabeled: 0, unsupported: 0, omitted: 0 } } };
return this.inventory;
}
async startProducer(recording) {
this.capacity();
this.assertOpen();
if (this.producer || !this.node)
throw new Error("The previous Linux producer must close before another can start.");
const node = this.node, fd = await this.io.open(node);
let resolveFirst;
const first = new Promise(resolve => { resolveFirst = resolve; });
let producer;
const preview = new LinuxFrameStream(frame => {
if (this.producer !== producer)
return;
producer.frameCount++;
producer.firstPTS ??= frame.sourcePTS;
producer.lastFrameAt = Date.now();
this.frameId++;
this.latest = frame;
resolveFirst();
});
const args = [...prefix, "-copyts", "-start_at_zero", "-f", "v4l2", "-i", "/proc/self/fd/6"];
const pipes = { 1: bytes => preview.pushBytes(bytes), 3: bytes => preview.pushMetadata(bytes) };
if (recording) {
args.push("-map", "0:v:0", "-an", "-vf", "scale=iw*8:ih*8:flags=neighbor", "-c:v", "libvpx", "-deadline", "realtime", "-cpu-used", "8", "-b:v", "8M", "-enc_time_base", "1:1000000", "-fps_mode", "passthrough", "-f", "tee", "[f=webm]pipe:4|[f=framehash:hash=sha256:format_version=1:flush_packets=1]pipe:5");
pipes[4] = async (bytes) => {
if (recording.bytes + bytes.length > NATIVE_LIMITS.maxRecordingBytes) {
recording.reason = "limit";
throw new Error("Linux recording reached its byte limit; its partial file is retained.");
}
let offset = 0;
while (offset < bytes.length) {
const result = await recording.file.write(bytes, offset, bytes.length - offset);
if (!result.bytesWritten)
throw new Error("Linux recording write stopped.");
recording.digest.update(bytes.subarray(offset, offset + result.bytesWritten));
recording.bytes += result.bytesWritten;
offset += result.bytesWritten;
}
};
pipes[5] = bytes => recording.stats.push(bytes);
}
args.push("-map", "0:v:0", "-an", "-c:v", "mjpeg", "-q:v", "3", "-pix_fmt", "yuvj420p", "-enc_time_base", "1:1000000", "-fps_mode", "passthrough", "-f", "tee", "[f=image2pipe:flush_packets=1]pipe:1|[f=framehash:hash=sha256:format_version=1:flush_packets=1]pipe:3");
try {
this.assertOpen();
if (recording && performance.now() >= recording.expires)
throw new Error("Recording deadline expired while opening the selected device.");
const process = this.remember(this.io.launch(this.executable, args, recording ? "recording" : "preview", pipes, fd.fd, this.diagnostics));
producer = { process, preview, first, recording, frameCount: 0, lastFrameAt: Date.now(), generation: ++this.sourceGeneration };
this.producer = producer;
if (recording)
recording.generation = producer.generation;
producer.watchdog = setInterval(() => { if (Date.now() - producer.lastFrameAt >= 10000)
process.fail(new Error("The selected Chromatic stopped providing complete video frames.")); }, 1000);
producer.watchdog.unref?.();
}
finally {
try {
await fd.close();
}
catch {
throw new NativeHelperError("Linux video descriptor closure is unconfirmed.", true, "LINUX_CLOSE_UNKNOWN");
}
}
void producer.process.finished.then(() => { if (!producer.stop)
void this.finishProducer(producer, "error").catch(error => { if (error instanceof NativeHelperError && error.uncertain)
this.fatal(error); }); });
let timer;
try {
await Promise.race([first, producer.process.finished.then(() => { throw new Error("Linux capture exited before its first complete frame."); }), new Promise((_, reject) => { timer = setTimeout(() => reject(new Error("No frame arrived from the selected Chromatic.")), 10000); })]);
this.assertOpen();
if (this.producer !== producer)
throw new Error("Linux producer closed during connection.");
this.connection = "connected";
}
catch (error) {
producer.process.fail(error instanceof Error ? error : new Error(String(error)));
await this.finishProducer(producer, "error");
throw error;
}
finally {
clearTimeout(timer);
}
}
fatal(error) {
const message = error instanceof Error ? error.message : String(error);
this.cancelled = true;
this.unresolved = error instanceof Error ? error : new Error(message);
this.lastError = { code: "LINUX_CLOSE_UNKNOWN", message: message.slice(0, 2000) };
this.connection = "failed";
this.events({ event: "fatal", sessionId: this.ready.sessionId, error: this.lastError, devicesReleased: false });
}
finishProducer(producer, reason) {
if (producer.stop)
return producer.stop;
producer.stop = (async () => {
clearInterval(producer.watchdog);
const result = await producer.process.stop();
let error = result.error ?? (reason === "error" ? new Error("Linux capture ended before its requested stop or deadline.") : undefined);
if (!result.record.spawnObserved || result.record.exitCode !== 0 || result.record.signal !== null)
error ??= new Error(`Linux capture exited with ${result.record.signal ?? result.record.exitCode ?? "unknown status"}.`);
try {
producer.preview.finish();
producer.recording?.stats.finish();
}
catch (failure) {
error ??= failure;
}
if (this.producer === producer) {
this.producer = undefined;
this.latest = undefined;
this.latestImage = undefined;
this.connection = error ? "failed" : "disconnected";
}
const r = producer.recording;
if (r) {
clearTimeout(r.timer);
r.error ??= error;
try {
await r.file.sync();
}
catch (failure) {
r.error ??= failure;
}
await r.file.close();
const durationMs = Math.min(Math.max(0, performance.now() - r.started), r.requested + 15000);
r.final = { captureId: r.id, state: r.error ? r.bytes ? "partial" : "failed" : "complete", basename: r.bytes ? r.basename : null, format: "webm", bytes: r.bytes,
sha256: r.bytes ? r.digest.digest("hex") : null, durationMs, requestedDurationMs: r.requested, videoSamples: r.stats.count, audioSamples: 0,
droppedVideoSamples: null, droppedAudioSamples: 0, reason: r.reason === "limit" ? "limit" : reason,
error: r.error ? { code: "LINUX_RECORDING_FAILED", message: r.error.message.slice(0, 2000) } : null, devicesReleased: true,
processes: this.processes.map(record => ({ ...record })), sampleEvidence: "encoded-output",
encodedDurationMs: Math.max(0, ((r.stats.lastPTS ?? 0) - (r.stats.firstPTS ?? 0) + r.stats.lastDuration) * 1000), sourceGeneration: r.generation, sourceClock: "ffmpeg-relative-input" };
this.events({ event: "recording_finished", ...this.next(), recording: r.final });
}
if (error)
this.lastError = { code: "LINUX_CAPTURE_FAILED", message: error.message.slice(0, 2000) };
this.events({ event: "status", status: this.status() });
if (error && !r)
throw new NativeHelperError(error.message, false, "LINUX_CAPTURE_FAILED");
if (r && !error && !r.error && !this.cancelled && ["stop", "duration"].includes(reason))
await this.startProducer();
})();
return producer.stop;
}
async writeImage(frame, png = false) {
const frameId = this.frameId, sourceGeneration = this.sourceGeneration;
let bytes = frame.bytes;
if (png) {
this.capacity();
if (++this.screenshots > NATIVE_LIMITS.maxScreenshots)
throw new Error("Screenshot limit reached.");
const chunks = [];
let length = 0;
const process = this.remember(this.io.launch(this.executable, [...prefix, "-f", "image2pipe", "-c:v", "mjpeg", "-i", "pipe:0", "-frames:v", "1", "-c:v", "png", "-f", "image2pipe", "pipe:1"], "screenshot", { 1: chunk => { length += chunk.length; if (length > NATIVE_LIMITS.maxScreenshotBytes)
throw new Error("PNG exceeds its byte limit."); chunks.push(chunk); } }, undefined, this.diagnostics));
process.input(bytes);
const result = await process.wait(5000);
if (result.error || result.record.exitCode !== 0 || result.record.signal !== null)
throw result.error ?? new Error("Linux PNG encoding failed.");
bytes = Buffer.concat(chunks);
}
const basename = `${png ? "snapshot" : "preview"}-${randomUUID()}.${png ? "png" : "jpeg"}`;
const file = await open(path.join(this.ready.sessionDir, basename), constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL | constants.O_NOFOLLOW, 0o600);
let stat;
try {
await file.writeFile(bytes);
stat = await file.stat();
}
finally {
await file.close();
}
if (!png) {
this.previewFiles.push({ name: basename, stat });
while (this.previewFiles.length > 2) {
const old = this.previewFiles.shift(), filename = path.join(this.ready.sessionDir, old.name), current = await lstat(filename);
if (current.ino !== old.stat.ino || current.dev !== old.stat.dev || !current.isFile() || current.nlink !== 1)
throw new Error("Owned preview file changed before retirement.");
await unlink(filename);
}
}
const image = { ...this.next(), frameId, basename, bytes: bytes.length, sha256: createHash("sha256").update(bytes).digest("hex"), width: frame.width, height: frame.height, sourcePTS: frame.sourcePTS, format: png ? "png" : "jpeg", sourceGeneration, sourceClock: "ffmpeg-relative-input" };
if (!png)
this.latestImage = image;
return image;
}
request(action, fields = {}) {
if (action === "close")
return this.close();
if (this.operation)
return Promise.reject(new Error("A Linux capture operation is still finishing."));
const operation = this.perform(action, fields);
this.operation = operation;
return operation.catch(error => { if (error instanceof NativeHelperError && error.uncertain) {
this.cancelled = true;
this.fatal(error);
} throw error; })
.finally(() => { if (this.operation === operation)
this.operation = undefined; });
}
async perform(action, fields) {
this.assertOpen();
if (action === "list")
return this.list();
if (action === "status" || action === "permission_status") {
if (this.latest && this.latestImage?.frameId !== this.frameId)
await this.writeImage(this.latest);
return this.status();
}
if (action === "request_permission")
throw new Error("Linux capture does not request macOS permissions.");
if (action === "connect") {
const selection = fields.selection;
if (selection?.audioDeviceId)
throw new Error("Linux USB audio is not supported in this build.");
const node = selection?.videoDeviceId ? this.nodes.get(selection.videoDeviceId) : undefined;
if (!node || fields.inventoryId !== this.inventory?.inventoryId || this.producer)
throw new Error("Refresh and select the exact Linux Chromatic before connecting.");
this.node = node;
this.connection = "connecting";
this.lastError = null;
await this.startProducer();
return this.status();
}
if (action === "live_frame" || action === "screenshot") {
if (!this.latest || this.connection !== "connected")
throw new NativeHelperError("No complete Linux video frame has arrived.", false, "NO_VIDEO_FRAME");
return this.writeImage(this.latest, action === "screenshot");
}
if (action === "start_recording") {
const duration = fields.durationMs, id = fields.captureId;
if (!Number.isInteger(duration) || Number(duration) < 1000 || Number(duration) > NATIVE_LIMITS.maxDurationMs || typeof id !== "string" || !/^[a-f0-9-]{36}$/.test(id))
throw new Error("Invalid Linux recording request.");
if (!this.producer || this.connection !== "connected" || this.producer.recording || this.recordings >= 2)
throw new Error("Connect before recording; at most two recordings are allowed per session.");
// Start the one deadline BEFORE preview closure and re-open, never reset it.
const started = performance.now(), deadline = Date.now() + Number(duration), expires = started + Number(duration);
await this.finishProducer(this.producer, "stop");
this.assertOpen();
if (performance.now() >= expires)
throw new Error("Recording deadline expired during preview closure.");
const basename = `capture-${id}.webm`, file = await open(path.join(this.ready.sessionDir, basename), constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL | constants.O_NOFOLLOW, 0o600);
const r = { id, basename, file, bytes: 0, digest: createHash("sha256"), stats: new LinuxPacketStats(), requested: Number(duration), started, deadline, expires, reason: "stop" };
this.recording = r;
this.recordings++;
r.timer = setTimeout(() => { if (this.producer?.recording === r)
void this.finishProducer(this.producer, "duration").catch(error => this.fatal(error)); }, expires - performance.now());
try {
await this.startProducer(r);
}
catch (error) {
clearTimeout(r.timer);
if (!r.final && !this.producer && !(error instanceof NativeHelperError && error.uncertain)) {
await file.close();
r.final = { captureId: r.id, state: "failed", basename: null, format: "webm", bytes: 0, sha256: null, durationMs: Math.min(performance.now() - started, r.requested + 15000), requestedDurationMs: r.requested,
videoSamples: 0, audioSamples: 0, droppedVideoSamples: null, droppedAudioSamples: 0, reason: "error", error: { code: "LINUX_RECORDING_NOT_STARTED", message: String(error instanceof Error ? error.message : error).slice(0, 2000) }, devicesReleased: true, sampleEvidence: "encoded-output", processes: this.processes.map(record => ({ ...record })) };
this.events({ event: "recording_finished", ...this.next(), recording: r.final });
}
throw error;
}
return this.status();
}
if (action === "stop_recording" || action === "disconnect") {
if (this.producer)
await this.finishProducer(this.producer, action === "disconnect" ? "disconnect" : "stop");
return this.status();
}
throw new Error("Unknown Linux capture action.");
}
close() {
if (this.closing)
return this.closing;
this.cancelled = true;
this.closing = (async () => {
await this.operation?.catch(() => { });
if (this.producer)
await this.finishProducer(this.producer, "close").catch(error => { if (!(error instanceof NativeHelperError) || error.uncertain)
throw error; });
for (const process of this.owned)
if (!process.record.closedAt || !process.record.streamsSettled)
await process.stop();
if (this.recording && !this.recording.final)
await this.recording.file.close();
if (this.unresolved)
throw this.unresolved;
this.connection = "closed";
this.events({ event: "closed", ...this.next(), devicesReleased: true });
})();
return this.closing;
}
}
//# sourceMappingURL=linux-helper.js.mapSHA-256: 747bbdd7264946d915c72d0630ef998ade06e1deaf01289445ca934970b1709c