Files
sundynix-agentix/sundynix-dispatcher/internal/llm/pool.go
T
Blizzard ecf4a80466 feat(llm): T2.1 模型路由 + Fallback —— 单 provider 抖动不再整体宕
现状:单 provider,一家 API 抖动/挂掉全平台不可用。加主备 failover:主模型调用失败
自动按序切备用,compose/ReAct/Chat 全路径透明白嫖。

- llm/failover.go: failoverModel 把多个 ToolCallingChatModel 串成主备链,按序调用、
  遇错切下一个;它本身是 model.ToolCallingChatModel 故全路径透明。调用方主动取消
  (ctx.Err()!=nil) 不切;模型自身请求超时走内部 ctx、不污染父 ctx 故仍 failover。
  局限(v1):Stream 仅建流同步报错时切(已回流 token 的中途失败不切)。
- llm/pool.go: SetConfig 用激活配置(含 Fallbacks)重建——主+可用备用串成 failover 链,
  无备用则直接用主;备用单个构建失败跳过不影响主链。
- contract.ModelConfig: 加 Fallbacks 字段(骑在主配置里下发,不改任何 bus/ServeConfig 签名)。
- gateway store.ActiveConfig: chat 把"其它已登记 chat 模型"按序填进 Fallbacks;
  provide(main) + broadcast(admin) 共用 → 注册多个 chat 模型即自动成主备。
- bus.decryptConfig: 备用模型的 api_key(密文)一并解密。

测试:failover 单测(主可用不调备/主挂切备/全挂报错/取消不切/Stream 切备/WithTools 链)。
live 验证:active=死 ollama 主 + deepseek 备 → 任务连主拒连→自动切 deepseek→4s 完成。
DEPTH_ROADMAP T2.1(admin 注册多模型即主备,无需新 UI)。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-30 09:26:18 +08:00

