This repository has no description
1package xrpc
2
3import (
4 "context"
5 "encoding/json"
6 "fmt"
7 "net/http"
8 "sync"
9 "time"
10
11 "github.com/bluesky-social/indigo/atproto/atclient"
12 "github.com/bluesky-social/indigo/atproto/syntax"
13 "github.com/gorilla/websocket"
14 "tangled.org/core/api/tangled"
15 "tangled.org/core/spindle/logview"
16 "tangled.org/core/spindle/models"
17)
18
19func (x *Xrpc) HandleCiSubscribePipelineLogs(w http.ResponseWriter, r *http.Request) {
20 var (
21 pipelineQuery = r.URL.Query().Get("pipeline")
22 workflows = r.URL.Query()["workflows"]
23 )
24
25 pipeline, err := syntax.ParseTID(pipelineQuery)
26 if err != nil {
27 writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "BadRequest", Message: fmt.Sprintf("pipeline parameter invalid: %s", pipelineQuery)})
28 return
29 }
30
31 x.handleSubscribeLogs(w, r, pipeline, workflows)
32}
33
34var wsUpgrader = websocket.Upgrader{
35 ReadBufferSize: 10_000,
36 WriteBufferSize: 10_000,
37}
38
39func (x *Xrpc) handleSubscribeLogs(w http.ResponseWriter, r *http.Request, pipeline syntax.TID, workflows []string) {
40 l := x.Logger.With("pipeline", pipeline, "workflows", workflows)
41
42 // 1. query the event from database to get the knot
43 var eventJson string
44 err := x.Db.QueryRow(
45 `select event from events where nsid = ? and rkey = ?`,
46 tangled.PipelineNSID,
47 pipeline.String(),
48 ).Scan(&eventJson)
49 if err != nil {
50 l.Error("failed to find pipeline event", "err", err)
51 writeJson(w, http.StatusNotFound, atclient.ErrorBody{Name: "NotFound", Message: fmt.Sprintf("pipeline not found: %s", pipeline.String())})
52 return
53 }
54
55 var tpl tangled.Pipeline
56 if err := json.Unmarshal([]byte(eventJson), &tpl); err != nil {
57 l.Error("failed to unmarshal pipeline event", "err", err)
58 writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalError", Message: "failed to parse pipeline event"})
59 return
60 }
61
62 if tpl.TriggerMetadata == nil || tpl.TriggerMetadata.Repo == nil {
63 l.Error("pipeline event trigger metadata is incomplete")
64 writeJson(w, http.StatusInternalServerError, atclient.ErrorBody{Name: "InternalError", Message: "pipeline event trigger metadata is incomplete"})
65 return
66 }
67 knot := tpl.TriggerMetadata.Repo.Knot
68
69 // 2. if workflows is empty, default to all workflows defined in the pipeline
70 if len(workflows) == 0 {
71 for _, wf := range tpl.Workflows {
72 if wf != nil && wf.Name != "" {
73 workflows = append(workflows, wf.Name)
74 }
75 }
76 }
77
78 if len(workflows) == 0 {
79 writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "BadRequest", Message: "no workflows specified or found"})
80 return
81 }
82
83 // 3. upgrade to websocket
84 ctx, cancel := context.WithCancel(r.Context())
85 defer cancel()
86
87 conn, err := wsUpgrader.Upgrade(w, r, w.Header())
88 if err != nil {
89 l.Error("websocket upgrade failed", "err", err)
90 return
91 }
92 defer conn.Close()
93
94 lastWriteLk := sync.Mutex{}
95 lastWrite := time.Now()
96
97 // Ping loop
98 go func() {
99 ticker := time.NewTicker(30 * time.Second)
100 defer ticker.Stop()
101
102 for {
103 select {
104 case <-ticker.C:
105 lastWriteLk.Lock()
106 lw := lastWrite
107 lastWriteLk.Unlock()
108
109 if time.Since(lw) < 30*time.Second {
110 continue
111 }
112
113 if err := conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(5*time.Second)); err != nil {
114 l.Warn("failed to ping client", "err", err)
115 cancel()
116 return
117 }
118 case <-ctx.Done():
119 return
120 }
121 }
122 }()
123
124 conn.SetPingHandler(func(message string) error {
125 err := conn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(time.Second*60))
126 if err == websocket.ErrCloseSent {
127 return nil
128 }
129 return err
130 })
131
132 // Read discard loop
133 go func() {
134 for {
135 _, _, err := conn.ReadMessage()
136 if err != nil {
137 l.Warn("failed to read message from client", "err", err)
138 cancel()
139 return
140 }
141 }
142 }()
143
144 eventsChan := make(chan tangled.CiSubscribePipelineLogs_Event, 128)
145 wg := sync.WaitGroup{}
146
147 // 4. start a tail reader goroutine for each workflow
148 for _, wf := range workflows {
149 wg.Add(1)
150 go func(wfName string) {
151 defer wg.Done()
152
153 wid := models.WorkflowId{
154 PipelineId: models.PipelineId{
155 Knot: knot,
156 Rkey: pipeline.String(),
157 },
158 Name: wfName,
159 }
160
161 // check if finished, but poll database to know when it finishes
162 var isFinished bool
163 status, err := x.Db.GetStatus(wid)
164 if err == nil {
165 isFinished = models.StatusKind(status.Status).IsFinish()
166 }
167
168 lines, stop, err := logview.Follow(ctx, x.Db, x.ArtifactReader, x.Config.Server.LogDir, wid, isFinished)
169 if err != nil {
170 l.Error("failed to follow workflow log", "workflow", wfName, "err", err)
171 return
172 }
173 defer stop()
174
175 // if we are following, poll status in database to stop when finished
176 if !isFinished {
177 go func() {
178 ticker := time.NewTicker(2 * time.Second)
179 defer ticker.Stop()
180 for {
181 select {
182 case <-ctx.Done():
183 return
184 case <-ticker.C:
185 status, err := x.Db.GetStatus(wid)
186 if err == nil && models.StatusKind(status.Status).IsFinish() {
187 stop()
188 return
189 }
190 }
191 }
192 }()
193 }
194
195 for {
196 select {
197 case <-ctx.Done():
198 return
199 case line, ok := <-lines:
200 if !ok || line == nil {
201 return
202 }
203
204 if line.Err != nil {
205 l.Warn("error tailing log file", "workflow", wfName, "err", line.Err)
206 return
207 }
208
209 var logLine models.LogLine
210 if err := json.Unmarshal([]byte(line.Text), &logLine); err != nil {
211 // if it's not JSON, treat it as a raw data line
212 logLine = models.NewDataLogLine(0, line.Text, "stdout")
213 }
214
215 var ev tangled.CiSubscribePipelineLogs_Event
216 timeStr := logLine.Time.Format(time.RFC3339)
217 if logLine.Time.IsZero() {
218 timeStr = time.Now().Format(time.RFC3339)
219 }
220
221 if logLine.Kind == models.LogKindControl {
222 stepKindStr := "user"
223 if logLine.StepKind == models.StepKindSystem {
224 stepKindStr = "system"
225 }
226 ev = tangled.CiSubscribePipelineLogs_Event{Control: &tangled.CiSubscribePipelineLogs_Control{
227 Time: timeStr,
228 Workflow: wfName,
229 Step: int64(logLine.StepId),
230 Content: logLine.Content,
231 Command: strptrOrNil(logLine.StepCommand),
232 Status: strptrOrNil(string(logLine.StepStatus)),
233 Kind: strptrOrNil(stepKindStr),
234 }}
235 } else {
236 streamType := logLine.Stream
237 if streamType != "stdout" && streamType != "stderr" {
238 streamType = "stdout"
239 }
240 ev = tangled.CiSubscribePipelineLogs_Event{Data: &tangled.CiSubscribePipelineLogs_Data{
241 Time: timeStr,
242 Workflow: wfName,
243 Step: int64(logLine.StepId),
244 Content: logLine.Content + "\n", // Append newline back since logger trims it
245 Stream: streamType,
246 }}
247 }
248
249 select {
250 case eventsChan <- ev:
251 case <-ctx.Done():
252 return
253 }
254 }
255 }
256 }(wf)
257 }
258
259 // Closer goroutine for eventsChan
260 go func() {
261 wg.Wait()
262 close(eventsChan)
263 }()
264
265 // Main writer loop
266 for {
267 select {
268 case <-ctx.Done():
269 return
270 case evt, ok := <-eventsChan:
271 if !ok {
272 return
273 }
274
275 wc, err := conn.NextWriter(websocket.BinaryMessage)
276 if err != nil {
277 l.Error("failed to get next writer", "err", err)
278 return
279 }
280
281 err = evt.Serialize(wc)
282 if err != nil {
283 l.Error("failed to serialize event", "err", err)
284 wc.Close()
285 return
286 }
287
288 if err := wc.Close(); err != nil {
289 l.Warn("failed to flush-close event write", "err", err)
290 return
291 }
292
293 lastWriteLk.Lock()
294 lastWrite = time.Now()
295 lastWriteLk.Unlock()
296 }
297 }
298}
299
300func strptr(s string) *string { return &s }
301
302func strptrOrNil(s string) *string {
303 if s == "" {
304 return nil
305 }
306 return &s
307}