(job: Job<GenerationInput>)
| 29 | } |
| 30 | |
| 31 | async process(job: Job<GenerationInput>): Promise<void> { |
| 32 | const { projectId, runId, threadId } = job.data; |
| 33 | console.log(`[processor] === Starting generation job ===`); |
| 34 | console.log(`[processor] projectId=${projectId} runId=${runId} threadId=${threadId}`); |
| 35 | |
| 36 | const [run] = await this.db |
| 37 | .select({ status: schema.agentRuns.status }) |
| 38 | .from(schema.agentRuns) |
| 39 | .where(eq(schema.agentRuns.id, runId)); |
| 40 | if (run && run.status !== "running" && run.status !== "queued") { |
| 41 | console.log(`[processor] Skipping run=${runId}; status=${run.status}`); |
| 42 | return; |
| 43 | } |
| 44 | |
| 45 | // Build the graph with checkpointer |
| 46 | console.log(`[processor] Building graph...`); |
| 47 | const graph = buildPhase2Graph(); |
| 48 | const checkpointer = new MemorySaver(); |
| 49 | const compiled = graph.compile({ checkpointer }); |
| 50 | console.log(`[processor] Graph compiled with ${Object.keys(graph).length} nodes`); |
| 51 | |
| 52 | const config = { configurable: { thread_id: threadId } }; |
| 53 | const [project] = await this.db |
| 54 | .select({ |
| 55 | title: schema.projects.title, |
| 56 | story: schema.projects.story, |
| 57 | style: schema.projects.style, |
| 58 | targetShotCount: schema.projects.targetShotCount, |
| 59 | storyOutline: schema.projects.storyOutline, |
| 60 | visualBible: schema.projects.visualBible, |
| 61 | }) |
| 62 | .from(schema.projects) |
| 63 | .where(eq(schema.projects.id, projectId)); |
| 64 | |
| 65 | const projectContext = buildProjectContext(projectId, project, job.data); |
| 66 | let latestOutline: Record<string, unknown> | null = project?.storyOutline || null; |
| 67 | let latestVisualBible: string | null = project?.visualBible || null; |
| 68 | |
| 69 | // ---- Build the REAL agent context ---- |
| 70 | const ctx = { |
| 71 | projectId, |
| 72 | runId, |
| 73 | threadId, |
| 74 | sendMessage: async (content: string, opts?: { summary?: string; progress?: number; isLoading?: boolean; stage?: string }) => { |
| 75 | console.log(`[processor:msg] content="${content.substring(0, 80)}..." progress=${opts?.progress ?? "N/A"} stage=${opts?.stage ?? "N/A"}`); |
| 76 | try { |
| 77 | await this.db.insert(schema.messages).values({ |
| 78 | projectId, |
| 79 | runId, |
| 80 | agent: "system", |
| 81 | role: "assistant", |
| 82 | content, |
| 83 | summary: opts?.summary || null, |
| 84 | progress: opts?.progress || null, |
| 85 | isLoading: opts?.isLoading || false, |
| 86 | createdAt: new Date(), |
| 87 | }); |
| 88 | } catch (dbErr) { |
nothing calls this directly
no test coverage detected