返回 oh-my-ppt
job-manager.test.ts
根目录 / tests / unit / generation / job-manager.test.ts
1 import { beforeEach, describe, expect, it, vi } from 'vitest'
2
3 const logMocks = vi.hoisted(() => ({
4 info: vi.fn(),
5 warn: vi.fn(),
6 error: vi.fn()
7 }))
8 const finalizeGenerationFailureMock = vi.hoisted(() => vi.fn())
9
10 vi.mock('electron-log/main.js', () => ({ default: logMocks }))
11 vi.mock('../../../src/main/generation/finalization', () => ({
12 finalizeGenerationFailure: finalizeGenerationFailureMock,
13 resolveGenerationFailureSessionStatus: () => 'failed'
14 }))
15
16 import { GenerateJobManager } from '../../../src/main/generation/job-manager'
17 import { JobCoordinator } from '../../../src/main/agent-runtime/job/coordinator'
18
19 const createGenerationJobContext = (ctx: Record<string, any>) => ({
20 ...ctx,
21 sessionRuns: {
22 sessionRunStates: ctx.sessionRunStates,
23 beginSessionRunState: ctx.beginSessionRunState,
24 pruneFinishedSessionRunStates: () => undefined,
25 trackSessionRunChunk: () => undefined
26 },
27 runtimeEmitters: {
28 emitSessionRunLifecycle: ctx.emitSessionRunLifecycle || (() => undefined),
29 emitGenerateChunk: ctx.emitGenerateChunk,
30 emitRuntimeJobStarted: ctx.emitRuntimeJobStarted || (() => undefined),
31 emitRuntimeJobTerminal: ctx.emitRuntimeJobTerminal,
32 createDeckProgressEmitter: () => () => undefined
33 }
34 })
35
36 describe('GenerateJobManager', () => {
37 beforeEach(() => {
38 vi.clearAllMocks()
39 finalizeGenerationFailureMock.mockReset()
40 })
41
42 it('persists a background generation as a unified session job', async () => {
43 let resolveExecution: (() => void) | undefined
44 const execution = new Promise<void>((resolve) => {
45 resolveExecution = resolve
46 })
47 const beginSessionRunState = vi.fn()
48 const emitRuntimeJobTerminal = vi.fn()
49 const ctx = {
50 db: {
51 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
52 updateSessionJobStatus: vi.fn().mockResolvedValue(undefined)
53 },
54 sessionRunStates: new Map(),
55 beginSessionRunState,
56 emitGenerateChunk: vi.fn(),
57 emitRuntimeJobTerminal,
58 agentManager: {
59 removeSession: vi.fn(),
60 cancelSession: vi.fn()
61 }
62 }
63 const coordinator = new JobCoordinator()
64 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never, coordinator)
65 const reserved = await manager.reserve('generate:start', 'session-1', 'run-generate-1')
66 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
67
68 const result = await manager.enqueue({
69 reservation: reserved.reservation,
70 kind: 'standard',
71 context: {
72 sessionId: 'session-1',
73 runId: 'run-generate-1',
74 styleId: 'style-1',
75 previousSessionStatus: 'completed',
76 effectiveMode: 'generate',
77 messageScope: 'main',
78 projectId: 'project-1'
79 },
80 totalPages: 1,
81 execute: async () => execution
82 })
83
84 expect(result).toEqual({ runId: 'run-generate-1', queued: false })
85 expect(ctx.db.createGenerationRunWithSessionJob).toHaveBeenCalledWith(
86 expect.objectContaining({
87 run: expect.objectContaining({ id: 'run-generate-1', mode: 'generate', totalPages: 1 }),
88 job: expect.objectContaining({
89 id: 'run-generate-1',
90 kind: 'standard',
91 previousSessionStatus: 'completed',
92 totalPages: 1
93 })
94 })
95 )
96 expect(beginSessionRunState).toHaveBeenCalledWith(
97 expect.objectContaining({
98 kind: 'standard'
99 })
100 )
101
102 resolveExecution?.()
103 await vi.waitFor(() => {
104 expect(ctx.db.updateSessionJobStatus).toHaveBeenCalledWith('run-generate-1', 'finished')
105 })
106 expect(emitRuntimeJobTerminal).toHaveBeenCalledWith({
107 sessionId: 'session-1',
108 jobId: 'run-generate-1',
109 domain: 'generation',
110 status: 'completed'
111 })
112 const sessionJobFinishedCall = ctx.db.updateSessionJobStatus.mock.calls.findIndex(
113 ([runId, status]) => runId === 'run-generate-1' && status === 'finished'
114 )
115 expect(sessionJobFinishedCall).toBeGreaterThanOrEqual(0)
116 expect(
117 ctx.db.updateSessionJobStatus.mock.invocationCallOrder[sessionJobFinishedCall]
118 ).toBeLessThan(emitRuntimeJobTerminal.mock.invocationCallOrder[0])
119 expect(coordinator.getByOwner({ kind: 'session', id: 'session-1' })).toBeNull()
120 })
121
122 it('does not leave a run behind when atomic job creation fails', async () => {
123 const ctx = {
124 db: {
125 createGenerationRunWithSessionJob: vi.fn().mockRejectedValue(new Error('job insert failed')),
126 updateSessionJobStatus: vi.fn(),
127 updateGenerationRunStatus: vi.fn().mockResolvedValue(undefined),
128 updateSessionStatus: vi.fn().mockResolvedValue(undefined)
129 },
130 sessionRunStates: new Map(),
131 beginSessionRunState: vi.fn(),
132 emitGenerateChunk: vi.fn(),
133 emitRuntimeJobTerminal: vi.fn(),
134 agentManager: {
135 removeSession: vi.fn(),
136 cancelSession: vi.fn()
137 }
138 }
139 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never)
140 const reserved = await manager.reserve('generate:start', 'session-2', 'run-generate-2')
141 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
142
143 await expect(
144 manager.enqueue({
145 reservation: reserved.reservation,
146 kind: 'standard',
147 context: {
148 sessionId: 'session-2',
149 runId: 'run-generate-2',
150 styleId: 'style-1',
151 previousSessionStatus: 'completed',
152 effectiveMode: 'generate',
153 messageScope: 'main',
154 projectId: 'project-1'
155 },
156 totalPages: 1,
157 execute: vi.fn()
158 })
159 ).rejects.toThrow('job insert failed')
160
161 expect(ctx.db.updateGenerationRunStatus).not.toHaveBeenCalled()
162 expect(ctx.db.updateSessionStatus).not.toHaveBeenCalled()
163 })
164
165 it('aborts the persisted session job when setup fails after it has been created', async () => {
166 const ctx = {
167 db: {
168 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
169 updateSessionJobStatus: vi.fn().mockResolvedValue(undefined),
170 updateGenerationRunStatus: vi.fn().mockResolvedValue(undefined),
171 updateSessionStatus: vi.fn().mockResolvedValue(undefined)
172 },
173 sessionRunStates: new Map(),
174 beginSessionRunState: vi.fn(() => {
175 throw new Error('state initialization failed')
176 }),
177 emitGenerateChunk: vi.fn(),
178 emitRuntimeJobTerminal: vi.fn(),
179 agentManager: {
180 removeSession: vi.fn(),
181 cancelSession: vi.fn()
182 }
183 }
184 const coordinator = new JobCoordinator()
185 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never, coordinator)
186 const reserved = await manager.reserve(
187 'generate:start',
188 'session-setup-failure',
189 'run-generate-setup-failure'
190 )
191 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
192
193 await expect(
194 manager.enqueue({
195 reservation: reserved.reservation,
196 kind: 'standard',
197 context: {
198 sessionId: 'session-setup-failure',
199 runId: 'run-generate-setup-failure',
200 styleId: 'style-1',
201 previousSessionStatus: 'completed',
202 effectiveMode: 'generate',
203 messageScope: 'main',
204 projectId: 'project-1'
205 },
206 totalPages: 1,
207 execute: vi.fn()
208 })
209 ).rejects.toThrow('state initialization failed')
210
211 expect(ctx.db.updateSessionJobStatus).toHaveBeenCalledWith(
212 'run-generate-setup-failure',
213 'aborted',
214 { abortReason: 'setup_failed' }
215 )
216 expect(ctx.db.updateGenerationRunStatus).toHaveBeenCalledWith(
217 'run-generate-setup-failure',
218 'failed',
219 'state initialization failed'
220 )
221 expect(coordinator.getByOwner({ kind: 'session', id: 'session-setup-failure' })).toBeNull()
222 })
223
224 it('restores session status after an interrupted persisted job', async () => {
225 const ctx = {
226 db: {
227 listActiveSessionJobs: vi.fn().mockResolvedValue([
228 {
229 id: 'run-generate-3',
230 session_id: 'session-3',
231 previous_session_status: 'completed'
232 }
233 ]),
234 updateSessionJobStatus: vi.fn().mockResolvedValue(undefined),
235 updateGenerationRunStatus: vi.fn().mockResolvedValue(undefined),
236 updateSessionStatus: vi.fn().mockResolvedValue(undefined),
237 getGenerationRun: vi.fn().mockResolvedValue({ status: 'running' })
238 },
239 sessionRunStates: new Map(),
240 beginSessionRunState: vi.fn(),
241 emitGenerateChunk: vi.fn(),
242 emitRuntimeJobTerminal: vi.fn(),
243 agentManager: {
244 removeSession: vi.fn(),
245 cancelSession: vi.fn()
246 }
247 }
248 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never)
249
250 await manager.abortInterruptedJobs('应用退出导致生成中断')
251
252 expect(ctx.db.updateSessionJobStatus).toHaveBeenCalledWith('run-generate-3', 'aborted', {
253 abortReason: '应用退出导致生成中断'
254 })
255 expect(ctx.db.updateGenerationRunStatus).toHaveBeenCalledWith(
256 'run-generate-3',
257 'failed',
258 '应用退出导致生成中断'
259 )
260 expect(ctx.db.updateSessionStatus).toHaveBeenCalledWith('session-3', 'completed')
261 })
262
263 it('settles an interrupted job that already persisted a successful generation', async () => {
264 const ctx = {
265 db: {
266 listActiveSessionJobs: vi.fn().mockResolvedValue([
267 {
268 id: 'run-completed-before-crash',
269 session_id: 'session-completed-before-crash',
270 previous_session_status: 'active'
271 }
272 ]),
273 getGenerationRun: vi.fn().mockResolvedValue({ status: 'completed' }),
274 updateSessionJobStatus: vi.fn().mockResolvedValue(undefined),
275 updateGenerationRunStatus: vi.fn(),
276 updateSessionStatus: vi.fn()
277 },
278 sessionRunStates: new Map(),
279 beginSessionRunState: vi.fn(),
280 emitGenerateChunk: vi.fn(),
281 emitRuntimeJobTerminal: vi.fn(),
282 agentManager: { removeSession: vi.fn() }
283 }
284 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never)
285
286 await manager.abortInterruptedJobs('应用退出导致生成中断')
287
288 expect(ctx.db.updateSessionJobStatus).toHaveBeenCalledWith(
289 'run-completed-before-crash',
290 'finished'
291 )
292 expect(ctx.db.updateGenerationRunStatus).not.toHaveBeenCalled()
293 expect(ctx.db.updateSessionStatus).not.toHaveBeenCalled()
294 })
295
296 it('uses the active JobCoordinator lease signal for cancellation and terminal persistence', async () => {
297 let executionSignal: AbortSignal | undefined
298 let executionStarted!: () => void
299 const started = new Promise<void>((resolve) => {
300 executionStarted = resolve
301 })
302 const emitRuntimeJobTerminal = vi.fn()
303 const ctx = {
304 db: {
305 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
306 updateSessionJobStatus: vi.fn().mockResolvedValue(undefined),
307 updateGenerationRunStatus: vi.fn().mockResolvedValue(undefined),
308 updateSessionStatus: vi.fn().mockResolvedValue(undefined),
309 getGenerationRun: vi.fn().mockResolvedValue(null),
310 addMessage: vi.fn().mockResolvedValue(undefined)
311 },
312 sessionRunStates: new Map(),
313 beginSessionRunState: vi.fn(),
314 emitGenerateChunk: vi.fn(),
315 emitRuntimeJobTerminal,
316 agentManager: { removeSession: vi.fn() }
317 }
318 const coordinator = new JobCoordinator()
319 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never, coordinator)
320 const reserved = await manager.reserve('generate:start', 'session-active', 'run-active')
321 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
322
323 await manager.enqueue({
324 reservation: reserved.reservation,
325 kind: 'standard',
326 context: {
327 sessionId: 'session-active',
328 runId: 'run-active',
329 styleId: 'style-1',
330 previousSessionStatus: 'completed',
331 effectiveMode: 'generate',
332 messageScope: 'main',
333 projectId: 'project-1',
334 // Generation handlers resolve expensive context only after reserve(),
335 // so their execution context carries this exact JobLease signal.
336 abortSignal: reserved.reservation.signal
337 },
338 totalPages: 1,
339 execute: async (context) => {
340 executionSignal = context.abortSignal
341 executionStarted()
342 await new Promise<void>((_resolve, reject) => {
343 context.abortSignal.addEventListener(
344 'abort',
345 () => reject(new Error('生成已取消')),
346 { once: true }
347 )
348 })
349 }
350 })
351 await started
352
353 expect(executionSignal).toBe(reserved.reservation.signal)
354 await expect(manager.cancel('session-active')).resolves.toBe(true)
355 await vi.waitFor(() => {
356 expect(ctx.db.updateSessionJobStatus).toHaveBeenCalledWith('run-active', 'aborted', {
357 abortReason: 'cancelled'
358 })
359 })
360 expect(emitRuntimeJobTerminal).toHaveBeenCalledWith({
361 sessionId: 'session-active',
362 jobId: 'run-active',
363 domain: 'generation',
364 status: 'cancelled',
365 errorCode: undefined,
366 errorMessage: undefined
367 })
368 expect(coordinator.getByOwner({ kind: 'session', id: 'session-active' })).toBeNull()
369 })
370
371 it('settles and publishes a failed job when generation finalization fails', async () => {
372 const emitRuntimeJobTerminal = vi.fn()
373 const ctx = {
374 db: {
375 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
376 updateSessionJobStatus: vi.fn().mockResolvedValue(undefined),
377 updateGenerationRunStatus: vi.fn().mockResolvedValue(undefined),
378 updateSessionStatus: vi.fn().mockResolvedValue(undefined)
379 },
380 sessionRunStates: new Map(),
381 beginSessionRunState: vi.fn(),
382 emitGenerateChunk: vi.fn(),
383 emitRuntimeJobTerminal,
384 agentManager: { removeSession: vi.fn() }
385 }
386 const coordinator = new JobCoordinator()
387 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never, coordinator)
388 const reserved = await manager.reserve('generate:start', 'session-finalize-failure', 'run-finalize-failure')
389 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
390 finalizeGenerationFailureMock.mockRejectedValueOnce(new Error('database temporarily unavailable'))
391
392 await manager.enqueue({
393 reservation: reserved.reservation,
394 kind: 'standard',
395 context: {
396 sessionId: 'session-finalize-failure',
397 runId: 'run-finalize-failure',
398 styleId: 'style-1',
399 previousSessionStatus: 'completed',
400 effectiveMode: 'generate',
401 messageScope: 'main',
402 projectId: 'project-1'
403 },
404 totalPages: 1,
405 execute: async () => {
406 throw new Error('generation failed')
407 }
408 })
409
410 await vi.waitFor(() => {
411 expect(ctx.db.updateSessionJobStatus).toHaveBeenCalledWith('run-finalize-failure', 'finished')
412 })
413 expect(ctx.db.updateGenerationRunStatus).toHaveBeenCalledWith(
414 'run-finalize-failure',
415 'failed',
416 'generation failed'
417 )
418 expect(ctx.db.updateSessionStatus).toHaveBeenCalledWith('session-finalize-failure', 'failed')
419 expect(ctx.emitGenerateChunk).toHaveBeenCalledWith('session-finalize-failure', {
420 type: 'run_error',
421 payload: {
422 runId: 'run-finalize-failure',
423 message: 'generation failed',
424 cancelled: false
425 }
426 })
427 expect(logMocks.error).toHaveBeenCalledWith(
428 '[generate:job] failed to finalize generation',
429 expect.objectContaining({ runId: 'run-finalize-failure' })
430 )
431 expect(emitRuntimeJobTerminal).toHaveBeenCalledWith({
432 sessionId: 'session-finalize-failure',
433 jobId: 'run-finalize-failure',
434 domain: 'generation',
435 status: 'failed',
436 errorCode: 'generation_failed',
437 errorMessage: 'generation failed'
438 })
439 const sessionJobFinishedCall = ctx.db.updateSessionJobStatus.mock.calls.findIndex(
440 ([runId, status]) => runId === 'run-finalize-failure' && status === 'finished'
441 )
442 expect(
443 ctx.db.updateSessionJobStatus.mock.invocationCallOrder[sessionJobFinishedCall]
444 ).toBeLessThan(emitRuntimeJobTerminal.mock.invocationCallOrder[0])
445 const fallbackRunFailureCall = ctx.db.updateGenerationRunStatus.mock.calls.findIndex(
446 ([runId, status]) => runId === 'run-finalize-failure' && status === 'failed'
447 )
448 expect(
449 ctx.db.updateGenerationRunStatus.mock.invocationCallOrder[fallbackRunFailureCall]
450 ).toBeLessThan(ctx.db.updateSessionJobStatus.mock.invocationCallOrder[sessionJobFinishedCall])
451 expect(coordinator.getByOwner({ kind: 'session', id: 'session-finalize-failure' })).toBeNull()
452 })
453
454 it('leaves the session job recoverable when finalization cannot restore the session', async () => {
455 const emitRuntimeJobTerminal = vi.fn()
456 const ctx = {
457 db: {
458 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
459 updateGenerationRunStatus: vi.fn().mockResolvedValue(undefined),
460 updateSessionStatus: vi.fn().mockRejectedValue(new Error('session database unavailable')),
461 updateSessionJobStatus: vi.fn().mockResolvedValue(undefined)
462 },
463 sessionRunStates: new Map(),
464 beginSessionRunState: vi.fn(),
465 emitGenerateChunk: vi.fn(),
466 emitRuntimeJobTerminal,
467 agentManager: { removeSession: vi.fn() }
468 }
469 const coordinator = new JobCoordinator()
470 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never, coordinator)
471 const reserved = await manager.reserve('generate:start', 'session-unrecoverable', 'run-unrecoverable')
472 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
473 finalizeGenerationFailureMock.mockRejectedValueOnce(new Error('finalization database unavailable'))
474
475 await manager.enqueue({
476 reservation: reserved.reservation,
477 kind: 'standard',
478 context: {
479 sessionId: 'session-unrecoverable',
480 runId: 'run-unrecoverable',
481 styleId: 'style-1',
482 previousSessionStatus: 'completed',
483 effectiveMode: 'generate',
484 messageScope: 'main',
485 projectId: 'project-1'
486 },
487 totalPages: 1,
488 execute: async () => {
489 throw new Error('generation failed')
490 }
491 })
492
493 await vi.waitFor(() => {
494 expect(logMocks.error).toHaveBeenCalledWith(
495 '[generate:job] failed to persist fallback generation terminal state',
496 expect.objectContaining({ runId: 'run-unrecoverable' })
497 )
498 })
499 expect(ctx.db.updateSessionJobStatus).not.toHaveBeenCalledWith('run-unrecoverable', 'finished')
500 expect(emitRuntimeJobTerminal).not.toHaveBeenCalled()
501 expect(coordinator.getByOwner({ kind: 'session', id: 'session-unrecoverable' })).toBeNull()
502 })
503
504 it('does not execute or publish started when activation cannot be persisted', async () => {
505 const execute = vi.fn()
506 const emitRuntimeJobStarted = vi.fn()
507 const emitRuntimeJobTerminal = vi.fn()
508 const ctx = {
509 db: {
510 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
511 updateSessionJobStatus: vi
512 .fn()
513 .mockRejectedValueOnce(new Error('activation database unavailable'))
514 .mockResolvedValue(undefined)
515 },
516 sessionRunStates: new Map(),
517 beginSessionRunState: vi.fn(),
518 emitGenerateChunk: vi.fn(),
519 emitRuntimeJobStarted,
520 emitRuntimeJobTerminal,
521 agentManager: { removeSession: vi.fn() }
522 }
523 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never)
524 const reserved = await manager.reserve('generate:start', 'session-activation-failure', 'run-activation-failure')
525 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
526 finalizeGenerationFailureMock.mockResolvedValueOnce(undefined)
527
528 await manager.enqueue({
529 reservation: reserved.reservation,
530 kind: 'standard',
531 context: {
532 sessionId: 'session-activation-failure',
533 runId: 'run-activation-failure',
534 styleId: 'style-1',
535 previousSessionStatus: 'completed',
536 effectiveMode: 'generate',
537 messageScope: 'main',
538 projectId: 'project-1'
539 },
540 totalPages: 1,
541 execute
542 })
543
544 await vi.waitFor(() => {
545 expect(emitRuntimeJobTerminal).toHaveBeenCalledWith({
546 sessionId: 'session-activation-failure',
547 jobId: 'run-activation-failure',
548 domain: 'generation',
549 status: 'failed',
550 errorCode: 'generation_failed',
551 errorMessage: 'activation database unavailable'
552 })
553 })
554 expect(execute).not.toHaveBeenCalled()
555 expect(emitRuntimeJobStarted).not.toHaveBeenCalled()
556 })
557
558 it('does not turn completed generation into failure when only the session-job terminal write fails', async () => {
559 const emitRuntimeJobTerminal = vi.fn()
560 const ctx = {
561 db: {
562 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
563 updateSessionJobStatus: vi
564 .fn()
565 .mockResolvedValueOnce(undefined)
566 .mockRejectedValueOnce(new Error('session job database unavailable'))
567 },
568 sessionRunStates: new Map(),
569 beginSessionRunState: vi.fn(),
570 emitGenerateChunk: vi.fn(),
571 emitRuntimeJobTerminal,
572 agentManager: { removeSession: vi.fn() }
573 }
574 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never)
575 const reserved = await manager.reserve('generate:start', 'session-completed-write-failure', 'run-completed-write-failure')
576 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
577
578 await manager.enqueue({
579 reservation: reserved.reservation,
580 kind: 'standard',
581 context: {
582 sessionId: 'session-completed-write-failure',
583 runId: 'run-completed-write-failure',
584 styleId: 'style-1',
585 previousSessionStatus: 'completed',
586 effectiveMode: 'generate',
587 messageScope: 'main',
588 projectId: 'project-1'
589 },
590 totalPages: 1,
591 execute: vi.fn().mockResolvedValue(undefined)
592 })
593
594 await vi.waitFor(() => {
595 expect(logMocks.error).toHaveBeenCalledWith(
596 '[generate:job] failed to settle completed session job',
597 expect.objectContaining({ runId: 'run-completed-write-failure' })
598 )
599 })
600 expect(finalizeGenerationFailureMock).not.toHaveBeenCalled()
601 expect(emitRuntimeJobTerminal).not.toHaveBeenCalled()
602 })
603
604 it('does not publish a terminal event until the failed job status is persisted', async () => {
605 const emitRuntimeJobTerminal = vi.fn()
606 const ctx = {
607 db: {
608 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
609 updateSessionJobStatus: vi
610 .fn()
611 .mockResolvedValueOnce(undefined)
612 .mockRejectedValueOnce(new Error('session job database unavailable'))
613 },
614 sessionRunStates: new Map(),
615 beginSessionRunState: vi.fn(),
616 emitGenerateChunk: vi.fn(),
617 emitRuntimeJobTerminal,
618 agentManager: { removeSession: vi.fn() }
619 }
620 const coordinator = new JobCoordinator()
621 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never, coordinator)
622 const reserved = await manager.reserve('generate:start', 'session-status-failure', 'run-status-failure')
623 if (reserved.alreadyRunning) throw new Error('expected available job reservation')
624
625 await manager.enqueue({
626 reservation: reserved.reservation,
627 kind: 'standard',
628 context: {
629 sessionId: 'session-status-failure',
630 runId: 'run-status-failure',
631 styleId: 'style-1',
632 previousSessionStatus: 'completed',
633 effectiveMode: 'generate',
634 messageScope: 'main',
635 projectId: 'project-1'
636 },
637 totalPages: 1,
638 execute: async () => {
639 throw new Error('generation failed')
640 }
641 })
642
643 await vi.waitFor(() => {
644 expect(logMocks.error).toHaveBeenCalledWith(
645 '[generate:job] failed to settle session job',
646 expect.objectContaining({ runId: 'run-status-failure' })
647 )
648 })
649 expect(emitRuntimeJobTerminal).not.toHaveBeenCalled()
650 expect(coordinator.getByOwner({ kind: 'session', id: 'session-status-failure' })).toBeNull()
651 })
652
653 it('settles a queued job cancelled directly through JobCoordinator and starts the next job in FIFO order', async () => {
654 const executions: string[] = []
655 const completions = new Map<string, { promise: Promise<void>; resolve: () => void }>()
656 const createCompletion = (runId: string): void => {
657 let resolve!: () => void
658 const promise = new Promise<void>((complete) => {
659 resolve = complete
660 })
661 completions.set(runId, { promise, resolve })
662 }
663 for (const runId of ['run-1', 'run-2', 'run-3', 'run-4']) createCompletion(runId)
664
665 const emitRuntimeJobStarted = vi.fn()
666 const ctx = {
667 db: {
668 createGenerationRunWithSessionJob: vi.fn().mockResolvedValue(undefined),
669 updateSessionJobStatus: vi.fn().mockResolvedValue(undefined),
670 updateGenerationRunStatus: vi.fn().mockResolvedValue(undefined),
671 updateSessionStatus: vi.fn().mockResolvedValue(undefined)
672 },
673 sessionRunStates: new Map(),
674 beginSessionRunState: vi.fn(),
675 emitGenerateChunk: vi.fn(),
676 emitRuntimeJobStarted,
677 emitRuntimeJobTerminal: vi.fn(),
678 agentManager: { removeSession: vi.fn() }
679 }
680 const coordinator = new JobCoordinator()
681 const manager = new GenerateJobManager(createGenerationJobContext(ctx) as never, coordinator)
682
683 const enqueue = async (sessionId: string, runId: string): Promise<void> => {
684 const reservation = await manager.reserve('generate:start', sessionId, runId)
685 if (reservation.alreadyRunning) throw new Error(`unexpected busy reservation for ${runId}`)
686 await manager.enqueue({
687 reservation: reservation.reservation,
688 kind: 'standard',
689 context: {
690 sessionId,
691 runId,
692 styleId: 'style-1',
693 previousSessionStatus: 'completed',
694 effectiveMode: 'generate',
695 messageScope: 'main',
696 projectId: 'project-1'
697 },
698 totalPages: 1,
699 execute: async () => {
700 executions.push(runId)
701 await completions.get(runId)?.promise
702 }
703 })
704 }
705
706 await enqueue('session-1', 'run-1')
707 await enqueue('session-2', 'run-2')
708 await vi.waitFor(() => expect(executions).toEqual(['run-1', 'run-2']))
709
710 await enqueue('session-3', 'run-3')
711 await enqueue('session-4', 'run-4')
712 expect(executions).toEqual(['run-1', 'run-2'])
713
714 // A capacity-queued generation keeps its session write lease. Every outer
715 // writer must observe the queued run as busy rather than slipping between
716 // resource and capacity scheduling.
717 const queuedGenerationConflicts = [
718 { name: 'retry-failed-pages', domain: 'generation' as const },
719 { name: 'add-page', domain: 'generation' as const },
720 { name: 'single-page-retry', domain: 'generation' as const },
721 { name: 'page-edit', domain: 'edit' as const },
722 { name: 'page-beautify', domain: 'edit' as const },
723 { name: 'deck-edit', domain: 'edit' as const },
724 { name: 'style-switch', domain: 'style' as const }
725 ]
726 for (const operation of queuedGenerationConflicts) {
727 await expect(
728 coordinator.reserve({
729 jobId: `${operation.name}-while-queued`,
730 domain: operation.domain,
731 owner: { kind: 'session', id: 'session-3' },
732 claims: { write: ['session:session-3'] },
733 wait: 'fail'
734 })
735 ).resolves.toEqual({ status: 'busy', conflictingJobId: 'run-3' })
736 }
737
738 expect(coordinator.cancel('run-3')).toBe(true)
739 await vi.waitFor(() => {
740 expect(ctx.db.updateSessionJobStatus).toHaveBeenCalledWith('run-3', 'aborted', {
741 abortReason: 'cancelled'
742 })
743 })
744 await vi.waitFor(() => {
745 expect(coordinator.getByOwner({ kind: 'session', id: 'session-3' })).toBeNull()
746 })
747
748 const retryAfterQueuedGeneration = await manager.reserve(
749 'generate:retryFailedPages',
750 'session-3',
751 'run-3-retry'
752 )
753 expect(retryAfterQueuedGeneration).toMatchObject({ alreadyRunning: false })
754 if (!retryAfterQueuedGeneration.alreadyRunning) retryAfterQueuedGeneration.reservation.release()
755
756 completions.get('run-1')?.resolve()
757 await vi.waitFor(() => expect(executions).toEqual(['run-1', 'run-2', 'run-4']))
758 expect(emitRuntimeJobStarted).toHaveBeenCalledWith({
759 sessionId: 'session-4',
760 jobId: 'run-4',
761 domain: 'generation'
762 })
763 expect(executions).not.toContain('run-3')
764
765 completions.get('run-2')?.resolve()
766 completions.get('run-4')?.resolve()
767 await vi.waitFor(() => {
768 expect(coordinator.getByOwner({ kind: 'session', id: 'session-2' })).toBeNull()
769 expect(coordinator.getByOwner({ kind: 'session', id: 'session-4' })).toBeNull()
770 })
771 })
772 })
773
773 lines TYPESCRIPT