This repository has no description
0

Configure Feed

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

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