Files
sundynix-agentix/sundynix-dispatcher/internal/llm/pool.go
T
Blizzard 55d50417a9 feat: 模型健康/熔断态 surface 到管理端(T4.F 可观测)
failover/熔断的运行时态原来只在 dispatcher 日志、admin 看不到 —— 本次接到概览可见:
- harness: CircuitBreaker.Snapshot() 只读观测访问器(state + fails,不动状态机)
- llm: Pool.ModelHealth() 上报主备链每模型 {provider,model,role,state,fails};
  buildWithFallbacks 把模型名↔breaker 配对(同包直接读 failoverModel.breakers);
  newFailoverModel 改返回具体类型以便读 breakers
- dispatcher 心跳 payload 加 models[]
- gateway /admin/overview 独立超时 Ping dispatcher,合并进 models.health
- admin 概览「模型路由」新增「运行时链路态(实时)」:逐模型状态点
  (🟢在线/🔴熔断中+失败数/🟡半开探测/单点)
- 单测:Snapshot、Pool.ModelHealth(名字↔态配对/单点/空)
- live:配坏主→提交任务打熔断→概览显示 broken-demo「熔断中·失败3」,备用在线

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-06 12:00:13 +08:00

356 lines
12 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-dispatcher/internal/harness"
"github.com/sundynix/sundynix-shared/contract"
"github.com/sundynix/sundynix-shared/otelx"
)
// ModelHealth 是主备链上单个模型的实时健康态(供管理端展示 failover/熔断可见)。
type ModelHealth struct {
Provider string `json:"provider"`
Model string `json:"model"`
Role string `json:"role"` // primary / fallback
State string `json:"state"` // closed(在线) / open(熔断中) / half-open(半开探测) / single(无备用链)
Fails int `json:"fails"` // 闭合态连续失败计数
}
// 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
health []ModelHealth // 主备链各模型静态信息(provider/model/role),与 breakers 同序
breakers []*harness.CircuitBreaker // 与 health 一一对应;无 failover 链时为对应 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
var health []ModelHealth
var breakers []*harness.CircuitBreaker
if cfg != nil && cfg.Ready() && !forceStub() {
cm, health, breakers = buildWithFallbacks(cfg)
}
p.mu.Lock()
p.cfg = cfg
p.cm = cm
p.health = health
p.breakers = breakers
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(降级桩);备用单个失败 → 跳过该备用,不影响主链。
// 第二/三返回值为主备链的健康态元信息(provider/model/role)与对应熔断器(无 failover 链时为 nil),
// 供 Pool.ModelHealth 上报管理端。
func buildWithFallbacks(cfg *contract.ModelConfig) (model.BaseChatModel, []ModelHealth, []*harness.CircuitBreaker) {
primary, err := buildChatModel(cfg)
if err != nil {
fmt.Printf("[llm] 构建主 ChatModel 失败(降级桩运行): %v\n", err)
return nil, nil, nil
}
ptcm, ok := primary.(model.ToolCallingChatModel)
if !ok {
// 不支持 WithTools(无法包 failover/cache)→ 单模型,无独立熔断器。
return primary,
[]ModelHealth{{Provider: cfg.Provider, Model: cfg.Model, Role: "primary"}},
[]*harness.CircuitBreaker{nil}
}
// 主链:主模型 +(可用的)备用模型串成 failover。metas 与 models 同序。
models := []model.ToolCallingChatModel{ptcm}
metas := []ModelHealth{{Provider: cfg.Provider, Model: cfg.Model, Role: "primary"}}
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)
metas = append(metas, ModelHealth{Provider: fb.Provider, Model: fb.Model, Role: "fallback"})
}
}
var chain model.ToolCallingChatModel = ptcm
breakers := make([]*harness.CircuitBreaker, len(models)) // 无 failover 时全 nil
if len(models) > 1 {
fmt.Printf("[llm] 启用模型 failover:主 %s + %d 个备用\n", cfg.Model, len(models)-1)
fm := newFailoverModel(models, func(idx int, ferr error) {
fmt.Printf("[llm] 模型 failover:第 %d 个模型失败(%v),切下一个\n", idx, ferr)
})
copy(breakers, fm.breakers) // 每模型熔断器(同序),供上报态
chain = fm
}
// 缓存包在最外层:命中直接跳过整条 failover 链(省成本+提速)。键含模型名 → 换模型自然失效。
return withCache(chain, cfg.Model), metas, breakers
}
// ModelHealth 返回当前主备链各模型的实时健康态(含每模型熔断状态),供管理端展示。
func (p *Pool) ModelHealth() []ModelHealth {
p.mu.RLock()
defer p.mu.RUnlock()
out := make([]ModelHealth, len(p.health))
for i, h := range p.health {
out[i] = h
if i < len(p.breakers) && p.breakers[i] != nil {
snap := p.breakers[i].Snapshot()
out[i].State = snap.State.String()
out[i].Fails = snap.Fails
} else {
out[i].State = "single" // 无备用链 → 该模型无独立熔断器
}
}
return out
}
// 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
}