MCPcopy Create free account
hub / github.com/astercloud/aster / Stream

Method Stream

pkg/agent/streaming.go:55–161  ·  view source on GitHub ↗

Stream 实现流式执行接口 返回 stream.Reader,支持: - 流式生成事件 - 客户端控制的取消 - 与 LLM 流式 API 无缝集成 - 多消费者支持(Copy, Merge, Transform)

(ctx context.Context, message string, opts ...Option)

Source from the content-addressed store, hash-verified

53// - 与 LLM 流式 API 无缝集成
54// - 多消费者支持(Copy, Merge, Transform)
55func (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

Callers 1

handleChatMethod · 0.95

Calls 13

validateMessageMethod · 0.95
GetContentMethod · 0.95
persistMessageMethod · 0.95
runModelStepStreamingMethod · 0.95
EnhanceSystemPromptMethod · 0.80
truncateFunction · 0.70
generateEventIDFunction · 0.70
CloseMethod · 0.65
InfoMethod · 0.65
WarnMethod · 0.65
DebugMethod · 0.65

Tested by

no test coverage detected