309 lines
9.9 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 llm 抽象 LLM Pool(第三方在线 API / 自部署)的配置与流式推理。
// 底层用 Eino 的 ChatModel 组件(eino-ext/openaiOpenAI 兼容);本层只负责
// 配置热更新、降级桩与对外的稳定签名(Chat / ChatStream / StreamText)。
package llm
import (
"context"
"fmt"
"io"
"os"
"strconv"
"strings"
"sync"
"time"
"github.com/cloudwego/eino-ext/components/model/openai"
"github.com/cloudwego/eino/components/model"
"github.com/cloudwego/eino/schema"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
"github.com/sundynix/sundynix-shared/contract"
"github.com/sundynix/sundynix-shared/otelx"
)
// requestTimeout 是单次推理请求的上限。
const requestTimeout = 120 * time.Second
// ChatMessage 是一条对话消息(role: system/user/assistant)。
type ChatMessage struct {
Role string `json:"role"`
Content string `json:"content"`
}
// Pool 维护当前激活的后端配置 + 据此构建的 Eino ChatModel(控制面经 NATS 下发,可热更新)。
type Pool struct {
mu sync.RWMutex
cfg *contract.ModelConfig
cm model.BaseChatModel // 由 SetConfig 用激活配置构建;未配置时为 nil
}
func NewPool() *Pool { return &Pool{} }
// forceStub 报告是否强制走降级桩(LLM_FORCE_STUB=1)——压测平台自身吞吐时用,
// 绕开真实 LLM 推理与计费,只压全链路 plumbing(网关→NATS→调度→工具RTT→回流)。
func forceStub() bool { return os.Getenv("LLM_FORCE_STUB") == "1" }
// SetConfig 热更新后端配置:用激活配置(含备用模型)重建 ChatModel(控制面变更时调用)。
func (p *Pool) SetConfig(cfg *contract.ModelConfig) {
var cm model.BaseChatModel
if cfg != nil && cfg.Ready() && !forceStub() {
cm = buildWithFallbacks(cfg)
}
p.mu.Lock()
p.cfg = cfg
p.cm = cm
p.mu.Unlock()
if cfg != nil {
// 不打印 api_key。
fmt.Printf("[llm] model config set: provider=%s base=%s model=%s fallbacks=%d\n",
cfg.Provider, normalizeBaseURL(cfg), cfg.Model, len(cfg.Fallbacks))
}
}
// buildWithFallbacks 构建主模型,并把可用的备用模型串成 failover 链(无备用则直接返回主模型)。
// 主模型构建失败 → 返回 nil(降级桩);备用单个失败 → 跳过该备用,不影响主链。
func buildWithFallbacks(cfg *contract.ModelConfig) model.BaseChatModel {
primary, err := buildChatModel(cfg)
if err != nil {
fmt.Printf("[llm] 构建主 ChatModel 失败(降级桩运行): %v\n", err)
return nil
}
ptcm, ok := primary.(model.ToolCallingChatModel)
if !ok || len(cfg.Fallbacks) == 0 {
return primary // 不支持 WithTools 或无备用 → 直接用主模型
}
models := []model.ToolCallingChatModel{ptcm}
for i := range cfg.Fallbacks {
fb := cfg.Fallbacks[i]
if !fb.Ready() {
continue
}
fbm, ferr := buildChatModel(&fb)
if ferr != nil {
fmt.Printf("[llm] 构建备用模型 %s 失败,跳过: %v\n", fb.Model, ferr)
continue
}
if t, ok := fbm.(model.ToolCallingChatModel); ok {
models = append(models, t)
}
}
if len(models) == 1 {
return primary // 无有效备用
}
fmt.Printf("[llm] 启用模型 failover:主 %s + %d 个备用\n", cfg.Model, len(models)-1)
return newFailoverModel(models, func(idx int, ferr error) {
fmt.Printf("[llm] 模型 failover:第 %d 个模型失败(%v),切下一个\n", idx, ferr)
})
}
// buildChatModel 据 provider 归一化连接参数后构建 OpenAI 兼容 ChatModel。
// vLLM 与 Ollama 都暴露 OpenAI 兼容 API(底层 go-openai 请求 {base}/chat/completions),
// 故统一走 openai 客户端,仅差在 BaseURL 是否带 /v1 与是否需要占位 key:
// - Ollama:端点固定为 {host}/v1,且不校验 api_key → 缺省补 /v1 + 占位 key "ollama"。
// - vLLMopenai server 在 /v1,默认不校验 key(除非 --api-key 启动)→ 补 /v1 + 占位 "EMPTY"。
// - openai-compatibleDeepSeek/OpenAI 等在线):BaseURL 原样(DeepSeek 两种都收,OpenAI 默认带 /v1)。
func buildChatModel(cfg *contract.ModelConfig) (model.BaseChatModel, error) {
return openai.NewChatModel(context.Background(), &openai.ChatModelConfig{
APIKey: apiKeyOrPlaceholder(cfg),
BaseURL: normalizeBaseURL(cfg),
Model: cfg.Model,
Timeout: requestTimeout,
})
}
// 本地后端 provider 标识(与控制台下拉、contract.ModelConfig.Provider 对齐)。
const (
providerOllama = "ollama"
providerVLLM = "vllm"
)
// normalizeBaseURL 为本地后端补全 /v1 端点(Ollama/vLLM 的 OpenAI 兼容 API 均在 /v1
// 漏写会打到 /chat/completions 而 404);在线 provider 原样返回。
func normalizeBaseURL(cfg *contract.ModelConfig) string {
base := strings.TrimRight(cfg.BaseURL, "/")
switch strings.ToLower(cfg.Provider) {
case providerOllama, providerVLLM:
if base != "" && !strings.HasSuffix(base, "/v1") {
base += "/v1"
}
}
return base
}
// apiKeyOrPlaceholder 为不校验 key 的本地后端补占位符(openai 客户端要求 api_key 非空)。
func apiKeyOrPlaceholder(cfg *contract.ModelConfig) string {
if cfg.APIKey != "" {
return cfg.APIKey
}
switch strings.ToLower(cfg.Provider) {
case providerOllama:
return "ollama"
case providerVLLM:
return "EMPTY"
}
return cfg.APIKey
}
func (p *Pool) config() *contract.ModelConfig {
p.mu.RLock()
defer p.mu.RUnlock()
return p.cfg
}
func (p *Pool) model() model.BaseChatModel {
p.mu.RLock()
defer p.mu.RUnlock()
return p.cm
}
// Ready 报告是否已配置可用后端(且 ChatModel 构建成功)。
func (p *Pool) Ready() bool { return p.model() != nil }
// ChatModel 返回当前 Eino ChatModel 组件(用于 compose.Graph 编排);未就绪则 nil。
func (p *Pool) ChatModel() model.BaseChatModel { return p.model() }
// ToolCallingModel 返回支持函数调用的模型(用于 ReAct agent);未就绪 / 不支持则 nil。
func (p *Pool) ToolCallingModel() model.ToolCallingChatModel {
if tcm, ok := p.model().(model.ToolCallingChatModel); ok {
return tcm
}
return nil
}
// ModelName 返回当前激活的对话模型名(未配置则空)—— 供服务状态面板展示。
func (p *Pool) ModelName() string {
if cfg := p.config(); cfg != nil {
return cfg.Model
}
return ""
}
// ChatStream 流式推理,逐 token 回调 onToken(经 Eino ChatModel.Stream)。
// onReasoning(可空)接收 reasoning 模型(DeepSeek-R1 / Qwen3 思考 / QwQ 等)的「思考过程」分片:
// 思考阶段的分片只有 ReasoningContent、Content 为空,本就不会污染答案;onReasoning 让调用方
// 可把思考流单独 surface 到观测轨迹。仅在 Ready() 时可用(调用方据此决定真实推理或降级桩)。
func (p *Pool) ChatStream(ctx context.Context, msgs []ChatMessage, onToken func(string), onReasoning func(string)) error {
cm := p.model()
if cm == nil {
return fmt.Errorf("no model configured")
}
ctx, span := otelx.Tracer().Start(ctx, "llm.stream",
trace.WithAttributes(
attribute.String("llm.model", p.ModelName()),
attribute.Int("llm.messages", len(msgs)),
))
defer span.End()
sr, err := cm.Stream(ctx, toSchema(msgs))
if err != nil {
span.RecordError(err)
return fmt.Errorf("llm stream: %w", err)
}
defer sr.Close()
for {
chunk, rerr := sr.Recv()
if rerr == io.EOF {
return nil
}
if rerr != nil {
return fmt.Errorf("llm stream recv: %w", rerr)
}
if chunk.ReasoningContent != "" && onReasoning != nil {
onReasoning(chunk.ReasoningContent)
}
if chunk.Content != "" {
onToken(chunk.Content)
}
}
}
// Chat 非流式:经 Eino ChatModel.Generate 拿到整段文本。
// 报告生成的「规划大纲 / 撰写章节」等需要拿到完整结果再继续,用它而非流式。
func (p *Pool) Chat(ctx context.Context, msgs []ChatMessage) (string, error) {
cm := p.model()
if cm == nil {
return "", fmt.Errorf("no model configured")
}
ctx, span := otelx.Tracer().Start(ctx, "llm.generate",
trace.WithAttributes(
attribute.String("llm.model", p.ModelName()),
attribute.Int("llm.messages", len(msgs)),
))
defer span.End()
out, err := cm.Generate(ctx, toSchema(msgs))
if err != nil {
span.RecordError(err)
return "", fmt.Errorf("llm generate: %w", err)
}
return out.Content, nil
}
// toSchema 把内部 ChatMessage 转为 Eino schema.Message。
func toSchema(msgs []ChatMessage) []*schema.Message {
out := make([]*schema.Message, 0, len(msgs))
for _, m := range msgs {
var role schema.RoleType
switch m.Role {
case "system":
role = schema.System
case "assistant":
role = schema.Assistant
default:
role = schema.User
}
out = append(out, &schema.Message{Role: role, Content: m.Content})
}
return out
}
// ---- 占位降级(未配置后端时)----
// 占位参数:模拟真实后端的 TTFT(首 token 延迟) 与逐 token 间隔。
// 可经 env 调整(压测平台吞吐时设 0 → 近瞬时桩,测出 plumbing 天花板而非被桩节奏掩盖)。
var (
timeToFirstToken = envDuration("LLM_STUB_TTFT_MS", 700*time.Millisecond)
interTokenDelay = envDuration("LLM_STUB_INTERTOKEN_MS", 60*time.Millisecond)
)
// envDuration 读毫秒 env(允许 0),缺省回退 def。
func envDuration(key string, def time.Duration) time.Duration {
if v := os.Getenv(key); v != "" {
if n, err := strconv.Atoi(v); err == nil && n >= 0 {
return time.Duration(n) * time.Millisecond
}
}
return def
}
// StreamText 按节奏把给定文本流式回调(未配置真实后端时的降级桩)。
func (p *Pool) StreamText(ctx context.Context, text string, onToken func([]byte)) error {
if timeToFirstToken > 0 {
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(timeToFirstToken):
}
}
for _, tok := range tokenize(text) {
select {
case <-ctx.Done():
return ctx.Err()
default:
}
onToken([]byte(tok))
if interTokenDelay > 0 {
time.Sleep(interTokenDelay)
}
}
return nil
}
func tokenize(s string) []string {
out := make([]string, 0, len(s))
for _, r := range s {
out = append(out, string(r))
}
return out
}