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