package eino import ( "context" "io" "strings" "github.com/cloudwego/eino/compose" "github.com/cloudwego/eino/schema" ) // sentenceDelimiters 句子分隔符集合。 var sentenceDelimiters = map[rune]bool{ '。': true, '!': true, '?': true, '\n': true, '.': true, '!': true, '?': true, } // NewMessageToStringLambda 创建 Message → String 转换 Lambda 节点。 // 输入: *schema.Message → 输出: string // // 提取 Message.Content 文本,供 Splitter 节点消费。 func NewMessageToStringLambda() *compose.Lambda { return compose.TransformableLambda(func(ctx context.Context, input *schema.StreamReader[*schema.Message]) (*schema.StreamReader[string], error) { sr, sw := schema.Pipe[string](8) go func() { defer sw.Close() defer input.Close() for { msg, err := input.Recv() if err != nil { if err == io.EOF { return } sw.Send("", err) return } if msg != nil && msg.Content != "" { sw.Send(msg.Content, nil) } } }() return sr, nil }) } // NewSplitterLambda 创建句子分割 Transform Lambda 节点。 // 输入: StreamReader[string](LLM token 流)→ 输出: StreamReader[string](完整句子流) // // 逐字符累积,按句子分隔符切分。每切出一个完整句子就输出一次, // 供下游 TTS 节点实时合成。 func NewSplitterLambda() *compose.Lambda { return compose.TransformableLambda(func(ctx context.Context, input *schema.StreamReader[string]) (*schema.StreamReader[string], error) { sr, sw := schema.Pipe[string](8) go func() { defer sw.Close() defer input.Close() var buffer strings.Builder for { chunk, err := input.Recv() if err != nil { if err == io.EOF { // 流结束,flush 剩余缓冲 if buffer.Len() > 0 { text := strings.TrimSpace(buffer.String()) if text != "" { sw.Send(text, nil) } } return } sw.Send("", err) return } // 逐字符累积,按句子分隔符切分 for _, r := range chunk { buffer.WriteRune(r) if sentenceDelimiters[r] { text := strings.TrimSpace(buffer.String()) if text != "" { sw.Send(text, nil) } buffer.Reset() } } } }() return sr, nil }) }