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
17 kB 657 lines
1package pipelines 2 3import ( 4 "bytes" 5 "context" 6 "fmt" 7 "log/slog" 8 "net/http" 9 "sync" 10 "time" 11 12 "tangled.org/core/api/tangled" 13 "tangled.org/core/appview/config" 14 "tangled.org/core/appview/db" 15 "tangled.org/core/appview/middleware" 16 "tangled.org/core/appview/oauth" 17 "tangled.org/core/appview/pages" 18 "tangled.org/core/appview/reporesolver" 19 "tangled.org/core/hostutil" 20 "tangled.org/core/idresolver" 21 "tangled.org/core/lexutil" 22 "tangled.org/core/rbac" 23 "tangled.org/core/types" 24 25 "github.com/bluesky-social/indigo/atproto/syntax" 26 indigoxrpc "github.com/bluesky-social/indigo/xrpc" 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 db *db.DB 38 enforcer *rbac.Enforcer 39 logger *slog.Logger 40} 41 42func (p *Pipelines) Router(mw *middleware.Middleware) http.Handler { 43 r := chi.NewRouter() 44 r.Get("/", p.Index) 45 r.Get("/{pipeline}/workflow/{workflow}", p.Workflow) 46 r.Get("/{pipeline}/workflow/{workflow}/logs", p.Logs) 47 r.Group(func(r chi.Router) { 48 r.Use(mw.RepoPermissionMiddleware("repo:push")) 49 r.Post("/{pipeline}/workflow/{workflow}/cancel", p.CancelWorkflow) 50 r.Post("/{pipeline}/retry", p.RetryPipeline) 51 r.Post("/{pipeline}/workflow/{workflow}/retry", p.RetryWorkflow) 52 }) 53 54 return r 55} 56 57func New( 58 oauth *oauth.OAuth, 59 repoResolver *reporesolver.RepoResolver, 60 pages *pages.Pages, 61 idResolver *idresolver.Resolver, 62 db *db.DB, 63 config *config.Config, 64 enforcer *rbac.Enforcer, 65 logger *slog.Logger, 66) *Pipelines { 67 return &Pipelines{ 68 oauth: oauth, 69 repoResolver: repoResolver, 70 pages: pages, 71 idResolver: idResolver, 72 config: config, 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 var kinds []string 91 switch filterKind { 92 case "push", "pull_request": 93 kinds = []string{filterKind} 94 default: 95 filterKind = "all" 96 } 97 98 if f.Spindle == "" { 99 p.pages.Pipelines(w, pages.PipelinesParams{ 100 BaseParams: pages.BaseParamsFromContext(r.Context()), 101 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 102 Pipelines: nil, 103 FilterKind: filterKind, 104 Total: 0, 105 }) 106 return 107 } 108 109 spindleUrl, err := hostutil.EnsureHttpScheme(f.Spindle) 110 if err != nil { 111 l.Error("invalid spindle host", "host", f.Spindle, "err", err) 112 p.pages.Pipelines(w, pages.PipelinesParams{ 113 BaseParams: pages.BaseParamsFromContext(r.Context()), 114 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 115 Pipelines: nil, 116 FilterKind: filterKind, 117 Total: 0, 118 }) 119 return 120 } 121 122 xrpcc := indigoxrpc.Client{Host: spindleUrl} 123 out, err := tangled.CiQueryPipelines(r.Context(), &xrpcc, nil, "", kinds, 30, f.RepoDid) 124 if err != nil { 125 l.Error("failed to fetch pipelines", "err", err) 126 p.pages.Pipelines(w, pages.PipelinesParams{ 127 BaseParams: pages.BaseParamsFromContext(r.Context()), 128 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 129 Pipelines: nil, 130 FilterKind: filterKind, 131 Total: 0, 132 }) 133 return 134 } 135 136 var pipelines []types.Pipeline 137 for _, pipeline := range out.Pipelines { 138 if pipeline.Repo == nil || *pipeline.Repo != f.RepoDid { 139 l.Warn("spindle returned pipeline for unexpected repo", 140 "want", f.RepoDid, "got", pipeline.Repo, "spindle", f.Spindle) 141 continue 142 } 143 pipelines = append(pipelines, types.Pipeline{CiPipeline: pipeline}) 144 } 145 146 p.pages.Pipelines(w, pages.PipelinesParams{ 147 BaseParams: pages.BaseParamsFromContext(r.Context()), 148 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 149 Pipelines: pipelines, 150 FilterKind: filterKind, 151 Total: out.Total, 152 }) 153} 154 155func (p *Pipelines) Workflow(w http.ResponseWriter, r *http.Request) { 156 user := p.oauth.GetMultiAccountUser(r) 157 l := p.logger.With("handler", "Workflow") 158 159 f, err := p.repoResolver.Resolve(r) 160 if err != nil { 161 l.Error("failed to get repo and knot", "err", err) 162 p.pages.Error404(w) 163 return 164 } 165 166 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) 167 if err != nil { 168 l.Debug("invalid pipeline id", "id", pipelineId) 169 p.pages.Error404(w) 170 return 171 } 172 173 workflowName := chi.URLParam(r, "workflow") 174 if workflowName == "" { 175 l.Debug("empty workflow name") 176 p.pages.Error404(w) 177 return 178 } 179 180 l = l.With("pipeline", pipelineId, "workflow", workflowName) 181 182 // TODO: change url path to: 183 // /{owner}/{slug}/pipelines/{spindle-did}/{pipeline-id}/workflow/{workflow-id} 184 185 if f.Spindle == "" { 186 p.pages.Error404(w) 187 return 188 } 189 190 spindleUrl, err := hostutil.EnsureHttpScheme(f.Spindle) 191 if err != nil { 192 l.Error("invalid spindle host", "host", f.Spindle, "err", err) 193 p.pages.Error404(w) 194 return 195 } 196 197 xrpcc := &indigoxrpc.Client{Host: spindleUrl} 198 out, err := tangled.CiGetPipeline(r.Context(), xrpcc, pipelineId.String()) 199 if err != nil { 200 // TODO(boltless): change behavior based on error 201 l.Debug("failed to get pipeline", "err", err) 202 p.pages.Error404(w) 203 return 204 } 205 206 if out.Repo == nil || *out.Repo != f.RepoDid { 207 l.Debug("spindle returned pipeline for unexpected repo", "want", f.RepoDid, "got", out.Repo) 208 p.pages.Error404(w) 209 return 210 } 211 // ensure workflow exists 212 exist := false 213 for _, workflow := range out.Workflows { 214 if workflow.Name == workflowName { 215 exist = true 216 break 217 } 218 } 219 if !exist { 220 l.Debug("workflow doesn't exist in pipeline") 221 p.pages.Error404(w) 222 return 223 } 224 225 p.pages.Workflow(w, pages.WorkflowParams{ 226 BaseParams: pages.BaseParamsFromContext(r.Context()), 227 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 228 Pipeline: types.Pipeline{CiPipeline: out}, 229 Workflow: workflowName, 230 }) 231} 232 233var upgrader = websocket.Upgrader{ 234 ReadBufferSize: 1024, 235 WriteBufferSize: 1024, 236} 237 238type webLogScheduler struct { 239 ch chan *tangled.CiSubscribePipelineLogs_Event 240} 241 242var _ lexutil.Scheduler[tangled.CiSubscribePipelineLogs_Event] = (*webLogScheduler)(nil) 243 244// AddWork implements [lexutil.Scheduler]. 245func (w *webLogScheduler) AddWork(ctx context.Context, _ string, val *tangled.CiSubscribePipelineLogs_Event) error { 246 select { 247 case w.ch <- val: 248 return nil 249 case <-ctx.Done(): 250 return ctx.Err() 251 } 252} 253 254// Shutdown implements [lexutil.Scheduler]. 255func (w *webLogScheduler) Shutdown() { close(w.ch) } 256 257func retryPipelineTrigger(orig *tangled.CiPipeline) *tangled.CiTriggerPipeline_Input_Trigger { 258 if orig.Trigger != nil && orig.Trigger.CiTrigger_PullRequest != nil { 259 pr := orig.Trigger.CiTrigger_PullRequest 260 sourceSha := pr.SourceSha 261 if sourceSha == "" { 262 sourceSha = orig.Commit 263 } 264 sourceRepo := pr.SourceRepo 265 if sourceRepo == nil { 266 sourceRepo = orig.SourceRepo 267 } 268 269 return &tangled.CiTriggerPipeline_Input_Trigger{ 270 CiTrigger_PullRequest: &tangled.CiTrigger_PullRequest{ 271 Pull: pr.Pull, 272 SourceBranch: pr.SourceBranch, 273 SourceRepo: sourceRepo, 274 SourceSha: sourceSha, 275 TargetBranch: pr.TargetBranch, 276 }, 277 } 278 } 279 280 manual := &tangled.CiTrigger_Manual{ 281 Sha: orig.Commit, 282 SourceRepo: orig.SourceRepo, 283 } 284 if orig.Trigger != nil && orig.Trigger.CiTrigger_Manual != nil { 285 origManual := orig.Trigger.CiTrigger_Manual 286 if origManual.Sha != "" { 287 manual.Sha = origManual.Sha 288 } 289 manual.Ref = origManual.Ref 290 manual.Inputs = origManual.Inputs 291 if origManual.SourceRepo != nil { 292 manual.SourceRepo = origManual.SourceRepo 293 } 294 } 295 296 return &tangled.CiTriggerPipeline_Input_Trigger{ 297 CiTrigger_Manual: manual, 298 } 299} 300 301func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { 302 l := p.logger.With("handler", "logs") 303 304 f, err := p.repoResolver.Resolve(r) 305 if err != nil { 306 l.Error("failed to get repo and knot", "err", err) 307 http.Error(w, "bad repo/knot", http.StatusBadRequest) 308 return 309 } 310 311 if f.Spindle == "" { 312 http.Error(w, "invalid repo info", http.StatusBadRequest) 313 return 314 } 315 316 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) 317 if err != nil { 318 l.Debug("invalid pipeline id", "id", pipelineId) 319 http.Error(w, "invalid pipeline id", http.StatusBadRequest) 320 return 321 } 322 323 workflowName := chi.URLParam(r, "workflow") 324 if workflowName == "" { 325 l.Debug("empty workflow name") 326 http.Error(w, "invalid workflow name", http.StatusBadRequest) 327 return 328 } 329 330 clientConn, err := upgrader.Upgrade(w, r, nil) 331 if err != nil { 332 l.Error("websocket upgrade failed", "err", err) 333 return 334 } 335 defer clientConn.Close() 336 337 ctx, cancel := context.WithCancel(r.Context()) 338 defer cancel() 339 340 spindleUrl, err := hostutil.EnsureHttpScheme(f.Spindle) 341 if err != nil { 342 l.Error("invalid spindle host", "host", f.Spindle, "err", err) 343 return 344 } 345 346 evChan := make(chan *tangled.CiSubscribePipelineLogs_Event, 100) 347 done := make(chan error, 1) 348 sched := &webLogScheduler{ch: evChan} 349 xrpcc := &lexutil.Client{Client: indigoxrpc.Client{Host: spindleUrl}} 350 go func() { 351 done <- tangled.CiSubscribePipelineLogs(ctx, xrpcc, pipelineId.String(), []string{workflowName}, sched) 352 }() 353 354 var lastWriteLk sync.Mutex 355 lastWrite := time.Now() 356 357 // Start a goroutine to ping the client periodically to check if it's still 358 // alive. If the client doesn't respond to a ping within 5 seconds, we'll 359 // close the connection and teardown the consumer. 360 go func() { 361 ticker := time.NewTicker(30 * time.Second) 362 defer ticker.Stop() 363 for { 364 select { 365 case <-ticker.C: 366 lastWriteLk.Lock() 367 lw := lastWrite 368 lastWriteLk.Unlock() 369 if time.Since(lw) < 30*time.Second { 370 continue 371 } 372 if err := clientConn.WriteControl(websocket.PingMessage, nil, time.Now().Add(5*time.Second)); err != nil { 373 l.Warn("failed to ping client", "err", err) 374 cancel() 375 return 376 } 377 case <-ctx.Done(): 378 return 379 } 380 } 381 }() 382 383 clientConn.SetPingHandler(func(message string) error { 384 err := clientConn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(60*time.Second)) 385 if err == websocket.ErrCloseSent { 386 return nil 387 } 388 return err 389 }) 390 391 // Start a goroutine to read messages from the client and discard them. 392 go func() { 393 for { 394 if _, _, err := clientConn.ReadMessage(); err != nil { 395 cancel() 396 return 397 } 398 } 399 }() 400 401 // Main loop: sole writer of data frames to the client. 402 stepStartTimes := make(map[int]time.Time) 403 stepAnsi := make(map[int]*ansiState) 404 var fragment bytes.Buffer 405 for { 406 select { 407 case <-ctx.Done(): 408 l.Info("client disconnected") 409 return 410 411 case ev, ok := <-evChan: 412 if !ok { 413 // Stream ended: Shutdown closed the upstream channel. 414 if err := <-done; !isExpectedClose(err) { 415 l.Error("spindle stream error", "err", err) 416 } 417 msg := websocket.FormatCloseMessage(websocket.CloseNormalClosure, "finished") 418 _ = clientConn.WriteMessage(websocket.CloseMessage, msg) 419 return 420 } 421 422 fragment.Reset() 423 424 switch { 425 case ev.Error != nil: 426 l.Error("spindle error frame", "err", ev.Error.Error, "msg", ev.Error.Message) 427 return 428 429 case ev.Control != nil: 430 c := ev.Control 431 step := int(c.Step) 432 switch derefStr(c.Status) { 433 case "start": 434 t := parseRFC3339(c.Time) 435 stepStartTimes[step] = t 436 // "system" steps are injected by the CI runner; collapse them. 437 collapsed := derefStr(c.Kind) == "system" 438 err = p.pages.LogBlock(&fragment, pages.LogBlockParams{ 439 Id: step, 440 Name: c.Content, 441 Command: derefStr(c.Command), 442 Collapsed: collapsed, 443 StartTime: t, 444 }) 445 case "end": 446 err = p.pages.LogBlockEnd(&fragment, pages.LogBlockEndParams{ 447 Id: step, 448 StartTime: stepStartTimes[step], 449 EndTime: parseRFC3339(c.Time), 450 }) 451 } 452 453 case ev.Data != nil: 454 d := ev.Data 455 step := int(d.Step) 456 ansi, ok := stepAnsi[step] 457 if !ok { 458 ansi = NewAnsiState() 459 stepAnsi[step] = ansi 460 } 461 err = p.pages.LogLine(&fragment, pages.LogLineParams{ 462 Id: step, 463 Content: ansi.Render(d.Content), 464 }) 465 } 466 if err != nil { 467 l.Error("failed to render log line", "err", err) 468 return 469 } 470 471 if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil { 472 l.Error("error writing to client", "err", err) 473 return 474 } 475 lastWriteLk.Lock() 476 lastWrite = time.Now() 477 lastWriteLk.Unlock() 478 } 479 } 480} 481 482func (p *Pipelines) CancelWorkflow(w http.ResponseWriter, r *http.Request) { 483 l := p.logger.With("handler", "CancelWorkflow") 484 errorId := "workflow-error" 485 486 f, err := p.repoResolver.Resolve(r) 487 if err != nil { 488 l.Error("failed to get repo and knot", "err", err) 489 p.pages.Notice(w, errorId, "Failed to cancel workflow") 490 return 491 } 492 l = l.With("repo", f.RepoDid) 493 494 if f.Spindle == "" { 495 l.Debug("spindle is empty") 496 p.pages.Notice(w, errorId, "Failed to cancel workflow") 497 return 498 } 499 500 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) 501 if err != nil { 502 l.Debug("invalid pipeline id", "id", pipelineId) 503 p.pages.Error404(w) 504 return 505 } 506 507 workflowName := chi.URLParam(r, "workflow") 508 if workflowName == "" { 509 l.Debug("empty workflow name") 510 p.pages.Error404(w) 511 return 512 } 513 514 l = l.With("pipeline", pipelineId, "workflow", workflowName) 515 516 spindleClient, err := p.oauth.SpindleServiceClient(r, f.Spindle, tangled.CiCancelPipelineNSID) 517 if err != nil { 518 l.Error("failed to prepare spindle client", "err", err) 519 p.pages.Notice(w, errorId, "Failed to cancel workflow") 520 return 521 } 522 523 if err := tangled.CiCancelPipeline( 524 r.Context(), 525 spindleClient, 526 &tangled.CiCancelPipeline_Input{ 527 Repo: f.RepoDid, 528 Pipeline: pipelineId.String(), 529 Workflows: []string{workflowName}, 530 }, 531 ); err != nil { 532 l.Error("failed to cancel workflow", "err", err) 533 p.pages.Notice(w, errorId, "Failed to cancel workflow") 534 return 535 } 536 l.Debug("canceled workflow") 537} 538 539// RetryPipeline retries all workflows in a pipeline 540func (p *Pipelines) RetryPipeline(w http.ResponseWriter, r *http.Request) { 541 p.retry(w, r, "") 542} 543 544// RetryWorkflow retries a single workflow in a pipeline 545func (p *Pipelines) RetryWorkflow(w http.ResponseWriter, r *http.Request) { 546 p.retry(w, r, chi.URLParam(r, "workflow")) 547} 548 549// retry triggers a new pipeline run for the original commit, either for all or a single workflow 550func (p *Pipelines) retry(w http.ResponseWriter, r *http.Request, only string) { 551 user := p.oauth.GetMultiAccountUser(r) 552 l := p.logger.With("handler", "retry", "only", only) 553 errorId := "workflow-error" 554 555 // fail logs the error and shows a notice to the user 556 fail := func(msg string, err error) { 557 if err != nil { 558 l.Error(msg, "err", err) 559 p.pages.Notice(w, errorId, fmt.Sprintf("%s: %v", msg, err)) 560 } else { 561 l.Error(msg) 562 p.pages.Notice(w, errorId, msg) 563 } 564 } 565 566 f, err := p.repoResolver.Resolve(r) 567 if err != nil { 568 fail("failed to resolve repository", err) 569 return 570 } 571 l = l.With("repo", f.RepoDid) 572 573 if f.Spindle == "" { 574 fail("this repository has no spindle configured", nil) 575 return 576 } 577 578 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) 579 if err != nil { 580 l.Debug("invalid pipeline id", "id", pipelineId) 581 p.pages.Error404(w) 582 return 583 } 584 l = l.With("pipeline", pipelineId) 585 586 spindleUrl, err := hostutil.EnsureHttpScheme(f.Spindle) 587 if err != nil { 588 fail("invalid spindle host", err) 589 return 590 } 591 592 // fetch the original pipeline to replay the same commit and workflows 593 queryClient := &indigoxrpc.Client{Host: spindleUrl} 594 orig, err := tangled.CiGetPipeline(r.Context(), queryClient, pipelineId.String()) 595 if err != nil { 596 fail("failed to load the original pipeline", err) 597 return 598 } 599 if orig.Commit == "" { 600 fail("cannot retry: the original pipeline has no commit", nil) 601 return 602 } 603 604 // figure out which workflows to run and where to redirect 605 var workflows []string 606 if only != "" { 607 workflows = []string{only} 608 } else { 609 for _, wf := range orig.Workflows { 610 workflows = append(workflows, wf.Name) 611 } 612 } 613 if len(workflows) == 0 { 614 fail("cannot retry: the original pipeline has no workflows", nil) 615 return 616 } 617 redirectWf := workflows[0] 618 619 spindleClient, err := p.oauth.SpindleServiceClient(r, f.Spindle, tangled.CiTriggerPipelineNSID) 620 if err != nil { 621 fail("failed to authorize with spindle", err) 622 return 623 } 624 625 out, err := tangled.CiTriggerPipeline( 626 r.Context(), 627 spindleClient, 628 &tangled.CiTriggerPipeline_Input{ 629 Repo: f.RepoDid, 630 Trigger: retryPipelineTrigger(orig), 631 Workflows: workflows, 632 }, 633 ) 634 if err != nil { 635 fail("spindle rejected the trigger", err) 636 return 637 } 638 639 newAt, err := syntax.ParseATURI(out.Pipeline) 640 if err != nil { 641 fail("pipeline triggered, but the response was malformed", err) 642 return 643 } 644 newId := newAt.RecordKey().String() 645 l = l.With("new", newId) 646 l.Info("pipeline retried") 647 648 repoInfo := p.repoResolver.GetRepoInfo(r, user) 649 dest := fmt.Sprintf("/%s/pipelines/%s/workflow/%s", repoInfo.FullName(), newId, redirectWf) 650 651 if r.Header.Get("HX-Request") == "true" { 652 w.Header().Set("HX-Redirect", dest) 653 w.WriteHeader(http.StatusOK) 654 return 655 } 656 http.Redirect(w, r, dest, http.StatusSeeOther) 657}