Stream 实现流式执行接口 返回 stream.Reader,支持: - 流式生成事件 - 客户端控制的取消 - 与 LLM 流式 API 无缝集成 - 多消费者支持(Copy, Merge, Transform)
(ctx context.Context, message string, opts ...Option)
| 53 | // - 与 LLM 流式 API 无缝集成 |
| 54 | // - 多消费者支持(Copy, Merge, Transform) |
| 55 | func (a *Agent) Stream(ctx context.Context, message string, opts ...Option) *stream.Reader[*session.Event] { |
| 56 | reader, writer := stream.Pipe[*session.Event](10) |
| 57 | |
| 58 | go func() { |
| 59 | defer writer.Close() |
| 60 | |
| 61 | // 应用选项 |
| 62 | config := &streamConfig{} |
| 63 | for _, opt := range opts { |
| 64 | opt(config) |
| 65 | } |
| 66 | |
| 67 | streamLog.Info(ctx, "starting stream", map[string]any{"message": truncate(message, 50)}) |
| 68 | |
| 69 | // 1. 前置验证 |
| 70 | if err := a.validateMessage(message); err != nil { |
| 71 | writer.Send(nil, fmt.Errorf("validate message: %w", err)) |
| 72 | return |
| 73 | } |
| 74 | |
| 75 | // 2. 创建用户消息 |
| 76 | userMsg := types.Message{ |
| 77 | Role: types.RoleUser, |
| 78 | Content: message, |
| 79 | ContentBlocks: []types.ContentBlock{ |
| 80 | &types.TextBlock{Text: message}, |
| 81 | }, |
| 82 | } |
| 83 | |
| 84 | // 3. 检查 Slash Commands |
| 85 | if a.commandExecutor != nil { |
| 86 | if handled, err := a.handleSlashCommandForStream(ctx, &userMsg, writer); handled { |
| 87 | if err != nil { |
| 88 | writer.Send(nil, fmt.Errorf("slash command: %w", err)) |
| 89 | } |
| 90 | return |
| 91 | } |
| 92 | } |
| 93 | |
| 94 | // 4. 应用 Skills 增强 |
| 95 | if a.skillInjector != nil { |
| 96 | skillContext := skills.SkillContext{ |
| 97 | UserMessage: userMsg.GetContent(), |
| 98 | } |
| 99 | enhancedPrompt := a.skillInjector.EnhanceSystemPrompt(ctx, a.template.SystemPrompt, skillContext) |
| 100 | a.template.SystemPrompt = enhancedPrompt |
| 101 | } |
| 102 | |
| 103 | // 5. 入队消息 |
| 104 | a.mu.Lock() |
| 105 | a.messages = append(a.messages, userMsg) |
| 106 | a.mu.Unlock() |
| 107 | |
| 108 | // 6. 持久化消息 |
| 109 | if err := a.persistMessage(ctx, &userMsg); err != nil { |
| 110 | streamLog.Warn(ctx, "failed to persist message", map[string]any{"error": err}) |
| 111 | } |
| 112 |
no test coverage detected