|
|
@@ -7,6 +7,7 @@ import type {
|
|
|
ExportFile,
|
|
|
} from "@pipeline/shared";
|
|
|
import { PLATFORM_PRESETS } from "@pipeline/shared";
|
|
|
+import { createLogger } from "@pipeline/shared/node";
|
|
|
import { parseText } from "./stages/parse.js";
|
|
|
import { generateTTS } from "./stages/tts.js";
|
|
|
import { resolveAssets } from "./stages/assets.js";
|
|
|
@@ -73,13 +74,18 @@ export async function runPipeline(
|
|
|
callbacks?: PipelineCallbacks
|
|
|
): Promise<PipelineJob> {
|
|
|
const jobId = randomUUID();
|
|
|
+ const log = createLogger("pipeline");
|
|
|
+ const t0 = Date.now();
|
|
|
const workDir = join(config.output.dir, "tmp", jobId);
|
|
|
await mkdir(workDir, { recursive: true });
|
|
|
|
|
|
// Opportunistic cache sweep: prune outputs older than retentionDays.
|
|
|
// Runs on every job (service jobs are low-frequency); never throws.
|
|
|
if (config.output.retentionDays && config.output.retentionDays > 0) {
|
|
|
- await cleanupExpiredOutput(config.output.dir, config.output.retentionDays);
|
|
|
+ const cleaned = await cleanupExpiredOutput(config.output.dir, config.output.retentionDays);
|
|
|
+ log.debug(
|
|
|
+ `cleanup scanned=${cleaned.scanned} removed=${cleaned.removed} freed=${Math.round(cleaned.bytesFreed / 1024 / 1024)}MB`
|
|
|
+ );
|
|
|
}
|
|
|
|
|
|
const platforms = input.platforms;
|
|
|
@@ -91,19 +97,27 @@ export async function runPipeline(
|
|
|
updatedAt: Date.now(),
|
|
|
};
|
|
|
|
|
|
+ log.info(
|
|
|
+ `job ${jobId} start template=${input.template} platforms=${platforms.join(",")} ` +
|
|
|
+ `chars=${input.text.length} skipLlm=${!!config.skipLlm} skipTts=${!!config.skipTts}`
|
|
|
+ );
|
|
|
+
|
|
|
try {
|
|
|
// Stage 1: Parse (shared across all platforms)
|
|
|
callbacks?.onStageStart?.("parse");
|
|
|
+ const tParse = Date.now();
|
|
|
const parsed = await parseText(input.text, input.template, {
|
|
|
llm: config.llm,
|
|
|
skipLlm: config.skipLlm,
|
|
|
source: input.source,
|
|
|
});
|
|
|
job.parsed = parsed;
|
|
|
+ log.info(`parse ok scenes=${parsed.scenes?.length ?? 0} (${Date.now() - tParse}ms)`);
|
|
|
callbacks?.onStageComplete?.("parse");
|
|
|
|
|
|
// Stage 2: TTS (shared across all platforms)
|
|
|
callbacks?.onStageStart?.("tts");
|
|
|
+ const tTts = Date.now();
|
|
|
const tts = await generateTTS(parsed, workDir, {
|
|
|
provider: config.tts.provider,
|
|
|
voiceId: config.tts.voiceId || input.voiceId,
|
|
|
@@ -114,6 +128,10 @@ export async function runPipeline(
|
|
|
alignment: config.alignment,
|
|
|
});
|
|
|
job.tts = tts;
|
|
|
+ log.info(
|
|
|
+ `tts ok provider=${config.tts.provider} scenes=${tts.scenes?.length ?? 0} ` +
|
|
|
+ `duration=${(tts.totalDurationSeconds ?? 0).toFixed(1)}s (${Date.now() - tTts}ms)`
|
|
|
+ );
|
|
|
callbacks?.onStageComplete?.("tts");
|
|
|
|
|
|
// Stages 3-6: Per-platform loop
|
|
|
@@ -138,7 +156,9 @@ export async function runPipeline(
|
|
|
// Stage 5: Render
|
|
|
callbacks?.onStageStart?.(`render:${platform}`);
|
|
|
const renderOutput = join(workDir, `render-${platform}.mp4`);
|
|
|
+ const tRender = Date.now();
|
|
|
await renderVideo(composed, renderOutput, config.templates.entryPoint, config.assets.root);
|
|
|
+ log.info(`render ${platform} ok -> ${renderOutput} (${Date.now() - tRender}ms)`);
|
|
|
callbacks?.onStageComplete?.(`render:${platform}`);
|
|
|
|
|
|
// Stage 6: Export
|
|
|
@@ -149,6 +169,11 @@ export async function runPipeline(
|
|
|
config.output.dir,
|
|
|
jobId
|
|
|
);
|
|
|
+ const expFile = exported.files[0];
|
|
|
+ log.info(
|
|
|
+ `export ${platform} -> ${expFile?.filePath} ` +
|
|
|
+ `(${((expFile?.fileSizeBytes ?? 0) / 1024 / 1024).toFixed(1)}MB)`
|
|
|
+ );
|
|
|
allExportFiles.push(...exported.files);
|
|
|
callbacks?.onStageComplete?.(`export:${platform}`);
|
|
|
}
|
|
|
@@ -158,6 +183,9 @@ export async function runPipeline(
|
|
|
|
|
|
// Stage 7: Publish — upload to OSS + notify Feishu. Opt-in via config.
|
|
|
if (!config.skipPublish && config.publish) {
|
|
|
+ log.info(
|
|
|
+ `publish oss=${config.publish.oss ? "on" : "off"} feishu=${config.publish.feishu ? "on" : "off"}`
|
|
|
+ );
|
|
|
callbacks?.onStageStart?.("publish");
|
|
|
const publishCtx = {
|
|
|
jobId,
|
|
|
@@ -176,6 +204,7 @@ export async function runPipeline(
|
|
|
};
|
|
|
if (!outcome.ok && outcome.error) {
|
|
|
job.publishError = outcome.error;
|
|
|
+ log.error(`publish failed: ${outcome.error}`);
|
|
|
callbacks?.onError?.("publish", new Error(outcome.error));
|
|
|
}
|
|
|
callbacks?.onStageComplete?.("publish");
|
|
|
@@ -183,12 +212,14 @@ export async function runPipeline(
|
|
|
// Publish must never mask a successful render.
|
|
|
const msg = err instanceof Error ? err.message : String(err);
|
|
|
job.publishError = msg;
|
|
|
+ log.error(`publish error: ${msg}`);
|
|
|
callbacks?.onError?.("publish", err instanceof Error ? err : new Error(msg));
|
|
|
}
|
|
|
}
|
|
|
} catch (err) {
|
|
|
job.status = "failed";
|
|
|
job.error = err instanceof Error ? err.message : String(err);
|
|
|
+ log.error(`job ${jobId} failed: ${job.error}`);
|
|
|
callbacks?.onError?.("pipeline", err instanceof Error ? err : new Error(String(err)));
|
|
|
|
|
|
// Render failed before any file existed — push failure info to Feishu.
|
|
|
@@ -201,6 +232,11 @@ export async function runPipeline(
|
|
|
}
|
|
|
}
|
|
|
|
|
|
+ const elapsed = ((Date.now() - t0) / 1000).toFixed(1);
|
|
|
+ log.info(
|
|
|
+ `job ${jobId} ${job.status} in ${elapsed}s` +
|
|
|
+ (job.publishError ? ` publishError=${job.publishError.slice(0, 160)}` : "")
|
|
|
+ );
|
|
|
job.updatedAt = Date.now();
|
|
|
return job;
|
|
|
}
|