Monorepo for Tangled tangled.org
1

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / xrpc / ci_pipeline_subscribe_logs.go
7.9 kB 320 lines
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}