Repository navigation
Expand file tree
/
Copy pathWorkflowEngine.ts
More file actions
123 lines (106 loc) · 4.6 KB
/
Copy pathWorkflowEngine.ts
File metadata and controls
123 lines (106 loc) · 4.6 KB
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
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
import { randomUUID } from 'node:crypto'
import type { ExecutionContext } from '../core/ExecutionContext.ts'
import type { IProtocol } from './protocols/IProtocol.ts'
import type { WorkflowStep, WorkflowContext, WorkflowDefinition } from './types.ts'
import { UserProtocol } from './protocols/UserProtocol.ts'
import { AgentProtocol } from './protocols/AgentProtocol.ts'
import { ToolProtocol } from './protocols/ToolProtocol.ts'
import { WorkflowCheckpoint } from './WorkflowCheckpoint.ts'
export class WorkflowEngine {
private protocols: Map<string, IProtocol<unknown, unknown>> = new Map()
private checkpoint: WorkflowCheckpoint
constructor() {
this.protocols.set('user-protocol', new UserProtocol())
this.protocols.set('agent-protocol', new AgentProtocol())
this.protocols.set('tool-protocol', new ToolProtocol())
this.checkpoint = new WorkflowCheckpoint()
}
registerProtocol(name: string, protocol: IProtocol<unknown, unknown>): void {
this.protocols.set(name, protocol)
}
async execute(definition: WorkflowDefinition, context: ExecutionContext): Promise<WorkflowContext> {
// A new correlationId is minted per top-level workflow run; nested
// workflows inherit the parent's so all step traces chain together.
const correlationId = context.correlationId ?? randomUUID()
const stepContext: ExecutionContext = { ...context, correlationId }
const wfContext: WorkflowContext = {
executionId: context.executionId,
artifacts: new Map(),
stepResults: new Map(),
metadata: { workflowName: definition.name, correlationId },
}
console.log(`\n[WorkflowEngine] Starting "${definition.name}" (${definition.steps.length} steps)\n`)
for (const step of definition.steps) {
if (step.type === 'agent' || step.type === 'tool') {
const cpId = this.checkpoint.save(step.name, wfContext)
console.log(`[WorkflowEngine] Checkpoint saved for step "${step.name}" (${cpId})`)
}
try {
const result = await this.executeStep(step, stepContext, wfContext)
wfContext.stepResults.set(step.name, result)
} catch (err) {
if (step.type === 'human-in-loop') {
const latest = this.checkpoint.getLatest()
if (latest) {
console.log(`[WorkflowEngine] Human rejected step "${step.name}" — restoring checkpoint "${latest.id}"`)
this.checkpoint.restore(latest.id, wfContext)
this.checkpoint.clear(latest.id)
}
}
throw err
}
if (step.type === 'human-in-loop') {
this.checkpoint.clearAll()
console.log(`[WorkflowEngine] Human approved step "${step.name}" — checkpoints cleared`)
}
}
return wfContext
}
private async executeStep(step: WorkflowStep, context: ExecutionContext, _wfContext: WorkflowContext): Promise<unknown> {
const protocol = step.protocol ?? this.resolveProtocol(step)
if (!protocol) {
console.log(`[WorkflowEngine] Step "${step.name}" has no protocol — skipping`)
return null
}
const maxRetries = protocol.config.retryPolicy.maxRetries
const backoffMs = protocol.config.retryPolicy.backoffMs
for (let attempt = 0; attempt <= maxRetries; attempt++) {
try {
const result = await protocol.execute(step.input ?? {}, context)
return result
} catch (err) {
const error = err instanceof Error ? err : new Error(String(err))
const action = protocol.onFailure(error, attempt)
console.log(`[WorkflowEngine] Step "${step.name}" failed (attempt ${attempt + 1}/${maxRetries + 1}): ${error.message}`)
if (action === 'abort') {
throw error
}
if (action === 'escalate') {
console.log(`[WorkflowEngine] Escalating failure for step "${step.name}"`)
throw error
}
if (action === 'retry' && attempt < maxRetries) {
const delay = backoffMs * Math.pow(2, attempt)
console.log(`[WorkflowEngine] Retrying step "${step.name}" in ${delay}ms`)
await new Promise((resolve) => setTimeout(resolve, delay))
}
}
}
throw new Error(`Step "${step.name}" failed after ${maxRetries + 1} attempts`)
}
private resolveProtocol(step: WorkflowStep): IProtocol<unknown, unknown> | undefined {
switch (step.type) {
case 'agent':
case 'goal':
return this.protocols.get('agent-protocol')
case 'human-in-loop':
return this.protocols.get('user-protocol')
case 'tool':
return this.protocols.get('tool-protocol')
case 'parallel':
return undefined
default:
return undefined
}
}
}