// Package nats 是调度器对共享 bus 的薄封装(消费任务 / 回写 Token)。 package nats import ( "context" "log" sharedbus "github.com/sundynix/sundynix-shared/bus" "github.com/sundynix/sundynix-shared/contract" ) // TaskHandler 处理单个任务。 type TaskHandler func(ctx context.Context, t *contract.Task) error // Subscriber 包装共享 bus,向调度器暴露消费能力。 type Subscriber struct { inner *sharedbus.Bus } // MustConnect 接入 NATS 并确保任务流存在(消费者声明在 Consume 时完成)。 func MustConnect(url string) *Subscriber { inner, err := sharedbus.Connect(url) if err != nil { log.Fatalf("[dispatcher/nats] connect: %v", err) } if err := inner.EnsureTaskStream(context.Background()); err != nil { log.Fatalf("[dispatcher/nats] ensure stream: %v", err) } log.Printf("[dispatcher/nats] connected %s", url) return &Subscriber{inner: inner} } // ConsumeTasks 从 sundynix.tasks.* 持续消费任务(队列组负载均衡),阻塞至 ctx 取消。 func (s *Subscriber) ConsumeTasks(ctx context.Context, h TaskHandler) error { stop, err := s.inner.ConsumeTasks(ctx, func(c context.Context, t *contract.Task) error { return h(c, t) }) if err != nil { return err } defer stop() <-ctx.Done() return ctx.Err() } // PublishToken / CompleteStream 让 Subscriber 满足 eino.TokenSink, // 把推理 Token 回流到 sundynix.streams.。 func (s *Subscriber) PublishToken(taskID string, token []byte) error { return s.inner.PublishToken(taskID, token) } func (s *Subscriber) CompleteStream(taskID string) error { return s.inner.CompleteStream(taskID) } // CallTool 让 Subscriber 满足 eino.ToolCaller,经 NATS request-reply 调起第 5 层 MCP 工具。 func (s *Subscriber) CallTool(ctx context.Context, subject string, call *contract.ToolCall) (*contract.ToolResult, error) { return s.inner.CallTool(ctx, subject, call) } // RequestModelConfig 向控制面(Gateway)取当前激活的模型配置。 func (s *Subscriber) RequestModelConfig(ctx context.Context) (*contract.ModelConfig, error) { return s.inner.RequestModelConfig(ctx) } // SubscribeModelConfigUpdated 订阅模型配置热更新。 func (s *Subscriber) SubscribeModelConfigUpdated(onUpdate func(*contract.ModelConfig)) (func() error, error) { return s.inner.SubscribeModelConfigUpdated(onUpdate) } func (s *Subscriber) Close() { s.inner.Close() }