Monorepo for Tangled tangled.org
1

Configure Feed

Select the types of activity you want to include in your feed.

core / appview / pipelines / pipelines.go
17 kB 654 lines
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}