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