From 71102d2424b9480efc2b61a7b0a4694c6126b095 Mon Sep 17 00:00:00 2001 From: Blizzard Date: Tue, 23 Jun 2026 09:48:44 +0800 Subject: [PATCH] =?UTF-8?q?feat(dispatcher,gateway):=20Eino=20=E9=87=87?= =?UTF-8?q?=E7=BA=B3=20Phase=20D=20=E2=80=94=E2=80=94=20=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E7=94=9F=E5=91=BD=E5=91=A8=E6=9C=9F=E7=8A=B6=E6=80=81=E6=9C=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- EINO_ADOPTION.md | 12 ++--- sundynix-dispatcher/cmd/dispatcher/main.go | 4 +- .../internal/eino/orchestrator.go | 53 +++++++++++++++++-- .../internal/nats/subscriber.go | 8 +++ sundynix-gateway/cmd/server/main.go | 9 ++++ .../internal/handler/task_handler.go | 11 ++++ sundynix-gateway/internal/nats/publisher.go | 5 ++ sundynix-gateway/internal/router/router.go | 1 + sundynix-gateway/internal/store/models.go | 3 +- sundynix-gateway/internal/store/pgsql.go | 26 ++++++++- sundynix-shared/bus/bus.go | 25 +++++++++ sundynix-shared/contract/task.go | 24 +++++++++ 12 files changed, 166 insertions(+), 15 deletions(-) diff --git a/EINO_ADOPTION.md b/EINO_ADOPTION.md index 2895ff4..3acb352 100644 --- a/EINO_ADOPTION.md +++ b/EINO_ADOPTION.md @@ -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`):出现真实多角色编排需求时再上,现在无用例。 diff --git a/sundynix-dispatcher/cmd/dispatcher/main.go b/sundynix-dispatcher/cmd/dispatcher/main.go index 65a0620..c0caa6f 100644 --- a/sundynix-dispatcher/cmd/dispatcher/main.go +++ b/sundynix-dispatcher/cmd/dispatcher/main.go @@ -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) } diff --git a/sundynix-dispatcher/internal/eino/orchestrator.go b/sundynix-dispatcher/internal/eino/orchestrator.go index 5e6e197..1a5428d 100644 --- a/sundynix-dispatcher/internal/eino/orchestrator.go +++ b/sundynix-dispatcher/internal/eino/orchestrator.go @@ -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.。 @@ -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) diff --git a/sundynix-dispatcher/internal/nats/subscriber.go b/sundynix-dispatcher/internal/nats/subscriber.go index 99f3db9..9353e1a 100644 --- a/sundynix-dispatcher/internal/nats/subscriber.go +++ b/sundynix-dispatcher/internal/nats/subscriber.go @@ -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) diff --git a/sundynix-gateway/cmd/server/main.go b/sundynix-gateway/cmd/server/main.go index 7652a0a..53c52bd 100644 --- a/sundynix-gateway/cmd/server/main.go +++ b/sundynix-gateway/cmd/server/main.go @@ -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) diff --git a/sundynix-gateway/internal/handler/task_handler.go b/sundynix-gateway/internal/handler/task_handler.go index 28eca24..43c98f3 100644 --- a/sundynix-gateway/internal/handler/task_handler.go +++ b/sundynix-gateway/internal/handler/task_handler.go @@ -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.,以 SSE 把零拷贝 Token Stream 推给客户端。 func (h *Handler) StreamTask(c *gin.Context) { taskID := c.Param("id") diff --git a/sundynix-gateway/internal/nats/publisher.go b/sundynix-gateway/internal/nats/publisher.go index cd328f8..181a8ea 100644 --- a/sundynix-gateway/internal/nats/publisher.go +++ b/sundynix-gateway/internal/nats/publisher.go @@ -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) diff --git a/sundynix-gateway/internal/router/router.go b/sundynix-gateway/internal/router/router.go index 0bb7275..1bff4f2 100644 --- a/sundynix-gateway/internal/router/router.go +++ b/sundynix-gateway/internal/router/router.go @@ -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=) diff --git a/sundynix-gateway/internal/store/models.go b/sundynix-gateway/internal/store/models.go index fbe4ae1..81092f3 100644 --- a/sundynix-gateway/internal/store/models.go +++ b/sundynix-gateway/internal/store/models.go @@ -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"` // 失败/超时原因等(状态机回写) } diff --git a/sundynix-gateway/internal/store/pgsql.go b/sundynix-gateway/internal/store/pgsql.go index 0e4c751..8b011f7 100644 --- a/sundynix-gateway/internal/store/pgsql.go +++ b/sundynix-gateway/internal/store/pgsql.go @@ -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)。 diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index 63631cf..dca7377 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -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)。 diff --git a/sundynix-shared/contract/task.go b/sundynix-shared/contract/task.go index 73c7614..8965a90 100644 --- a/sundynix-shared/contract/task.go +++ b/sundynix-shared/contract/task.go @@ -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 中承载会话标识的键(用于短期多轮历史)。