Files
sundynix-agentix/sundynix-dispatcher/internal/llm/pool.go
T
Blizzard 2ee16d1f99 feat(dispatcher): Eino 采纳 Phase C —— 对话主流程跑 compose.Graph + callbacks 归一
按"并存 + 等价回归"策略,对话主流程可跑在 Eino compose.Graph 上:
- compose_graph.go:runConversation 按 EINO_COMPOSE 灰度开关分流;
  runComposeConversation 建图 START→ChatModel→END,Compile→Stream 回流 token;
  模型未就绪/编译失败降级回 runAgent。默认关,graph.go 仍是默认且权威。
- compose_callbacks.go:composeTracer 用 utils/callbacks 把 ChatModel/Tool 的
  start/end/error 桥到 ExecEvent(可观测归一,不再各处手写 emit)。
- LLM 接口 + Pool 增 ChatModel();fakeLLM 加 cm 字段 + stub Eino 模型。
- 测试:compose 图编译运行 / compose 对话流式 / 开关关→走 runAgent。

顺带修真 bug:SubjectTaskStatus 原 sundynix.tasks.status 落在任务流通配
sundynix.tasks.> 内 → 状态事件被当成"幽灵任务"自我放大(实测污染 2300+ 条)。
挪到 sundynix.status.task + dispatcher 加空任务护栏。

验收:make test-go 全绿;live compose 路径 54 字答复 eval 1.00、默认路径
eval 1.00、幽灵任务 0 复发。剩余 branch/map/render 等节点逐步迁移后 graph.go 退役。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-23 10:42:50 +08:00

190 lines
5.1 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"
"sync"
"time"
"github.com/cloudwego/eino-ext/components/model/openai"
"github.com/cloudwego/eino/components/model"
"github.com/cloudwego/eino/schema"
"github.com/sundynix/sundynix-shared/contract"
)
// 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{} }
// SetConfig 热更新后端配置:重建 ChatModel 实例(控制面变更时调用)。
func (p *Pool) SetConfig(cfg *contract.ModelConfig) {
var cm model.BaseChatModel
if cfg != nil && cfg.Ready() {
built, err := openai.NewChatModel(context.Background(), &openai.ChatModelConfig{
APIKey: cfg.APIKey,
BaseURL: cfg.BaseURL,
Model: cfg.Model,
Timeout: requestTimeout,
})
if err != nil {
fmt.Printf("[llm] 构建 ChatModel 失败(降级桩运行): %v\n", err)
} else {
cm = built
}
}
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\n", cfg.Provider, cfg.BaseURL, cfg.Model)
}
}
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)。
// 仅在 Ready() 时可用(调用方据此决定真实推理或降级桩)。
func (p *Pool) ChatStream(ctx context.Context, msgs []ChatMessage, onToken func(string)) error {
cm := p.model()
if cm == nil {
return fmt.Errorf("no model configured")
}
sr, err := cm.Stream(ctx, toSchema(msgs))
if err != nil {
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.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")
}
out, err := cm.Generate(ctx, toSchema(msgs))
if err != nil {
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 间隔。
const (
timeToFirstToken = 700 * time.Millisecond
interTokenDelay = 60 * time.Millisecond
)
// StreamText 按节奏把给定文本流式回调(未配置真实后端时的降级桩)。
func (p *Pool) StreamText(ctx context.Context, text string, onToken func([]byte)) error {
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))
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
}