Skip to content

Commit b1cb6c5

Browse files
fix(operation-queue): guard queue ownership by operation id
1 parent 64280c9 commit b1cb6c5

3 files changed

Lines changed: 290 additions & 71 deletions

File tree

apps/sim/stores/operation-queue/store.test.ts

Lines changed: 221 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,43 @@
11
/**
22
* @vitest-environment node
33
*/
4+
import Module from 'node:module'
45
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
56
import { registerEmitFunctions, useOperationQueueStore } from '@/stores/operation-queue/store'
67

8+
function addWorkflowOperation(id: string) {
9+
useOperationQueueStore.getState().addToQueue({
10+
id,
11+
workflowId: 'workflow-a',
12+
userId: 'user-1',
13+
operation: {
14+
operation: 'replace-state',
15+
target: 'workflow',
16+
payload: { state: { operationId: id } },
17+
},
18+
})
19+
}
20+
21+
function addBlockOperation(id: string, blockId: string) {
22+
useOperationQueueStore.getState().addToQueue({
23+
id,
24+
workflowId: 'workflow-a',
25+
userId: 'user-1',
26+
operation: {
27+
operation: 'block-update',
28+
target: 'block',
29+
payload: { id: blockId },
30+
},
31+
})
32+
}
33+
734
describe('operation queue room gating', () => {
835
beforeEach(() => {
936
vi.clearAllMocks()
1037
useOperationQueueStore.setState({
1138
operations: [],
1239
workflowOperationVersions: {},
13-
isProcessing: false,
40+
processingOperationId: null,
1441
hasOperationError: false,
1542
})
1643
registerEmitFunctions(vi.fn(), vi.fn(), vi.fn(), null)
@@ -20,7 +47,7 @@ describe('operation queue room gating', () => {
2047
useOperationQueueStore.setState({
2148
operations: [],
2249
workflowOperationVersions: {},
23-
isProcessing: false,
50+
processingOperationId: null,
2451
hasOperationError: false,
2552
})
2653
registerEmitFunctions(vi.fn(), vi.fn(), vi.fn(), null)
@@ -92,7 +119,7 @@ describe('operation queue room gating', () => {
92119
expect(skippingEmit).toHaveBeenCalledTimes(1)
93120

94121
const state = useOperationQueueStore.getState()
95-
expect(state.isProcessing).toBe(false)
122+
expect(state.processingOperationId).toBeNull()
96123
expect(state.hasOperationError).toBe(false)
97124
expect(state.operations).toEqual([
98125
expect.objectContaining({ id: 'op-1', status: 'pending', retryCount: 0 }),
@@ -133,63 +160,214 @@ describe('operation queue room gating', () => {
133160
useOperationQueueStore.getState().confirmOperation('op-1')
134161
})
135162

136-
it('triggers offline mode for a non-retryable failure and recovers via clearError', () => {
137-
registerEmitFunctions(
138-
vi.fn(() => true),
139-
vi.fn(),
140-
vi.fn(),
141-
'workflow-a'
142-
)
163+
it('advances exactly once after a successful acknowledgement', () => {
164+
const workflowEmit = vi.fn(() => true)
165+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
143166

144-
useOperationQueueStore.getState().addToQueue({
145-
id: 'op-1',
146-
workflowId: 'workflow-a',
147-
userId: 'user-1',
148-
operation: {
149-
operation: 'replace-state',
150-
target: 'workflow',
151-
payload: { state: {} },
152-
},
167+
addWorkflowOperation('op-1')
168+
addWorkflowOperation('op-2')
169+
170+
expect(workflowEmit).toHaveBeenCalledTimes(1)
171+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-1')
172+
173+
useOperationQueueStore.getState().confirmOperation('op-1')
174+
175+
expect(workflowEmit).toHaveBeenCalledTimes(2)
176+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-2')
177+
178+
useOperationQueueStore.getState().confirmOperation('op-2')
179+
180+
expect(workflowEmit).toHaveBeenCalledTimes(2)
181+
expect(useOperationQueueStore.getState().processingOperationId).toBeNull()
182+
expect(useOperationQueueStore.getState().operations).toEqual([])
183+
})
184+
185+
it('does not release the active operation for an unknown acknowledgement', () => {
186+
const workflowEmit = vi.fn(() => true)
187+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
188+
189+
addWorkflowOperation('op-1')
190+
addWorkflowOperation('op-2')
191+
192+
useOperationQueueStore.getState().confirmOperation('unknown-operation')
193+
194+
expect(workflowEmit).toHaveBeenCalledTimes(1)
195+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-1')
196+
expect(useOperationQueueStore.getState().operations.map((operation) => operation.id)).toEqual([
197+
'op-1',
198+
'op-2',
199+
])
200+
201+
useOperationQueueStore.getState().confirmOperation('op-1')
202+
useOperationQueueStore.getState().confirmOperation('op-2')
203+
})
204+
205+
it('does not release the active operation for a stale acknowledgement', () => {
206+
const workflowEmit = vi.fn(() => true)
207+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
208+
209+
addWorkflowOperation('op-1')
210+
addWorkflowOperation('op-2')
211+
addWorkflowOperation('op-3')
212+
213+
useOperationQueueStore.getState().confirmOperation('op-2')
214+
215+
expect(workflowEmit).toHaveBeenCalledTimes(1)
216+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-1')
217+
expect(useOperationQueueStore.getState().operations.map((operation) => operation.id)).toEqual([
218+
'op-1',
219+
'op-3',
220+
])
221+
222+
useOperationQueueStore.getState().confirmOperation('op-1')
223+
224+
expect(workflowEmit).toHaveBeenCalledTimes(2)
225+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-3')
226+
useOperationQueueStore.getState().confirmOperation('op-3')
227+
})
228+
229+
it('ignores a stale failure for a pending operation', () => {
230+
const workflowEmit = vi.fn(() => true)
231+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
232+
233+
addWorkflowOperation('op-1')
234+
addWorkflowOperation('op-2')
235+
236+
useOperationQueueStore.getState().failOperation('op-2')
237+
238+
expect(workflowEmit).toHaveBeenCalledTimes(1)
239+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-1')
240+
expect(useOperationQueueStore.getState().operations).toEqual([
241+
expect.objectContaining({ id: 'op-1', status: 'processing', retryCount: 0 }),
242+
expect.objectContaining({ id: 'op-2', status: 'pending', retryCount: 0 }),
243+
])
244+
245+
useOperationQueueStore.getState().confirmOperation('op-1')
246+
useOperationQueueStore.getState().confirmOperation('op-2')
247+
})
248+
249+
it('does not advance when cancelling an unrelated pending operation', () => {
250+
const workflowEmit = vi.fn(() => true)
251+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
252+
253+
addWorkflowOperation('op-1')
254+
addBlockOperation('op-2', 'block-2')
255+
addWorkflowOperation('op-3')
256+
257+
useOperationQueueStore.getState().cancelOperationsForBlock('block-2')
258+
259+
expect(workflowEmit).toHaveBeenCalledTimes(1)
260+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-1')
261+
expect(useOperationQueueStore.getState().operations.map((operation) => operation.id)).toEqual([
262+
'op-1',
263+
'op-3',
264+
])
265+
266+
useOperationQueueStore.getState().confirmOperation('op-1')
267+
268+
expect(workflowEmit).toHaveBeenCalledTimes(2)
269+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-3')
270+
useOperationQueueStore.getState().confirmOperation('op-3')
271+
})
272+
273+
it('releases ownership and advances once when cancelling the active operation', () => {
274+
const workflowEmit = vi.fn(() => true)
275+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
276+
277+
addBlockOperation('op-1', 'block-1')
278+
addWorkflowOperation('op-2')
279+
280+
useOperationQueueStore.getState().cancelOperationsForBlock('block-1')
281+
282+
expect(workflowEmit).toHaveBeenCalledTimes(2)
283+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-2')
284+
expect(useOperationQueueStore.getState().operations).toEqual([
285+
expect.objectContaining({ id: 'op-2', status: 'processing' }),
286+
])
287+
288+
useOperationQueueStore.getState().confirmOperation('op-2')
289+
expect(workflowEmit).toHaveBeenCalledTimes(2)
290+
expect(useOperationQueueStore.getState().processingOperationId).toBeNull()
291+
})
292+
293+
it('drops a failed operation with a missing target and advances once', () => {
294+
const originalRequire = Module.prototype.require
295+
const requireSpy = vi.spyOn(Module.prototype, 'require').mockImplementation(function (
296+
moduleId: string
297+
) {
298+
if (moduleId === '@/stores/workflows/workflow/store') {
299+
return { useWorkflowStore: { getState: () => ({ blocks: {} }) } }
300+
}
301+
return originalRequire.call(this, moduleId)
153302
})
154303

304+
try {
305+
const workflowEmit = vi.fn(() => true)
306+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
307+
308+
addBlockOperation('op-1', 'missing-block')
309+
addWorkflowOperation('op-2')
310+
311+
useOperationQueueStore.getState().failOperation('op-1', false)
312+
313+
expect(workflowEmit).toHaveBeenCalledTimes(2)
314+
expect(useOperationQueueStore.getState().hasOperationError).toBe(false)
315+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-2')
316+
expect(useOperationQueueStore.getState().operations).toEqual([
317+
expect.objectContaining({ id: 'op-2', status: 'processing' }),
318+
])
319+
320+
useOperationQueueStore.getState().confirmOperation('op-2')
321+
} finally {
322+
requireSpy.mockRestore()
323+
}
324+
})
325+
326+
it('triggers offline mode for a non-retryable failure and recovers via clearError', () => {
327+
const workflowEmit = vi.fn(() => true)
328+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
329+
330+
addWorkflowOperation('op-1')
331+
addWorkflowOperation('op-2')
332+
155333
useOperationQueueStore.getState().failOperation('op-1', false)
156334

335+
expect(workflowEmit).toHaveBeenCalledTimes(1)
157336
expect(useOperationQueueStore.getState().hasOperationError).toBe(true)
158337
expect(useOperationQueueStore.getState().operations).toEqual([])
338+
expect(useOperationQueueStore.getState().processingOperationId).toBeNull()
159339

160340
useOperationQueueStore.getState().clearError()
161341

162342
expect(useOperationQueueStore.getState().hasOperationError).toBe(false)
163343
})
164344

165-
it('triggers offline mode once retries exhaust for retryable failures', () => {
166-
registerEmitFunctions(
167-
vi.fn(() => true),
168-
vi.fn(),
169-
vi.fn(),
170-
'workflow-a'
171-
)
345+
it('retries only after ownership is reacquired and enters offline mode after max retries', async () => {
346+
vi.useFakeTimers()
347+
try {
348+
const workflowEmit = vi.fn(() => true)
349+
registerEmitFunctions(workflowEmit, vi.fn(), vi.fn(), 'workflow-a')
350+
addWorkflowOperation('op-1')
172351

173-
useOperationQueueStore.getState().addToQueue({
174-
id: 'op-1',
175-
workflowId: 'workflow-a',
176-
userId: 'user-1',
177-
operation: {
178-
operation: 'replace-state',
179-
target: 'workflow',
180-
payload: { state: {} },
181-
},
182-
})
352+
for (const delay of [2000, 4000, 8000]) {
353+
useOperationQueueStore.getState().failOperation('op-1', true)
183354

184-
useOperationQueueStore.getState().failOperation('op-1', true)
185-
useOperationQueueStore.getState().failOperation('op-1', true)
186-
useOperationQueueStore.getState().failOperation('op-1', true)
187-
expect(useOperationQueueStore.getState().hasOperationError).toBe(false)
355+
expect(useOperationQueueStore.getState().processingOperationId).toBeNull()
356+
expect(useOperationQueueStore.getState().hasOperationError).toBe(false)
188357

189-
useOperationQueueStore.getState().failOperation('op-1', true)
358+
await vi.advanceTimersByTimeAsync(delay)
359+
expect(useOperationQueueStore.getState().processingOperationId).toBe('op-1')
360+
}
190361

191-
expect(useOperationQueueStore.getState().hasOperationError).toBe(true)
192-
expect(useOperationQueueStore.getState().operations).toEqual([])
362+
useOperationQueueStore.getState().failOperation('op-1', true)
363+
364+
expect(workflowEmit).toHaveBeenCalledTimes(4)
365+
expect(useOperationQueueStore.getState().hasOperationError).toBe(true)
366+
expect(useOperationQueueStore.getState().operations).toEqual([])
367+
expect(useOperationQueueStore.getState().processingOperationId).toBeNull()
368+
} finally {
369+
vi.useRealTimers()
370+
}
193371
})
194372

195373
it('reports pending operations per workflow', () => {

0 commit comments

Comments
 (0)