Repository navigation
Expand file tree
/
Copy pathcloudflow_build.go
More file actions
115 lines (107 loc) · 3.34 KB
/
Copy pathcloudflow_build.go
File metadata and controls
115 lines (107 loc) · 3.34 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
// CloudFlow NL-builder stream handling: build-cloud-flow and
// refine-cloud-flow answer with a text/event-stream of build events — tool
// lifecycle markers, token-by-token assistant text, and custom events carrying
// the created flow's ID. Raw, that stream reaches the output formatter as one
// quoted string: unreadable for humans and unparseable-without-work for
// agents. This chapter parses it into the result the caller actually needs:
// the flow ID, the conversation ID (to continue with refine-cloud-flow), the
// assistant's answer text, and the build steps that ran.
package main
import (
"encoding/json"
"strings"
)
var cloudflowBuilderOperations = map[string]bool{
"build-cloud-flow": true,
"refine-cloud-flow": true,
}
// transformCloudflowBuildStream parses the NL flow builder's SSE body into a
// structured result. Anything that isn't recognizably that stream — another
// command, a non-string body, no data: lines — passes through untouched, so a
// server-side format change degrades to today's raw output rather than an
// empty one.
func transformCloudflowBuildStream(body interface{}) interface{} {
if !cloudflowBuilderOperations[invokedCommandName] {
return body
}
stream, ok := body.(string)
if !ok || !strings.Contains(stream, "data: ") {
return body
}
var answer strings.Builder
var steps []string
result := map[string]interface{}{}
for _, line := range strings.Split(stream, "\n") {
payload, ok := strings.CutPrefix(strings.TrimSpace(line), "data: ")
if !ok {
continue
}
var event map[string]interface{}
if err := json.Unmarshal([]byte(payload), &event); err != nil {
continue
}
if id, ok := event["conversationId"].(string); ok && id != "" {
result["conversationId"] = id
}
token, ok := event["answer"].(string)
if !ok {
continue
}
lifecycle := parseBuilderLifecycleEvent(token)
if lifecycle == nil {
answer.WriteString(token)
continue
}
if step, ok := lifecycle["toolStart"].(string); ok && step != "" {
steps = append(steps, step)
}
if flowID := builderCreatedFlowID(lifecycle); flowID != "" {
result["flowId"] = flowID
}
}
text := strings.TrimSpace(answer.String())
if text == "" && result["conversationId"] == nil && result["flowId"] == nil {
return body
}
if text != "" {
result["answer"] = text
}
if len(steps) > 0 {
result["steps"] = steps
}
return result
}
// parseBuilderLifecycleEvent distinguishes the stream's embedded lifecycle
// JSON (toolStart/toolEnd/llmStart/llmEnd/customEvent) from assistant text.
// A text token that merely looks like JSON stays text unless it carries one
// of the lifecycle keys.
func parseBuilderLifecycleEvent(token string) map[string]interface{} {
if !strings.HasPrefix(token, "{") {
return nil
}
var event map[string]interface{}
if err := json.Unmarshal([]byte(token), &event); err != nil {
return nil
}
for _, key := range []string{"toolStart", "toolEnd", "llmStart", "llmEnd", "customEvent"} {
if _, ok := event[key]; ok {
return event
}
}
return nil
}
func builderCreatedFlowID(lifecycle map[string]interface{}) string {
custom, ok := lifecycle["customEvent"].(map[string]interface{})
if !ok {
return ""
}
if custom["messageId"] != "cloudflow_created" {
return ""
}
data, ok := custom["data"].(map[string]interface{})
if !ok {
return ""
}
flowID, _ := data["flowId"].(string)
return flowID
}