Files
sundynix-agentix/sundynix-shared/bus/bus.go
T
Blizzard 16a6c4b1aa feat(hitl): 人工审批中断(Eino Phase D)—— 审批节点暂停→批准续跑/拒绝中止
专用「审批」节点方案,全栈打通。

后端:
- contract:TaskWaiting/TaskRejected 状态 + ApprovalSubject/ApprovalDecision。
- bus:PublishApproval + WaitApproval(订阅决定主题,带超时);消费者 AckWait
  30s→15min(阻塞等人审期间消息未 ack,否则重投成重复任务)。
- dispatcher:ApprovalWaiter 接口 + approvalNode——执行到审批节点发 await 事件 +
  置 waiting,阻塞等决定。批准→回 running 放行下游;拒绝/超时→errRejected 哨兵→
  剪下游→rejected,优雅收尾不计熔断。超时安全默认拒绝。
  taskExecTimeout 3→10min(审批5 < 执行10 < AckWait15)。
- 网关:POST /tasks/:id/approve(仅 waiting 态受理,幂等)。

桌面端:
- Studio 新增「人工审批」节点(nodeCatalog,可填标题/说明)。
- run.ts pendingApproval() 从 exec 流派生待审中断 + waiting 节点状态。
- BottomDrawer ApprovalBar:琥珀审批条(摘要 + 批准/拒绝 + 备注,调 api.approveTask)。
- ExecTrace waiting 状态灯。

验证:后端 curl 实测 waiting→批准→running→done;waiting→拒绝→rejected(下游未跑)。
前端 tsc + vitest 36 过(pendingApproval 4 例 + waiting 状态)。全模块 build+vet+test 全绿。
EINO_ADOPTION Phase D 标记 HITL 完成。

未覆盖:compose 路径(EINO_COMPOSE,默认关)的 approval 节点;桌面端审批条未在 GUI 实点。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-24 13:18:47 +08:00

