Monorepo for Tangled
tangled.org
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}