classroom-job-store.ts 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226
  1. import { promises as fs } from 'fs';
  2. import path from 'path';
  3. import type {
  4. ClassroomGenerationProgress,
  5. ClassroomGenerationStep,
  6. GenerateClassroomInput,
  7. GenerateClassroomResult,
  8. } from '@/lib/server/classroom-generation';
  9. import {
  10. CLASSROOM_JOBS_DIR,
  11. ensureClassroomJobsDir,
  12. writeJsonFileAtomic,
  13. } from '@/lib/server/classroom-storage';
  14. export type ClassroomGenerationJobStatus = 'queued' | 'running' | 'succeeded' | 'failed';
  15. export interface ClassroomGenerationJob {
  16. id: string;
  17. status: ClassroomGenerationJobStatus;
  18. step: ClassroomGenerationStep | 'queued' | 'failed';
  19. progress: number;
  20. message: string;
  21. createdAt: string;
  22. updatedAt: string;
  23. startedAt?: string;
  24. completedAt?: string;
  25. inputSummary: {
  26. requirementPreview: string;
  27. language: string;
  28. hasPdf: boolean;
  29. pdfTextLength: number;
  30. pdfImageCount: number;
  31. };
  32. scenesGenerated: number;
  33. totalScenes?: number;
  34. result?: {
  35. classroomId: string;
  36. url: string;
  37. scenesCount: number;
  38. };
  39. error?: string;
  40. }
  41. function jobFilePath(jobId: string) {
  42. return path.join(CLASSROOM_JOBS_DIR, `${jobId}.json`);
  43. }
  44. function buildInputSummary(input: GenerateClassroomInput): ClassroomGenerationJob['inputSummary'] {
  45. return {
  46. requirementPreview:
  47. input.requirement.length > 200 ? `${input.requirement.slice(0, 197)}...` : input.requirement,
  48. language: input.language || 'zh-CN',
  49. hasPdf: !!input.pdfContent,
  50. pdfTextLength: input.pdfContent?.text.length || 0,
  51. pdfImageCount: input.pdfContent?.images.length || 0,
  52. };
  53. }
  54. /** Simple per-job mutex to serialize read-modify-write on the same job file. */
  55. const jobLocks = new Map<string, Promise<void>>();
  56. async function withJobLock<T>(jobId: string, fn: () => Promise<T>): Promise<T> {
  57. const prev = jobLocks.get(jobId) ?? Promise.resolve();
  58. let resolve: () => void;
  59. const next = new Promise<void>((r) => {
  60. resolve = r;
  61. });
  62. jobLocks.set(jobId, next);
  63. try {
  64. await prev;
  65. return await fn();
  66. } finally {
  67. resolve!();
  68. if (jobLocks.get(jobId) === next) jobLocks.delete(jobId);
  69. }
  70. }
  71. /** Max age (ms) before a "running" job without an active runner is considered stale. */
  72. const STALE_JOB_TIMEOUT_MS = 30 * 60 * 1000; // 30 minutes
  73. function markStaleIfNeeded(job: ClassroomGenerationJob): ClassroomGenerationJob {
  74. if (job.status !== 'running') return job;
  75. const updatedAt = new Date(job.updatedAt).getTime();
  76. if (Date.now() - updatedAt > STALE_JOB_TIMEOUT_MS) {
  77. return {
  78. ...job,
  79. status: 'failed',
  80. step: 'failed',
  81. message: 'Job appears stale (no progress update for 30 minutes)',
  82. error: 'Stale job: process may have restarted during generation',
  83. completedAt: new Date().toISOString(),
  84. updatedAt: new Date().toISOString(),
  85. };
  86. }
  87. return job;
  88. }
  89. export function isValidClassroomJobId(jobId: string): boolean {
  90. return /^[a-zA-Z0-9_-]+$/.test(jobId);
  91. }
  92. export async function createClassroomGenerationJob(
  93. jobId: string,
  94. input: GenerateClassroomInput,
  95. ): Promise<ClassroomGenerationJob> {
  96. const now = new Date().toISOString();
  97. const job: ClassroomGenerationJob = {
  98. id: jobId,
  99. status: 'queued',
  100. step: 'queued',
  101. progress: 0,
  102. message: 'Classroom generation job queued',
  103. createdAt: now,
  104. updatedAt: now,
  105. inputSummary: buildInputSummary(input),
  106. scenesGenerated: 0,
  107. };
  108. await ensureClassroomJobsDir();
  109. await writeJsonFileAtomic(jobFilePath(jobId), job);
  110. return job;
  111. }
  112. export async function readClassroomGenerationJob(
  113. jobId: string,
  114. ): Promise<ClassroomGenerationJob | null> {
  115. try {
  116. const content = await fs.readFile(jobFilePath(jobId), 'utf-8');
  117. const job = JSON.parse(content) as ClassroomGenerationJob;
  118. return markStaleIfNeeded(job);
  119. } catch (error) {
  120. if ((error as NodeJS.ErrnoException).code === 'ENOENT') {
  121. return null;
  122. }
  123. throw error;
  124. }
  125. }
  126. export async function updateClassroomGenerationJob(
  127. jobId: string,
  128. patch: Partial<ClassroomGenerationJob>,
  129. ): Promise<ClassroomGenerationJob> {
  130. return withJobLock(jobId, async () => {
  131. const existing = await readClassroomGenerationJob(jobId);
  132. if (!existing) {
  133. throw new Error(`Classroom generation job not found: ${jobId}`);
  134. }
  135. const updated: ClassroomGenerationJob = {
  136. ...existing,
  137. ...patch,
  138. updatedAt: new Date().toISOString(),
  139. };
  140. await writeJsonFileAtomic(jobFilePath(jobId), updated);
  141. return updated;
  142. });
  143. }
  144. export async function markClassroomGenerationJobRunning(
  145. jobId: string,
  146. ): Promise<ClassroomGenerationJob> {
  147. return withJobLock(jobId, async () => {
  148. const existing = await readClassroomGenerationJob(jobId);
  149. if (!existing) {
  150. throw new Error(`Classroom generation job not found: ${jobId}`);
  151. }
  152. const updated: ClassroomGenerationJob = {
  153. ...existing,
  154. status: 'running',
  155. startedAt: existing.startedAt || new Date().toISOString(),
  156. message: 'Classroom generation started',
  157. updatedAt: new Date().toISOString(),
  158. };
  159. await writeJsonFileAtomic(jobFilePath(jobId), updated);
  160. return updated;
  161. });
  162. }
  163. export async function updateClassroomGenerationJobProgress(
  164. jobId: string,
  165. progress: ClassroomGenerationProgress,
  166. ): Promise<ClassroomGenerationJob> {
  167. return updateClassroomGenerationJob(jobId, {
  168. status: 'running',
  169. step: progress.step,
  170. progress: progress.progress,
  171. message: progress.message,
  172. scenesGenerated: progress.scenesGenerated,
  173. totalScenes: progress.totalScenes,
  174. });
  175. }
  176. export async function markClassroomGenerationJobSucceeded(
  177. jobId: string,
  178. result: GenerateClassroomResult,
  179. ): Promise<ClassroomGenerationJob> {
  180. return updateClassroomGenerationJob(jobId, {
  181. status: 'succeeded',
  182. step: 'completed',
  183. progress: 100,
  184. message: 'Classroom generation completed',
  185. completedAt: new Date().toISOString(),
  186. scenesGenerated: result.scenesCount,
  187. result: {
  188. classroomId: result.id,
  189. url: result.url,
  190. scenesCount: result.scenesCount,
  191. },
  192. });
  193. }
  194. export async function markClassroomGenerationJobFailed(
  195. jobId: string,
  196. error: string,
  197. ): Promise<ClassroomGenerationJob> {
  198. return updateClassroomGenerationJob(jobId, {
  199. status: 'failed',
  200. step: 'failed',
  201. message: 'Classroom generation failed',
  202. completedAt: new Date().toISOString(),
  203. error,
  204. });
  205. }