This repository has no description
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}