Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions apps/web/src/lib/flows/__tests__/editor-graph.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,30 @@ describe('flow editor graph helpers', () => {
})
})

it('remaps a fork join when the join node is renamed', () => {
const flow: FlowDefinition = {
edges: [
{ id: 'edge-1', sourceNodeId: 'fork-1', targetNodeId: 'agent-2' },
{ id: 'edge-2', sourceNodeId: 'agent-2', targetNodeId: 'merge-1' },
],
layout: undefined,
nodes: [
{ id: 'fork-1', joinNodeId: 'merge-1', name: 'Fan out', type: 'fork' },
{ compactOutput: false, id: 'agent-2', name: 'Agent', promptTemplate: 'Work', targetAgentId: null, type: 'agent' },
{ id: 'merge-1', name: 'Collect', type: 'merge' },
],
startNodeId: 'fork-1',
version: 1,
}
const mergeNode = flow.nodes[2]!

const result = updateFlowDefinitionNode(flow, { ...mergeNode, name: 'Collect results' })

expect(result?.nodeId).toBe('collect-results')
const fork = result?.definition.nodes[0]
expect(fork?.type === 'fork' && fork.joinNodeId).toBe('collect-results')
})

it('renames condition targets and template references across node types', () => {
const flow: FlowDefinition = {
edges: [
Expand Down
92 changes: 92 additions & 0 deletions apps/web/src/lib/flows/__tests__/runner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -351,6 +351,98 @@ describe('triggerFlowNow', () => {
expect(mocks.markRunSucceeded).toHaveBeenCalledWith('run-1', expect.objectContaining({ openCodeSessionId: 'session-1' }))
})

describe('fork execution', () => {
function createForkFlowDefinition() {
return {
edges: [
{ id: 'edge-1', sourceNodeId: 'agent-1', targetNodeId: 'fork-1' },
{ id: 'edge-2', sourceNodeId: 'fork-1', targetNodeId: 'agent-2' },
{ id: 'edge-3', sourceNodeId: 'fork-1', targetNodeId: 'agent-3' },
{ id: 'edge-4', sourceNodeId: 'agent-2', targetNodeId: 'merge-1' },
{ id: 'edge-5', sourceNodeId: 'agent-3', targetNodeId: 'merge-1' },
{ id: 'edge-6', sourceNodeId: 'merge-1', targetNodeId: 'agent-4' },
],
nodes: [
{ compactOutput: false, id: 'agent-1', name: 'Orient', promptTemplate: 'Start', targetAgentId: null, type: 'agent' },
{ id: 'fork-1', joinNodeId: 'merge-1', name: 'Fan out', type: 'fork' },
{ compactOutput: false, id: 'agent-2', name: 'Hunt bugs', promptTemplate: 'Branch A', targetAgentId: null, type: 'agent' },
{ compactOutput: false, id: 'agent-3', name: 'Hunt perf', promptTemplate: 'Branch B', targetAgentId: null, type: 'agent' },
{ id: 'merge-1', name: 'Collect', type: 'merge' },
{ compactOutput: false, id: 'agent-4', name: 'Verify', promptTemplate: 'Merged: {{steps.agent-2.output}} / {{steps.agent-3.output}}', targetAgentId: null, type: 'agent' },
],
startNodeId: 'agent-1',
version: 1,
}
}

function createForkFlow() {
const flow = createClaimedFlow()
flow.definition = createForkFlowDefinition()
return flow
}

function promptedPrompts() {
return mocks.runFlowPromptAndReadOutput.mock.calls.map((call) => (call[0] as { prompt: string }).prompt)
}

it('runs fork branches and exposes branch outputs after the join', async () => {
mocks.userFindByIdSelect.mockResolvedValue({ slug: 'alice' })
mocks.createRun.mockResolvedValue(createRunRecord())
mocks.runFlowPromptAndReadOutput.mockImplementation(async (params: { prompt: string }) => ({
ok: true,
output: params.prompt,
}))

await runClaimedFlow(createForkFlow(), FlowRunTrigger.manual)

const prompts = promptedPrompts()
expect(prompts).toContain('Branch A')
expect(prompts).toContain('Branch B')
expect(prompts.at(-1)).toBe('Merged: Branch A / Branch B')
expect(mocks.markRunSucceeded).toHaveBeenCalledWith('run-1', expect.objectContaining({ openCodeSessionId: 'session-1' }))
expect(mocks.markRunFailed).not.toHaveBeenCalled()
})

it('creates a separate session per fork branch', async () => {
mocks.userFindByIdSelect.mockResolvedValue({ slug: 'alice' })
mocks.createRun.mockResolvedValue(createRunRecord())
mocks.runFlowPromptAndReadOutput.mockImplementation(async (params: { prompt: string }) => ({
ok: true,
output: params.prompt,
}))
const createSession = vi.fn()
.mockResolvedValueOnce({ data: { id: 'session-primary' } })
.mockResolvedValueOnce({ data: { id: 'session-branch-a' } })
.mockResolvedValueOnce({ data: { id: 'session-branch-b' } })
mocks.createInstanceClient.mockResolvedValue({ session: { create: createSession } })

await runClaimedFlow(createForkFlow(), FlowRunTrigger.manual)

expect(mocks.markRunSucceeded).toHaveBeenCalled()
const titles = createSession.mock.calls.map((call) => (call[0] as { title: string }).title)
expect(titles).toHaveLength(3)
expect(titles.some((title) => title.includes('Hunt bugs'))).toBe(true)
expect(titles.some((title) => title.includes('Hunt perf'))).toBe(true)
})

it('fails the run without executing the join when a branch fails', async () => {
mocks.userFindByIdSelect.mockResolvedValue({ slug: 'alice' })
mocks.createRun.mockResolvedValue(createRunRecord())
mocks.runFlowPromptAndReadOutput.mockImplementation(async (params: { prompt: string }) => {
if (params.prompt === 'Branch B') {
return { ok: false, type: 'failed', error: 'flow_step_failed' }
}
return { ok: true, output: params.prompt }
})

await runClaimedFlow(createForkFlow(), FlowRunTrigger.manual)

expect(promptedPrompts()).not.toContain('Merged: Branch A / Branch B')
expect(mocks.markRunFailed).toHaveBeenCalledWith('run-1', expect.objectContaining({ error: 'flow_step_failed' }))
expect(mocks.markRunSucceeded).not.toHaveBeenCalled()
})
})

it('keeps the flow run and lease active when runtime termination is unconfirmed', async () => {
mocks.userFindByIdSelect.mockResolvedValue({ slug: 'alice' })
mocks.runFlowPromptAndReadOutput.mockResolvedValue({
Expand Down
139 changes: 139 additions & 0 deletions apps/web/src/lib/flows/__tests__/validation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -259,3 +259,142 @@ describe('validateFlowDefinition', () => {
}
})
})

describe('fork topology validation', () => {
function agentNode(id: string, name: string, promptTemplate = 'Work') {
return { compactOutput: false, id, name, promptTemplate, targetAgentId: null, type: 'agent' as const }
}

function forkNode(id: string, joinNodeId: string, name = 'Fan out') {
return { id, joinNodeId, name, type: 'fork' as const }
}

function createForkDefinition(): FlowDefinition {
return {
edges: [
{ id: 'edge-1', sourceNodeId: 'agent-1', targetNodeId: 'fork-1' },
{ id: 'edge-2', sourceNodeId: 'fork-1', targetNodeId: 'agent-2' },
{ id: 'edge-3', sourceNodeId: 'fork-1', targetNodeId: 'agent-3' },
{ id: 'edge-4', sourceNodeId: 'agent-2', targetNodeId: 'merge-1' },
{ id: 'edge-5', sourceNodeId: 'agent-3', targetNodeId: 'merge-1' },
],
nodes: [
agentNode('agent-1', 'Orient'),
forkNode('fork-1', 'merge-1'),
agentNode('agent-2', 'Hunt bugs'),
agentNode('agent-3', 'Hunt perf'),
{ id: 'merge-1', name: 'Collect', type: 'merge' },
],
startNodeId: 'agent-1',
version: 1,
}
}

it('accepts a well-formed fork/join flow and keeps the join reference', () => {
const result = validateFlowDefinition(createForkDefinition())

expect(result.ok).toBe(true)
if (!result.ok) return
const fork = result.definition.nodes.find((node) => node.id === 'fork-1')
expect(fork?.type === 'fork' && fork.joinNodeId).toBe('merge-1')
})

it('rejects fork nodes without a join reference', () => {
const cases: unknown[] = [
{ ...forkNode('fork-1', 'merge-1'), joinNodeId: '' },
{ ...forkNode('fork-1', 'merge-1'), joinNodeId: ' ' },
{ id: 'fork-1', name: 'Fan out', type: 'fork' },
]

for (const fork of cases) {
const definition = createForkDefinition()
definition.nodes = [fork as FlowDefinition['nodes'][number], ...definition.nodes.slice(1)]
expect(validateFlowDefinition(definition)).toEqual({ ok: false, error: 'invalid_flow_nodes' })
}
})

it('rejects non-condition nodes with more than one outgoing edge', () => {
const definition = createForkDefinition()
definition.edges.push({ id: 'edge-6', sourceNodeId: 'agent-1', targetNodeId: 'merge-1' })

expect(validateFlowDefinition(definition)).toEqual({ ok: false, error: 'multiple_outgoing_edges:agent-1' })
})

it('rejects fork joins that are missing or not merge nodes', () => {
const missing = createForkDefinition()
missing.nodes[1] = forkNode('fork-1', 'missing')
expect(validateFlowDefinition(missing)).toEqual({ ok: false, error: 'fork_unknown_join:fork-1' })

const notMerge = createForkDefinition()
notMerge.nodes[1] = forkNode('fork-1', 'agent-2')
expect(validateFlowDefinition(notMerge)).toEqual({ ok: false, error: 'fork_join_not_merge:fork-1' })
})

it('rejects two forks declaring the same join', () => {
const definition = createForkDefinition()
definition.nodes.push(forkNode('fork-2', 'merge-1', 'Fan out again'))

expect(validateFlowDefinition(definition)).toEqual({ ok: false, error: 'fork_join_shared' })
})

it('rejects forks with fewer than two branches or a direct fork-to-join edge', () => {
const single = createForkDefinition()
single.edges = single.edges.filter((edge) => edge.id !== 'edge-3' && edge.id !== 'edge-5')
expect(validateFlowDefinition(single)).toEqual({ ok: false, error: 'fork_without_branches:fork-1' })

const direct = createForkDefinition()
direct.edges.push({ id: 'edge-6', sourceNodeId: 'fork-1', targetNodeId: 'merge-1' })
expect(validateFlowDefinition(direct)).toEqual({ ok: false, error: 'fork_branch_empty:fork-1' })
})

it('rejects branches that dead-end before reaching the join', () => {
const deadEnd = createForkDefinition()
deadEnd.edges = deadEnd.edges.filter((edge) => edge.id !== 'edge-4')
expect(validateFlowDefinition(deadEnd)).toEqual({ ok: false, error: 'fork_branch_dead_end:agent-2' })
})

it('rejects human and slack nodes inside a branch region', () => {
const definition = createForkDefinition()
definition.nodes[2] = { id: 'agent-2', instructions: 'Review', name: 'Hunt bugs', required: true, type: 'human' }

expect(validateFlowDefinition(definition)).toEqual({ ok: false, error: 'fork_branch_unsupported_node:agent-2' })
})

it('rejects join inputs from outside the branch region', () => {
const definition = createForkDefinition()
definition.nodes.push(agentNode('agent-4', 'Side chain'))
definition.edges.push({ id: 'edge-6', sourceNodeId: 'agent-4', targetNodeId: 'merge-1' })

expect(validateFlowDefinition(definition)).toEqual({ ok: false, error: 'fork_join_external_input:fork-1' })
})

it('accepts nested forks with their own joins', () => {
const definition: FlowDefinition = {
edges: [
{ id: 'edge-1', sourceNodeId: 'agent-1', targetNodeId: 'fork-1' },
{ id: 'edge-2', sourceNodeId: 'fork-1', targetNodeId: 'agent-2' },
{ id: 'edge-3', sourceNodeId: 'fork-1', targetNodeId: 'fork-2' },
{ id: 'edge-4', sourceNodeId: 'agent-2', targetNodeId: 'merge-1' },
{ id: 'edge-5', sourceNodeId: 'fork-2', targetNodeId: 'agent-4' },
{ id: 'edge-6', sourceNodeId: 'fork-2', targetNodeId: 'agent-5' },
{ id: 'edge-7', sourceNodeId: 'agent-4', targetNodeId: 'merge-2' },
{ id: 'edge-8', sourceNodeId: 'agent-5', targetNodeId: 'merge-2' },
{ id: 'edge-9', sourceNodeId: 'merge-2', targetNodeId: 'merge-1' },
],
nodes: [
agentNode('agent-1', 'Orient'),
forkNode('fork-1', 'merge-1'),
agentNode('agent-2', 'Hunt bugs'),
forkNode('fork-2', 'merge-2', 'Fan out inner'),
agentNode('agent-4', 'Hunt perf'),
agentNode('agent-5', 'Hunt improve'),
{ id: 'merge-2', name: 'Collect inner', type: 'merge' },
{ id: 'merge-1', name: 'Collect', type: 'merge' },
],
startNodeId: 'agent-1',
version: 1,
}

expect(validateFlowDefinition(definition).ok).toBe(true)
})
})
4 changes: 4 additions & 0 deletions apps/web/src/lib/flows/editor-graph.ts
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,10 @@ function updateNodeReferences(node: FlowNode, previousId: string, nextId: string
return { ...node, promptTemplate: replaceNodeVariableReferences(node.promptTemplate, previousId, nextId) }
}

if (node.type === 'fork') {
return { ...node, joinNodeId: node.joinNodeId === previousId ? nextId : node.joinNodeId }
}

return node
}

Expand Down
8 changes: 6 additions & 2 deletions apps/web/src/lib/flows/node-executor-utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,14 @@ import { FlowNodeType as PrismaFlowNodeType } from '@prisma/client'
import {
FLOW_RUN_CANCELLED_ERROR,
} from '@/lib/flows/session-executor'
import type { FlowNode } from '@/lib/flows/types'
import type { FlowNode, ForkFlowNode } from '@/lib/flows/types'
import type { FlowRunStepRecord } from '@/lib/services/flow'

export function nodeTypeToPrisma(node: FlowNode): PrismaFlowNodeType {
// Fork nodes never record run steps (they are handled by the runner loop),
// so the step-record node mapping excludes them.
export type StepRecordFlowNode = Exclude<FlowNode, ForkFlowNode>

export function nodeTypeToPrisma(node: StepRecordFlowNode): PrismaFlowNodeType {
switch (node.type) {
case 'agent':
return PrismaFlowNodeType.agent
Expand Down
Loading
Loading