|
| 1 | +/** |
| 2 | + * @vitest-environment node |
| 3 | + */ |
| 4 | +import type { SessionPrincipal } from '@sim/auth/principal' |
| 5 | +import { createSerializedBlock, createSerializedWorkflow } from '@sim/testing' |
| 6 | +import { beforeEach, describe, expect, it, vi } from 'vitest' |
| 7 | + |
| 8 | +const { executed, gateChoice } = vi.hoisted(() => ({ |
| 9 | + executed: [] as string[], |
| 10 | + gateChoice: { value: 'if' as 'if' | 'else' }, |
| 11 | +})) |
| 12 | + |
| 13 | +vi.mock('@/executor/handlers/registry', () => ({ |
| 14 | + createBlockHandlers: () => [ |
| 15 | + { |
| 16 | + canHandle: () => true, |
| 17 | + execute: async ( |
| 18 | + _ctx: unknown, |
| 19 | + block: { id: string; metadata?: { id?: string; name?: string } } |
| 20 | + ) => { |
| 21 | + executed.push(block.id) |
| 22 | + if (block.metadata?.id === 'condition') { |
| 23 | + return { selectedOption: `${block.metadata.name}-${gateChoice.value}` } |
| 24 | + } |
| 25 | + return { ok: true } |
| 26 | + }, |
| 27 | + }, |
| 28 | + ], |
| 29 | +})) |
| 30 | + |
| 31 | +import { BlockType } from '@/executor/constants' |
| 32 | +import { DAGExecutor } from '@/executor/execution/executor' |
| 33 | +import { stripCloneSuffixes } from '@/executor/utils/subflow-utils' |
| 34 | +import type { SerializedWorkflow } from '@/serializer/types' |
| 35 | + |
| 36 | +const PRINCIPAL: SessionPrincipal = { kind: 'session', userId: 'user-1', sessionId: 'session-1' } |
| 37 | + |
| 38 | +/** Block ids double as their type for the special blocks; every other block is a function. */ |
| 39 | +const BLOCK_TYPES: Record<string, string> = { |
| 40 | + start: BlockType.STARTER, |
| 41 | + gate: BlockType.CONDITION, |
| 42 | + loop: BlockType.LOOP, |
| 43 | + fanOut: BlockType.PARALLEL, |
| 44 | +} |
| 45 | + |
| 46 | +function workflow( |
| 47 | + connections: SerializedWorkflow['connections'], |
| 48 | + loops: SerializedWorkflow['loops'] = {}, |
| 49 | + parallels: SerializedWorkflow['parallels'] = {} |
| 50 | +): SerializedWorkflow { |
| 51 | + const ids = new Set(connections.flatMap((c) => [c.source, c.target])) |
| 52 | + const blocks = [...ids].map((id) => |
| 53 | + createSerializedBlock({ id, name: id, type: BLOCK_TYPES[id] ?? BlockType.FUNCTION }) |
| 54 | + ) |
| 55 | + return { ...createSerializedWorkflow(blocks, connections), loops, parallels } |
| 56 | +} |
| 57 | + |
| 58 | +function createExecutor(wf: SerializedWorkflow, executionId: string): DAGExecutor { |
| 59 | + return new DAGExecutor({ |
| 60 | + workflow: wf, |
| 61 | + contextExtensions: { workspaceId: 'ws', executionId, principal: PRINCIPAL }, |
| 62 | + }) |
| 63 | +} |
| 64 | + |
| 65 | +/** A full run with one gate decision, then a run-from-block seeded from that run's snapshot. */ |
| 66 | +async function rerunFromBlock( |
| 67 | + wf: SerializedWorkflow, |
| 68 | + startBlockId: string, |
| 69 | + first: 'if' | 'else', |
| 70 | + second: 'if' | 'else' |
| 71 | +): Promise<string[]> { |
| 72 | + gateChoice.value = first |
| 73 | + const run1 = await createExecutor(wf, 'exec-1').execute('wf') |
| 74 | + expect(run1.success).toBe(true) |
| 75 | + expect(run1.executionState).toBeDefined() |
| 76 | + |
| 77 | + executed.length = 0 |
| 78 | + gateChoice.value = second |
| 79 | + const run2 = await createExecutor(wf, 'exec-2').executeFromBlock( |
| 80 | + 'wf', |
| 81 | + startBlockId, |
| 82 | + run1.executionState! |
| 83 | + ) |
| 84 | + expect(run2.success).toBe(true) |
| 85 | + return [...executed] |
| 86 | +} |
| 87 | + |
| 88 | +const exclusive = workflow([ |
| 89 | + { source: 'start', target: 'prep' }, |
| 90 | + { source: 'prep', target: 'gate' }, |
| 91 | + { source: 'gate', target: 'readDoc', sourceHandle: 'condition-gate-if' }, |
| 92 | + { source: 'gate', target: 'noDoc', sourceHandle: 'condition-gate-else' }, |
| 93 | +]) |
| 94 | + |
| 95 | +const diamond = workflow([ |
| 96 | + { source: 'start', target: 'prep' }, |
| 97 | + { source: 'prep', target: 'gate' }, |
| 98 | + { source: 'gate', target: 'readDoc', sourceHandle: 'condition-gate-if' }, |
| 99 | + { source: 'gate', target: 'docPlan', sourceHandle: 'condition-gate-else' }, |
| 100 | + { source: 'readDoc', target: 'docPlan' }, |
| 101 | +]) |
| 102 | + |
| 103 | +const siblings = workflow([ |
| 104 | + { source: 'start', target: 'a' }, |
| 105 | + { source: 'start', target: 'b' }, |
| 106 | + { source: 'a', target: 'join' }, |
| 107 | + { source: 'b', target: 'join' }, |
| 108 | +]) |
| 109 | + |
| 110 | +const looped = workflow( |
| 111 | + [ |
| 112 | + { source: 'start', target: 'prep' }, |
| 113 | + { source: 'prep', target: 'loop' }, |
| 114 | + { source: 'loop', target: 'gate', sourceHandle: 'loop-start-source' }, |
| 115 | + { source: 'gate', target: 'readDoc', sourceHandle: 'condition-gate-if' }, |
| 116 | + { source: 'gate', target: 'noDoc', sourceHandle: 'condition-gate-else' }, |
| 117 | + ], |
| 118 | + { loop: { id: 'loop', nodes: ['gate', 'readDoc', 'noDoc'], iterations: 1, loopType: 'for' } } |
| 119 | +) |
| 120 | + |
| 121 | +describe('DAGExecutor run-from-block edge state', () => { |
| 122 | + beforeEach(() => { |
| 123 | + executed.length = 0 |
| 124 | + }) |
| 125 | + |
| 126 | + it.each([ |
| 127 | + { |
| 128 | + title: 'does not run an unselected branch the source execution had activated', |
| 129 | + wf: exclusive, |
| 130 | + start: 'prep', |
| 131 | + first: 'if', |
| 132 | + second: 'else', |
| 133 | + expected: ['prep', 'gate', 'noDoc'], |
| 134 | + }, |
| 135 | + { |
| 136 | + title: 'runs the branch the source execution had deactivated when it is selected', |
| 137 | + wf: exclusive, |
| 138 | + start: 'prep', |
| 139 | + first: 'else', |
| 140 | + second: 'if', |
| 141 | + expected: ['prep', 'gate', 'readDoc'], |
| 142 | + }, |
| 143 | + { |
| 144 | + title: 'waits for the live input of a join the source execution had released early', |
| 145 | + wf: diamond, |
| 146 | + start: 'prep', |
| 147 | + first: 'else', |
| 148 | + second: 'if', |
| 149 | + expected: ['prep', 'gate', 'readDoc', 'docPlan'], |
| 150 | + }, |
| 151 | + { |
| 152 | + title: 'still runs the join directly when its other input is deselected', |
| 153 | + wf: diamond, |
| 154 | + start: 'prep', |
| 155 | + first: 'if', |
| 156 | + second: 'else', |
| 157 | + expected: ['prep', 'gate', 'docPlan'], |
| 158 | + }, |
| 159 | + { |
| 160 | + title: 'does not wait on a cached sibling input outside the re-run region', |
| 161 | + wf: siblings, |
| 162 | + start: 'a', |
| 163 | + first: 'if', |
| 164 | + second: 'if', |
| 165 | + expected: ['a', 'join'], |
| 166 | + }, |
| 167 | + { |
| 168 | + title: 'does not run a stale branch inside a loop', |
| 169 | + wf: looped, |
| 170 | + start: 'prep', |
| 171 | + first: 'if', |
| 172 | + second: 'else', |
| 173 | + expected: ['prep', 'gate', 'noDoc'], |
| 174 | + }, |
| 175 | + ] as const)('$title', async ({ wf, start, first, second, expected }) => { |
| 176 | + expect(await rerunFromBlock(wf, start, first, second)).toEqual(expected) |
| 177 | + }) |
| 178 | + |
| 179 | + it('runs a join once after every re-run input completes', async () => { |
| 180 | + const fanIn = workflow([ |
| 181 | + { source: 'start', target: 'prep' }, |
| 182 | + { source: 'prep', target: 'a' }, |
| 183 | + { source: 'prep', target: 'b' }, |
| 184 | + { source: 'a', target: 'join' }, |
| 185 | + { source: 'b', target: 'join' }, |
| 186 | + ]) |
| 187 | + const order = await rerunFromBlock(fanIn, 'prep', 'if', 'if') |
| 188 | + expect(order.filter((id) => id === 'join')).toHaveLength(1) |
| 189 | + expect(order.indexOf('join')).toBeGreaterThan(Math.max(order.indexOf('a'), order.indexOf('b'))) |
| 190 | + }) |
| 191 | + |
| 192 | + it('does not run a stale branch in any parallel branch copy', async () => { |
| 193 | + const parallel = workflow( |
| 194 | + [ |
| 195 | + { source: 'start', target: 'prep' }, |
| 196 | + { source: 'prep', target: 'fanOut' }, |
| 197 | + { source: 'fanOut', target: 'gate', sourceHandle: 'parallel-start-source' }, |
| 198 | + { source: 'gate', target: 'readDoc', sourceHandle: 'condition-gate-if' }, |
| 199 | + { source: 'gate', target: 'noDoc', sourceHandle: 'condition-gate-else' }, |
| 200 | + ], |
| 201 | + {}, |
| 202 | + { |
| 203 | + fanOut: { |
| 204 | + id: 'fanOut', |
| 205 | + nodes: ['gate', 'readDoc', 'noDoc'], |
| 206 | + count: 3, |
| 207 | + parallelType: 'count', |
| 208 | + }, |
| 209 | + } |
| 210 | + ) |
| 211 | + const ids = (await rerunFromBlock(parallel, 'prep', 'if', 'else')).map(stripCloneSuffixes) |
| 212 | + expect(ids).not.toContain('readDoc') |
| 213 | + expect(ids.filter((id) => id === 'noDoc')).toHaveLength(3) |
| 214 | + }) |
| 215 | +}) |
0 commit comments