Monorepo for Tangled
tangled.org
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 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 // ensure workflow exists
209 exist := false
210 for _, workflow := range out.Workflows {
211 if workflow.Name == workflowName {
212 exist = true
213 break
214 }
215 }
216 if !exist {
217 l.Debug("workflow doesn't exist in pipeline")
218 p.pages.Error404(w)
219 return
220 }
221
222 p.pages.Workflow(w, pages.WorkflowParams{
223 BaseParams: pages.BaseParamsFromContext(r.Context()),
224 RepoInfo: p.repoResolver.GetRepoInfo(r, user),
225 Pipeline: types.Pipeline{CiPipeline: out},
226 Workflow: workflowName,
227 })
228}
229
230var upgrader = websocket.Upgrader{
231 ReadBufferSize: 1024,
232 WriteBufferSize: 1024,
233}
234
235type webLogScheduler struct {
236 ch chan *tangled.CiSubscribePipelineLogs_Event
237}
238
239var _ lexutil.Scheduler[tangled.CiSubscribePipelineLogs_Event] = (*webLogScheduler)(nil)
240
241// AddWork implements [lexutil.Scheduler].
242func (w *webLogScheduler) AddWork(ctx context.Context, _ string, val *tangled.CiSubscribePipelineLogs_Event) error {
243 select {
244 case w.ch <- val:
245 return nil
246 case <-ctx.Done():
247 return ctx.Err()
248 }
249}
250
251// Shutdown implements [lexutil.Scheduler].
252func (w *webLogScheduler) Shutdown() { close(w.ch) }
253
254func retryPipelineTrigger(orig *tangled.CiPipeline) *tangled.CiTriggerPipeline_Input_Trigger {
255 if orig.Trigger != nil && orig.Trigger.CiTrigger_PullRequest != nil {
256 pr := orig.Trigger.CiTrigger_PullRequest
257 sourceSha := pr.SourceSha
258 if sourceSha == "" {
259 sourceSha = orig.Commit
260 }
261 sourceRepo := pr.SourceRepo
262 if sourceRepo == nil {
263 sourceRepo = orig.SourceRepo
264 }
265
266 return &tangled.CiTriggerPipeline_Input_Trigger{
267 CiTrigger_PullRequest: &tangled.CiTrigger_PullRequest{
268 Pull: pr.Pull,
269 SourceBranch: pr.SourceBranch,
270 SourceRepo: sourceRepo,
271 SourceSha: sourceSha,
272 TargetBranch: pr.TargetBranch,
273 },
274 }
275 }
276
277 manual := &tangled.CiTrigger_Manual{
278 Sha: orig.Commit,
279 SourceRepo: orig.SourceRepo,
280 }
281 if orig.Trigger != nil && orig.Trigger.CiTrigger_Manual != nil {
282 origManual := orig.Trigger.CiTrigger_Manual
283 if origManual.Sha != "" {
284 manual.Sha = origManual.Sha
285 }
286 manual.Ref = origManual.Ref
287 manual.Inputs = origManual.Inputs
288 if origManual.SourceRepo != nil {
289 manual.SourceRepo = origManual.SourceRepo
290 }
291 }
292
293 return &tangled.CiTriggerPipeline_Input_Trigger{
294 CiTrigger_Manual: manual,
295 }
296}
297
298func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) {
299 l := p.logger.With("handler", "logs")
300
301 f, err := p.repoResolver.Resolve(r)
302 if err != nil {
303 l.Error("failed to get repo and knot", "err", err)
304 http.Error(w, "bad repo/knot", http.StatusBadRequest)
305 return
306 }
307
308 if f.Spindle == "" {
309 http.Error(w, "invalid repo info", http.StatusBadRequest)
310 return
311 }
312
313 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline"))
314 if err != nil {
315 l.Debug("invalid pipeline id", "id", pipelineId)
316 http.Error(w, "invalid pipeline id", http.StatusBadRequest)
317 return
318 }
319
320 workflowName := chi.URLParam(r, "workflow")
321 if workflowName == "" {
322 l.Debug("empty workflow name")
323 http.Error(w, "invalid workflow name", http.StatusBadRequest)
324 return
325 }
326
327 clientConn, err := upgrader.Upgrade(w, r, nil)
328 if err != nil {
329 l.Error("websocket upgrade failed", "err", err)
330 return
331 }
332 defer clientConn.Close()
333
334 ctx, cancel := context.WithCancel(r.Context())
335 defer cancel()
336
337 spindleUrl, err := hostutil.EnsureHttpScheme(f.Spindle)
338 if err != nil {
339 l.Error("invalid spindle host", "host", f.Spindle, "err", err)
340 return
341 }
342
343 evChan := make(chan *tangled.CiSubscribePipelineLogs_Event, 100)
344 done := make(chan error, 1)
345 sched := &webLogScheduler{ch: evChan}
346 xrpcc := &lexutil.Client{Client: indigoxrpc.Client{Host: spindleUrl}}
347 go func() {
348 done <- tangled.CiSubscribePipelineLogs(ctx, xrpcc, pipelineId.String(), []string{workflowName}, sched)
349 }()
350
351 var lastWriteLk sync.Mutex
352 lastWrite := time.Now()
353
354 // Start a goroutine to ping the client periodically to check if it's still
355 // alive. If the client doesn't respond to a ping within 5 seconds, we'll
356 // close the connection and teardown the consumer.
357 go func() {
358 ticker := time.NewTicker(30 * time.Second)
359 defer ticker.Stop()
360 for {
361 select {
362 case <-ticker.C:
363 lastWriteLk.Lock()
364 lw := lastWrite
365 lastWriteLk.Unlock()
366 if time.Since(lw) < 30*time.Second {
367 continue
368 }
369 if err := clientConn.WriteControl(websocket.PingMessage, nil, time.Now().Add(5*time.Second)); err != nil {
370 l.Warn("failed to ping client", "err", err)
371 cancel()
372 return
373 }
374 case <-ctx.Done():
375 return
376 }
377 }
378 }()
379
380 clientConn.SetPingHandler(func(message string) error {
381 err := clientConn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(60*time.Second))
382 if err == websocket.ErrCloseSent {
383 return nil
384 }
385 return err
386 })
387
388 // Start a goroutine to read messages from the client and discard them.
389 go func() {
390 for {
391 if _, _, err := clientConn.ReadMessage(); err != nil {
392 cancel()
393 return
394 }
395 }
396 }()
397
398 // Main loop: sole writer of data frames to the client.
399 stepStartTimes := make(map[int]time.Time)
400 stepAnsi := make(map[int]*ansiState)
401 var fragment bytes.Buffer
402 for {
403 select {
404 case <-ctx.Done():
405 l.Info("client disconnected")
406 return
407
408 case ev, ok := <-evChan:
409 if !ok {
410 // Stream ended: Shutdown closed the upstream channel.
411 if err := <-done; !isExpectedClose(err) {
412 l.Error("spindle stream error", "err", err)
413 }
414 msg := websocket.FormatCloseMessage(websocket.CloseNormalClosure, "finished")
415 _ = clientConn.WriteMessage(websocket.CloseMessage, msg)
416 return
417 }
418
419 fragment.Reset()
420
421 switch {
422 case ev.Error != nil:
423 l.Error("spindle error frame", "err", ev.Error.Error, "msg", ev.Error.Message)
424 return
425
426 case ev.Control != nil:
427 c := ev.Control
428 step := int(c.Step)
429 switch derefStr(c.Status) {
430 case "start":
431 t := parseRFC3339(c.Time)
432 stepStartTimes[step] = t
433 // "system" steps are injected by the CI runner; collapse them.
434 collapsed := derefStr(c.Kind) == "system"
435 err = p.pages.LogBlock(&fragment, pages.LogBlockParams{
436 Id: step,
437 Name: c.Content,
438 Command: derefStr(c.Command),
439 Collapsed: collapsed,
440 StartTime: t,
441 })
442 case "end":
443 err = p.pages.LogBlockEnd(&fragment, pages.LogBlockEndParams{
444 Id: step,
445 StartTime: stepStartTimes[step],
446 EndTime: parseRFC3339(c.Time),
447 })
448 }
449
450 case ev.Data != nil:
451 d := ev.Data
452 step := int(d.Step)
453 ansi, ok := stepAnsi[step]
454 if !ok {
455 ansi = NewAnsiState()
456 stepAnsi[step] = ansi
457 }
458 err = p.pages.LogLine(&fragment, pages.LogLineParams{
459 Id: step,
460 Content: ansi.Render(d.Content),
461 })
462 }
463 if err != nil {
464 l.Error("failed to render log line", "err", err)
465 return
466 }
467
468 if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil {
469 l.Error("error writing to client", "err", err)
470 return
471 }
472 lastWriteLk.Lock()
473 lastWrite = time.Now()
474 lastWriteLk.Unlock()
475 }
476 }
477}
478
479func (p *Pipelines) CancelWorkflow(w http.ResponseWriter, r *http.Request) {
480 l := p.logger.With("handler", "CancelWorkflow")
481 errorId := "workflow-error"
482
483 f, err := p.repoResolver.Resolve(r)
484 if err != nil {
485 l.Error("failed to get repo and knot", "err", err)
486 p.pages.Notice(w, errorId, "Failed to cancel workflow")
487 return
488 }
489 l = l.With("repo", f.RepoDid)
490
491 if f.Spindle == "" {
492 l.Debug("spindle is empty")
493 p.pages.Notice(w, errorId, "Failed to cancel workflow")
494 return
495 }
496
497 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline"))
498 if err != nil {
499 l.Debug("invalid pipeline id", "id", pipelineId)
500 p.pages.Error404(w)
501 return
502 }
503
504 workflowName := chi.URLParam(r, "workflow")
505 if workflowName == "" {
506 l.Debug("empty workflow name")
507 p.pages.Error404(w)
508 return
509 }
510
511 l = l.With("pipeline", pipelineId, "workflow", workflowName)
512
513 spindleClient, err := p.oauth.SpindleServiceClient(r, f.Spindle, tangled.CiCancelPipelineNSID)
514 if err != nil {
515 l.Error("failed to prepare spindle client", "err", err)
516 p.pages.Notice(w, errorId, "Failed to cancel workflow")
517 return
518 }
519
520 if err := tangled.CiCancelPipeline(
521 r.Context(),
522 spindleClient,
523 &tangled.CiCancelPipeline_Input{
524 Repo: f.RepoDid,
525 Pipeline: pipelineId.String(),
526 Workflows: []string{workflowName},
527 },
528 ); err != nil {
529 l.Error("failed to cancel workflow", "err", err)
530 p.pages.Notice(w, errorId, "Failed to cancel workflow")
531 return
532 }
533 l.Debug("canceled workflow")
534}
535
536// RetryPipeline retries all workflows in a pipeline
537func (p *Pipelines) RetryPipeline(w http.ResponseWriter, r *http.Request) {
538 p.retry(w, r, "")
539}
540
541// RetryWorkflow retries a single workflow in a pipeline
542func (p *Pipelines) RetryWorkflow(w http.ResponseWriter, r *http.Request) {
543 p.retry(w, r, chi.URLParam(r, "workflow"))
544}
545
546// retry triggers a new pipeline run for the original commit, either for all or a single workflow
547func (p *Pipelines) retry(w http.ResponseWriter, r *http.Request, only string) {
548 user := p.oauth.GetMultiAccountUser(r)
549 l := p.logger.With("handler", "retry", "only", only)
550 errorId := "workflow-error"
551
552 // fail logs the error and shows a notice to the user
553 fail := func(msg string, err error) {
554 if err != nil {
555 l.Error(msg, "err", err)
556 p.pages.Notice(w, errorId, fmt.Sprintf("%s: %v", msg, err))
557 } else {
558 l.Error(msg)
559 p.pages.Notice(w, errorId, msg)
560 }
561 }
562
563 f, err := p.repoResolver.Resolve(r)
564 if err != nil {
565 fail("failed to resolve repository", err)
566 return
567 }
568 l = l.With("repo", f.RepoDid)
569
570 if f.Spindle == "" {
571 fail("this repository has no spindle configured", nil)
572 return
573 }
574
575 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline"))
576 if err != nil {
577 l.Debug("invalid pipeline id", "id", pipelineId)
578 p.pages.Error404(w)
579 return
580 }
581 l = l.With("pipeline", pipelineId)
582
583 spindleUrl, err := hostutil.EnsureHttpScheme(f.Spindle)
584 if err != nil {
585 fail("invalid spindle host", err)
586 return
587 }
588
589 // fetch the original pipeline to replay the same commit and workflows
590 queryClient := &indigoxrpc.Client{Host: spindleUrl}
591 orig, err := tangled.CiGetPipeline(r.Context(), queryClient, pipelineId.String())
592 if err != nil {
593 fail("failed to load the original pipeline", err)
594 return
595 }
596 if orig.Commit == "" {
597 fail("cannot retry: the original pipeline has no commit", nil)
598 return
599 }
600
601 // figure out which workflows to run and where to redirect
602 var workflows []string
603 if only != "" {
604 workflows = []string{only}
605 } else {
606 for _, wf := range orig.Workflows {
607 workflows = append(workflows, wf.Name)
608 }
609 }
610 if len(workflows) == 0 {
611 fail("cannot retry: the original pipeline has no workflows", nil)
612 return
613 }
614 redirectWf := workflows[0]
615
616 spindleClient, err := p.oauth.SpindleServiceClient(r, f.Spindle, tangled.CiTriggerPipelineNSID)
617 if err != nil {
618 fail("failed to authorize with spindle", err)
619 return
620 }
621
622 out, err := tangled.CiTriggerPipeline(
623 r.Context(),
624 spindleClient,
625 &tangled.CiTriggerPipeline_Input{
626 Repo: f.RepoDid,
627 Trigger: retryPipelineTrigger(orig),
628 Workflows: workflows,
629 },
630 )
631 if err != nil {
632 fail("spindle rejected the trigger", err)
633 return
634 }
635
636 newAt, err := syntax.ParseATURI(out.Pipeline)
637 if err != nil {
638 fail("pipeline triggered, but the response was malformed", err)
639 return
640 }
641 newId := newAt.RecordKey().String()
642 l = l.With("new", newId)
643 l.Info("pipeline retried")
644
645 repoInfo := p.repoResolver.GetRepoInfo(r, user)
646 dest := fmt.Sprintf("/%s/pipelines/%s/workflow/%s", repoInfo.FullName(), newId, redirectWf)
647
648 if r.Header.Get("HX-Request") == "true" {
649 w.Header().Set("HX-Redirect", dest)
650 w.WriteHeader(http.StatusOK)
651 return
652 }
653 http.Redirect(w, r, dest, http.StatusSeeOther)
654}