feat: 引入 Eino 框架,实现四阶段生成管线
基于 compose.Graph 编排 PromptBuilder → AssetGenerator → QualitySupervisor → FormatAdapter, 含质检重试分支(最多 3 次)和降级输出。推理层提供 mock 模式,可替换为真实 API。
This commit is contained in:
@@ -0,0 +1,102 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/cloudwego/eino/compose"
|
||||
)
|
||||
|
||||
const (
|
||||
nodePromptBuilder = "prompt_builder"
|
||||
nodeAssetGenerator = "asset_generator"
|
||||
nodeQualitySupervisor = "quality_supervisor"
|
||||
nodeFormatAdapter = "format_adapter"
|
||||
)
|
||||
|
||||
// NewGenerateGraph 创建四阶段生成管线 Graph。
|
||||
//
|
||||
// START → PromptBuilder → AssetGenerator → QualitySupervisor
|
||||
// ├── pass → FormatAdapter → END
|
||||
// └── fail, retry<3 → PromptBuilder
|
||||
// └── fail, retry>=3 → FormatAdapter (降级)
|
||||
func NewGenerateGraph() (*compose.Graph[PipelineInput, PipelineOutput], error) {
|
||||
g := compose.NewGraph[PipelineInput, PipelineOutput](
|
||||
compose.WithGenLocalState(func(ctx context.Context) *PipelineState {
|
||||
return &PipelineState{}
|
||||
}),
|
||||
)
|
||||
|
||||
// 添加节点
|
||||
if err := g.AddLambdaNode(nodePromptBuilder, promptBuilderNode,
|
||||
compose.WithStatePreHandler(promptBuilderPreHandler),
|
||||
compose.WithStatePostHandler(promptBuilderPostHandler),
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("add %s node: %w", nodePromptBuilder, err)
|
||||
}
|
||||
|
||||
if err := g.AddLambdaNode(nodeAssetGenerator, assetGeneratorNode,
|
||||
compose.WithStatePostHandler(assetGeneratorPostHandler),
|
||||
); err != nil {
|
||||
return nil, fmt.Errorf("add %s node: %w", nodeAssetGenerator, err)
|
||||
}
|
||||
|
||||
if err := g.AddLambdaNode(nodeQualitySupervisor, qualitySupervisorNode); err != nil {
|
||||
return nil, fmt.Errorf("add %s node: %w", nodeQualitySupervisor, err)
|
||||
}
|
||||
|
||||
if err := g.AddLambdaNode(nodeFormatAdapter, formatAdapterNode); err != nil {
|
||||
return nil, fmt.Errorf("add %s node: %w", nodeFormatAdapter, err)
|
||||
}
|
||||
|
||||
// 连线:正常路径
|
||||
if err := g.AddEdge(compose.START, nodePromptBuilder); err != nil {
|
||||
return nil, fmt.Errorf("add edge START->%s: %w", nodePromptBuilder, err)
|
||||
}
|
||||
if err := g.AddEdge(nodePromptBuilder, nodeAssetGenerator); err != nil {
|
||||
return nil, fmt.Errorf("add edge %s->%s: %w", nodePromptBuilder, nodeAssetGenerator, err)
|
||||
}
|
||||
if err := g.AddEdge(nodeAssetGenerator, nodeQualitySupervisor); err != nil {
|
||||
return nil, fmt.Errorf("add edge %s->%s: %w", nodeAssetGenerator, nodeQualitySupervisor, err)
|
||||
}
|
||||
if err := g.AddEdge(nodeFormatAdapter, compose.END); err != nil {
|
||||
return nil, fmt.Errorf("add edge %s->END: %w", nodeFormatAdapter, err)
|
||||
}
|
||||
|
||||
// 连线:质检分支(从 state.NextNode 读取路由目标)
|
||||
if err := g.AddBranch(nodeQualitySupervisor, compose.NewGraphBranch(
|
||||
func(ctx context.Context, _ PipelineInput) (string, error) {
|
||||
var next string
|
||||
_ = compose.ProcessState[*PipelineState](ctx, func(_ context.Context, state *PipelineState) error {
|
||||
next = state.NextNode
|
||||
return nil
|
||||
})
|
||||
return next, nil
|
||||
},
|
||||
map[string]bool{nodePromptBuilder: true, nodeFormatAdapter: true},
|
||||
)); err != nil {
|
||||
return nil, fmt.Errorf("add branch at %s: %w", nodeQualitySupervisor, err)
|
||||
}
|
||||
|
||||
return g, nil
|
||||
}
|
||||
|
||||
// RunPipeline 编译并执行生成管线
|
||||
func RunPipeline(ctx context.Context, in PipelineInput) (*PipelineOutput, error) {
|
||||
g, err := NewGenerateGraph()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create graph: %w", err)
|
||||
}
|
||||
|
||||
r, err := g.Compile(ctx, compose.WithMaxRunSteps(20))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("compile graph: %w", err)
|
||||
}
|
||||
|
||||
output, err := r.Invoke(ctx, in)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("invoke pipeline: %w", err)
|
||||
}
|
||||
|
||||
return &output, nil
|
||||
}
|
||||
Reference in New Issue
Block a user