1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
| // 工作流定义
interface Workflow {
id: string;
states: WorkflowState[];
startAt: string;
}
type WorkflowState =
| TaskState
| ParallelState
| ChoiceState
| WaitState
| MapState;
// 工作流执行引擎
class WorkflowEngine {
async execute(workflow: Workflow, input: unknown): Promise<WorkflowResult> {
const execution: WorkflowExecution = {
id: this.generateId(),
workflowId: workflow.id,
status: 'running',
history: [],
startedAt: Date.now()
};
try {
let currentState = workflow.startAt;
let currentInput = input;
while (currentState) {
const state = workflow.states.find(s => s.name === currentState);
// 记录执行历史
execution.history.push({
state: currentState,
timestamp: Date.now(),
input: currentInput
});
// 执行状态
const result = await this.executeState(state, currentInput);
// 持久化中间状态
await this.stateStore.save(execution.id, {
state: currentState,
result
});
// 转移到下一个状态
currentState = state.next;
currentInput = result.output;
}
execution.status = 'succeeded';
execution.completedAt = Date.now();
execution.output = currentInput;
return execution;
} catch (error) {
execution.status = 'failed';
execution.error = error.message;
throw error;
}
}
private async executeState(
state: WorkflowState,
input: unknown
): Promise<StateExecutionResult> {
switch (state.type) {
case 'Task':
return this.executeTask(state, input);
case 'Parallel':
return this.executeParallel(state, input);
case 'Choice':
return this.executeChoice(state, input);
case 'Map':
return this.executeMap(state, input);
default:
throw new Error(`Unknown state type: ${state.type}`);
}
}
private async executeMap(
state: MapState,
input: unknown[]
): Promise<StateExecutionResult> {
const items = Array.isArray(input) ? input : [input];
// 并行处理数组
const results = await Promise.all(
items.map(item =>
this.executeState(state.iteration, item)
)
);
return {
output: results.map(r => r.output),
status: 'succeeded'
};
}
}
|