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