Files
Blizzard aec7ad949c feat(model): 工作模型与 JARVIS 语音模型分开配置 + 语音任务路由到快模型
模型配置加"用途"维度:工作主力(chat,要强)与 JARVIS 语音(voice,要快/低时延)各配各激活,
语音任务走语音模型池,未配置则透明回落工作模型——不影响现有功能。

- contract: ConfigKindVoice="voice" + Meta[model_profile]=voice(与 intent==report 同类路由)
- gateway: ServeConfig/broadcastActive 循环纳入 voice;submitVoiceTask 打 model_profile=voice 标记
- dispatcher: 第二个 llm.Pool(voicePool)吃 voice 配置热更新;board.useVoice 从 Meta 派生(含快照);
  Orchestrator.agentPool(b) 按黑板选池——语音且语音池就绪→语音池,否则工作池;
  agent 生成路径(graph/react/coordinator/compose)全改走 agentPool(b),报告/护栏/记忆固定工作池
- admin: 模型页三 Tab(工作主力/JARVIS语音/向量化),复用 ModelManager;api Kind 加 voice

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-22 11:39:17 +08:00

170 lines
7.3 KiB
Go
Raw Permalink 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
}
// IsConnected 报告 NATS 此刻是否可用(供 readiness 探针)。
func (s *Subscriber) IsConnected() bool { return s.inner.IsConnected() }
// 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)
}
// 状态/用量回写流:dispatcher 是发布方,js.Publish 需流已存在(幂等,网关也会 ensure)。
if err := inner.EnsureStatusStream(context.Background()); err != nil {
log.Fatalf("[dispatcher/nats] ensure status stream: %v", err)
}
if err := inner.EnsureUsageStream(context.Background()); err != nil {
log.Fatalf("[dispatcher/nats] ensure usage stream: %v", err)
}
if err := inner.EnsureEvalStream(context.Background()); err != nil {
log.Fatalf("[dispatcher/nats] ensure eval 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 {
drain, err := s.inner.ConsumeTasks(ctx, func(c context.Context, t *contract.Task) error {
return h(c, t)
})
if err != nil {
return err
}
<-ctx.Done() // 收到停机信号
// 优雅停机:停止消费新任务 + 等在途任务跑完(至多 DrainTimeout),再退出。
to := sharedbus.DrainTimeout()
log.Printf("[dispatcher] 收到停机信号,drain 在途任务(≤%s)…", to)
dctx, cancel := context.WithTimeout(context.Background(), to)
defer cancel()
drain(dctx)
log.Printf("[dispatcher] drain 完成,退出")
return nil
}
// 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(),
})
}
// PublishEval 让 Subscriber 满足 eino.EvalSink,把自动化评测结果回写给网关落库。
func (s *Subscriber) PublishEval(ev *contract.EvalEvent) error {
return s.inner.PublishEval(ev)
}
// PublishUsage 让 Subscriber 满足 eino.UsageSink,把任务 token 用量回写给网关累计/计费。
func (s *Subscriber) PublishUsage(ev *contract.UsageEvent) error {
return s.inner.PublishUsage(ev)
}
// WaitApproval 让 Subscriber 满足 eino.ApprovalWaiter,阻塞等待审批节点的人工决定。
func (s *Subscriber) WaitApproval(ctx context.Context, taskID string, timeout time.Duration) (*contract.ApprovalDecision, error) {
return s.inner.WaitApproval(ctx, taskID, timeout)
}
// Checkpoints 打开 HITL 持久化中断的 KV 桶(compose checkpoint + resume 记录)。
// 返回的 *KVHandle 结构化满足 eino.CheckpointKVGet/Put/Delete),可直接 SetCheckpoints。
func (s *Subscriber) Checkpoints(ctx context.Context, bucket string, ttl time.Duration) (*sharedbus.KVHandle, error) {
return s.inner.Checkpoints(ctx, bucket, ttl)
}
// EnsureApprovalStream 确保审批决定流存在(中断/恢复模型用,决定持久抗 dispatcher 离线)。
func (s *Subscriber) EnsureApprovalStream(ctx context.Context) error {
return s.inner.EnsureApprovalStream(ctx)
}
// ConsumeApprovals 持久消费审批决定,每条交给 h 驱动 resume(队列组多副本安全)。
func (s *Subscriber) ConsumeApprovals(ctx context.Context, h func(context.Context, *contract.ApprovalDecision)) (func(context.Context), error) {
return s.inner.ConsumeApprovals(ctx, h)
}
// 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)
}
// SubscribeVoiceConfigUpdated 订阅 JARVIS 语音模型配置热更新(可空——未配置语音模型时不生效)。
func (s *Subscriber) SubscribeVoiceConfigUpdated(onUpdate func(*contract.ModelConfig)) (func() error, error) {
return s.inner.SubscribeConfigUpdated(contract.ConfigKindVoice, onUpdate)
}
// FetchVoiceConfigWithRetry 后台重试拉取初始语音模型配置(未配置则一直拿不到,语音任务回落工作模型)。
func (s *Subscriber) FetchVoiceConfigWithRetry(ctx context.Context, apply func(*contract.ModelConfig)) {
s.inner.RequestConfigWithRetry(ctx, contract.ConfigKindVoice, apply)
}
// FetchPromptsWithRetry 后台重试拉取初始激活 prompt 集(覆盖内置默认)。
func (s *Subscriber) FetchPromptsWithRetry(ctx context.Context, apply func(map[string]string)) {
s.inner.RequestPromptsWithRetry(ctx, apply)
}
// SubscribePromptsUpdated 订阅 prompt 激活集热更新。
func (s *Subscriber) SubscribePromptsUpdated(onUpdate func(map[string]string)) (func() error, error) {
return s.inner.SubscribePromptsUpdated(onUpdate)
}
func (s *Subscriber) Close() { s.inner.Close() }