This repository has no description
0

Configure Feed

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

core / appview / pipelines / pipelines.go
11 kB 448 lines
1package pipelines 2 3import ( 4 "bytes" 5 "context" 6 "encoding/json" 7 "fmt" 8 "log/slog" 9 "net/http" 10 "strings" 11 "time" 12 13 "tangled.org/core/api/tangled" 14 "tangled.org/core/appview/config" 15 "tangled.org/core/appview/db" 16 "tangled.org/core/appview/middleware" 17 "tangled.org/core/appview/models" 18 "tangled.org/core/appview/oauth" 19 "tangled.org/core/appview/pages" 20 "tangled.org/core/appview/reporesolver" 21 "tangled.org/core/eventconsumer" 22 "tangled.org/core/idresolver" 23 "tangled.org/core/orm" 24 "tangled.org/core/rbac" 25 spindlemodel "tangled.org/core/spindle/models" 26 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 spindlestream *eventconsumer.Consumer 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. 49 With(mw.RepoPermissionMiddleware("repo:owner")). 50 Post("/{pipeline}/workflow/{workflow}/cancel", p.Cancel) 51 52 return r 53} 54 55func New( 56 oauth *oauth.OAuth, 57 repoResolver *reporesolver.RepoResolver, 58 pages *pages.Pages, 59 spindlestream *eventconsumer.Consumer, 60 idResolver *idresolver.Resolver, 61 db *db.DB, 62 config *config.Config, 63 enforcer *rbac.Enforcer, 64 logger *slog.Logger, 65) *Pipelines { 66 return &Pipelines{ 67 oauth: oauth, 68 repoResolver: repoResolver, 69 pages: pages, 70 idResolver: idResolver, 71 config: config, 72 spindlestream: spindlestream, 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 filters := []orm.Filter{ 91 orm.FilterEq("p.repo_did", f.RepoDid), 92 } 93 switch filterKind { 94 case "push": 95 filters = append(filters, orm.FilterEq("t.kind", "push")) 96 case "pull_request": 97 filters = append(filters, orm.FilterEq("t.kind", "pull_request")) 98 default: 99 // no filters otherwise, default to "all" 100 filterKind = "all" 101 } 102 103 ps, err := db.GetPipelineStatuses( 104 p.db, 105 30, 106 filters..., 107 ) 108 if err != nil { 109 l.Error("failed to query db", "err", err) 110 return 111 } 112 113 total, err := db.GetTotalPipelineStatuses(p.db, filters...) 114 if err != nil { 115 l.Error("failed to query db", "err", err) 116 return 117 } 118 119 p.pages.Pipelines(w, pages.PipelinesParams{ 120 LoggedInUser: user, 121 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 122 Pipelines: ps, 123 FilterKind: filterKind, 124 Total: total, 125 }) 126} 127 128func (p *Pipelines) Workflow(w http.ResponseWriter, r *http.Request) { 129 user := p.oauth.GetMultiAccountUser(r) 130 l := p.logger.With("handler", "Workflow") 131 132 f, err := p.repoResolver.Resolve(r) 133 if err != nil { 134 l.Error("failed to get repo and knot", "err", err) 135 return 136 } 137 138 pipelineId := chi.URLParam(r, "pipeline") 139 if pipelineId == "" { 140 l.Error("empty pipeline ID") 141 return 142 } 143 144 workflow := chi.URLParam(r, "workflow") 145 if workflow == "" { 146 l.Error("empty workflow name") 147 return 148 } 149 150 ps, err := db.GetPipelineStatuses( 151 p.db, 152 1, 153 orm.FilterEq("p.repo_did", f.RepoDid), 154 orm.FilterEq("p.id", pipelineId), 155 ) 156 if err != nil { 157 l.Error("failed to query db", "err", err) 158 return 159 } 160 161 if len(ps) != 1 { 162 l.Error("invalid number of pipelines", "len", len(ps)) 163 return 164 } 165 166 singlePipeline := ps[0] 167 168 p.pages.Workflow(w, pages.WorkflowParams{ 169 LoggedInUser: user, 170 RepoInfo: p.repoResolver.GetRepoInfo(r, user), 171 Pipeline: singlePipeline, 172 Workflow: workflow, 173 }) 174} 175 176var upgrader = websocket.Upgrader{ 177 ReadBufferSize: 1024, 178 WriteBufferSize: 1024, 179} 180 181func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) { 182 l := p.logger.With("handler", "logs") 183 184 clientConn, err := upgrader.Upgrade(w, r, nil) 185 if err != nil { 186 l.Error("websocket upgrade failed", "err", err) 187 return 188 } 189 defer func() { 190 _ = clientConn.WriteControl( 191 websocket.CloseMessage, 192 websocket.FormatCloseMessage(websocket.CloseNormalClosure, "log stream complete"), 193 time.Now().Add(time.Second), 194 ) 195 clientConn.Close() 196 }() 197 198 ctx, cancel := context.WithCancel(r.Context()) 199 defer cancel() 200 201 f, err := p.repoResolver.Resolve(r) 202 if err != nil { 203 l.Error("failed to get repo and knot", "err", err) 204 http.Error(w, "bad repo/knot", http.StatusBadRequest) 205 return 206 } 207 208 pipelineId := chi.URLParam(r, "pipeline") 209 workflow := chi.URLParam(r, "workflow") 210 if pipelineId == "" || workflow == "" { 211 http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest) 212 return 213 } 214 215 ps, err := db.GetPipelineStatuses( 216 p.db, 217 1, 218 orm.FilterEq("p.repo_did", f.RepoDid), 219 orm.FilterEq("p.id", pipelineId), 220 ) 221 if err != nil || len(ps) != 1 { 222 l.Error("pipeline query failed", "err", err, "count", len(ps)) 223 http.Error(w, "pipeline not found", http.StatusNotFound) 224 return 225 } 226 227 singlePipeline := ps[0] 228 spindle := f.Spindle 229 knot := f.Knot 230 rkey := singlePipeline.Rkey 231 232 if spindle == "" || knot == "" || rkey == "" { 233 http.Error(w, "invalid repo info", http.StatusBadRequest) 234 return 235 } 236 237 scheme := "wss" 238 if p.config.Core.Dev { 239 scheme = "ws" 240 } 241 242 url := scheme + "://" + strings.Join([]string{spindle, "logs", knot, rkey, workflow}, "/") 243 l = l.With("url", url) 244 l.Info("logs endpoint hit") 245 246 spindleConn, _, err := websocket.DefaultDialer.Dial(url, nil) 247 if err != nil { 248 l.Error("websocket dial failed", "err", err) 249 http.Error(w, "failed to connect to log stream", http.StatusBadGateway) 250 return 251 } 252 defer spindleConn.Close() 253 254 // create a channel for incoming messages 255 evChan := make(chan logEvent, 100) 256 // start a goroutine to read from spindle 257 go readLogs(spindleConn, evChan) 258 259 stepStartTimes := make(map[int]time.Time) 260 var fragment bytes.Buffer 261 for { 262 select { 263 case <-ctx.Done(): 264 l.Info("client disconnected") 265 return 266 267 case ev, ok := <-evChan: 268 if !ok { 269 continue 270 } 271 272 if ev.err != nil && ev.isCloseError() { 273 l.Debug("graceful shutdown, tail complete", "err", err) 274 return 275 } 276 if ev.err != nil { 277 l.Error("error reading from spindle", "err", err) 278 return 279 } 280 281 var logLine spindlemodel.LogLine 282 if err = json.Unmarshal(ev.msg, &logLine); err != nil { 283 l.Error("failed to parse logline", "err", err) 284 continue 285 } 286 287 fragment.Reset() 288 289 switch logLine.Kind { 290 case spindlemodel.LogKindControl: 291 switch logLine.StepStatus { 292 case spindlemodel.StepStatusStart: 293 stepStartTimes[logLine.StepId] = logLine.Time 294 collapsed := false 295 if logLine.StepKind == spindlemodel.StepKindSystem { 296 collapsed = true 297 } 298 err = p.pages.LogBlock(&fragment, pages.LogBlockParams{ 299 Id: logLine.StepId, 300 Name: logLine.Content, 301 Command: logLine.StepCommand, 302 Collapsed: collapsed, 303 StartTime: logLine.Time, 304 }) 305 case spindlemodel.StepStatusEnd: 306 startTime := stepStartTimes[logLine.StepId] 307 endTime := logLine.Time 308 err = p.pages.LogBlockEnd(&fragment, pages.LogBlockEndParams{ 309 Id: logLine.StepId, 310 StartTime: startTime, 311 EndTime: endTime, 312 }) 313 } 314 315 case spindlemodel.LogKindData: 316 // data messages simply insert new log lines into current step 317 err = p.pages.LogLine(&fragment, pages.LogLineParams{ 318 Id: logLine.StepId, 319 Content: logLine.Content, 320 }) 321 } 322 if err != nil { 323 l.Error("failed to render log line", "err", err) 324 return 325 } 326 327 if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil { 328 l.Error("error writing to client", "err", err) 329 return 330 } 331 332 case <-time.After(30 * time.Second): 333 l.Debug("sent keepalive") 334 if err = clientConn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil { 335 l.Error("failed to write control", "err", err) 336 return 337 } 338 } 339 } 340} 341 342func (p *Pipelines) Cancel(w http.ResponseWriter, r *http.Request) { 343 l := p.logger.With("handler", "Cancel") 344 345 var ( 346 pipelineId = chi.URLParam(r, "pipeline") 347 workflow = chi.URLParam(r, "workflow") 348 ) 349 if pipelineId == "" || workflow == "" { 350 http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest) 351 return 352 } 353 354 f, err := p.repoResolver.Resolve(r) 355 if err != nil { 356 l.Error("failed to get repo and knot", "err", err) 357 http.Error(w, "bad repo/knot", http.StatusBadRequest) 358 return 359 } 360 361 pipeline, err := func() (models.Pipeline, error) { 362 ps, err := db.GetPipelineStatuses( 363 p.db, 364 1, 365 orm.FilterEq("p.repo_did", f.RepoDid), 366 orm.FilterEq("p.id", pipelineId), 367 ) 368 if err != nil { 369 return models.Pipeline{}, err 370 } 371 if len(ps) != 1 { 372 return models.Pipeline{}, fmt.Errorf("wrong pipeline count %d", len(ps)) 373 } 374 return ps[0], nil 375 }() 376 if err != nil { 377 l.Error("pipeline query failed", "err", err) 378 http.Error(w, "pipeline not found", http.StatusNotFound) 379 } 380 var ( 381 spindle = f.Spindle 382 knot = f.Knot 383 rkey = pipeline.Rkey 384 ) 385 386 if spindle == "" || knot == "" || rkey == "" { 387 http.Error(w, "invalid repo info", http.StatusBadRequest) 388 return 389 } 390 391 spindleClient, err := p.oauth.ServiceClient( 392 r, 393 oauth.WithService(f.Spindle), 394 oauth.WithLxm(tangled.PipelineCancelPipelineNSID), 395 oauth.WithDev(p.config.Core.Dev), 396 oauth.WithTimeout(time.Second*30), // workflow cleanup usually takes time 397 ) 398 399 err = tangled.PipelineCancelPipeline( 400 r.Context(), 401 spindleClient, 402 &tangled.PipelineCancelPipeline_Input{ 403 Repo: string(f.RepoAt()), 404 Pipeline: pipeline.AtUri().String(), 405 Workflow: workflow, 406 }, 407 ) 408 errorId := "workflow-error" 409 if err != nil { 410 l.Error("failed to cancel workflow", "err", err) 411 p.pages.Notice(w, errorId, "Failed to cancel workflow") 412 return 413 } 414 l.Debug("canceled pipeline", "uri", pipeline.AtUri()) 415} 416 417// either a message or an error 418type logEvent struct { 419 msg []byte 420 err error 421} 422 423func (ev *logEvent) isCloseError() bool { 424 return websocket.IsCloseError( 425 ev.err, 426 websocket.CloseNormalClosure, 427 websocket.CloseGoingAway, 428 websocket.CloseAbnormalClosure, 429 ) 430} 431 432// read logs from spindle and pass through to chan 433func readLogs(conn *websocket.Conn, ch chan logEvent) { 434 defer close(ch) 435 436 for { 437 if conn == nil { 438 return 439 } 440 441 _, msg, err := conn.ReadMessage() 442 if err != nil { 443 ch <- logEvent{err: err} 444 return 445 } 446 ch <- logEvent{msg: msg} 447 } 448}