pipeline.ts 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147
  1. import { randomUUID } from "node:crypto";
  2. import { mkdir } from "node:fs/promises";
  3. import { join } from "node:path";
  4. import type {
  5. PipelineInput,
  6. PipelineJob,
  7. ExportFile,
  8. } from "@pipeline/shared";
  9. import { PLATFORM_PRESETS } from "@pipeline/shared";
  10. import { parseText } from "./stages/parse.js";
  11. import { generateTTS } from "./stages/tts.js";
  12. import { resolveAssets } from "./stages/assets.js";
  13. import { composeProject } from "./stages/compose.js";
  14. import { renderVideo } from "./stages/render.js";
  15. import { exportVideo } from "./stages/export.js";
  16. export interface PipelineConfig {
  17. branding: {
  18. channelName: string;
  19. };
  20. llm: {
  21. baseURL?: string;
  22. apiKey?: string;
  23. model: string;
  24. };
  25. tts: {
  26. provider: string;
  27. voiceId?: string;
  28. model?: string;
  29. format?: "mp3" | "wav" | "pcm";
  30. speed?: number;
  31. };
  32. alignment?: {
  33. provider: "whisper" | "native";
  34. whisperModel?: string;
  35. language?: string;
  36. };
  37. output: {
  38. dir: string;
  39. };
  40. assets: {
  41. root: string;
  42. inputDir: string;
  43. };
  44. templates: {
  45. entryPoint: string;
  46. };
  47. skipTts?: boolean;
  48. skipLlm?: boolean;
  49. }
  50. export interface PipelineCallbacks {
  51. onStageStart?: (stage: string) => void;
  52. onStageComplete?: (stage: string) => void;
  53. onError?: (stage: string, error: Error) => void;
  54. }
  55. export async function runPipeline(
  56. input: PipelineInput,
  57. config: PipelineConfig,
  58. callbacks?: PipelineCallbacks
  59. ): Promise<PipelineJob> {
  60. const jobId = randomUUID();
  61. const workDir = join(config.output.dir, "tmp", jobId);
  62. await mkdir(workDir, { recursive: true });
  63. const platforms = input.platforms;
  64. const job: PipelineJob = {
  65. id: jobId,
  66. input,
  67. status: "running",
  68. createdAt: Date.now(),
  69. updatedAt: Date.now(),
  70. };
  71. try {
  72. // Stage 1: Parse (shared across all platforms)
  73. callbacks?.onStageStart?.("parse");
  74. const parsed = await parseText(input.text, input.template, {
  75. llm: config.llm,
  76. skipLlm: config.skipLlm,
  77. });
  78. job.parsed = parsed;
  79. callbacks?.onStageComplete?.("parse");
  80. // Stage 2: TTS (shared across all platforms)
  81. callbacks?.onStageStart?.("tts");
  82. const tts = await generateTTS(parsed, workDir, {
  83. provider: config.tts.provider,
  84. voiceId: config.tts.voiceId || input.voiceId,
  85. model: config.tts.model,
  86. format: config.tts.format,
  87. speed: config.tts.speed,
  88. skip: config.skipTts,
  89. alignment: config.alignment,
  90. });
  91. job.tts = tts;
  92. callbacks?.onStageComplete?.("tts");
  93. // Stages 3-6: Per-platform loop
  94. const allExportFiles: ExportFile[] = [];
  95. job.composed = {};
  96. for (const platform of platforms) {
  97. // Stage 3: Assets
  98. callbacks?.onStageStart?.(`assets:${platform}`);
  99. const assets = await resolveAssets(
  100. parsed, workDir, config.assets.root, config.assets.inputDir,
  101. { template: input.template, aspect: PLATFORM_PRESETS[platform].aspect }
  102. );
  103. callbacks?.onStageComplete?.(`assets:${platform}`);
  104. // Stage 4: Compose
  105. callbacks?.onStageStart?.(`compose:${platform}`);
  106. const composed = composeProject(parsed, tts, assets, platform, input.template, config.branding.channelName);
  107. job.composed[platform] = composed;
  108. callbacks?.onStageComplete?.(`compose:${platform}`);
  109. // Stage 5: Render
  110. callbacks?.onStageStart?.(`render:${platform}`);
  111. const renderOutput = join(workDir, `render-${platform}.mp4`);
  112. await renderVideo(composed, renderOutput, config.templates.entryPoint);
  113. callbacks?.onStageComplete?.(`render:${platform}`);
  114. // Stage 6: Export
  115. callbacks?.onStageStart?.(`export:${platform}`);
  116. const exported = await exportVideo(
  117. renderOutput,
  118. composed,
  119. config.output.dir,
  120. jobId
  121. );
  122. allExportFiles.push(...exported.files);
  123. callbacks?.onStageComplete?.(`export:${platform}`);
  124. }
  125. job.exported = { jobId, createdAt: Date.now(), files: allExportFiles };
  126. job.status = "completed";
  127. } catch (err) {
  128. job.status = "failed";
  129. job.error = err instanceof Error ? err.message : String(err);
  130. callbacks?.onError?.("pipeline", err instanceof Error ? err : new Error(String(err)));
  131. }
  132. job.updatedAt = Date.now();
  133. return job;
  134. }