返回 oh-my-ppt
job-coordinator.test.ts
根目录 / tests / unit / agent-runtime / job-coordinator.test.ts
1 import { describe, expect, it } from 'vitest'
2 import { JobCoordinator } from '../../../src/main/agent-runtime/job/coordinator'
3 import { sessionLockKey } from '../../../src/main/agent-runtime/lock/keys'
4
5 const reserve = (
6 coordinator: JobCoordinator,
7 args: Partial<Parameters<JobCoordinator['reserve']>[0]> & {
8 jobId: string
9 owner: { kind: 'session' | 'style' | 'image-history'; id: string }
10 claims: { read?: string[]; write?: string[] }
11 }
12 ) =>
13 coordinator.reserve({
14 domain: 'generation',
15 wait: 'block',
16 ...args
17 })
18
19 describe('JobCoordinator', () => {
20 it('returns the existing job for the same owner and conflicting fail-fast claim', async () => {
21 const coordinator = new JobCoordinator()
22 const first = await reserve(coordinator, {
23 jobId: 'job-1',
24 owner: { kind: 'session', id: 'session-1' },
25 claims: { write: ['session:session-1'] }
26 })
27 if (first.status !== 'acquired') throw new Error('expected first job to acquire')
28
29 await expect(
30 reserve(coordinator, {
31 jobId: 'job-2',
32 owner: { kind: 'session', id: 'session-1' },
33 claims: { write: ['session:session-1'] }
34 })
35 ).resolves.toEqual({ status: 'busy', conflictingJobId: 'job-1' })
36
37 await expect(
38 reserve(coordinator, {
39 jobId: 'job-3',
40 owner: { kind: 'style', id: 'style-1' },
41 claims: { write: ['session:session-1'] },
42 wait: 'fail'
43 })
44 ).resolves.toEqual({ status: 'busy', conflictingJobId: 'job-1' })
45
46 first.lease.release()
47 })
48
49 it('keeps every outer session writer busy while generation owns the session', async () => {
50 const coordinator = new JobCoordinator()
51 const sessionId = 'session-1'
52 const generation = await reserve(coordinator, {
53 jobId: 'generation-run-1',
54 owner: { kind: 'session', id: sessionId },
55 claims: { write: [sessionLockKey(sessionId)] }
56 })
57 if (generation.status !== 'acquired') throw new Error('expected generation job')
58
59 const outerWriters = [
60 { name: 'retry-failed-pages', domain: 'generation' as const },
61 { name: 'add-page', domain: 'generation' as const },
62 { name: 'single-page-retry', domain: 'generation' as const },
63 { name: 'page-edit', domain: 'edit' as const },
64 { name: 'deck-edit', domain: 'edit' as const },
65 { name: 'style-switch', domain: 'style' as const }
66 ]
67 for (const operation of outerWriters) {
68 await expect(
69 coordinator.reserve({
70 jobId: `${operation.name}-job`,
71 domain: operation.domain,
72 owner: { kind: 'session', id: sessionId },
73 claims: { write: [sessionLockKey(sessionId)] },
74 wait: 'fail'
75 })
76 ).resolves.toEqual({ status: 'busy', conflictingJobId: 'generation-run-1' })
77 }
78
79 generation.lease.release()
80 const retry = await reserve(coordinator, {
81 jobId: 'retry-run-2',
82 owner: { kind: 'session', id: sessionId },
83 claims: { write: [sessionLockKey(sessionId)] }
84 })
85 expect(retry.status).toBe('acquired')
86 if (retry.status === 'acquired') retry.lease.release()
87 })
88
89 it('keeps queued jobs cancellable through the same signal and removes them', async () => {
90 const coordinator = new JobCoordinator()
91 const active = await reserve(coordinator, {
92 jobId: 'job-active',
93 owner: { kind: 'session', id: 'session-1' },
94 claims: { write: ['session:session-1'] }
95 })
96 if (active.status !== 'acquired') throw new Error('expected active job')
97
98 const waiting = reserve(coordinator, {
99 jobId: 'job-waiting',
100 owner: { kind: 'style', id: 'style-1' },
101 claims: { write: ['session:session-1'] }
102 })
103 expect(coordinator.getByOwner({ kind: 'style', id: 'style-1' })).toMatchObject({
104 jobId: 'job-waiting',
105 state: 'waiting'
106 })
107
108 expect(coordinator.cancel('job-waiting')).toBe(true)
109 await expect(waiting).rejects.toMatchObject({ name: 'AbortError' })
110 expect(coordinator.getByOwner({ kind: 'style', id: 'style-1' })).toBeNull()
111
112 active.lease.release()
113 })
114
115 it('relays external cancellation, cancels active jobs once, and releases idempotently', async () => {
116 const coordinator = new JobCoordinator()
117 const external = new AbortController()
118 const acquired = await reserve(coordinator, {
119 jobId: 'job-1',
120 owner: { kind: 'session', id: 'session-1' },
121 claims: { write: ['session:session-1'] },
122 signal: external.signal
123 })
124 if (acquired.status !== 'acquired') throw new Error('expected acquired job')
125
126 external.abort()
127 expect(acquired.lease.signal.aborted).toBe(true)
128 expect(coordinator.cancel('job-1')).toBe(false)
129
130 acquired.lease.release()
131 acquired.lease.release()
132
133 const next = await reserve(coordinator, {
134 jobId: 'job-2',
135 owner: { kind: 'session', id: 'session-1' },
136 claims: { write: ['session:session-1'] }
137 })
138 expect(next.status).toBe('acquired')
139 if (next.status === 'acquired') next.lease.release()
140 })
141
142 it('cancels every job owned by the requested owner', async () => {
143 const coordinator = new JobCoordinator()
144 const acquired = await reserve(coordinator, {
145 jobId: 'job-1',
146 owner: { kind: 'session', id: 'session-1' },
147 claims: { write: ['session:session-1'] }
148 })
149 if (acquired.status !== 'acquired') throw new Error('expected acquired job')
150
151 expect(coordinator.cancelOwner({ kind: 'session', id: 'session-1' })).toBe(1)
152 expect(acquired.lease.signal.aborted).toBe(true)
153 acquired.lease.release()
154 })
155 })
156
156 lines TYPESCRIPT