| 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 |