This repository has no description
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.Hostname(), 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}