158fe094ab
现象:已完成的审批任务,复盘轨迹停在「已中断等待审批」、审批节点永远转圈,看着 像卡住(实际任务已 done、有输出)。 根因:任务中断时 Handle 的 defer tr.done() 给 exec 流发了 CompleteExec → 网关 exec 录制器关闭。resume 是另一次独立调用,其 exec 事件(审批通过/拒绝、续跑节点)发到 同一主题时录制器已关 → 没录进去。(token 流中断时特意没关,所以输出录到了;exec 流却被关了,不对称。) - exec.go: execTracer 加 suspended 标志,done() 挂起时不关流。 - orchestrator.go: 中断分支置 tr.suspended=true(与 token 流一致)。 - ExecTrace.tsx: 前端兜底——运行已完成(phase done)时把悬挂的 waiting/running 节点 收敛为 done,修旧任务(已 done、轨迹缺 resume 事件)的转圈假象。 live 验证:审批任务批准后复盘轨迹完整呈现 approval await→end「批准:放行」→ agent start→推理→end,审批节点不再转圈。 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
72 lines
2.3 KiB
Go
72 lines
2.3 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
|
|
suspended bool // HITL 中断挂起:done() 不关 exec 流,留给 resume 续录(对齐 token 流的中断处理)
|
|
}
|
|
|
|
// 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 客户端收尾)。
|
|
// 中断挂起时不关:任务停在审批、resume 是另一次调用,过早关流会让录制器错过 resume 的轨迹事件
|
|
// (审批通过/拒绝、续跑的节点),导致复盘里审批节点永远停在「等待中」。
|
|
func (e *execTracer) done() {
|
|
if e == nil || e.sink == nil || e.suspended {
|
|
return
|
|
}
|
|
_ = e.sink.CompleteExec(e.task)
|
|
}
|