feat(dispatcher,gateway): Eino 采纳 Phase D —— 任务生命周期状态机
submitted → running → done / failed / timeout 显式状态机,根治"卡运行中 看不出来"(此前 Task.Status 只写死 submitted、从不流转)。 - contract:新增 SubjectTaskStatus 主题 + TaskStatusEvent + 状态常量 - shared/bus:PublishTaskStatus / SubscribeTaskStatus(core NATS pub-sub) - dispatcher:StatusSink 接口 + Orchestrator.Handle 状态钩子——进入执行 →running、收尾 finishStatus→done/failed、整体超时上限 taskExecTimeout =3min→timeout;经 sub 回写 - gateway:SubscribeTaskStatus 落 PG(Task 增 Detail 字段,AutoMigrate); 新增 GET /api/v1/tasks/:id 供 UI 轮询状态 验收:实测 submitted→running→done 流转 + PG 持久化 + 端点查询闭环; make test-go 全绿。HITL 中断恢复 / 多智能体仍按场景待做。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
+6
-6
@@ -10,7 +10,7 @@
|
||||
- [x] **Phase A · 地基**:`llm.Pool` 换 Eino ChatModel 组件(commit d84b1ec,验收通过)
|
||||
- [x] **Phase B · 质变**:MCP 工具→`InvokableTool` + ReAct agent(模型自主调工具,验收 7/7 命中)
|
||||
- [ ] **Phase C · 编排归一**:`Flow→compose.Graph` + callbacks 桥接
|
||||
- [ ] **Phase D · 状态化执行**:任务生命周期 FSM / HITL 中断恢复 / 多智能体
|
||||
- [~] **Phase D · 状态化执行**:✅ 任务生命周期 FSM(已完成)/ ⬜ HITL 中断恢复 / ⬜ 多智能体(按场景)
|
||||
|
||||
---
|
||||
|
||||
@@ -139,11 +139,11 @@ github.com/cloudwego/eino-ext/... # ⚠️ 官方组件实现(open
|
||||
|
||||
主题:把"执行"从一次性 DAG 升级为**可持久化、可恢复的状态机**。三件事同一条线,一起做。
|
||||
|
||||
- **任务生命周期 FSM** 🆕:现在 `Task.Status` 只写死 `submitted`、全仓从不流转(`store/models.go` + `pgsql.go:96`),是个摆设——这正是"卡运行中看不出来"的根因。
|
||||
- 设计:`submitted → running → done / failed / timeout` 显式状态机。
|
||||
- dispatcher 开跑/跑完/出错 经 NATS 回写状态(新增 `sundynix.tasks.status` 或复用 exec 流),网关落 PG 并推给 UI。
|
||||
- 收益:管理端「服务状态」/ 桌面端能看到任务真实进度,超时自动翻红,无需人工猜。
|
||||
- 与下面同源:Eino compose 的 graph state + 节点级状态正好承载它。
|
||||
- **任务生命周期 FSM** ✅ 已完成:`submitted → running → done / failed / timeout` 显式状态机。
|
||||
- dispatcher 经新主题 `sundynix.tasks.status` 回写(`TaskStatusEvent`):进入执行→running、收尾→done/failed、整体超时上限 `taskExecTimeout=3min`→timeout;网关 `SubscribeTaskStatus` 落 PG(`Task.Status/Detail`)。
|
||||
- UI 轮询:`GET /api/v1/tasks/:id` 返回 `{status, detail}`。
|
||||
- 验收:submitted→running→done 实测流转、PG 持久化、3min 超时兜底——根治"卡运行中看不出来"。
|
||||
- 桌面端轮询接线仍可补(当前后端 + 端点已就绪)。
|
||||
- **中断/恢复(HITL)**:审批型工业流程(生成中途人工确认)。需 checkpoint 持久化(PG/Redis)。等有具体审批用例再做。
|
||||
- **多智能体协同**(`flow/agent/multiagent`):出现真实多角色编排需求时再上,现在无用例。
|
||||
|
||||
|
||||
@@ -41,8 +41,8 @@ func main() {
|
||||
log.Printf("[dispatcher] subscribe model config: %v", err)
|
||||
}
|
||||
|
||||
// sub 同时作为 Token 回流出口(TokenSink)、MCP 工具调用出口(ToolCaller)与执行事件出口(ExecSink)。
|
||||
orch, err := eino.NewOrchestrator(pool, breaker, eval, sub, sub, sub)
|
||||
// sub 同时作为 Token 回流(TokenSink)、MCP 工具调用(ToolCaller)、执行事件(ExecSink)与任务状态回写(StatusSink)出口。
|
||||
orch, err := eino.NewOrchestrator(pool, breaker, eval, sub, sub, sub, sub)
|
||||
if err != nil {
|
||||
log.Fatalf("[dispatcher] build eino graph: %v", err)
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ package eino
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"sync"
|
||||
@@ -29,6 +30,11 @@ type ToolCaller interface {
|
||||
CallTool(ctx context.Context, subject string, call *contract.ToolCall) (*contract.ToolResult, error)
|
||||
}
|
||||
|
||||
// StatusSink 回写任务生命周期状态(由 NATS bus 实现;可为 nil → 不回写)。
|
||||
type StatusSink interface {
|
||||
PublishTaskStatus(taskID, status, detail string) error
|
||||
}
|
||||
|
||||
// LLM 是编排所需的语言模型能力(生产由 *llm.Pool 实现)。抽成接口便于测试注入假模型。
|
||||
type LLM interface {
|
||||
Ready() bool
|
||||
@@ -42,6 +48,9 @@ type LLM interface {
|
||||
// 工具调用超时;超时即降级(不带工具上下文继续推理)。
|
||||
const toolCallTimeout = 3 * time.Second
|
||||
|
||||
// taskExecTimeout 是单个任务整体执行上限;超时即判 timeout(状态机),避免无限期"运行中"。
|
||||
const taskExecTimeout = 3 * time.Minute
|
||||
|
||||
// Orchestrator 把每个 DSL 任务动态编译为 Eino 图并执行(记忆召回 → 工具节点 → 注入 → 流式)。
|
||||
type Orchestrator struct {
|
||||
pool LLM
|
||||
@@ -50,15 +59,39 @@ type Orchestrator struct {
|
||||
sink TokenSink
|
||||
tools ToolCaller
|
||||
exec ExecSink
|
||||
status StatusSink // 任务生命周期状态回写(可为 nil)
|
||||
|
||||
turnMu sync.Mutex // 保护 turns(攒批计数,多任务 goroutine 共享)
|
||||
turns map[string]int // sessionID → 累计轮次,用于每 N 轮触发 consolidate
|
||||
}
|
||||
|
||||
// NewOrchestrator 持有依赖;图按任务的 DSL 在 Handle 内动态编译。
|
||||
// exec 为执行可视化事件出口(可为 nil,则不发轨迹事件);eval 为自动化评测(可为 nil)。
|
||||
func NewOrchestrator(pool LLM, breaker *harness.CircuitBreaker, eval *harness.Evaluator, sink TokenSink, tools ToolCaller, exec ExecSink) (*Orchestrator, error) {
|
||||
return &Orchestrator{pool: pool, breaker: breaker, eval: eval, sink: sink, tools: tools, exec: exec}, nil
|
||||
// exec 为执行可视化事件出口(可为 nil,则不发轨迹事件);eval 为自动化评测(可为 nil);
|
||||
// status 为任务生命周期状态回写出口(可为 nil)。
|
||||
func NewOrchestrator(pool LLM, breaker *harness.CircuitBreaker, eval *harness.Evaluator, sink TokenSink, tools ToolCaller, exec ExecSink, status StatusSink) (*Orchestrator, error) {
|
||||
return &Orchestrator{pool: pool, breaker: breaker, eval: eval, sink: sink, tools: tools, exec: exec, status: status}, nil
|
||||
}
|
||||
|
||||
// setStatus 回写一次任务状态流转(status 为 nil 时静默跳过)。
|
||||
func (o *Orchestrator) setStatus(taskID, status, detail string) {
|
||||
if o.status == nil {
|
||||
return
|
||||
}
|
||||
if err := o.status.PublishTaskStatus(taskID, status, detail); err != nil {
|
||||
log.Printf("[eino] 回写任务状态 %s=%s 失败: %v", taskID, status, err)
|
||||
}
|
||||
}
|
||||
|
||||
// finishStatus 据收尾错误把任务置为 done / timeout / failed。
|
||||
func (o *Orchestrator) finishStatus(taskID string, err error) {
|
||||
switch {
|
||||
case err == nil:
|
||||
o.setStatus(taskID, contract.TaskDone, "")
|
||||
case errors.Is(err, context.DeadlineExceeded):
|
||||
o.setStatus(taskID, contract.TaskTimeout, "执行超时")
|
||||
default:
|
||||
o.setStatus(taskID, contract.TaskFailed, truncate(err.Error(), 200))
|
||||
}
|
||||
}
|
||||
|
||||
// Handle 消费一个任务:按 DSL 编译 Eino 图并执行,把 Token 流回流到 sundynix.streams.<id>。
|
||||
@@ -72,22 +105,31 @@ func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error {
|
||||
tr.info("task", "system", "服务熔断", "后端连续失败,暂时拒绝新任务,请稍后重试")
|
||||
_ = o.sink.PublishToken(t.ID, []byte("⚠️ 服务繁忙(已触发熔断保护),请稍后重试。"))
|
||||
_ = o.sink.CompleteStream(t.ID)
|
||||
o.setStatus(t.ID, contract.TaskFailed, "服务熔断")
|
||||
return nil
|
||||
}
|
||||
|
||||
// 任务状态机:进入执行 → running;整体加超时上限,超时判 timeout(杜绝无限期"运行中")。
|
||||
o.setStatus(t.ID, contract.TaskRunning, "")
|
||||
tctx, cancel := context.WithTimeout(ctx, taskExecTimeout)
|
||||
defer cancel()
|
||||
|
||||
// 报告生成走专用多步编排(规划→分章并行检索撰写→汇聚→渲染 Word),而非通用对话图。
|
||||
if intent, _ := t.Meta[contract.MetaIntent].(string); intent == contract.IntentReport {
|
||||
return o.handleReport(ctx, t, tr)
|
||||
err := o.handleReport(tctx, t, tr)
|
||||
o.finishStatus(t.ID, err)
|
||||
return err
|
||||
}
|
||||
log.Printf("[eino] task %s received (graph=%d bytes), 按图执行(拓扑+连线+分支)...", t.ID, len(t.Graph))
|
||||
tr.info("task", "system", "任务受理", fmt.Sprintf("DSL %d 字节,按图执行", len(t.Graph)))
|
||||
|
||||
// 按 DSL 图的真实拓扑/连线/分支执行(graph.go 解释器),agent 节点流式回流 token。
|
||||
answer, err := o.runGraph(ctx, t, tr)
|
||||
answer, err := o.runGraph(tctx, t, tr)
|
||||
if err != nil {
|
||||
log.Printf("[eino] task %s graph error: %v", t.ID, err)
|
||||
_ = o.sink.CompleteStream(t.ID)
|
||||
o.breaker.Report(false)
|
||||
o.finishStatus(t.ID, err)
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -96,6 +138,7 @@ func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error {
|
||||
}
|
||||
log.Printf("[eino] task %s done (%d 字答复)", t.ID, len([]rune(answer)))
|
||||
o.breaker.Report(true)
|
||||
o.finishStatus(t.ID, nil)
|
||||
|
||||
// 写回阶段:离开热路径、异步落历史 + (TODO)抽取记忆。
|
||||
go o.memorize(t, answer)
|
||||
|
||||
@@ -4,6 +4,7 @@ package nats
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
sharedbus "github.com/sundynix/sundynix-shared/bus"
|
||||
"github.com/sundynix/sundynix-shared/contract"
|
||||
@@ -73,6 +74,13 @@ func (s *Subscriber) ServeHealth(provide func() []byte) (func() error, error) {
|
||||
return s.inner.ServeHealth(contract.SubjectHealthDispatcher, provide)
|
||||
}
|
||||
|
||||
// PublishTaskStatus 让 Subscriber 满足 eino.StatusSink,把任务状态流转回写给网关。
|
||||
func (s *Subscriber) PublishTaskStatus(taskID, status, detail string) error {
|
||||
return s.inner.PublishTaskStatus(&contract.TaskStatusEvent{
|
||||
TaskID: taskID, Status: status, Detail: detail, TS: time.Now().UnixMilli(),
|
||||
})
|
||||
}
|
||||
|
||||
// RequestModelConfig 向控制面(Gateway)取当前激活的对话模型配置。
|
||||
func (s *Subscriber) RequestModelConfig(ctx context.Context) (*contract.ModelConfig, error) {
|
||||
return s.inner.RequestConfig(ctx, contract.ConfigKindChat)
|
||||
|
||||
@@ -46,6 +46,15 @@ func main() {
|
||||
}
|
||||
}
|
||||
|
||||
// 任务生命周期:订阅 dispatcher 回写的状态流转(running/done/failed/timeout),落 PG 供 UI 查询。
|
||||
if _, err := bus.SubscribeTaskStatus(func(ev *contract.TaskStatusEvent) {
|
||||
if err := db.UpdateTaskStatus(context.Background(), ev.TaskID, ev.Status, ev.Detail); err != nil {
|
||||
log.Printf("[gateway] 更新任务状态 %s=%s 失败: %v", ev.TaskID, ev.Status, err)
|
||||
}
|
||||
}); err != nil {
|
||||
log.Printf("[gateway] subscribe task status: %v", err)
|
||||
}
|
||||
|
||||
r := router.New(db, cache, bus, blobStore)
|
||||
addr := envOr("GATEWAY_ADDR", ":8080")
|
||||
log.Printf("[gateway] listening on %s", addr)
|
||||
|
||||
@@ -56,6 +56,17 @@ func (h *Handler) SubmitTask(c *gin.Context) {
|
||||
c.JSON(http.StatusAccepted, gin.H{"task_id": task.ID})
|
||||
}
|
||||
|
||||
// TaskStatus: GET /api/v1/tasks/:id —— 返回任务生命周期状态(供 UI 轮询,根治"卡运行中看不出来")。
|
||||
func (h *Handler) TaskStatus(c *gin.Context) {
|
||||
id := c.Param("id")
|
||||
status, detail := h.db.GetTaskStatus(c.Request.Context(), id)
|
||||
if status == "" {
|
||||
c.JSON(http.StatusNotFound, gin.H{"error": "task not found"})
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"task_id": id, "status": status, "detail": detail})
|
||||
}
|
||||
|
||||
// StreamTask: 订阅 sundynix.streams.<task_id>,以 SSE 把零拷贝 Token Stream 推给客户端。
|
||||
func (h *Handler) StreamTask(c *gin.Context) {
|
||||
taskID := c.Param("id")
|
||||
|
||||
@@ -59,6 +59,11 @@ func (b *Bus) Ping(ctx context.Context, subject string) ([]byte, error) {
|
||||
return b.inner.Ping(ctx, subject)
|
||||
}
|
||||
|
||||
// SubscribeTaskStatus 订阅 dispatcher 回写的任务生命周期状态(落 PG)。
|
||||
func (b *Bus) SubscribeTaskStatus(onEvent func(*contract.TaskStatusEvent)) (func() error, error) {
|
||||
return b.inner.SubscribeTaskStatus(onEvent)
|
||||
}
|
||||
|
||||
// ServeConfig 让网关作为配置控制面,响应某 kind 的配置请求。
|
||||
func (b *Bus) ServeConfig(kind string, provide func() *contract.ModelConfig) (func() error, error) {
|
||||
return b.inner.ServeConfig(kind, provide)
|
||||
|
||||
@@ -49,6 +49,7 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus, blobStore *blob.
|
||||
p := api.Group("", middleware.RequireAuth())
|
||||
{
|
||||
p.POST("/tasks", h.SubmitTask) // 解析 DSL 并 Publish 到 NATS(带已验证 uid)
|
||||
p.GET("/tasks/:id", h.TaskStatus) // 任务生命周期状态(UI 轮询 submitted/running/done/failed/timeout)
|
||||
p.PUT("/memory", h.SetMemory) // 偏好记忆登记(→ mcp-go memory_upsert)
|
||||
p.GET("/memory", h.ListMemory) // 列出当前用户偏好(记忆面板)
|
||||
p.DELETE("/memory", h.DeleteMemory) // 软删一条偏好(?key=)
|
||||
|
||||
@@ -19,5 +19,6 @@ type Task struct {
|
||||
BaseModel
|
||||
TaskID string `gorm:"uniqueIndex;size:64"` // task_xxx
|
||||
Graph string `gorm:"type:jsonb"` // React Flow 导出的 DSL 原文
|
||||
Status string `gorm:"size:32"` // submitted / done / failed
|
||||
Status string `gorm:"size:32"` // submitted / running / done / failed / timeout
|
||||
Detail string `gorm:"type:text"` // 失败/超时原因等(状态机回写)
|
||||
}
|
||||
|
||||
@@ -9,6 +9,8 @@ import (
|
||||
"gorm.io/driver/postgres"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/schema"
|
||||
|
||||
"github.com/sundynix/sundynix-shared/contract"
|
||||
)
|
||||
|
||||
// errStoreDisabled 表示 Postgres 处于降级(未连接)模式,写操作无法进行。
|
||||
@@ -93,7 +95,29 @@ func (p *Postgres) SaveTask(ctx context.Context, id, graph string) error {
|
||||
if p.db == nil {
|
||||
return nil
|
||||
}
|
||||
return p.db.WithContext(ctx).Create(&Task{TaskID: id, Graph: graph, Status: "submitted"}).Error
|
||||
return p.db.WithContext(ctx).Create(&Task{TaskID: id, Graph: graph, Status: contract.TaskSubmitted}).Error
|
||||
}
|
||||
|
||||
// UpdateTaskStatus 流转任务状态(running/done/failed/timeout),由 dispatcher 经 NATS 回写驱动。
|
||||
func (p *Postgres) UpdateTaskStatus(ctx context.Context, id, status, detail string) error {
|
||||
if p.db == nil {
|
||||
return nil
|
||||
}
|
||||
return p.db.WithContext(ctx).Model(&Task{}).
|
||||
Where("task_id = ?", id).
|
||||
Updates(map[string]any{"status": status, "detail": detail}).Error
|
||||
}
|
||||
|
||||
// GetTaskStatus 取一条任务的当前状态(供 UI 轮询;不存在返回空串)。
|
||||
func (p *Postgres) GetTaskStatus(ctx context.Context, id string) (status, detail string) {
|
||||
if p.db == nil {
|
||||
return "", ""
|
||||
}
|
||||
var t Task
|
||||
if err := p.db.WithContext(ctx).Select("status", "detail").Where("task_id = ?", id).First(&t).Error; err != nil {
|
||||
return "", ""
|
||||
}
|
||||
return t.Status, t.Detail
|
||||
}
|
||||
|
||||
// CountTasks 返回已提交任务数(降级模式返回 0)。
|
||||
|
||||
@@ -237,6 +237,31 @@ func (b *Bus) Ping(ctx context.Context, subject string) ([]byte, error) {
|
||||
return msg.Data, nil
|
||||
}
|
||||
|
||||
// ---- 任务生命周期状态回写(core NATS pub-sub)----
|
||||
|
||||
// PublishTaskStatus 广播一次任务状态流转(dispatcher 调用)。
|
||||
func (b *Bus) PublishTaskStatus(ev *contract.TaskStatusEvent) error {
|
||||
data, err := json.Marshal(ev)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return b.nc.Publish(contract.SubjectTaskStatus, data)
|
||||
}
|
||||
|
||||
// SubscribeTaskStatus 订阅任务状态流转(网关调用,落 PG + 推 UI)。
|
||||
func (b *Bus) SubscribeTaskStatus(onEvent func(*contract.TaskStatusEvent)) (unsub func() error, err error) {
|
||||
sub, err := b.nc.Subscribe(contract.SubjectTaskStatus, func(m *nats.Msg) {
|
||||
var ev contract.TaskStatusEvent
|
||||
if json.Unmarshal(m.Data, &ev) == nil {
|
||||
onEvent(&ev)
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("subscribe task status: %w", err)
|
||||
}
|
||||
return sub.Unsubscribe, nil
|
||||
}
|
||||
|
||||
// ---- 配置控制面(core NATS request-reply + broadcast)----
|
||||
|
||||
// RequestConfig 向控制面(Gateway)请求某 kind 当前激活配置(chat/embedding)。
|
||||
|
||||
@@ -29,6 +29,30 @@ const (
|
||||
// request-reply 心跳主题让控制面(管理端「服务状态」)能判定它在不在线。
|
||||
SubjectHealthDispatcher = "sundynix.health.dispatcher"
|
||||
|
||||
// 任务生命周期状态回写:dispatcher 开跑/跑完/出错经此主题广播,网关订阅落 PG 并推 UI。
|
||||
// core NATS pub-sub(状态是幂等覆盖,丢一条由下一条纠正,无需持久化)。
|
||||
SubjectTaskStatus = "sundynix.tasks.status"
|
||||
)
|
||||
|
||||
// 任务生命周期状态机:submitted(网关建任务)→ running(dispatcher 开跑)
|
||||
// → done / failed / timeout(dispatcher 收尾)。
|
||||
const (
|
||||
TaskSubmitted = "submitted"
|
||||
TaskRunning = "running"
|
||||
TaskDone = "done"
|
||||
TaskFailed = "failed"
|
||||
TaskTimeout = "timeout"
|
||||
)
|
||||
|
||||
// TaskStatusEvent 是一次任务状态流转事件(经 SubjectTaskStatus 回流给网关)。
|
||||
type TaskStatusEvent struct {
|
||||
TaskID string `json:"task_id"`
|
||||
Status string `json:"status"` // running / done / failed / timeout
|
||||
Detail string `json:"detail,omitempty"` // 失败原因等
|
||||
TS int64 `json:"ts"` // unix 毫秒
|
||||
}
|
||||
|
||||
const (
|
||||
// MetaUserID 是 Task.Meta 中承载已登录用户标识的键(用于偏好记忆召回)。
|
||||
MetaUserID = "user_id"
|
||||
// MetaSessionID 是 Task.Meta 中承载会话标识的键(用于短期多轮历史)。
|
||||
|
||||
Reference in New Issue
Block a user