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/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}