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
13 kB 493 lines
1package pipelines 2 3import ( 4 "bytes" 5 "context" 6 "log/slog" 7 "net/http" 8 "sync" 9 "time" 10 11 "tangled.org/core/api/tangled" 12 "tangled.org/core/appview/config" 13 "tangled.org/core/appview/db" 14 "tangled.org/core/appview/middleware" 15 "tangled.org/core/appview/oauth" 16 "tangled.org/core/appview/pages" 17 "tangled.org/core/appview/reporesolver" 18 "tangled.org/core/hostutil" 19 "tangled.org/core/idresolver" 20 "tangled.org/core/lexutil" 21 "tangled.org/core/orm" 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. 48 With(mw.RepoPermissionMiddleware("repo:owner")). 49 Post("/{pipeline}/workflow/{workflow}/cancel", p.CancelWorkflow) 50 51 return r 52} 53 54func New( 55 oauth *oauth.OAuth, 56 repoResolver *reporesolver.RepoResolver, 57 pages *pages.Pages, 58 idResolver *idresolver.Resolver, 59 db *db.DB, 60 config *config.Config, 61 enforcer *rbac.Enforcer, 62 logger *slog.Logger, 63) *Pipelines { 64 return &Pipelines{ 65 oauth: oauth, 66 repoResolver: repoResolver, 67 pages: pages, 68 idResolver: idResolver, 69 config: config, 70 db: db, 71 enforcer: enforcer, 72 logger: logger, 73 } 74} 75 76func (p *Pipelines) Index(w http.ResponseWriter, r *http.Request) { 77 user := p.oauth.GetMultiAccountUser(r) 78 l := p.logger.With("handler", "Index") 79 80 f, err := p.repoResolver.Resolve(r) 81 if err != nil { 82 l.Error("failed to get repo and knot", "err", err) 83 return 84 } 85 86 filterKind := r.URL.Query().Get("trigger") 87 filters := []orm.Filter{ 88 orm.FilterEq("p.repo_did", f.RepoDid), 89 } 90 switch filterKind { 91 case "push": 92 filters = append(filters, orm.FilterEq("t.kind", "push")) 93 case "pull_request": 94 filters = append(filters, orm.FilterEq("t.kind", "pull_request")) 95 default: 96 // no filters otherwise, default to "all" 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 // sh.tangled.ci.queryPipelines(repo, kind, limit=30) 125 xrpcc := indigoxrpc.Client{Host: spindleUrl} 126 out, err := tangled.CiQueryPipelines(r.Context(), &xrpcc, nil, "", 30, f.RepoDid) 127 if err != nil { 128 l.Error("failed to fetch pipelines", "err", err) 129 p.pages.Pipelines(w, pages.PipelinesParams{ 130 BaseParams: pages.BaseParamsFromContext(r.Context()), 131 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 132 Pipelines: nil, 133 FilterKind: filterKind, 134 Total: 0, 135 }) 136 return 137 } 138 139 var pipelines []types.Pipeline 140 for _, pipeline := range out.Pipelines { 141 pipelines = append(pipelines, types.Pipeline{CiDefs_Pipeline: pipeline}) 142 } 143 144 p.pages.Pipelines(w, pages.PipelinesParams{ 145 BaseParams: pages.BaseParamsFromContext(r.Context()), 146 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 147 Pipelines: pipelines, 148 FilterKind: filterKind, 149 Total: out.Total, 150 }) 151} 152 153func (p *Pipelines) Workflow(w http.ResponseWriter, r *http.Request) { 154 user := p.oauth.GetMultiAccountUser(r) 155 l := p.logger.With("handler", "Workflow") 156 157 f, err := p.repoResolver.Resolve(r) 158 if err != nil { 159 l.Error("failed to get repo and knot", "err", err) 160 p.pages.Error404(w) 161 return 162 } 163 164 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) 165 if err != nil { 166 l.Debug("invalid pipeline id", "id", pipelineId) 167 p.pages.Error404(w) 168 return 169 } 170 171 workflowName := chi.URLParam(r, "workflow") 172 if workflowName == "" { 173 l.Debug("empty workflow name") 174 p.pages.Error404(w) 175 return 176 } 177 178 l = l.With("pipeline", pipelineId, "workflow", workflowName) 179 180 // TODO: change url path to: 181 // /{owner}/{slug}/pipelines/{spindle-did}/{pipeline-id}/workflow/{workflow-id} 182 183 if f.Spindle == "" { 184 p.pages.Error404(w) 185 return 186 } 187 188 spindleUrl, err := hostutil.EnsureHttpScheme(f.Spindle) 189 if err != nil { 190 l.Error("invalid spindle host", "host", f.Spindle, "err", err) 191 p.pages.Error404(w) 192 return 193 } 194 195 xrpcc := &indigoxrpc.Client{Host: spindleUrl} 196 out, err := tangled.CiGetPipeline(r.Context(), xrpcc, pipelineId.String()) 197 if err != nil { 198 // TODO(boltless): change behavior based on error 199 l.Debug("failed to get pipeline", "err", err) 200 p.pages.Error404(w) 201 return 202 } 203 204 // ensure workflow exists 205 exist := false 206 for _, workflow := range out.Workflows { 207 if workflow.Name == workflowName { 208 exist = true 209 break 210 } 211 } 212 if !exist { 213 l.Debug("workflow doesn't exist in pipeline") 214 p.pages.Error404(w) 215 return 216 } 217 218 p.pages.Workflow(w, pages.WorkflowParams{ 219 BaseParams: pages.BaseParamsFromContext(r.Context()), 220 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 221 Pipeline: types.Pipeline{CiDefs_Pipeline: out}, 222 Workflow: workflowName, 223 }) 224} 225 226var upgrader = websocket.Upgrader{ 227 ReadBufferSize: 1024, 228 WriteBufferSize: 1024, 229} 230 231type webLogScheduler struct { 232 ch chan *tangled.CiPipelineSubscribeLogs_Event 233} 234 235var _ lexutil.Scheduler[tangled.CiPipelineSubscribeLogs_Event] = (*webLogScheduler)(nil) 236 237// AddWork implements [lexutil.Scheduler]. 238func (w *webLogScheduler) AddWork(ctx context.Context, _ string, val *tangled.CiPipelineSubscribeLogs_Event) error { 239 select { 240 case w.ch <- val: 241 return nil 242 case <-ctx.Done(): 243 return ctx.Err() 244 } 245} 246 247// Shutdown implements [lexutil.Scheduler]. 248func (w *webLogScheduler) Shutdown() { close(w.ch) } 249 250func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { 251 l := p.logger.With("handler", "logs") 252 253 f, err := p.repoResolver.Resolve(r) 254 if err != nil { 255 l.Error("failed to get repo and knot", "err", err) 256 http.Error(w, "bad repo/knot", http.StatusBadRequest) 257 return 258 } 259 260 if f.Spindle == "" { 261 http.Error(w, "invalid repo info", http.StatusBadRequest) 262 return 263 } 264 265 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) 266 if err != nil { 267 l.Debug("invalid pipeline id", "id", pipelineId) 268 http.Error(w, "invalid pipeline id", http.StatusBadRequest) 269 return 270 } 271 272 workflowName := chi.URLParam(r, "workflow") 273 if workflowName == "" { 274 l.Debug("empty workflow name") 275 http.Error(w, "invalid workflow name", http.StatusBadRequest) 276 return 277 } 278 279 clientConn, err := upgrader.Upgrade(w, r, nil) 280 if err != nil { 281 l.Error("websocket upgrade failed", "err", err) 282 return 283 } 284 defer clientConn.Close() 285 286 ctx, cancel := context.WithCancel(r.Context()) 287 defer cancel() 288 289 spindleUrl, err := hostutil.EnsureHttpScheme(f.Spindle) 290 if err != nil { 291 l.Error("invalid spindle host", "host", f.Spindle, "err", err) 292 return 293 } 294 295 evChan := make(chan *tangled.CiPipelineSubscribeLogs_Event, 100) 296 done := make(chan error, 1) 297 sched := &webLogScheduler{ch: evChan} 298 xrpcc := &lexutil.Client{Client: indigoxrpc.Client{Host: spindleUrl}} 299 go func() { 300 done <- tangled.CiPipelineSubscribeLogs(ctx, xrpcc, pipelineId.String(), []string{workflowName}, sched) 301 }() 302 303 var lastWriteLk sync.Mutex 304 lastWrite := time.Now() 305 306 // Start a goroutine to ping the client periodically to check if it's still 307 // alive. If the client doesn't respond to a ping within 5 seconds, we'll 308 // close the connection and teardown the consumer. 309 go func() { 310 ticker := time.NewTicker(30 * time.Second) 311 defer ticker.Stop() 312 for { 313 select { 314 case <-ticker.C: 315 lastWriteLk.Lock() 316 lw := lastWrite 317 lastWriteLk.Unlock() 318 if time.Since(lw) < 30*time.Second { 319 continue 320 } 321 if err := clientConn.WriteControl(websocket.PingMessage, nil, time.Now().Add(5*time.Second)); err != nil { 322 l.Warn("failed to ping client", "err", err) 323 cancel() 324 return 325 } 326 case <-ctx.Done(): 327 return 328 } 329 } 330 }() 331 332 clientConn.SetPingHandler(func(message string) error { 333 err := clientConn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(60*time.Second)) 334 if err == websocket.ErrCloseSent { 335 return nil 336 } 337 return err 338 }) 339 340 // Start a goroutine to read messages from the client and discard them. 341 go func() { 342 for { 343 if _, _, err := clientConn.ReadMessage(); err != nil { 344 cancel() 345 return 346 } 347 } 348 }() 349 350 // Main loop: sole writer of data frames to the client. 351 stepStartTimes := make(map[int]time.Time) 352 stepAnsi := make(map[int]*ansiState) 353 var fragment bytes.Buffer 354 for { 355 select { 356 case <-ctx.Done(): 357 l.Info("client disconnected") 358 return 359 360 case ev, ok := <-evChan: 361 if !ok { 362 // Stream ended: Shutdown closed the upstream channel. 363 if err := <-done; !isExpectedClose(err) { 364 l.Error("spindle stream error", "err", err) 365 } 366 msg := websocket.FormatCloseMessage(websocket.CloseNormalClosure, "finished") 367 _ = clientConn.WriteMessage(websocket.CloseMessage, msg) 368 return 369 } 370 371 fragment.Reset() 372 373 switch { 374 case ev.Error != nil: 375 l.Error("spindle error frame", "err", ev.Error.Error, "msg", ev.Error.Message) 376 return 377 378 case ev.Control != nil: 379 c := ev.Control 380 step := int(c.Step) 381 switch derefStr(c.Status) { 382 case "start": 383 t := parseRFC3339(c.Time) 384 stepStartTimes[step] = t 385 // "system" steps are injected by the CI runner; collapse them. 386 collapsed := derefStr(c.Kind) == "system" 387 err = p.pages.LogBlock(&fragment, pages.LogBlockParams{ 388 Id: step, 389 Name: c.Content, 390 Command: derefStr(c.Command), 391 Collapsed: collapsed, 392 StartTime: t, 393 }) 394 case "end": 395 err = p.pages.LogBlockEnd(&fragment, pages.LogBlockEndParams{ 396 Id: step, 397 StartTime: stepStartTimes[step], 398 EndTime: parseRFC3339(c.Time), 399 }) 400 } 401 402 case ev.Data != nil: 403 d := ev.Data 404 step := int(d.Step) 405 ansi, ok := stepAnsi[step] 406 if !ok { 407 ansi = NewAnsiState() 408 stepAnsi[step] = ansi 409 } 410 err = p.pages.LogLine(&fragment, pages.LogLineParams{ 411 Id: step, 412 Content: ansi.Render(d.Content), 413 }) 414 } 415 if err != nil { 416 l.Error("failed to render log line", "err", err) 417 return 418 } 419 420 if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil { 421 l.Error("error writing to client", "err", err) 422 return 423 } 424 lastWriteLk.Lock() 425 lastWrite = time.Now() 426 lastWriteLk.Unlock() 427 } 428 } 429} 430 431func (p *Pipelines) CancelWorkflow(w http.ResponseWriter, r *http.Request) { 432 l := p.logger.With("handler", "CancelWorkflow") 433 errorId := "workflow-error" 434 435 f, err := p.repoResolver.Resolve(r) 436 if err != nil { 437 l.Error("failed to get repo and knot", "err", err) 438 p.pages.Notice(w, errorId, "Failed to cancel workflow") 439 return 440 } 441 l = l.With("repo", f.RepoDid) 442 443 if f.Spindle == "" { 444 l.Debug("spindle is empty") 445 p.pages.Notice(w, errorId, "Failed to cancel workflow") 446 return 447 } 448 449 pipelineId, err := syntax.ParseTID(chi.URLParam(r, "pipeline")) 450 if err != nil { 451 l.Debug("invalid pipeline id", "id", pipelineId) 452 p.pages.Error404(w) 453 return 454 } 455 456 workflowName := chi.URLParam(r, "workflow") 457 if workflowName == "" { 458 l.Debug("empty workflow name") 459 p.pages.Error404(w) 460 return 461 } 462 463 l = l.With("pipeline", pipelineId, "workflow", workflowName) 464 465 hostname, noTLS, err := hostutil.ParseHostname(f.Spindle) 466 if err != nil { 467 http.Error(w, "invalid spindle hostname", http.StatusBadRequest) 468 return 469 } 470 471 spindleClient, err := p.oauth.ServiceClient( 472 r, 473 oauth.WithService(hostname), 474 oauth.WithLxm(tangled.CiPipelineCancelPipelineNSID), 475 oauth.WithDev(noTLS), 476 oauth.WithTimeout(time.Second*30), // workflow cleanup usually takes time 477 ) 478 479 if err := tangled.CiPipelineCancelPipeline( 480 r.Context(), 481 spindleClient, 482 &tangled.CiPipelineCancelPipeline_Input{ 483 Repo: string(f.RepoAt()), 484 Pipeline: pipelineId.String(), 485 Workflows: []string{workflowName}, 486 }, 487 ); err != nil { 488 l.Error("failed to cancel workflow", "err", err) 489 p.pages.Notice(w, errorId, "Failed to cancel workflow") 490 return 491 } 492 l.Debug("canceled workflow") 493}