a2d184b7ec
新增 sundynix-shared/otelx:otelx.Init(ctx,服务名) 注册 W3C 传播器 + OTLP/HTTP 批量导出(默认 localhost:4318,docker 里的 Jaeger);OTEL_SDK_DISABLED=true 关导出。 Jaeger 不在线 / 导出器构建失败都不阻断启动(可观测性是增益而非依赖)。 - NATS 跨进程传播(无现成中间件):bus/trace.go 的 natsHeaderCarrier + inject/extract, PublishTask/CallTool 注入 traceparent,ConsumeTasks/ServeTool 抽出续上 → 链路跨总线连成一棵树。 - 埋点:gateway 挂 otelgin(HTTP server span,链路根);dispatcher task.execute → node.<kind>(每节点,nctx 下传使工具/LLM 挂到节点下)→ llm.stream/llm.generate; bus 自动出 tool.call(client)↔tool.serve(server) 成对跨服务 span。 - docker-compose 加 jaeger all-in-one(UI :16686,OTLP :4318)。 - 依赖修复:otlptracehttp 触发 genproto 单体(旧)vs 拆分模块 ambiguous import(milvus 拉旧版), pin genproto 至后拆分版(go.work 工作区全局生效)。 - production_readiness.md 1.1 更新为「已实现」。 验证:真实 input→retriever→agent 任务在 Jaeger 出 14 span / 3 服务的完整树, 跨 NATS(publish→consume)、跨服务(tool.call→tool.serve)均连通, 瓶颈 kb_search 692ms、llm 1597ms 一眼可见;四模块 build+vet+test 全绿。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
440 lines
15 KiB
Go
440 lines
15 KiB
Go
// 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
|
||
}
|
||
|
||
// ---- 配置控制面(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,
|
||
})
|
||
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
|
||
}
|