This repository has no description
0

Configure Feed

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

core / appview / pipelines / pipelines.go
11 kB 454 lines
1package pipelines 2 3import ( 4 "bytes" 5 "context" 6 "encoding/json" 7 "fmt" 8 "log/slog" 9 "net/http" 10 "strings" 11 "time" 12 13 "tangled.org/core/api/tangled" 14 "tangled.org/core/appview/config" 15 "tangled.org/core/appview/db" 16 "tangled.org/core/appview/middleware" 17 "tangled.org/core/appview/models" 18 "tangled.org/core/appview/oauth" 19 "tangled.org/core/appview/pages" 20 "tangled.org/core/appview/reporesolver" 21 "tangled.org/core/eventconsumer" 22 "tangled.org/core/idresolver" 23 "tangled.org/core/orm" 24 "tangled.org/core/rbac" 25 spindlemodel "tangled.org/core/spindle/models" 26 27 "github.com/go-chi/chi/v5" 28 "github.com/gorilla/websocket" 29) 30 31type Pipelines struct { 32 repoResolver *reporesolver.RepoResolver 33 idResolver *idresolver.Resolver 34 config *config.Config 35 oauth *oauth.OAuth 36 pages *pages.Pages 37 spindlestream *eventconsumer.Consumer 38 db *db.DB 39 enforcer *rbac.Enforcer 40 logger *slog.Logger 41} 42 43func (p *Pipelines) Router(mw *middleware.Middleware) http.Handler { 44 r := chi.NewRouter() 45 r.Get("/", p.Index) 46 r.Get("/{pipeline}/workflow/{workflow}", p.Workflow) 47 r.Get("/{pipeline}/workflow/{workflow}/logs", p.Logs) 48 r. 49 With(mw.RepoPermissionMiddleware("repo:owner")). 50 Post("/{pipeline}/workflow/{workflow}/cancel", p.Cancel) 51 52 return r 53} 54 55func New( 56 oauth *oauth.OAuth, 57 repoResolver *reporesolver.RepoResolver, 58 pages *pages.Pages, 59 spindlestream *eventconsumer.Consumer, 60 idResolver *idresolver.Resolver, 61 db *db.DB, 62 config *config.Config, 63 enforcer *rbac.Enforcer, 64 logger *slog.Logger, 65) *Pipelines { 66 return &Pipelines{ 67 oauth: oauth, 68 repoResolver: repoResolver, 69 pages: pages, 70 idResolver: idResolver, 71 config: config, 72 spindlestream: spindlestream, 73 db: db, 74 enforcer: enforcer, 75 logger: logger, 76 } 77} 78 79func (p *Pipelines) Index(w http.ResponseWriter, r *http.Request) { 80 user := p.oauth.GetMultiAccountUser(r) 81 l := p.logger.With("handler", "Index") 82 83 f, err := p.repoResolver.Resolve(r) 84 if err != nil { 85 l.Error("failed to get repo and knot", "err", err) 86 return 87 } 88 89 filterKind := r.URL.Query().Get("trigger") 90 filters := []orm.Filter{ 91 orm.FilterEq("p.repo_did", f.RepoDid), 92 } 93 switch filterKind { 94 case "push": 95 filters = append(filters, orm.FilterEq("t.kind", "push")) 96 case "pull_request": 97 filters = append(filters, orm.FilterEq("t.kind", "pull_request")) 98 default: 99 // no filters otherwise, default to "all" 100 filterKind = "all" 101 } 102 103 ps, err := db.GetPipelineStatuses( 104 p.db, 105 30, 106 filters..., 107 ) 108 if err != nil { 109 l.Error("failed to query db", "err", err) 110 return 111 } 112 113 total, err := db.GetTotalPipelineStatuses(p.db, filters...) 114 if err != nil { 115 l.Error("failed to query db", "err", err) 116 return 117 } 118 119 p.pages.Pipelines(w, pages.PipelinesParams{ 120 LoggedInUser: user, 121 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 122 Pipelines: ps, 123 FilterKind: filterKind, 124 Total: total, 125 }) 126} 127 128func (p *Pipelines) Workflow(w http.ResponseWriter, r *http.Request) { 129 user := p.oauth.GetMultiAccountUser(r) 130 l := p.logger.With("handler", "Workflow") 131 132 f, err := p.repoResolver.Resolve(r) 133 if err != nil { 134 l.Error("failed to get repo and knot", "err", err) 135 return 136 } 137 138 pipelineId := chi.URLParam(r, "pipeline") 139 if pipelineId == "" { 140 l.Error("empty pipeline ID") 141 return 142 } 143 144 workflow := chi.URLParam(r, "workflow") 145 if workflow == "" { 146 l.Error("empty workflow name") 147 return 148 } 149 150 ps, err := db.GetPipelineStatuses( 151 p.db, 152 1, 153 orm.FilterEq("p.repo_owner", f.Did), 154 orm.FilterEq("p.repo_name", f.Rkey), 155 orm.FilterEq("p.knot", f.Knot), 156 orm.FilterEq("p.id", pipelineId), 157 ) 158 if err != nil { 159 l.Error("failed to query db", "err", err) 160 return 161 } 162 163 if len(ps) != 1 { 164 l.Error("invalid number of pipelines", "len", len(ps)) 165 return 166 } 167 168 singlePipeline := ps[0] 169 170 p.pages.Workflow(w, pages.WorkflowParams{ 171 LoggedInUser: user, 172 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 173 Pipeline: singlePipeline, 174 Workflow: workflow, 175 }) 176} 177 178var upgrader = websocket.Upgrader{ 179 ReadBufferSize: 1024, 180 WriteBufferSize: 1024, 181} 182 183func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { 184 l := p.logger.With("handler", "logs") 185 186 clientConn, err := upgrader.Upgrade(w, r, nil) 187 if err != nil { 188 l.Error("websocket upgrade failed", "err", err) 189 return 190 } 191 defer func() { 192 _ = clientConn.WriteControl( 193 websocket.CloseMessage, 194 websocket.FormatCloseMessage(websocket.CloseNormalClosure, "log stream complete"), 195 time.Now().Add(time.Second), 196 ) 197 clientConn.Close() 198 }() 199 200 ctx, cancel := context.WithCancel(r.Context()) 201 defer cancel() 202 203 f, err := p.repoResolver.Resolve(r) 204 if err != nil { 205 l.Error("failed to get repo and knot", "err", err) 206 http.Error(w, "bad repo/knot", http.StatusBadRequest) 207 return 208 } 209 210 pipelineId := chi.URLParam(r, "pipeline") 211 workflow := chi.URLParam(r, "workflow") 212 if pipelineId == "" || workflow == "" { 213 http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest) 214 return 215 } 216 217 ps, err := db.GetPipelineStatuses( 218 p.db, 219 1, 220 orm.FilterEq("p.repo_owner", f.Did), 221 orm.FilterEq("p.repo_name", f.Rkey), 222 orm.FilterEq("p.knot", f.Knot), 223 orm.FilterEq("p.id", pipelineId), 224 ) 225 if err != nil || len(ps) != 1 { 226 l.Error("pipeline query failed", "err", err, "count", len(ps)) 227 http.Error(w, "pipeline not found", http.StatusNotFound) 228 return 229 } 230 231 singlePipeline := ps[0] 232 spindle := f.Spindle 233 knot := f.Knot 234 rkey := singlePipeline.Rkey 235 236 if spindle == "" || knot == "" || rkey == "" { 237 http.Error(w, "invalid repo info", http.StatusBadRequest) 238 return 239 } 240 241 scheme := "wss" 242 if p.config.Core.Dev { 243 scheme = "ws" 244 } 245 246 url := scheme + "://" + strings.Join([]string{spindle, "logs", knot, rkey, workflow}, "/") 247 l = l.With("url", url) 248 l.Info("logs endpoint hit") 249 250 spindleConn, _, err := websocket.DefaultDialer.Dial(url, nil) 251 if err != nil { 252 l.Error("websocket dial failed", "err", err) 253 http.Error(w, "failed to connect to log stream", http.StatusBadGateway) 254 return 255 } 256 defer spindleConn.Close() 257 258 // create a channel for incoming messages 259 evChan := make(chan logEvent, 100) 260 // start a goroutine to read from spindle 261 go readLogs(spindleConn, evChan) 262 263 stepStartTimes := make(map[int]time.Time) 264 var fragment bytes.Buffer 265 for { 266 select { 267 case <-ctx.Done(): 268 l.Info("client disconnected") 269 return 270 271 case ev, ok := <-evChan: 272 if !ok { 273 continue 274 } 275 276 if ev.err != nil && ev.isCloseError() { 277 l.Debug("graceful shutdown, tail complete", "err", err) 278 return 279 } 280 if ev.err != nil { 281 l.Error("error reading from spindle", "err", err) 282 return 283 } 284 285 var logLine spindlemodel.LogLine 286 if err = json.Unmarshal(ev.msg, &logLine); err != nil { 287 l.Error("failed to parse logline", "err", err) 288 continue 289 } 290 291 fragment.Reset() 292 293 switch logLine.Kind { 294 case spindlemodel.LogKindControl: 295 switch logLine.StepStatus { 296 case spindlemodel.StepStatusStart: 297 stepStartTimes[logLine.StepId] = logLine.Time 298 collapsed := false 299 if logLine.StepKind == spindlemodel.StepKindSystem { 300 collapsed = true 301 } 302 err = p.pages.LogBlock(&fragment, pages.LogBlockParams{ 303 Id: logLine.StepId, 304 Name: logLine.Content, 305 Command: logLine.StepCommand, 306 Collapsed: collapsed, 307 StartTime: logLine.Time, 308 }) 309 case spindlemodel.StepStatusEnd: 310 startTime := stepStartTimes[logLine.StepId] 311 endTime := logLine.Time 312 err = p.pages.LogBlockEnd(&fragment, pages.LogBlockEndParams{ 313 Id: logLine.StepId, 314 StartTime: startTime, 315 EndTime: endTime, 316 }) 317 } 318 319 case spindlemodel.LogKindData: 320 // data messages simply insert new log lines into current step 321 err = p.pages.LogLine(&fragment, pages.LogLineParams{ 322 Id: logLine.StepId, 323 Content: logLine.Content, 324 }) 325 } 326 if err != nil { 327 l.Error("failed to render log line", "err", err) 328 return 329 } 330 331 if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil { 332 l.Error("error writing to client", "err", err) 333 return 334 } 335 336 case <-time.After(30 * time.Second): 337 l.Debug("sent keepalive") 338 if err = clientConn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { 339 l.Error("failed to write control", "err", err) 340 return 341 } 342 } 343 } 344} 345 346func (p *Pipelines) Cancel(w http.ResponseWriter, r *http.Request) { 347 l := p.logger.With("handler", "Cancel") 348 349 var ( 350 pipelineId = chi.URLParam(r, "pipeline") 351 workflow = chi.URLParam(r, "workflow") 352 ) 353 if pipelineId == "" || workflow == "" { 354 http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest) 355 return 356 } 357 358 f, err := p.repoResolver.Resolve(r) 359 if err != nil { 360 l.Error("failed to get repo and knot", "err", err) 361 http.Error(w, "bad repo/knot", http.StatusBadRequest) 362 return 363 } 364 365 pipeline, err := func() (models.Pipeline, error) { 366 ps, err := db.GetPipelineStatuses( 367 p.db, 368 1, 369 orm.FilterEq("p.repo_owner", f.Did), 370 orm.FilterEq("p.repo_name", f.Rkey), 371 orm.FilterEq("p.knot", f.Knot), 372 orm.FilterEq("p.id", pipelineId), 373 ) 374 if err != nil { 375 return models.Pipeline{}, err 376 } 377 if len(ps) != 1 { 378 return models.Pipeline{}, fmt.Errorf("wrong pipeline count %d", len(ps)) 379 } 380 return ps[0], nil 381 }() 382 if err != nil { 383 l.Error("pipeline query failed", "err", err) 384 http.Error(w, "pipeline not found", http.StatusNotFound) 385 } 386 var ( 387 spindle = f.Spindle 388 knot = f.Knot 389 rkey = pipeline.Rkey 390 ) 391 392 if spindle == "" || knot == "" || rkey == "" { 393 http.Error(w, "invalid repo info", http.StatusBadRequest) 394 return 395 } 396 397 spindleClient, err := p.oauth.ServiceClient( 398 r, 399 oauth.WithService(f.Spindle), 400 oauth.WithLxm(tangled.PipelineCancelPipelineNSID), 401 oauth.WithDev(p.config.Core.Dev), 402 oauth.WithTimeout(time.Second*30), // workflow cleanup usually takes time 403 ) 404 405 err = tangled.PipelineCancelPipeline( 406 r.Context(), 407 spindleClient, 408 &tangled.PipelineCancelPipeline_Input{ 409 Repo: string(f.RepoAt()), 410 Pipeline: pipeline.AtUri().String(), 411 Workflow: workflow, 412 }, 413 ) 414 errorId := "workflow-error" 415 if err != nil { 416 l.Error("failed to cancel workflow", "err", err) 417 p.pages.Notice(w, errorId, "Failed to cancel workflow") 418 return 419 } 420 l.Debug("canceled pipeline", "uri", pipeline.AtUri()) 421} 422 423// either a message or an error 424type logEvent struct { 425 msg []byte 426 err error 427} 428 429func (ev *logEvent) isCloseError() bool { 430 return websocket.IsCloseError( 431 ev.err, 432 websocket.CloseNormalClosure, 433 websocket.CloseGoingAway, 434 websocket.CloseAbnormalClosure, 435 ) 436} 437 438// read logs from spindle and pass through to chan 439func readLogs(conn *websocket.Conn, ch chan logEvent) { 440 defer close(ch) 441 442 for { 443 if conn == nil { 444 return 445 } 446 447 _, msg, err := conn.ReadMessage() 448 if err != nil { 449 ch <- logEvent{err: err} 450 return 451 } 452 ch <- logEvent{msg: msg} 453 } 454}