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_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}