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