This repository has no description
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_owner", f.Did),
154 orm.FilterEq("p.repo_name", f.Rkey),
155 orm.FilterEq("p.knot", f.Knot),
156 orm.FilterEq("p.id", pipelineId),
157 )
158 if err != nil {
159 l.Error("failed to query db", "err", err)
160 return
161 }
162
163 if len(ps) != 1 {
164 l.Error("invalid number of pipelines", "len", len(ps))
165 return
166 }
167
168 singlePipeline := ps[0]
169
170 p.pages.Workflow(w, pages.WorkflowParams{
171 LoggedInUser: user,
172 RepoInfo: p.repoResolver.GetRepoInfo(r, user),
173 Pipeline: singlePipeline,
174 Workflow: workflow,
175 })
176}
177
178var upgrader = websocket.Upgrader{
179 ReadBufferSize: 1024,
180 WriteBufferSize: 1024,
181}
182
183func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) {
184 l := p.logger.With("handler", "logs")
185
186 clientConn, err := upgrader.Upgrade(w, r, nil)
187 if err != nil {
188 l.Error("websocket upgrade failed", "err", err)
189 return
190 }
191 defer func() {
192 _ = clientConn.WriteControl(
193 websocket.CloseMessage,
194 websocket.FormatCloseMessage(websocket.CloseNormalClosure, "log stream complete"),
195 time.Now().Add(time.Second),
196 )
197 clientConn.Close()
198 }()
199
200 ctx, cancel := context.WithCancel(r.Context())
201 defer cancel()
202
203 f, err := p.repoResolver.Resolve(r)
204 if err != nil {
205 l.Error("failed to get repo and knot", "err", err)
206 http.Error(w, "bad repo/knot", http.StatusBadRequest)
207 return
208 }
209
210 pipelineId := chi.URLParam(r, "pipeline")
211 workflow := chi.URLParam(r, "workflow")
212 if pipelineId == "" || workflow == "" {
213 http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest)
214 return
215 }
216
217 ps, err := db.GetPipelineStatuses(
218 p.db,
219 1,
220 orm.FilterEq("p.repo_owner", f.Did),
221 orm.FilterEq("p.repo_name", f.Rkey),
222 orm.FilterEq("p.knot", f.Knot),
223 orm.FilterEq("p.id", pipelineId),
224 )
225 if err != nil || len(ps) != 1 {
226 l.Error("pipeline query failed", "err", err, "count", len(ps))
227 http.Error(w, "pipeline not found", http.StatusNotFound)
228 return
229 }
230
231 singlePipeline := ps[0]
232 spindle := f.Spindle
233 knot := f.Knot
234 rkey := singlePipeline.Rkey
235
236 if spindle == "" || knot == "" || rkey == "" {
237 http.Error(w, "invalid repo info", http.StatusBadRequest)
238 return
239 }
240
241 scheme := "wss"
242 if p.config.Core.Dev {
243 scheme = "ws"
244 }
245
246 url := scheme + "://" + strings.Join([]string{spindle, "logs", knot, rkey, workflow}, "/")
247 l = l.With("url", url)
248 l.Info("logs endpoint hit")
249
250 spindleConn, _, err := websocket.DefaultDialer.Dial(url, nil)
251 if err != nil {
252 l.Error("websocket dial failed", "err", err)
253 http.Error(w, "failed to connect to log stream", http.StatusBadGateway)
254 return
255 }
256 defer spindleConn.Close()
257
258 // create a channel for incoming messages
259 evChan := make(chan logEvent, 100)
260 // start a goroutine to read from spindle
261 go readLogs(spindleConn, evChan)
262
263 stepStartTimes := make(map[int]time.Time)
264 var fragment bytes.Buffer
265 for {
266 select {
267 case <-ctx.Done():
268 l.Info("client disconnected")
269 return
270
271 case ev, ok := <-evChan:
272 if !ok {
273 continue
274 }
275
276 if ev.err != nil && ev.isCloseError() {
277 l.Debug("graceful shutdown, tail complete", "err", err)
278 return
279 }
280 if ev.err != nil {
281 l.Error("error reading from spindle", "err", err)
282 return
283 }
284
285 var logLine spindlemodel.LogLine
286 if err = json.Unmarshal(ev.msg, &logLine); err != nil {
287 l.Error("failed to parse logline", "err", err)
288 continue
289 }
290
291 fragment.Reset()
292
293 switch logLine.Kind {
294 case spindlemodel.LogKindControl:
295 switch logLine.StepStatus {
296 case spindlemodel.StepStatusStart:
297 stepStartTimes[logLine.StepId] = logLine.Time
298 collapsed := false
299 if logLine.StepKind == spindlemodel.StepKindSystem {
300 collapsed = true
301 }
302 err = p.pages.LogBlock(&fragment, pages.LogBlockParams{
303 Id: logLine.StepId,
304 Name: logLine.Content,
305 Command: logLine.StepCommand,
306 Collapsed: collapsed,
307 StartTime: logLine.Time,
308 })
309 case spindlemodel.StepStatusEnd:
310 startTime := stepStartTimes[logLine.StepId]
311 endTime := logLine.Time
312 err = p.pages.LogBlockEnd(&fragment, pages.LogBlockEndParams{
313 Id: logLine.StepId,
314 StartTime: startTime,
315 EndTime: endTime,
316 })
317 }
318
319 case spindlemodel.LogKindData:
320 // data messages simply insert new log lines into current step
321 err = p.pages.LogLine(&fragment, pages.LogLineParams{
322 Id: logLine.StepId,
323 Content: logLine.Content,
324 })
325 }
326 if err != nil {
327 l.Error("failed to render log line", "err", err)
328 return
329 }
330
331 if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil {
332 l.Error("error writing to client", "err", err)
333 return
334 }
335
336 case <-time.After(30 * time.Second):
337 l.Debug("sent keepalive")
338 if err = clientConn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil {
339 l.Error("failed to write control", "err", err)
340 return
341 }
342 }
343 }
344}
345
346func (p *Pipelines) Cancel(w http.ResponseWriter, r *http.Request) {
347 l := p.logger.With("handler", "Cancel")
348
349 var (
350 pipelineId = chi.URLParam(r, "pipeline")
351 workflow = chi.URLParam(r, "workflow")
352 )
353 if pipelineId == "" || workflow == "" {
354 http.Error(w, "missing pipeline ID or workflow", http.StatusBadRequest)
355 return
356 }
357
358 f, err := p.repoResolver.Resolve(r)
359 if err != nil {
360 l.Error("failed to get repo and knot", "err", err)
361 http.Error(w, "bad repo/knot", http.StatusBadRequest)
362 return
363 }
364
365 pipeline, err := func() (models.Pipeline, error) {
366 ps, err := db.GetPipelineStatuses(
367 p.db,
368 1,
369 orm.FilterEq("p.repo_owner", f.Did),
370 orm.FilterEq("p.repo_name", f.Rkey),
371 orm.FilterEq("p.knot", f.Knot),
372 orm.FilterEq("p.id", pipelineId),
373 )
374 if err != nil {
375 return models.Pipeline{}, err
376 }
377 if len(ps) != 1 {
378 return models.Pipeline{}, fmt.Errorf("wrong pipeline count %d", len(ps))
379 }
380 return ps[0], nil
381 }()
382 if err != nil {
383 l.Error("pipeline query failed", "err", err)
384 http.Error(w, "pipeline not found", http.StatusNotFound)
385 }
386 var (
387 spindle = f.Spindle
388 knot = f.Knot
389 rkey = pipeline.Rkey
390 )
391
392 if spindle == "" || knot == "" || rkey == "" {
393 http.Error(w, "invalid repo info", http.StatusBadRequest)
394 return
395 }
396
397 spindleClient, err := p.oauth.ServiceClient(
398 r,
399 oauth.WithService(f.Spindle),
400 oauth.WithLxm(tangled.PipelineCancelPipelineNSID),
401 oauth.WithDev(p.config.Core.Dev),
402 oauth.WithTimeout(time.Second*30), // workflow cleanup usually takes time
403 )
404
405 err = tangled.PipelineCancelPipeline(
406 r.Context(),
407 spindleClient,
408 &tangled.PipelineCancelPipeline_Input{
409 Repo: string(f.RepoAt()),
410 Pipeline: pipeline.AtUri().String(),
411 Workflow: workflow,
412 },
413 )
414 errorId := "workflow-error"
415 if err != nil {
416 l.Error("failed to cancel workflow", "err", err)
417 p.pages.Notice(w, errorId, "Failed to cancel workflow")
418 return
419 }
420 l.Debug("canceled pipeline", "uri", pipeline.AtUri())
421}
422
423// either a message or an error
424type logEvent struct {
425 msg []byte
426 err error
427}
428
429func (ev *logEvent) isCloseError() bool {
430 return websocket.IsCloseError(
431 ev.err,
432 websocket.CloseNormalClosure,
433 websocket.CloseGoingAway,
434 websocket.CloseAbnormalClosure,
435 )
436}
437
438// read logs from spindle and pass through to chan
439func readLogs(conn *websocket.Conn, ch chan logEvent) {
440 defer close(ch)
441
442 for {
443 if conn == nil {
444 return
445 }
446
447 _, msg, err := conn.ReadMessage()
448 if err != nil {
449 ch <- logEvent{err: err}
450 return
451 }
452 ch <- logEvent{msg: msg}
453 }
454}