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