cdc5b3a847
把任务执行做成可观测:Dispatcher 在每个节点/阶段发结构化 ExecEvent, 经独立 NATS 通道回流,前端逐节点点亮(状态/耗时/工具入参产出)。 - shared: contract.ExecEvent + ExecSubject(sundynix.exec.<id>,与 Token 流分流); bus.PublishExec/CompleteExec/SubscribeExec(core NATS,复用结束头) - dispatcher: execTracer(自增 Seq 保序 + span 自动计耗时); Orchestrator 加 ExecSink;通用图(init 召回 / 各 tool 入参→产出 / prompt / model 首token+token数)与报告编排(规划大纲 / 各章并行 start-end / 渲染)全程埋点 - gateway: SubscribeExec + GET /tasks/:id/exec SSE(与 token 流并行) - desktop: streamExec + deriveNodes(按 node 归并 start/end/error/info); 复用组件 ExecTrace(竖向轨道,按 kind 着色,运行中脉冲灯); 新 RunsView(运行·观测:轨迹+输出双栏);BottomDrawer 轨迹/工具调用 tab 接真实数据; ReportView 加执行轨迹栏;左导航「运行」置就绪 实测: - 报告任务 /exec:规划(2680ms,4章) → 4 章并行(seq 交错,各~7-8s 重叠=真并行, 每章带 docs 知识库检索预览+成稿字数) → 渲染(docx 落盘) - 通用图 /exec:tool:kb_search(678ms,入参→Milvus 产出) → prompt(2消息) → model(首token 860ms / 4 tokens) - 浏览器(Preview):报告页执行轨迹逐节点点亮、章节带耗时/字数/检索片段,完成后下载 Word Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
69 lines
1.9 KiB
Go
69 lines
1.9 KiB
Go
package eino
|
|
|
|
import (
|
|
"encoding/json"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/sundynix/sundynix-shared/contract"
|
|
)
|
|
|
|
// ExecSink 是执行可视化事件的回流出口(由 NATS bus 实现):
|
|
// 把节点/阶段生命周期事件发到 sundynix.exec.<id>,供"运行·观测"实时点亮轨迹。
|
|
type ExecSink interface {
|
|
PublishExec(taskID string, data []byte) error
|
|
CompleteExec(taskID string) error
|
|
}
|
|
|
|
// execTracer 为一个任务发结构化执行事件,自增 Seq 保序、span 自动计耗时。
|
|
type execTracer struct {
|
|
sink ExecSink
|
|
task string
|
|
seq int32
|
|
}
|
|
|
|
// tracer 为某任务建一个事件发射器(sink 为空时所有方法变空操作)。
|
|
func (o *Orchestrator) tracer(taskID string) *execTracer {
|
|
return &execTracer{sink: o.exec, task: taskID}
|
|
}
|
|
|
|
func (e *execTracer) emit(node, kind, phase, label, detail string, ms int64) {
|
|
if e == nil || e.sink == nil {
|
|
return
|
|
}
|
|
ev := contract.ExecEvent{
|
|
Seq: int(atomic.AddInt32(&e.seq, 1)), TS: time.Now().UnixMilli(),
|
|
Node: node, Kind: kind, Phase: phase, Label: label, Detail: detail, MS: ms,
|
|
}
|
|
if data, err := json.Marshal(&ev); err == nil {
|
|
_ = e.sink.PublishExec(e.task, data)
|
|
}
|
|
}
|
|
|
|
// info 发一条瞬时事件(无耗时)。
|
|
func (e *execTracer) info(node, kind, label, detail string) {
|
|
e.emit(node, kind, "info", label, detail, 0)
|
|
}
|
|
|
|
// span 发 start,并返回一个结束函数:调用时按 err 发 end / error,附带耗时。
|
|
func (e *execTracer) span(node, kind, label string) func(detail string, err error) {
|
|
e.emit(node, kind, "start", label, "", 0)
|
|
t0 := time.Now()
|
|
return func(detail string, err error) {
|
|
ms := time.Since(t0).Milliseconds()
|
|
if err != nil {
|
|
e.emit(node, kind, "error", label, err.Error(), ms)
|
|
return
|
|
}
|
|
e.emit(node, kind, "end", label, detail, ms)
|
|
}
|
|
}
|
|
|
|
// done 关闭该任务的执行事件流(让 SSE 客户端收尾)。
|
|
func (e *execTracer) done() {
|
|
if e == nil || e.sink == nil {
|
|
return
|
|
}
|
|
_ = e.sink.CompleteExec(e.task)
|
|
}
|