track?.@_currentExplodedTrackIndex is invalid JS syntax — @ is not a valid identifier character. Replaced with track?.['@_currentExplodedTrackIndex'] so the worker process no longer crashes on startup.
241 lines
9 KiB
JavaScript
241 lines
9 KiB
JavaScript
import { join } from 'path';
|
|
import { unlink, writeFile, mkdir, rm } from 'fs/promises';
|
|
import { tmpdir } from 'os';
|
|
import { query } from '../db/client.js';
|
|
import { downloadFromS3, uploadToS3 } from '../s3/client.js';
|
|
import { trimSegment, concatSegments, runFFmpeg } from '../ffmpeg/executor.js';
|
|
import { parseEDL } from '../edl/parser.js';
|
|
import { XMLParser } from 'fast-xml-parser';
|
|
|
|
const S3_BUCKET = process.env.S3_BUCKET || 'wild-dragon';
|
|
|
|
const xmlParser = new XMLParser({
|
|
ignoreAttributes: false,
|
|
attributeNamePrefix: '@_',
|
|
});
|
|
|
|
function parseFcpXml(xmlContent) {
|
|
const doc = xmlParser.parse(xmlContent);
|
|
const sequence = doc?.xmeml?.sequence;
|
|
if (!sequence) throw new Error('Invalid FCP XML: no sequence element');
|
|
|
|
const name = sequence.name || 'Untitled';
|
|
const rate = sequence?.rate?.timebase ? parseInt(sequence.rate.timebase, 10) : 29.97;
|
|
const width = parseInt(sequence?.media?.video?.format?.samplecharacteristics?.width || 1920, 10);
|
|
const height = parseInt(sequence?.media?.video?.format?.samplecharacteristics?.height || 1080, 10);
|
|
|
|
const clips = [];
|
|
const videoTracks = sequence?.media?.video?.track || [];
|
|
const tracks = Array.isArray(videoTracks) ? videoTracks : [videoTracks];
|
|
|
|
for (const track of tracks) {
|
|
const trackNum = parseInt(track?.['@_currentExplodedTrackIndex'] || 0, 10);
|
|
const trackItems = track?.clipitem || [];
|
|
const items = Array.isArray(trackItems) ? trackItems : [trackItems];
|
|
|
|
for (const item of items) {
|
|
if (!item) continue;
|
|
const fileUrl = item?.file?.name || item?.file?.pathurl || '';
|
|
const fileName = fileUrl.split('/').pop() || fileUrl.split('\\').pop() || 'unknown';
|
|
const srcIn = parseFrame(item?.in?.toString() || '0', rate);
|
|
const srcOut = parseFrame(item?.out?.toString() || '0', rate);
|
|
const recIn = parseFrame(item?.start?.toString() || '0', rate);
|
|
const recOut = parseFrame(item?.end?.toString() || '0', rate);
|
|
const duration = parseFrame(item?.duration?.toString() || '0', rate);
|
|
|
|
if (srcOut <= srcIn || recOut <= recIn) continue;
|
|
|
|
clips.push({
|
|
trackIndex: trackNum,
|
|
fileName,
|
|
fileUrl,
|
|
sourceInFrames: srcIn,
|
|
sourceOutFrames: srcOut,
|
|
timelineInFrames: recIn,
|
|
timelineOutFrames: recOut,
|
|
duration,
|
|
});
|
|
}
|
|
}
|
|
|
|
return { name, frameRate: rate, width, height, clips };
|
|
}
|
|
|
|
function parseFrame(value, fps) {
|
|
// FCP XML stores timecode or frame count
|
|
const trimmed = value.trim();
|
|
// If it's a plain number, return as-is
|
|
if (/^\d+$/.test(trimmed)) return parseInt(trimmed, 10);
|
|
// HH:MM:SS:FF or HH:MM:SS;FF
|
|
const parts = trimmed.split(/[:;]/);
|
|
if (parts.length === 4) {
|
|
const hh = parseInt(parts[0], 10);
|
|
const mm = parseInt(parts[1], 10);
|
|
const ss = parseInt(parts[2], 10);
|
|
const ff = parseInt(parts[3], 10);
|
|
return hh * 3600 * fps + mm * 60 * fps + ss * fps + ff;
|
|
}
|
|
return 0;
|
|
}
|
|
|
|
export const conformWorker = async (job) => {
|
|
const { edl, fcpXml, projectId, sequenceName, frameRate, codec, quality, resolution, audio } = job.data;
|
|
const jobId = job.id;
|
|
|
|
const tmpDir = tmpdir();
|
|
const segmentsDir = join(tmpDir, `segments-${jobId}`);
|
|
const segmentListPath = join(tmpDir, `segments-${jobId}.txt`);
|
|
const outputPath = join(tmpDir, `output-${jobId}.mp4`);
|
|
|
|
try {
|
|
let edits = [];
|
|
let seqName = sequenceName || 'Conformed';
|
|
let seqFps = parseFloat(frameRate) || 29.97;
|
|
|
|
// Parse input — accept EDL, FCP XML, or structured JSON
|
|
if (edl) {
|
|
await job.updateProgress(5);
|
|
console.log(`[conform] Parsing EDL for job ${jobId}`);
|
|
edits = parseEDL(edl).map((e, i) => ({
|
|
editNumber: e.editNumber || i + 1,
|
|
reelName: e.reelName,
|
|
sourceIn: e.sourceIn,
|
|
sourceOut: e.sourceOut,
|
|
}));
|
|
} else if (fcpXml) {
|
|
await job.updateProgress(5);
|
|
console.log(`[conform] Parsing FCP XML for job ${jobId}`);
|
|
const parsed = parseFcpXml(fcpXml);
|
|
seqName = parsed.name || seqName;
|
|
seqFps = parsed.frameRate || seqFps;
|
|
edits = parsed.clips.map((c, i) => ({
|
|
editNumber: i + 1,
|
|
reelName: c.fileName,
|
|
sourceIn: c.sourceInFrames,
|
|
sourceOut: c.sourceOutFrames,
|
|
}));
|
|
} else {
|
|
throw new Error('No input provided — expected edl or fcpXml in job data');
|
|
}
|
|
|
|
await mkdir(segmentsDir, { recursive: true });
|
|
|
|
let processedEdits = 0;
|
|
const concatList = [];
|
|
|
|
for (const edit of edits) {
|
|
await job.updateProgress(Math.min(5 + (processedEdits / edits.length) * 50, 55));
|
|
console.log(`[conform] Processing edit ${edit.editNumber}: ${edit.reelName}`);
|
|
|
|
// BUG FIX #9: Scope asset lookup by project_id to prevent cross-project
|
|
// collisions when two projects contain assets with the same filename.
|
|
let assetRes;
|
|
if (projectId) {
|
|
assetRes = await query(
|
|
`SELECT id, original_s3_key FROM assets
|
|
WHERE filename = $1 AND project_id = $2
|
|
LIMIT 1`,
|
|
[edit.reelName, projectId]
|
|
);
|
|
// Fall back to unscoped lookup if no match in the current project
|
|
// (EDL reel names may reference assets not yet assigned to a project)
|
|
if (assetRes.rows.length === 0) {
|
|
assetRes = await query(
|
|
'SELECT id, original_s3_key FROM assets WHERE filename = $1 LIMIT 1',
|
|
[edit.reelName]
|
|
);
|
|
}
|
|
} else {
|
|
assetRes = await query(
|
|
'SELECT id, original_s3_key FROM assets WHERE filename = $1 LIMIT 1',
|
|
[edit.reelName]
|
|
);
|
|
}
|
|
|
|
if (assetRes.rows.length === 0) {
|
|
throw new Error(`Asset not found for reel: ${edit.reelName}`);
|
|
}
|
|
|
|
const { original_s3_key: sourceKey } = assetRes.rows[0];
|
|
const segmentInputPath = join(segmentsDir, `segment-${edit.editNumber}-src`);
|
|
const segmentOutputPath = join(segmentsDir, `segment-${edit.editNumber}.mov`);
|
|
|
|
console.log(`[conform] Downloading segment ${edit.editNumber} from S3 (${sourceKey})`);
|
|
await downloadFromS3(S3_BUCKET, sourceKey, segmentInputPath);
|
|
|
|
console.log(`[conform] Trimming ${edit.editNumber}: ${edit.sourceIn} → ${edit.sourceOut}`);
|
|
await trimSegment(segmentInputPath, segmentOutputPath, edit.sourceIn, edit.sourceOut);
|
|
|
|
concatList.push(segmentOutputPath);
|
|
await unlink(segmentInputPath).catch(() => {});
|
|
processedEdits++;
|
|
}
|
|
|
|
await job.updateProgress(60);
|
|
console.log(`[conform] Writing concat list for ${concatList.length} segments`);
|
|
const concatContent = concatList.map(p => `file '${p}'`).join('\n');
|
|
await writeFile(segmentListPath, concatContent, 'utf-8');
|
|
|
|
await job.updateProgress(70);
|
|
console.log(`[conform] Concatenating segments for job ${jobId}`);
|
|
|
|
// Use re-encode instead of stream copy for consistent output
|
|
const audioFlag = audio === 'include' ? ['-c:a', 'aac'] : ['-an'];
|
|
await runFFmpeg([
|
|
'-f', 'concat',
|
|
'-safe', '0',
|
|
'-i', segmentListPath,
|
|
'-c:v', codec === 'prores' ? 'prores_ks' : codec === 'h265' ? 'libx265' : 'libx264',
|
|
'-preset', quality === 'high' ? 'slow' : quality === 'broadcast' ? 'veryslow' : 'fast',
|
|
'-crf', quality === 'broadcast' ? '18' : quality === 'high' ? '23' : '28',
|
|
...audioFlag,
|
|
'-y', outputPath,
|
|
]);
|
|
|
|
await job.updateProgress(85);
|
|
const outputKey = `jobs/${jobId}/conformed.mp4`;
|
|
console.log(`[conform] Uploading output to ${outputKey}`);
|
|
await uploadToS3(S3_BUCKET, outputKey, outputPath);
|
|
|
|
// Register the conformed output as a new asset
|
|
const assetRes = await query(
|
|
`INSERT INTO assets (project_id, filename, display_name, media_type, status, original_s3_key, codec, resolution, fps, duration_ms, conform_source_sequence_id)
|
|
VALUES ($1, $2, $3, 'video', 'ready', $4, $5, $6, $7, $8, $9) RETURNING id`,
|
|
[
|
|
projectId || null,
|
|
`conformed-${seqName.replace(/[^a-z0-9]/gi, '_')}.mp4`,
|
|
`Conformed: ${seqName}`,
|
|
outputKey,
|
|
codec === 'prores' ? 'prores' : codec === 'h265' ? 'hevc' : 'h264',
|
|
resolution !== 'match' ? resolution : '1920x1080',
|
|
seqFps,
|
|
null,
|
|
job.data.sequenceId || null,
|
|
]
|
|
);
|
|
|
|
await job.updateProgress(100);
|
|
console.log(`[conform] Job ${jobId} complete → asset ${assetRes.rows[0].id}`);
|
|
|
|
return { jobId, outputKey, assetId: assetRes.rows[0].id };
|
|
|
|
} catch (error) {
|
|
console.error(`[conform] Error in job ${jobId}:`, error);
|
|
// BUG FIX #1: Mark the output asset (if any) as 'error' so the UI doesn't
|
|
// show a perpetually-spinning 'processing' state when the conform fails.
|
|
// We don't have an assetId until the INSERT succeeds, so target by job key.
|
|
await query(
|
|
`UPDATE assets
|
|
SET status = 'error', updated_at = NOW()
|
|
WHERE original_s3_key = $1`,
|
|
[`jobs/${jobId}/conformed.mp4`]
|
|
).catch(e => console.error('[conform] Failed to mark asset error:', e.message));
|
|
throw error;
|
|
} finally {
|
|
await Promise.all([
|
|
unlink(segmentListPath).catch(() => {}),
|
|
unlink(outputPath).catch(() => {}),
|
|
rm(segmentsDir, { recursive: true, force: true }).catch(() => {}),
|
|
]);
|
|
}
|
|
};
|