Files
sundynix-agentix/sundynix-dispatcher/internal/nats/subscriber.go
T
Blizzard 5cfed55e91 fix(dispatcher,mcp-go): 配置拉取改为后台重试,根治启动竞态
此前 dispatcher(chat)/mcp-go(embedding) 启动时一次性请求控制面配置,3s 扑空即
降级,且只能干等热更新广播——若消费方早于 gateway 启动,会全程降级(LLM 跑桩、
RAG 无向量),必须手动重启才恢复。

改为:先订阅热更新,再后台 RequestConfigWithRetry(重试至拿到配置,容忍 gateway
晚启)。新增 shared/bus.RequestConfigWithRetry + dispatcher Subscriber 包装。

验收:故意先起 dispatcher/mcp-go、后起 gateway,二者自动重试拿到 chat/embedding
配置,无需手动重启;make test-go 全绿。

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

100 lines
3.7 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 nats 是调度器对共享 bus 的薄封装(消费任务 / 回写 Token)。
package nats
import (
"context"
"log"
"time"
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.<taskID>。
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)
}
// PublishExec / CompleteExec 让 Subscriber 满足 eino.ExecSink
// 把执行轨迹事件回流到 sundynix.exec.<taskID>(与 Token 流分开)。
func (s *Subscriber) PublishExec(taskID string, data []byte) error {
return s.inner.PublishExec(taskID, data)
}
func (s *Subscriber) CompleteExec(taskID string) error {
return s.inner.CompleteExec(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)
}
// ServeHealth 在 dispatcher 心跳主题上应答探活,让管理端「服务状态」判定其在线。
func (s *Subscriber) ServeHealth(provide func() []byte) (func() error, error) {
return s.inner.ServeHealth(contract.SubjectHealthDispatcher, provide)
}
// PublishTaskStatus 让 Subscriber 满足 eino.StatusSink,把任务状态流转回写给网关。
func (s *Subscriber) PublishTaskStatus(taskID, status, detail string) error {
return s.inner.PublishTaskStatus(&contract.TaskStatusEvent{
TaskID: taskID, Status: status, Detail: detail, TS: time.Now().UnixMilli(),
})
}
// RequestModelConfig 向控制面(Gateway)取当前激活的对话模型配置。
func (s *Subscriber) RequestModelConfig(ctx context.Context) (*contract.ModelConfig, error) {
return s.inner.RequestConfig(ctx, contract.ConfigKindChat)
}
// SubscribeModelConfigUpdated 订阅对话模型配置热更新。
func (s *Subscriber) SubscribeModelConfigUpdated(onUpdate func(*contract.ModelConfig)) (func() error, error) {
return s.inner.SubscribeConfigUpdated(contract.ConfigKindChat, onUpdate)
}
// FetchModelConfigWithRetry 后台重试拉取初始对话模型配置(容忍 gateway 晚于 dispatcher 启动)。
func (s *Subscriber) FetchModelConfigWithRetry(ctx context.Context, apply func(*contract.ModelConfig)) {
s.inner.RequestConfigWithRetry(ctx, contract.ConfigKindChat, apply)
}
func (s *Subscriber) Close() { s.inner.Close() }