Documentation
¶
Index ¶
- Constants
- Variables
- func ExtractExpression(label string) string
- type AIConfig
- type Adapter
- type Bridge
- type DGAEdge
- type DGAEvaluationContext
- func (c *DGAEvaluationContext) All() map[string]any
- func (c *DGAEvaluationContext) Get(key string) (any, bool)
- func (c *DGAEvaluationContext) WithNode(node Node) EvaluationContext
- func (c *DGAEvaluationContext) WithParams(params map[string]any) EvaluationContext
- func (c *DGAEvaluationContext) WithPipeline(pipeline Pipeline) EvaluationContext
- type DGAGraph
- func (dga *DGAGraph) AddEdge(edge Edge) error
- func (dga *DGAGraph) AddVertex(node Node)
- func (dga *DGAGraph) Edges() []Edge
- func (dga *DGAGraph) HasCycle() bool
- func (dga *DGAGraph) Nodes() map[string]Node
- func (dga *DGAGraph) Traversal(ctx context.Context, evalCtx EvaluationContext, fn TraversalFn) error
- type DGANode
- type DefaultMetadataStoreFactory
- type Edge
- type Entry
- type EvaluationContext
- type Event
- type Executor
- type ExecutorConfig
- type Graph
- type GraphReader
- type HTTPMetadataConfig
- type HTTPMetadataStore
- type InConfigMetadataStore
- type Level
- type Listener
- type ListeningFn
- type LoggingConfig
- type Metadata
- type MetadataConfig
- type MetadataStore
- type MetadataStoreFactory
- type Node
- type NodeConfig
- type Pipeline
- type PipelineConfig
- type PipelineImpl
- func (p *PipelineImpl) Cancel()
- func (p *PipelineImpl) Done() <-chan struct{}
- func (p *PipelineImpl) GetGraph() Graph
- func (p *PipelineImpl) Id() string
- func (p *PipelineImpl) Listening(fn Listener)
- func (p *PipelineImpl) Metadata() Metadata
- func (p *PipelineImpl) Notify()
- func (p *PipelineImpl) Run(ctx context.Context) error
- func (p *PipelineImpl) SetGraph(graph Graph)
- func (p *PipelineImpl) SetMetadata(store MetadataStore)
- func (p *PipelineImpl) Status() string
- type Pongo2TemplateEngine
- type Pusher
- type RedisMetadataConfig
- type RedisMetadataStore
- type Runtime
- type RuntimeImpl
- func (r *RuntimeImpl) BuildGraph(config *PipelineConfig) Graph
- func (r *RuntimeImpl) Cancel(ctx context.Context, id string) error
- func (r *RuntimeImpl) Ctx() context.Context
- func (r *RuntimeImpl) Done() chan struct{}
- func (r *RuntimeImpl) Get(id string) (Pipeline, error)
- func (r *RuntimeImpl) Notify(data interface{}) error
- func (r *RuntimeImpl) Rm(id string)
- func (r *RuntimeImpl) RunAsync(ctx context.Context, id string, config string, listener Listener) (Pipeline, error)
- func (r *RuntimeImpl) RunSync(ctx context.Context, id string, config string, listener Listener) (Pipeline, error)
- func (r *RuntimeImpl) SetPusher(pusher Pusher)
- func (r *RuntimeImpl) SetTemplateEngine(engine TemplateEngine)
- func (r *RuntimeImpl) StartBackground()
- func (r *RuntimeImpl) StopBackground()
- type Step
- type TemplateEngine
- type TraversalFn
Constants ¶
const ( // 流水线状态常量 StatusRunning = "RUNNING" StatusFailed = "FAILED" StatusSuccess = "SUCCESS" StatusTerminate = "ABORTED" StatusPaused = "PAUSED" StatusUnknown = "UNKNOWN" StatusCancelled = "CANCELLED" // 流水线事件常量 EventPipelineInit = "pipeline-init" EventPipelineStart = "pipeline-start" EventPipelineFinish = "pipeline-finish" EventPipelineExecutorPrepare = "pipeline-executor-prepare" EventPipelineExecutorPrepareDone = "pipeline-executor-prepare-done" EventPipelineNodeStart = "pipeline-node-start" EventPipelineNodeFinish = "pipeline-node-finish" EventPipelineCancelled = "pipeline-cancelled" EventPipelineStatusUpdate = "pipeline-status-update" )
Variables ¶
var ( ErrInvalidGraph = errors.New("invalid graph") ErrHasCycle = errors.New("has cycle") )
Functions ¶
func ExtractExpression ¶
ExtractExpression 从边标签中提取条件表达式(公共函数供测试使用) 使用模板引擎的 Validate 方法验证表达式语法 先检查是否包含模板标记 {{ 或 {%,再使用模板引擎验证
Types ¶
type AIConfig ¶
type AIConfig struct {
Intent string `yaml:"intent"` // 核心意图描述
Constraints []string `yaml:"constraints"` // 关键约束列表
Template string `yaml:"template"` // 模板标识
GeneratedAt string `yaml:"generatedAt"` // 生成时间
Version int `yaml:"version"` // 版本号
}
AIConfig AI配置结构
type DGAEdge ¶
type DGAEdge struct {
// contains filtered or unexported fields
}
DGAEdge 是Edge接口的实现
func (*DGAEdge) Evaluate ¶
func (e *DGAEdge) Evaluate(ctx EvaluationContext) (bool, error)
Evaluate 评估条件表达式 如果表达式为空,返回true(无条件边总是可以通过) 如果有表达式但没有设置模板引擎,使用默认的Pongo2模板引擎
func (*DGAEdge) SetEngine ¶
func (e *DGAEdge) SetEngine(engine TemplateEngine)
SetEngine 设置模板引擎(用于延迟初始化)
type DGAEvaluationContext ¶
type DGAEvaluationContext struct {
// contains filtered or unexported fields
}
DGAEvaluationContext 是EvaluationContext接口的实现
func (*DGAEvaluationContext) All ¶
func (c *DGAEvaluationContext) All() map[string]any
All 返回上下文中所有数据的副本 合并了:基础数据、节点数据、流水线数据
func (*DGAEvaluationContext) Get ¶
func (c *DGAEvaluationContext) Get(key string) (any, bool)
Get 从上下文中获取值
func (*DGAEvaluationContext) WithNode ¶
func (c *DGAEvaluationContext) WithNode(node Node) EvaluationContext
WithNode 设置当前节点并返回新的上下文(链式调用)
func (*DGAEvaluationContext) WithParams ¶
func (c *DGAEvaluationContext) WithParams(params map[string]any) EvaluationContext
WithParams 添加参数到上下文并返回新的上下文(链式调用)
func (*DGAEvaluationContext) WithPipeline ¶
func (c *DGAEvaluationContext) WithPipeline(pipeline Pipeline) EvaluationContext
WithPipeline 设置流水线并返回新的上下文(链式调用)
type DGAGraph ¶
type DGAGraph struct {
// contains filtered or unexported fields
}
保存了流水线的图结构
func NewDGAGraph ¶
func NewDGAGraph() *DGAGraph
func (*DGAGraph) Traversal ¶
func (dga *DGAGraph) Traversal(ctx context.Context, evalCtx EvaluationContext, fn TraversalFn) error
Traversal 对DAG执行广度优先遍历 为图中的每个节点执行提供的 TraversalFn 函数 支持多个起始节点并发执行 支持条件边:如果边有表达式,会评估表达式决定是否遍历该边
type DGANode ¶
type DGANode struct {
// contains filtered or unexported fields
}
func NewDGANode ¶
NewDGANode creates a new DGANode with the specified id and state, initializing an empty property map.
func (*DGANode) PipelineId ¶
type DefaultMetadataStoreFactory ¶
type DefaultMetadataStoreFactory struct{}
DefaultMetadataStoreFactory 默认的元数据存储工厂
func (*DefaultMetadataStoreFactory) Create ¶
func (f *DefaultMetadataStoreFactory) Create(config MetadataConfig) (MetadataStore, error)
Create 根据配置类型创建对应的MetadataStore实例
type Edge ¶
type Edge interface {
// Source 返回边的源节点
Source() Node
// Target 返回边的目标节点
Target() Node
// Expression 返回边的条件表达式(pongo2模板语法)
// 如果返回空字符串,表示无条件边,总是可以遍历
Expression() string
// Evaluate 评估条件表达式,返回bool表示是否通过
Evaluate(ctx EvaluationContext) (bool, error)
// ID 返回边的唯一标识符(格式:source->target)
ID() string
}
Edge 表示DAG中的边,支持条件表达式
func NewConditionalEdge ¶
NewConditionalEdge 创建一条条件边
func NewConditionalEdgeWithEngine ¶
func NewConditionalEdgeWithEngine(source, target Node, expression string, engine TemplateEngine) Edge
NewConditionalEdgeWithEngine 创建一条带有模板引擎的条件边
type Entry ¶
type Entry struct {
Pipeline string `json:"pipeline"`
BuildID string `json:"buildId"`
Node string `json:"node"`
Step string `json:"step"`
Timestamp time.Time `json:"timestamp"`
Level Level `json:"level"`
Message string `json:"message"`
Output string `json:"output"` // 命令标准输出/错误
}
Entry 单条日志
type EvaluationContext ¶
type EvaluationContext interface {
Get(key string) (any, bool)
All() map[string]any
WithNode(node Node) EvaluationContext
WithPipeline(pipeline Pipeline) EvaluationContext
WithParams(params map[string]any) EvaluationContext
}
EvaluationContext 表达式求值上下文
func NewEvaluationContext ¶
func NewEvaluationContext() EvaluationContext
NewEvaluationContext 创建一个新的求值上下文
type Event ¶
type Event string
流水线事件
var ( //监听事件 PipelineInit Event = EventPipelineInit // 流水线初始化 PipelineStart Event = EventPipelineStart // 流水线开始执行 PipelineFinish Event = EventPipelineFinish // 流水线完成 PipelineExecutorPrepare Event = EventPipelineExecutorPrepare // 流水线执行器开始准备 PipelineExecutorPrepareDone Event = EventPipelineExecutorPrepareDone // 流水线执行器准备完毕 PipelineNodeStart Event = EventPipelineNodeStart // 节点开始 PipelineNodeFinish Event = EventPipelineNodeFinish // 节点完成 )
type Executor ¶
type Executor interface {
// Prepare 准备环境
Prepare(ctx context.Context) error
// Destruction 销毁环境
Destruction(ctx context.Context) error
// Transfer 传输需要执行的数据,并且反回执行的结果
Transfer(ctx context.Context, in chan<- any, out <-chan any)
}
Executor 执行器
type ExecutorConfig ¶
type ExecutorConfig struct {
Type string `yaml:"type"`
Config map[string]interface{} `yaml:"config"`
}
ExecutorConfig 执行器配置结构
type Graph ¶
type Graph interface {
GraphReader
//AddVertex 添加顶点
AddVertex(node Node)
//AddEdge 添加边
AddEdge(edge Edge) error
}
type GraphReader ¶
type GraphReader interface {
//Nodes
Nodes() map[string]Node
//Edges 返回所有的边
Edges() []Edge
//Traversal 遍历图结构
Traversal(ctx context.Context, evalCtx EvaluationContext, fn TraversalFn) error
}
type HTTPMetadataConfig ¶
type HTTPMetadataConfig struct {
URL string `yaml:"url"`
Method string `yaml:"method"`
Headers map[string]string `yaml:"headers"`
Timeout string `yaml:"timeout"`
}
HTTPMetadataConfig HTTP元数据配置
type HTTPMetadataStore ¶
type HTTPMetadataStore struct {
// contains filtered or unexported fields
}
HTTPMetadataStore HTTP 元数据存储
func NewHTTPMetadataStore ¶
func NewHTTPMetadataStore(config MetadataConfig) (*HTTPMetadataStore, error)
NewHTTPMetadataStore 创建基于HTTP的元数据存储
func (*HTTPMetadataStore) Delete ¶
func (s *HTTPMetadataStore) Delete(ctx context.Context, key string) error
Delete 通过HTTP接口删除元数据
type InConfigMetadataStore ¶
type InConfigMetadataStore struct {
// contains filtered or unexported fields
}
InConfigMetadataStore 从配置中直接读取数据的元数据存储
func NewInConfigMetadataStore ¶
func NewInConfigMetadataStore(config MetadataConfig) (*InConfigMetadataStore, error)
NewInConfigMetadataStore 创建基于配置的元数据存储
func (*InConfigMetadataStore) Delete ¶
func (s *InConfigMetadataStore) Delete(ctx context.Context, key string) error
Delete 删除元数据(in-config 类型为只读,返回错误)
type Listener ¶
type Listener interface {
// 处理对应的事件将事件发生的对应的流水线和对应的事件作为参数传入
Handle(p Pipeline, event Event)
// 获取当前注册的Event
Events() []Event
}
我们将整个流水线的运行过程中的事件抽象成对应的Event 这样我们就能再外部监听Event
type LoggingConfig ¶
type LoggingConfig struct {
Endpoint string `yaml:"endpoint"`
Headers map[string]string `yaml:"headers"`
Timeout string `yaml:"timeout"`
Retry int `yaml:"retry"`
}
LoggingConfig 日志配置结构
type MetadataConfig ¶
MetadataConfig 元数据配置结构
type MetadataStore ¶
type MetadataStore interface {
// Get 获取元数据值
Get(ctx context.Context, key string) (string, error)
// Set 设置元数据值
Set(ctx context.Context, key string, value string) error
// Delete 删除元数据
Delete(ctx context.Context, key string) error
// Close 关闭元数据存储连接
Close() error
}
MetadataStore 元数据存储接口
type MetadataStoreFactory ¶
type MetadataStoreFactory interface {
// Create 根据配置创建MetadataStore实例
Create(config MetadataConfig) (MetadataStore, error)
}
MetadataStoreFactory 元数据存储工厂接口,用于创建MetadataStore实例
func NewMetadataStoreFactory ¶
func NewMetadataStoreFactory() MetadataStoreFactory
NewMetadataStoreFactory 创建默认的元数据存储工厂
type NodeConfig ¶
type NodeConfig struct {
Executor string `yaml:"executor"`
Image string `yaml:"image"`
Steps []Step `yaml:"steps"`
Config map[string]interface{} `yaml:"Config"`
}
NodeConfig 节点配置结构
type Pipeline ¶
type Pipeline interface {
//ID 流水线的id
Id() string
//GetGraph 返回图结构
GetGraph() Graph
//SetGraph 设置图结构
SetGraph(graph Graph)
//Status 返回流水线的整体状态
Status() string
//SetMetadata 设置元数据
SetMetadata(store MetadataStore)
//Metadata 获取元数据
Metadata() Metadata
//Listening 流水线执行事件监听设置
Listening(listener Listener)
//Done流水线是否执行完成
Done() <-chan struct{}
//Run执行流水线
Run(ctx context.Context) error
//Notify 执行的步骤通知流水线
Notify()
//Cancel 取消流水线
Cancel()
}
func NewPipeline ¶
type PipelineConfig ¶
type PipelineConfig struct {
Version string `yaml:"Version"`
Name string `yaml:"Name"`
Metadate MetadataConfig `yaml:"Metadate"`
AI AIConfig `yaml:"AI"`
Param map[string]interface{} `yaml:"Param"`
Executors map[string]ExecutorConfig `yaml:"Executors"`
Logging LoggingConfig `yaml:"Logging"`
Graph string `yaml:"Graph"`
Status map[string]string `yaml:"Status"`
Nodes map[string]NodeConfig `yaml:"Nodes"`
}
PipelineConfig 流水线配置结构
type PipelineImpl ¶
type PipelineImpl struct {
// contains filtered or unexported fields
}
func (*PipelineImpl) Listening ¶
func (p *PipelineImpl) Listening(fn Listener)
Listening 设置流水线执行事件监听器
func (*PipelineImpl) Notify ¶
func (p *PipelineImpl) Notify()
这个主要是在运行过程中节点状态或者流水线状态变化,就会触发这个函数 节点 我们就可以在这里做一些处理 执行ListeningFn函数
func (*PipelineImpl) SetMetadata ¶
func (p *PipelineImpl) SetMetadata(store MetadataStore)
SetMetadata 设置流水线的元数据存储
type Pongo2TemplateEngine ¶
type Pongo2TemplateEngine struct{}
Pongo2TemplateEngine 使用pongo2作为模板引擎的实现
func (*Pongo2TemplateEngine) EvaluateBool ¶
EvaluateBool 评估模板表达式,返回布尔值
func (*Pongo2TemplateEngine) EvaluateString ¶
func (e *Pongo2TemplateEngine) EvaluateString(expression string, ctx map[string]any) (string, error)
EvaluateString 评估模板表达式,返回字符串
func (*Pongo2TemplateEngine) Validate ¶
func (e *Pongo2TemplateEngine) Validate(expression string) error
Validate 验证表达式语法是否正确
type Pusher ¶
type Pusher interface {
// Push 推送单条日志
Push(ctx context.Context, entry Entry) error
// PushBatch 批量推送
PushBatch(ctx context.Context, entries []Entry) error
// Close 关闭连接,刷新缓冲
Close() error
}
Pusher 日志推送接口
type RedisMetadataConfig ¶
type RedisMetadataConfig struct {
Host string `yaml:"host"`
Port int `yaml:"port"`
DB int `yaml:"db"`
Username string `yaml:"username"`
Password string `yaml:"password"`
}
RedisMetadataConfig Redis元数据配置
type RedisMetadataStore ¶
type RedisMetadataStore struct {
// contains filtered or unexported fields
}
RedisMetadataStore Redis 元数据存储
func NewRedisMetadataStore ¶
func NewRedisMetadataStore(config MetadataConfig) (*RedisMetadataStore, error)
NewRedisMetadataStore 创建基于Redis的元数据存储
func (*RedisMetadataStore) Delete ¶
func (s *RedisMetadataStore) Delete(ctx context.Context, key string) error
Delete 删除Redis元数据
type Runtime ¶
type Runtime interface {
//获取流水线状态
Get(id string) (Pipeline, error)
//取消运行中的流水线
Cancel(ctx context.Context, id string) error
//执行异步流水线
RunAsync(ctx context.Context, id string, config string, listener Listener) (Pipeline, error)
//执行同步流水线
RunSync(ctx context.Context, id string, config string, listener Listener) (Pipeline, error)
//移除流水线记录
Rm(id string)
//runtime已经执行完成
Done() chan struct{}
//通知runtime
Notify(data interface{}) error
//反回runtime公共
Ctx() context.Context
//停止后台处理
StopBackground()
// 启动后台
StartBackground()
// 设置日志推送器
SetPusher(pusher Pusher)
// 设置模板引擎
SetTemplateEngine(engine TemplateEngine)
}
Runtime 运行时
type RuntimeImpl ¶
type RuntimeImpl struct {
// contains filtered or unexported fields
}
RuntimeImpl Runtime接口的实现
func (*RuntimeImpl) BuildGraph ¶
func (r *RuntimeImpl) BuildGraph(config *PipelineConfig) Graph
BuildGraph 构建图结构
func (*RuntimeImpl) Cancel ¶
func (r *RuntimeImpl) Cancel(ctx context.Context, id string) error
Cancel 取消运行中的流水线
func (*RuntimeImpl) RunAsync ¶
func (r *RuntimeImpl) RunAsync(ctx context.Context, id string, config string, listener Listener) (Pipeline, error)
RunAsync 执行异步流水线
func (*RuntimeImpl) RunSync ¶
func (r *RuntimeImpl) RunSync(ctx context.Context, id string, config string, listener Listener) (Pipeline, error)
RunSync 执行同步流水线
func (*RuntimeImpl) SetTemplateEngine ¶
func (r *RuntimeImpl) SetTemplateEngine(engine TemplateEngine)
SetTemplateEngine 设置模板引擎
func (*RuntimeImpl) StartBackground ¶
func (r *RuntimeImpl) StartBackground()
StartBackground 启动后台处理
type TemplateEngine ¶
type TemplateEngine interface {
// EvaluateBool 评估模板表达式,返回布尔值
EvaluateBool(expression string, ctx map[string]any) (bool, error)
// EvaluateString 评估模板表达式,返回字符串
EvaluateString(expression string, ctx map[string]any) (string, error)
// Validate 验证表达式语法是否正确
Validate(expression string) error
}
TemplateEngine 模板引擎接口,用于表达式求值
func NewPongo2TemplateEngine ¶
func NewPongo2TemplateEngine() TemplateEngine
NewPongo2TemplateEngine 创建一个新的Pongo2模板引擎实例
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
executor
|
|
|
kubenetes/apis/agentcontroller/v1alpha1
Package v1alpha1 is the v1alpha1 version of the API.
|
Package v1alpha1 is the v1alpha1 version of the API. |
|
kubenetes/generated/clientset/versioned/fake
This package has the automatically generated fake clientset.
|
This package has the automatically generated fake clientset. |
|
kubenetes/generated/clientset/versioned/scheme
This package contains the scheme of the automatically generated clientset.
|
This package contains the scheme of the automatically generated clientset. |
|
kubenetes/generated/clientset/versioned/typed/agentcontroller/v1alpha1
This package has the automatically generated typed clients.
|
This package has the automatically generated typed clients. |
|
kubenetes/generated/clientset/versioned/typed/agentcontroller/v1alpha1/fake
Package fake has the automatically generated clients.
|
Package fake has the automatically generated clients. |