485 lines
16 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package bus 封装 NATS JetStream 的连接、流声明、任务发布与消费。
// Gateway 与 Dispatcher 共用这套真实收发逻辑,e2e 测试也直接覆盖它。
package bus
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
"github.com/sundynix/sundynix-shared/contract"
"github.com/sundynix/sundynix-shared/secrets"
)
// decryptConfig 在消费侧把配置里的 api_key 从密文还原为明文(控制面以密文过线缆,见 secrets 包)。
// 失败(密钥不匹配 / 密文损坏)时清空 api_key 并不再降级阻断——调用方据 Ready() 判定。
func decryptConfig(cfg *contract.ModelConfig) {
if cfg == nil || cfg.APIKey == "" {
return
}
if plain, err := secrets.Decrypt(cfg.APIKey); err == nil {
cfg.APIKey = plain
}
}
// Bus 持有 NATS 连接与 JetStream 上下文。
type Bus struct {
nc *nats.Conn
js jetstream.JetStream
}
// Connect 接入 NATS 骨干网并初始化 JetStream,使用默认重试参数。
func Connect(url string) (*Bus, error) {
return ConnectWithRetry(url, 30, time.Second)
}
// ConnectWithRetry 在 NATS 暂不可用时按固定间隔重试,容忍服务先于 NATS 启动。
func ConnectWithRetry(url string, attempts int, interval time.Duration) (*Bus, error) {
var lastErr error
for i := 0; i < attempts; i++ {
nc, err := nats.Connect(url,
nats.Timeout(5*time.Second),
nats.RetryOnFailedConnect(true),
nats.MaxReconnects(-1),
nats.ReconnectWait(interval),
)
if err != nil {
lastErr = err
time.Sleep(interval)
continue
}
// RetryOnFailedConnect 下 Connect 可能立即返回但尚未连上,等待真正建立。
if nc.Status() != nats.CONNECTED {
if !waitConnected(nc, 5*time.Second) {
lastErr = fmt.Errorf("nats not connected within timeout")
nc.Close()
time.Sleep(interval)
continue
}
}
js, err := jetstream.New(nc)
if err != nil {
nc.Close()
return nil, fmt.Errorf("jetstream init: %w", err)
}
return &Bus{nc: nc, js: js}, nil
}
return nil, fmt.Errorf("nats connect after %d attempts: %w", attempts, lastErr)
}
func waitConnected(nc *nats.Conn, d time.Duration) bool {
deadline := time.Now().Add(d)
for time.Now().Before(deadline) {
if nc.Status() == nats.CONNECTED {
return true
}
time.Sleep(50 * time.Millisecond)
}
return nc.Status() == nats.CONNECTED
}
// Close 关闭底层连接。
func (b *Bus) Close() {
if b.nc != nil {
b.nc.Close()
}
}
// EnsureTaskStream 幂等地创建/更新任务流,捕获 sundynix.tasks.>。
func (b *Bus) EnsureTaskStream(ctx context.Context) error {
_, err := b.js.CreateOrUpdateStream(ctx, jetstream.StreamConfig{
Name: contract.StreamTasks,
Subjects: []string{contract.SubjectTasksAll},
Storage: jetstream.FileStorage,
})
return err
}
// PublishTask 把任务发布到 sundynix.tasks.<id>,返回序列号。
// 发布前把链路上下文注入消息头,使 dispatcher 消费时能续上同一条 trace。
func (b *Bus) PublishTask(ctx context.Context, t *contract.Task) (uint64, error) {
data, err := t.Marshal()
if err != nil {
return 0, err
}
ctx, span := tracer().Start(ctx, "nats.publish task",
trace.WithSpanKind(trace.SpanKindProducer),
trace.WithAttributes(
attribute.String("messaging.system", "nats"),
attribute.String("messaging.destination.name", contract.TaskSubject(t.ID)),
attribute.String("sundynix.task_id", t.ID),
))
defer span.End()
msg := nats.NewMsg(contract.TaskSubject(t.ID))
msg.Data = data
injectTrace(ctx, msg.Header) // 把 traceparent 写进消息头 → 跨总线传播
ack, err := b.js.PublishMsg(ctx, msg)
if err != nil {
span.RecordError(err)
return 0, fmt.Errorf("publish task: %w", err)
}
return ack.Sequence, nil
}
// ---- Token 流回流(core NATS 零拷贝字节管道)----
// PublishToken 把一个推理 Token 以 core NATS 写到 sundynix.streams.<taskID>。
func (b *Bus) PublishToken(taskID string, token []byte) error {
return b.nc.Publish(contract.StreamSubject(taskID), token)
}
// CompleteStream 发送 Token 流结束信号(空体 + 结束头)。
func (b *Bus) CompleteStream(taskID string) error {
msg := nats.NewMsg(contract.StreamSubject(taskID))
msg.Header.Set(contract.HeaderStreamEnd, "1")
return b.nc.PublishMsg(msg)
}
// SubscribeTokens 订阅某 task 的 Token 流。每个 Token 触发 onToken
// 收到结束信号后触发 onDone。返回的 unsub 用于退订。
// 注意:core NATS 无持久化,订阅须在 Token 产生前建立(SSE 客户端先连)。
func (b *Bus) SubscribeTokens(taskID string, onToken func([]byte), onDone func()) (unsub func() error, err error) {
sub, err := b.nc.Subscribe(contract.StreamSubject(taskID), func(m *nats.Msg) {
if m.Header.Get(contract.HeaderStreamEnd) == "1" {
onDone()
return
}
// 拷贝,避免 nats 复用底层 buffer。
tok := make([]byte, len(m.Data))
copy(tok, m.Data)
onToken(tok)
})
if err != nil {
return nil, fmt.Errorf("subscribe tokens: %w", err)
}
return sub.Unsubscribe, nil
}
// ---- 执行可视化事件(core NATS,与 Token 流分流)----
// PublishExec 把一条执行事件(JSON)发到 sundynix.exec.<taskID>。
func (b *Bus) PublishExec(taskID string, data []byte) error {
return b.nc.Publish(contract.ExecSubject(taskID), data)
}
// CompleteExec 发送执行事件流结束信号(空体 + 结束头)。
func (b *Bus) CompleteExec(taskID string) error {
msg := nats.NewMsg(contract.ExecSubject(taskID))
msg.Header.Set(contract.HeaderStreamEnd, "1")
return b.nc.PublishMsg(msg)
}
// SubscribeExec 订阅某 task 的执行事件流。每条事件触发 onEvent;结束触发 onDone。
func (b *Bus) SubscribeExec(taskID string, onEvent func([]byte), onDone func()) (unsub func() error, err error) {
sub, err := b.nc.Subscribe(contract.ExecSubject(taskID), func(m *nats.Msg) {
if m.Header.Get(contract.HeaderStreamEnd) == "1" {
onDone()
return
}
data := make([]byte, len(m.Data))
copy(data, m.Data)
onEvent(data)
})
if err != nil {
return nil, fmt.Errorf("subscribe exec: %w", err)
}
return sub.Unsubscribe, nil
}
// ---- MCP 工具调用(core NATS request-reply----
// CallTool 同步调用一个 MCP 工具:发到 subject,阻塞等待应答。
// ctx 超时即视为工具不可用,由调用方决定降级。
func (b *Bus) CallTool(ctx context.Context, subject string, call *contract.ToolCall) (*contract.ToolResult, error) {
data, err := json.Marshal(call)
if err != nil {
return nil, fmt.Errorf("marshal tool call: %w", err)
}
ctx, span := tracer().Start(ctx, "tool.call "+call.Tool,
trace.WithSpanKind(trace.SpanKindClient),
trace.WithAttributes(
attribute.String("sundynix.tool", call.Tool),
attribute.String("messaging.destination.name", subject),
))
defer span.End()
req := nats.NewMsg(subject)
req.Data = data
injectTrace(ctx, req.Header) // 把链路上下文带给第 5 层工具服务
msg, err := b.nc.RequestMsgWithContext(ctx, req)
if err != nil {
span.RecordError(err)
return nil, fmt.Errorf("call tool %s: %w", subject, err)
}
var res contract.ToolResult
if err := json.Unmarshal(msg.Data, &res); err != nil {
span.RecordError(err)
return nil, fmt.Errorf("unmarshal tool result: %w", err)
}
span.SetAttributes(attribute.Bool("sundynix.tool.ok", res.OK))
return &res, nil
}
// ToolHandler 处理一次工具调用并返回结果。
type ToolHandler func(ctx context.Context, call *contract.ToolCall) *contract.ToolResult
// ServeTool 以队列组订阅工具主题(可用通配 sundynix.tools.go.>),
// 对每个请求调用 h 并 Respond,队列组内多副本自动负载均衡。
// 返回的 unsub 用于退订。
func (b *Bus) ServeTool(subject, queue string, h ToolHandler) (unsub func() error, err error) {
sub, err := b.nc.QueueSubscribe(subject, queue, func(m *nats.Msg) {
var call contract.ToolCall
if err := json.Unmarshal(m.Data, &call); err != nil {
respond(m, &contract.ToolResult{OK: false, Error: "bad tool call: " + err.Error()})
return
}
// 还原上游链路并开服务端 span(成为 dispatcher tool.call span 的子节点)。
mctx := extractTrace(context.Background(), m.Header)
mctx, span := tracer().Start(mctx, "tool.serve "+call.Tool,
trace.WithSpanKind(trace.SpanKindServer),
trace.WithAttributes(attribute.String("sundynix.tool", call.Tool)))
res := h(mctx, &call)
if res != nil {
span.SetAttributes(attribute.Bool("sundynix.tool.ok", res.OK))
}
span.End()
respond(m, res)
})
if err != nil {
return nil, fmt.Errorf("serve tool %s: %w", subject, err)
}
return sub.Unsubscribe, nil
}
func respond(m *nats.Msg, res *contract.ToolResult) {
data, err := json.Marshal(res)
if err != nil {
data, _ = json.Marshal(&contract.ToolResult{OK: false, Error: "marshal result: " + err.Error()})
}
_ = m.Respond(data)
}
// ---- 服务探活(core NATS request-reply 心跳)----
// ServeHealth 在 subject 上应答健康探测,provide 返回本节点状态 JSON(可为空)。
// 用于无 HTTP/工具端点的节点(如 dispatcher)向控制面暴露存活。
func (b *Bus) ServeHealth(subject string, provide func() []byte) (unsub func() error, err error) {
sub, err := b.nc.Subscribe(subject, func(m *nats.Msg) {
_ = m.Respond(provide())
})
if err != nil {
return nil, fmt.Errorf("serve health %s: %w", subject, err)
}
return sub.Unsubscribe, nil
}
// Ping 同步探测某节点健康:发到 subject 等应答。无人应答 / 超时即返回错误(视为下线)。
func (b *Bus) Ping(ctx context.Context, subject string) ([]byte, error) {
msg, err := b.nc.RequestWithContext(ctx, subject, nil)
if err != nil {
return nil, err
}
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
}
// ---- 人工审批(HITLcore NATS pub-sub----
// PublishApproval 广播一次人工审批决定(网关在收到 UI 的批准/拒绝后调用)。
func (b *Bus) PublishApproval(dec *contract.ApprovalDecision) error {
data, err := json.Marshal(dec)
if err != nil {
return err
}
return b.nc.Publish(contract.ApprovalSubject(dec.TaskID), data)
}
// WaitApproval 阻塞等待某任务的人工审批决定,直到收到、ctx 取消或超时。
// dispatcher 在审批节点调用:先订阅再等待(订阅早于决定到达,避免错过)。
// 超时返回 (nil, error) —— 调用方据安全默认(拒绝)处理。
func (b *Bus) WaitApproval(ctx context.Context, taskID string, timeout time.Duration) (*contract.ApprovalDecision, error) {
ch := make(chan *contract.ApprovalDecision, 1)
sub, err := b.nc.Subscribe(contract.ApprovalSubject(taskID), func(m *nats.Msg) {
var dec contract.ApprovalDecision
if json.Unmarshal(m.Data, &dec) == nil {
select {
case ch <- &dec:
default: // 已收到一条,丢弃后续重复决定
}
}
})
if err != nil {
return nil, fmt.Errorf("subscribe approval: %w", err)
}
defer func() { _ = sub.Unsubscribe() }()
timer := time.NewTimer(timeout)
defer timer.Stop()
select {
case dec := <-ch:
return dec, nil
case <-timer.C:
return nil, fmt.Errorf("approval timeout after %s", timeout)
case <-ctx.Done():
return nil, ctx.Err()
}
}
// ---- 配置控制面(core NATS request-reply + broadcast----
// RequestConfig 向控制面(Gateway)请求某 kind 当前激活配置(chat/embedding)。
// 无人应答 / 无激活配置时返回 (nil, nil),由调用方降级。
func (b *Bus) RequestConfig(ctx context.Context, kind string) (*contract.ModelConfig, error) {
msg, err := b.nc.RequestWithContext(ctx, contract.ConfigGetSubject(kind), nil)
if err != nil {
return nil, nil // 控制面暂不可用,降级
}
if len(msg.Data) == 0 {
return nil, nil
}
var cfg contract.ModelConfig
if err := json.Unmarshal(msg.Data, &cfg); err != nil {
return nil, fmt.Errorf("unmarshal %s config: %w", kind, err)
}
if !cfg.Ready() {
return nil, nil
}
decryptConfig(&cfg) // 线缆上是密文,消费侧还原
return &cfg, nil
}
// RequestConfigWithRetry 后台重试拉取某 kind 的初始配置,直到成功或重试耗尽。
// 容忍消费方(dispatcher/mcp-go)早于控制面(gateway)启动——一次性请求扑空后不再干等热更新。
// 拿到即调 apply 并返回;ctx 取消或重试上限到则放弃(此后仍可由热更新广播兜底)。
func (b *Bus) RequestConfigWithRetry(ctx context.Context, kind string, apply func(*contract.ModelConfig)) {
for i := 0; i < 60; i++ {
cctx, cancel := context.WithTimeout(ctx, 3*time.Second)
cfg, _ := b.RequestConfig(cctx, kind)
cancel()
if cfg != nil {
apply(cfg)
return
}
select {
case <-ctx.Done():
return
case <-time.After(2 * time.Second):
}
}
}
// ServeConfig 让控制面响应某 kind 的配置请求;provide 返回当前激活配置(可为 nil)。
func (b *Bus) ServeConfig(kind string, provide func() *contract.ModelConfig) (unsub func() error, err error) {
sub, err := b.nc.Subscribe(contract.ConfigGetSubject(kind), func(m *nats.Msg) {
var data []byte
if cfg := provide(); cfg != nil {
data, _ = json.Marshal(cfg)
}
_ = m.Respond(data)
})
if err != nil {
return nil, fmt.Errorf("serve %s config: %w", kind, err)
}
return sub.Unsubscribe, nil
}
// PublishConfigUpdated 广播某 kind 的配置变更(消费方据此热更新)。
func (b *Bus) PublishConfigUpdated(kind string, cfg *contract.ModelConfig) error {
data, err := json.Marshal(cfg)
if err != nil {
return err
}
return b.nc.Publish(contract.ConfigUpdatedSubject(kind), data)
}
// SubscribeConfigUpdated 订阅某 kind 的配置变更。
func (b *Bus) SubscribeConfigUpdated(kind string, onUpdate func(*contract.ModelConfig)) (unsub func() error, err error) {
sub, err := b.nc.Subscribe(contract.ConfigUpdatedSubject(kind), func(m *nats.Msg) {
var cfg contract.ModelConfig
if json.Unmarshal(m.Data, &cfg) == nil {
decryptConfig(&cfg) // 线缆上是密文,消费侧还原
onUpdate(&cfg)
}
})
if err != nil {
return nil, fmt.Errorf("subscribe %s config: %w", kind, err)
}
return sub.Unsubscribe, nil
}
// TaskHandler 处理一个消费到的任务。
type TaskHandler func(ctx context.Context, t *contract.Task) error
// ConsumeTasks 在持久消费者上消费任务,队列组内负载均衡。
// 返回的 stop 函数用于优雅停止消费。
func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err error) {
cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamTasks, jetstream.ConsumerConfig{
Durable: contract.ConsumerDurable,
AckPolicy: jetstream.AckExplicitPolicy,
FilterSubject: contract.SubjectTasksAll,
// HITL:审批节点会让 Handle 阻塞等人工决定(最长约 5 分钟),
// AckWait 必须覆盖「审批等待 + 图执行」总时长,否则消息在途未 ack 会被重投成重复任务。
AckWait: 15 * time.Minute,
})
if err != nil {
return nil, fmt.Errorf("create consumer: %w", err)
}
cc, err := cons.Consume(func(msg jetstream.Msg) {
t, err := contract.Unmarshal(msg.Data())
if err != nil {
_ = msg.Term() // 脏数据,丢弃不重投
return
}
// 从消息头还原上游链路,开消费 span(成为 gateway 发布 span 的子节点)。
mctx := extractTrace(ctx, nats.Header(msg.Headers()))
mctx, span := tracer().Start(mctx, "nats.consume task",
trace.WithSpanKind(trace.SpanKindConsumer),
trace.WithAttributes(attribute.String("sundynix.task_id", t.ID)))
herr := h(mctx, t)
if herr != nil {
span.RecordError(herr)
}
span.End()
if herr != nil {
_ = msg.NakWithDelay(time.Second) // 处理失败,延迟重投
return
}
_ = msg.Ack()
})
if err != nil {
return nil, fmt.Errorf("consume: %w", err)
}
return cc.Stop, nil
}