Monorepo for Tangled
tangled.org
1package spindle
2
3import (
4 "context"
5 "database/sql"
6 _ "embed"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "log/slog"
11 "maps"
12 "net/http"
13 "path/filepath"
14 "sync"
15 "time"
16
17 "github.com/bluesky-social/indigo/atproto/syntax"
18 indigoxrpc "github.com/bluesky-social/indigo/xrpc"
19 "github.com/go-chi/chi/v5"
20 "github.com/go-git/go-git/v5/plumbing/object"
21 "github.com/hashicorp/go-version"
22 "tangled.org/core/api/tangled"
23 "tangled.org/core/eventconsumer"
24 "tangled.org/core/eventconsumer/cursor"
25 "tangled.org/core/eventstream"
26 "tangled.org/core/idresolver"
27 "tangled.org/core/jetstream"
28 knotdb "tangled.org/core/knotserver/db"
29 kgit "tangled.org/core/knotserver/git"
30 "tangled.org/core/log"
31 "tangled.org/core/notifier"
32 "tangled.org/core/rbac"
33 "tangled.org/core/repoident"
34 "tangled.org/core/repoverify"
35 "tangled.org/core/spindle/config"
36 "tangled.org/core/spindle/db"
37 "tangled.org/core/spindle/engine"
38 "tangled.org/core/spindle/engines/dummy"
39 "tangled.org/core/spindle/engines/microvm"
40 "tangled.org/core/spindle/engines/nixery"
41 "tangled.org/core/spindle/git"
42 "tangled.org/core/spindle/models"
43 "tangled.org/core/spindle/secrets"
44 "tangled.org/core/spindle/xrpc"
45 "tangled.org/core/tid"
46 "tangled.org/core/workflow"
47 "tangled.org/core/xrpc/serviceauth"
48)
49
50//go:embed motd
51var defaultMotd []byte
52
53const (
54 rbacDomain = "thisserver"
55)
56
57type Spindle struct {
58 jc *jetstream.JetstreamClient
59 tap *Tap
60 embedTap *embeddedTap
61 db *db.DB
62 e *rbac.Enforcer
63 l *slog.Logger
64 n *notifier.Notifier
65 engs map[string]models.Engine
66 cfg *config.Config
67 ks *eventconsumer.Consumer
68 res *idresolver.Resolver
69 verify repoverify.Verifier
70 vault secrets.Manager
71 motd []byte
72 motdMu sync.RWMutex
73 rootCtx context.Context
74 jobWake chan struct{}
75}
76
77// New creates a new Spindle server with the provided configuration and engines.
78func New(ctx context.Context, cfg *config.Config, d *db.DB, engines map[string]models.Engine) (*Spindle, error) {
79 logger := log.FromContext(ctx)
80
81 e, err := rbac.NewEnforcer(cfg.Server.DBPath)
82 if err != nil {
83 return nil, fmt.Errorf("failed to setup rbac enforcer: %w", err)
84 }
85 e.E.EnableAutoSave(true)
86
87 n := notifier.New()
88
89 var vault secrets.Manager
90 switch cfg.Server.Secrets.Provider {
91 case "openbao":
92 if cfg.Server.Secrets.OpenBao.ProxyAddr == "" {
93 return nil, fmt.Errorf("openbao proxy address is required when using openbao secrets provider")
94 }
95 vault, err = secrets.NewOpenBaoManager(
96 cfg.Server.Secrets.OpenBao.ProxyAddr,
97 logger,
98 secrets.WithMountPath(cfg.Server.Secrets.OpenBao.Mount),
99 )
100 if err != nil {
101 return nil, fmt.Errorf("failed to setup openbao secrets provider: %w", err)
102 }
103 logger.Info("using openbao secrets provider", "proxy_address", cfg.Server.Secrets.OpenBao.ProxyAddr, "mount", cfg.Server.Secrets.OpenBao.Mount)
104 case "sqlite", "":
105 vault, err = secrets.NewSQLiteManager(cfg.Server.DBPath, secrets.WithTableName("secrets"))
106 if err != nil {
107 return nil, fmt.Errorf("failed to setup sqlite secrets provider: %w", err)
108 }
109 logger.Info("using sqlite secrets provider", "path", cfg.Server.DBPath)
110 default:
111 return nil, fmt.Errorf("unknown secrets provider: %s", cfg.Server.Secrets.Provider)
112 }
113
114 if err := runStartupMigrations(ctx, d, cfg.Server.Tap.Embed, cfg.Server.Tap.DBPath, logger); err != nil {
115 return nil, fmt.Errorf("failed to run startup migrations: %w", err)
116 }
117
118 collections := []string{
119 tangled.SpindleMemberNSID,
120 tangled.RepoNSID,
121 tangled.RepoCollaboratorNSID,
122 tangled.RepoPullNSID,
123 }
124 jc, err := jetstream.NewJetstreamClient(cfg.Server.JetstreamEndpoint, "spindle", collections, nil, log.SubLogger(logger, "jetstream"), d, true, true)
125 if err != nil {
126 return nil, fmt.Errorf("failed to setup jetstream client: %w", err)
127 }
128 jc.AddDid(cfg.Server.Owner)
129 // pull records are created by arbitrary users too, same hack as in tap
130 jc.ExemptCollection(tangled.RepoPullNSID)
131
132 // Check if the spindle knows about any Dids;
133 dids, err := d.GetAllDids()
134 if err != nil {
135 return nil, fmt.Errorf("failed to get all dids: %w", err)
136 }
137 for _, d := range dids {
138 jc.AddDid(d)
139 }
140
141 knownRepos, err := d.AllRepos()
142 if err != nil {
143 return nil, fmt.Errorf("failed to get known repos: %w", err)
144 }
145 for _, r := range knownRepos {
146 if r.Owner != "" {
147 jc.AddDid(r.Owner.String())
148 }
149 }
150
151 resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl)
152
153 spindle := &Spindle{
154 jc: jc,
155 e: e,
156 db: d,
157 l: logger,
158 n: &n,
159 engs: engines,
160 cfg: cfg,
161 res: resolver,
162 verify: repoverify.New(resolver, cfg.Server.Dev),
163 vault: vault,
164 motd: defaultMotd,
165 rootCtx: ctx,
166 jobWake: make(chan struct{}, 1),
167 }
168
169 err = e.AddSpindle(rbacDomain)
170 if err != nil {
171 return nil, fmt.Errorf("failed to set rbac domain: %w", err)
172 }
173 err = spindle.configureOwner()
174 if err != nil {
175 return nil, err
176 }
177 logger.Info("owner set", "did", cfg.Server.Owner)
178
179 cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath)
180 if err != nil {
181 return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err)
182 }
183
184 err = jc.StartJetstream(ctx, spindle.ingest())
185 if err != nil {
186 return nil, fmt.Errorf("failed to start jetstream consumer: %w", err)
187 }
188
189 // spindle listen to knot stream for sh.tangled.git.refUpdate
190 // which will sync the local workflow files in spindle and enqueues the
191 // pipeline job for on-push workflows
192 ccfg := eventconsumer.NewConsumerConfig()
193 ccfg.Logger = log.SubLogger(logger, "eventconsumer")
194 ccfg.ProcessFunc = spindle.processKnotStream
195 ccfg.CursorStore = cursorStore
196 ccfg.WorkerCount = 16
197 ccfg.QueueSize = 200
198 if cfg.Server.Dev {
199 ccfg.RetryInterval = 5 * time.Second
200 ccfg.MaxRetryInterval = 10 * time.Second
201 } else {
202 ccfg.RetryInterval = 1 * time.Minute
203 ccfg.MaxRetryInterval = 10 * time.Minute
204 }
205 knownKnots, err := d.Knots()
206 if err != nil {
207 return nil, err
208 }
209 for _, knot := range knownKnots {
210 logger.Info("adding source start", "knot", knot)
211 src := eventconsumer.NewKnotSource(knot)
212 eventconsumer.MigrateLegacyCursor(cursorStore, src)
213 ccfg.Sources[src] = struct{}{}
214 }
215 spindle.ks = eventconsumer.NewConsumer(*ccfg)
216
217 if cfg.Server.Tap.Embed {
218 pw, err := randomAdminPassword()
219 if err != nil {
220 return nil, err
221 }
222 cfg.Server.Tap.AdminPassword = pw
223 logger.Info("embedded tap: using random admin password")
224 }
225 spindle.tap = NewTapClient(spindle)
226
227 return spindle, nil
228}
229
230// DB returns the database instance.
231func (s *Spindle) DB() *db.DB {
232 return s.db
233}
234
235// Engines returns the map of available engines.
236func (s *Spindle) Engines() map[string]models.Engine {
237 return s.engs
238}
239
240// Vault returns the secrets manager instance.
241func (s *Spindle) Vault() secrets.Manager {
242 return s.vault
243}
244
245// Notifier returns the notifier instance.
246func (s *Spindle) Notifier() *notifier.Notifier {
247 return s.n
248}
249
250// Enforcer returns the RBAC enforcer instance.
251func (s *Spindle) Enforcer() *rbac.Enforcer {
252 return s.e
253}
254
255// SetMotdContent sets custom MOTD content, replacing the embedded default.
256func (s *Spindle) SetMotdContent(content []byte) {
257 s.motdMu.Lock()
258 defer s.motdMu.Unlock()
259 s.motd = content
260}
261
262// GetMotdContent returns the current MOTD content.
263func (s *Spindle) GetMotdContent() []byte {
264 s.motdMu.RLock()
265 defer s.motdMu.RUnlock()
266 return s.motd
267}
268
269// Start starts the Spindle server (blocking).
270func (s *Spindle) Start(ctx context.Context) error {
271 // starts a job queue runner in the background
272 s.StartJobWorkers(ctx)
273
274 // Stop vault token renewal if it implements Stopper
275 if stopper, ok := s.vault.(secrets.Stopper); ok {
276 defer stopper.Stop()
277 }
278
279 tapCtx, tapCancel := context.WithCancel(ctx)
280
281 if s.cfg.Server.Tap.Embed {
282 emb, err := startEmbeddedTap(tapCtx, s.cfg, log.SubLogger(s.l, "embedtap"))
283 if err != nil {
284 tapCancel()
285 return fmt.Errorf("starting embedded tap: %w", err)
286 }
287 s.embedTap = emb
288 defer func() {
289 tapCancel()
290 s.embedTap.Shutdown()
291 }()
292
293 go s.watchTapDrain(tapCtx, tapCancel)
294 } else {
295 defer tapCancel()
296 }
297
298 go func() {
299 s.l.Info("starting knot event consumer")
300 s.ks.Start(ctx)
301 }()
302
303 s.l.Info("starting tap client", "url", s.cfg.Server.Tap.Url)
304 s.tap.Start(tapCtx)
305
306 s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr)
307 return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router())
308}
309
310func (s *Spindle) declareTapInterest(ctx context.Context) {
311 repos, err := s.db.AllRepos()
312 if err != nil {
313 s.l.Warn("tap declare: failed to load known repos", "err", err)
314 return
315 }
316 seen := make(map[syntax.DID]struct{}, len(repos))
317 dids := make([]syntax.DID, 0, len(repos))
318 for _, r := range repos {
319 if r.Owner == "" {
320 continue
321 }
322 if _, ok := seen[r.Owner]; ok {
323 continue
324 }
325 seen[r.Owner] = struct{}{}
326 dids = append(dids, r.Owner)
327 }
328 if err := s.tap.AddOwnerDIDs(ctx, dids); err != nil {
329 s.l.Warn("tap declare: AddRepos rejected", "count", len(dids), "err", err)
330 return
331 }
332 s.l.Info("tap declare: known owner DIDs registered", "count", len(dids))
333}
334
335func Run(ctx context.Context) error {
336 cfg, err := config.Load(ctx)
337 if err != nil {
338 return fmt.Errorf("failed to load config: %w", err)
339 }
340
341 if err := ensureGitVersion(); err != nil {
342 return fmt.Errorf("ensuring git version: %w", err)
343 }
344
345 d, err := db.Make(ctx, cfg.Server.DBPath)
346 if err != nil {
347 return fmt.Errorf("failed to setup db: %w", err)
348 }
349
350 nixeryEng, err := nixery.New(ctx, cfg)
351 if err != nil {
352 return err
353 }
354
355 microvmEng, err := microvm.New(ctx, cfg, d)
356 if err != nil {
357 return err
358 }
359
360 s, err := New(ctx, cfg, d, map[string]models.Engine{
361 "nixery": nixeryEng,
362 "microvm": microvmEng,
363 "dummy": dummy.New(log.FromContext(ctx)),
364 })
365 if err != nil {
366 return err
367 }
368
369 return s.Start(ctx)
370}
371
372func (s *Spindle) Router() http.Handler {
373 mux := chi.NewRouter()
374
375 mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
376 w.Write(s.GetMotdContent())
377 })
378 mux.HandleFunc("/events", s.Events)
379 mux.HandleFunc("/logs/{knot}/{rkey}/{name}", s.Logs)
380
381 mux.Mount("/xrpc", s.XrpcRouter())
382 return mux
383}
384
385func (s *Spindle) XrpcRouter() http.Handler {
386 serviceAuth := serviceauth.NewServiceAuth(s.l, s.res.Directory(), s.cfg.Server.Did().String())
387
388 l := log.SubLogger(s.l, "xrpc")
389
390 x := xrpc.Xrpc{
391 Logger: l,
392 Db: s.db,
393 Enforcer: s.e,
394 Engines: s.engs,
395 Config: s.cfg,
396 Resolver: s.res,
397 Vault: s.vault,
398 Notifier: s.Notifier(),
399 ServiceAuth: serviceAuth,
400 Trigger: s,
401 }
402
403 return x.Router()
404}
405
406func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) error {
407 l := log.FromContext(ctx).With("handler", "processKnotStream")
408 l = l.With("src", src.Key(), "msg.Nsid", msg.Nsid, "msg.Rkey", msg.Rkey)
409 if msg.Nsid == knotdb.RepoCollaboratorUpdateNSID {
410 return s.ingestKnotCollaborator(ctx, l, src, msg)
411 }
412 if msg.Nsid == tangled.GitRefUpdateNSID {
413 event := tangled.GitRefUpdate{}
414 if err := json.Unmarshal(msg.EventJson, &event); err != nil {
415 l.Error("error unmarshalling", "err", err)
416 return err
417 }
418 l = l.With("repo", event.Repo, "ref", event.Ref, "newSha", event.NewSha)
419 l.Debug("debug")
420
421 repoDid := syntax.DID(event.Repo)
422 repo, err := s.db.GetRepoByDid(repoDid)
423 if err != nil {
424 return fmt.Errorf("unknown repoDid %s: %w", repoDid, err)
425 }
426
427 if src.Host != repo.Knot {
428 return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, repo.Knot)
429 }
430
431 if kgit.HasSkipCIPushOption(event.PushOptions) {
432 l.Info("push event requested ci skip, skipping the event")
433 return nil
434 }
435
436 // NOTE: we are blindly trusting the knot that it will return only repos it own
437 repoCloneUri := s.newRepoCloneUrl(src.Host, repoDid)
438 repoPath := s.newRepoPath(repoDid)
439
440 triggerRepo, err := s.buildTriggerRepo(ctx, repo)
441 if err != nil {
442 return fmt.Errorf("building trigger repo: %w", err)
443 }
444
445 trigger := tangled.Pipeline_TriggerMetadata{
446 Kind: string(workflow.TriggerKindPush),
447 Push: &tangled.Pipeline_PushTriggerData{
448 Ref: event.Ref,
449 OldSha: event.OldSha,
450 NewSha: event.NewSha,
451 },
452 Repo: triggerRepo,
453 }
454
455 pipelineId, err := s.runPipeline(ctx, repoDid, trigger, event.ChangedFiles, repoCloneUri, repoPath, event.NewSha, nil, triggerRepo)
456 if err != nil {
457 return err
458 }
459 if pipelineId.Rkey == "" {
460 l.Info("no workflow matched 'push' trigger, skipping the event")
461 return nil
462 }
463 l.Info("pipeline triggered", "pipeline", pipelineId.AtUri())
464 }
465
466 return nil
467}
468
469func (s *Spindle) ingestKnotCollaborator(ctx context.Context, l *slog.Logger, src eventconsumer.Source, msg eventstream.Event) error {
470 var rec knotdb.RepoCollaboratorUpdate
471 if err := json.Unmarshal(msg.EventJson, &rec); err != nil {
472 l.Error("error unmarshalling collaboratorUpdate", "err", err)
473 return err
474 }
475
476 subject, err := syntax.ParseDID(rec.Subject)
477 if err != nil {
478 l.Info("skipping collaboratorUpdate with malformed subject", "subject", rec.Subject, "err", err)
479 return nil
480 }
481 repoDid, err := syntax.ParseDID(rec.Repo)
482 if err != nil {
483 l.Info("skipping collaboratorUpdate with malformed repo", "repo", rec.Repo, "err", err)
484 return nil
485 }
486
487 repo, err := s.db.GetRepoByDid(repoDid)
488 if errors.Is(err, sql.ErrNoRows) {
489 l.Info("skipping collaboratorUpdate for unknown repo", "repo", repoDid)
490 return nil
491 }
492 if err != nil {
493 return fmt.Errorf("lookup repo %s: %w", repoDid, err)
494 }
495 if src.Host != repo.Knot {
496 l.Warn("dropping collaboratorUpdate from non-owning knot", "src", src.Host, "repoKnot", repo.Knot)
497 return nil
498 }
499
500 switch rec.Op {
501 case knotdb.AclOpAdd:
502 if err := s.e.AddCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil {
503 return fmt.Errorf("add collaborator policy: %w", err)
504 }
505 if err := s.db.AddKnotCollaborator(repoDid, subject); err != nil {
506 return fmt.Errorf("track collaborator: %w", err)
507 }
508 l.Info("added knot-managed collaborator", "subject", subject, "repo", repoDid)
509 case knotdb.AclOpRemove:
510 if err := s.e.RemoveCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil {
511 return fmt.Errorf("remove collaborator policy: %w", err)
512 }
513 if err := s.db.DeleteRepoCollaboratorBySubjectRepo(subject, repoDid); err != nil {
514 return fmt.Errorf("delete collaborator row: %w", err)
515 }
516 l.Info("removed knot-managed collaborator", "subject", subject, "repo", repoDid)
517 default:
518 return fmt.Errorf("collaboratorUpdate unknown op %q", rec.Op)
519 }
520 return nil
521}
522
523// buildTriggerRepo gathers trigger metadata, resolving default branch from the knot
524func (s *Spindle) buildTriggerRepo(ctx context.Context, repo *db.Repo) (*tangled.Pipeline_TriggerRepo, error) {
525 rkey := string(repo.Rkey)
526 repoDid := repo.RepoDid.String()
527 return s.buildTriggerRepoFrom(ctx, repo.Knot, repo.Owner.String(), rkey, repoDid), nil
528}
529
530func (s *Spindle) buildTriggerRepoFrom(ctx context.Context, knot, did, rkey, repoDid string) *tangled.Pipeline_TriggerRepo {
531 scheme := "https"
532 if s.cfg.Server.Dev {
533 scheme = "http"
534 }
535 client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, knot)}
536
537 // this should maybe (?) be in the refUpdate event itself to save a roundtrip
538 defaultBranch := ""
539 if out, err := tangled.RepoGetDefaultBranch(ctx, client, repoDid); err == nil {
540 defaultBranch = out.Name
541 }
542
543 var rkeyPtr *string
544 if rkey != "" {
545 rkeyPtr = &rkey
546 }
547 return &tangled.Pipeline_TriggerRepo{
548 Did: did,
549 Knot: knot,
550 Repo: rkeyPtr,
551 RepoDid: &repoDid,
552 DefaultBranch: defaultBranch,
553 }
554}
555
556func (s *Spindle) resolvePipelineSourceRepo(ctx context.Context, trigger *tangled.Pipeline_TriggerMetadata) (*tangled.Pipeline_TriggerRepo, error) {
557 if trigger == nil {
558 return nil, nil
559 }
560 if trigger.SourceRepo == nil || *trigger.SourceRepo == "" {
561 return trigger.Repo, nil
562 }
563 repoDid, err := syntax.ParseDID(*trigger.SourceRepo)
564 if err != nil {
565 return nil, fmt.Errorf("parse sourceRepo %s: %w", *trigger.SourceRepo, err)
566 }
567 return s.resolveSourceRepoInfo(ctx, repoDid)
568}
569
570// resolveSourceRepoInfo resolves trigger-repo metadata for a source repo DID.
571func (s *Spindle) resolveSourceRepoInfo(ctx context.Context, repoDid syntax.DID) (*tangled.Pipeline_TriggerRepo, error) {
572 repo, err := s.db.GetRepoByDid(repoDid)
573 if err == nil {
574 return s.buildTriggerRepo(ctx, repo)
575 }
576
577 // verify repo, we don't want git sync to point to arbitrary endpoints
578 res, err := s.verify(ctx, repoident.RepoDid(repoDid))
579 if err != nil {
580 return nil, fmt.Errorf("verify sourceRepo %s: %w", repoDid, err)
581 }
582 return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host, res.OwnerDid.String(), res.Rkey, repoDid.String()), nil
583}
584
585// runPipeline compiles and enqueues the pipeline for the given revision.
586// sourceRepo is the resolved repo the code was checked out from, forwarded to
587// processPipeline for env vars.
588func (s *Spindle) runPipeline(ctx context.Context, repoDid syntax.DID, trigger tangled.Pipeline_TriggerMetadata, changedFiles []string, repoCloneUri, repoPath, rev string, only []string, sourceRepo *tangled.Pipeline_TriggerRepo) (models.PipelineId, error) {
589 l := log.FromContext(ctx)
590
591 compiler := workflow.Compiler{
592 ChangedFiles: changedFiles,
593 Trigger: trigger,
594 }
595
596 rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev)
597 if err != nil {
598 return models.PipelineId{}, fmt.Errorf("loading pipeline: %w", err)
599 }
600 if len(rawPipeline) == 0 {
601 return models.PipelineId{}, nil
602 }
603
604 tpl := compiler.Compile(compiler.Parse(rawPipeline))
605 // todo(dawn): pass compile error to workflow log
606 for _, w := range compiler.Diagnostics.Errors {
607 l.Error(w.String())
608 }
609 for _, w := range compiler.Diagnostics.Warnings {
610 l.Warn(w.String())
611 }
612
613 if len(only) > 0 {
614 tpl.Workflows = filterWorkflows(tpl.Workflows, only)
615 }
616 if len(tpl.Workflows) == 0 {
617 return models.PipelineId{}, nil
618 }
619
620 pipelineId := models.PipelineId{
621 Knot: trigger.Repo.Knot,
622 Rkey: tid.TID(),
623 }
624 if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil {
625 return models.PipelineId{}, fmt.Errorf("creating pipeline event: %w", err)
626 }
627 err = s.processPipeline(repoDid, tpl, pipelineId, sourceRepo)
628 return pipelineId, err
629}
630
631// filterWorkflows filters workflows to the requested names
632func filterWorkflows(workflows []*tangled.Pipeline_Workflow, only []string) []*tangled.Pipeline_Workflow {
633 allowed := make(map[string]struct{}, len(only))
634 for _, n := range only {
635 allowed[n] = struct{}{}
636 }
637 var filtered []*tangled.Pipeline_Workflow
638 for _, w := range workflows {
639 if w == nil {
640 continue
641 }
642 if _, ok := allowed[w.Name]; ok {
643 filtered = append(filtered, w)
644 }
645 }
646 return filtered
647}
648
649// TriggerManual dispatches a pipeline at sha, authorized against and recorded
650// under repoDid. sourceRepo, pull, and inputs are optional trigger payload.
651func (s *Spindle) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull xrpc.PullContext, inputs []*tangled.Pipeline_Pair) (syntax.ATURI, error) {
652 repo, err := s.db.GetRepoByDid(repoDid)
653 if err != nil {
654 return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err)
655 }
656
657 triggerRepo, err := s.buildTriggerRepo(ctx, repo)
658 if err != nil {
659 return "", fmt.Errorf("building trigger repo: %w", err)
660 }
661
662 trigger := tangled.Pipeline_TriggerMetadata{Repo: triggerRepo}
663 if pull.IsPullRequest {
664 var pullAt *string
665 if pull.Pull != "" {
666 pullAtStr := pull.Pull.String()
667 pullAt = &pullAtStr
668 }
669 trigger.Kind = string(workflow.TriggerKindPullRequest)
670 trigger.PullRequest = &tangled.Pipeline_PullRequestTriggerData{
671 SourceBranch: pull.SourceBranch,
672 TargetBranch: pull.TargetBranch,
673 SourceSha: sha,
674 Pull: pullAt,
675 }
676 } else {
677 var refPtr *string
678 if ref != "" {
679 refPtr = &ref
680 }
681 trigger.Kind = string(workflow.TriggerKindManual)
682 trigger.Manual = &tangled.Pipeline_ManualTriggerData{
683 Sha: sha,
684 Ref: refPtr,
685 Inputs: inputs,
686 }
687 }
688
689 repoCloneUri := s.newRepoCloneUrl(repo.Knot, repoDid)
690 repoPath := s.newRepoPath(repoDid)
691 sourceInfo := triggerRepo // default: code comes from the repo itself
692 if sourceRepo != "" && sourceRepo != repoDid {
693 sourceInfo, err = s.resolveSourceRepoInfo(ctx, sourceRepo)
694 if err != nil {
695 return "", err
696 }
697 sourceRepoStr := sourceRepo.String()
698 trigger.SourceRepo = &sourceRepoStr
699 repoCloneUri = models.BuildRepoURL(sourceInfo)
700 repoPath = s.newRepoPath(sourceRepo)
701 }
702
703 pipelineId, err := s.runPipeline(ctx, repoDid, trigger, nil, repoCloneUri, repoPath, sha, workflows, sourceInfo)
704 if err != nil {
705 return "", err
706 }
707 if pipelineId.Rkey == "" {
708 return "", xrpc.ErrNoMatchingWorkflows
709 }
710 return pipelineId.AtUri(), nil
711}
712
713func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, error) {
714 if err := git.SparseSyncGitRepo(ctx, repoUri, repoPath, rev); err != nil {
715 return nil, fmt.Errorf("syncing git repo: %w", err)
716 }
717 gr, err := kgit.Open(repoPath, rev)
718 if err != nil {
719 return nil, fmt.Errorf("opening git repo: %w", err)
720 }
721
722 workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir)
723 if errors.Is(err, object.ErrDirectoryNotFound) {
724 // return empty RawPipeline when directory doesn't exist
725 return nil, nil
726 } else if err != nil {
727 return nil, fmt.Errorf("loading file tree: %w", err)
728 }
729
730 var rawPipeline workflow.RawPipeline
731 for _, e := range workflowDir {
732 if !e.IsFile() {
733 continue
734 }
735
736 fpath := filepath.Join(workflow.WorkflowDir, e.Name)
737 contents, err := gr.RawContent(fpath)
738 if err != nil {
739 return nil, fmt.Errorf("reading raw content of '%s': %w", fpath, err)
740 }
741
742 rawPipeline = append(rawPipeline, workflow.RawWorkflow{
743 Name: e.Name,
744 Contents: contents,
745 })
746 }
747
748 return rawPipeline, nil
749}
750
751func (s *Spindle) StartJobWorkers(ctx context.Context) {
752 for range s.cfg.Server.MaxJobCount {
753 go func() {
754 for {
755 job, err := s.db.DequeueJob(ctx)
756 if err != nil {
757 s.l.Error("failed to dequeue job", "error", err)
758 }
759 if job == nil {
760 // sleep until a new job wakes us
761 select {
762 case <-ctx.Done():
763 return
764 case <-s.jobWake:
765 }
766 continue
767 }
768 s.runJob(ctx, job)
769 }
770 }()
771 }
772}
773
774func (s *Spindle) runJob(ctx context.Context, job *db.JobRow) {
775 pipelineId := models.PipelineId{
776 Knot: job.PipelineIdKnot,
777 Rkey: job.PipelineIdRkey,
778 }
779
780 pipelineEnv := models.PipelineEnvVarsForSource(job.Tpl.TriggerMetadata, pipelineId, job.SourceRepo)
781 trustedSource := true
782 if tm := job.Tpl.TriggerMetadata; tm != nil && tm.SourceRepo != nil &&
783 *tm.SourceRepo != "" && *tm.SourceRepo != job.RepoDid {
784 trustedSource = false
785 }
786
787 initTpl := job.Tpl
788 if job.SourceRepo != nil && job.Tpl.TriggerMetadata != nil {
789 tm := *job.Tpl.TriggerMetadata
790 tm.Repo = job.SourceRepo
791 initTpl.TriggerMetadata = &tm
792 }
793
794 workflows := make(map[models.Engine][]models.Workflow)
795 for _, w := range job.Tpl.Workflows {
796 if w == nil {
797 continue
798 }
799 eng, ok := s.engs[w.Engine]
800 if !ok {
801 _ = s.db.StatusFailed(models.WorkflowId{
802 PipelineId: pipelineId,
803 Name: w.Name,
804 }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n)
805 continue
806 }
807
808 ewf, err := eng.InitWorkflow(*w, initTpl)
809 if err != nil {
810 _ = s.db.StatusFailed(models.WorkflowId{
811 PipelineId: pipelineId,
812 Name: w.Name,
813 }, fmt.Sprintf("init workflow: %s", err), -1, s.n)
814 continue
815 }
816
817 if ewf.Environment == nil {
818 ewf.Environment = make(map[string]string)
819 }
820 maps.Copy(ewf.Environment, pipelineEnv)
821 workflows[eng] = append(workflows[eng], *ewf)
822 }
823
824 engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.rootCtx, &models.Pipeline{
825 RepoDid: syntax.DID(job.RepoDid),
826 Workflows: workflows,
827 TrustedSource: trustedSource,
828 }, pipelineId)
829}
830
831// enqueues the workflows in tpl.
832func (s *Spindle) processPipeline(repoDid syntax.DID, tpl tangled.Pipeline, pipelineId models.PipelineId, sourceRepo *tangled.Pipeline_TriggerRepo) error {
833 err := s.db.EnqueueJob(s.rootCtx, repoDid.String(), pipelineId, sourceRepo, tpl)
834 if err != nil {
835 return fmt.Errorf("failed to enqueue durable job: %w", err)
836 }
837 s.l.Info("pipeline enqueued successfully to db", "id", pipelineId)
838
839 // wake up an idle worker to pick up more jobs if any
840 select {
841 case s.jobWake <- struct{}{}:
842 default:
843 }
844
845 // pipelines visible from now on, they are sitting in queue
846 for _, w := range tpl.Workflows {
847 if w == nil {
848 continue
849 }
850 if err := s.db.StatusPending(models.WorkflowId{
851 PipelineId: pipelineId,
852 Name: w.Name,
853 }, s.n); err != nil {
854 return fmt.Errorf("db.StatusPending: %w", err)
855 }
856 }
857 return nil
858}
859
860// newRepoPath creates a path to store repository by its did and rkey.
861// The path format would be: `/data/repos/did:plc:foo/sh.tangled.repo/repo-rkey
862func (s *Spindle) newRepoPath(repo syntax.DID) string {
863 return filepath.Join(s.cfg.Server.RepoDir, repo.String())
864}
865
866func (s *Spindle) newRepoCloneUrl(knot string, did syntax.DID) string {
867 scheme := "https://"
868 if s.cfg.Server.Dev {
869 scheme = "http://"
870 }
871 return fmt.Sprintf("%s%s/%s", scheme, knot, did)
872}
873
874const RequiredVersion = "2.49.0"
875
876func ensureGitVersion() error {
877 v, err := git.Version()
878 if err != nil {
879 return fmt.Errorf("fetching git version: %w", err)
880 }
881 if v.LessThan(version.Must(version.NewVersion(RequiredVersion))) {
882 return fmt.Errorf("installed git version %q is not supported, Spindle requires git version >= %q", v, RequiredVersion)
883 }
884 return nil
885}
886
887func (s *Spindle) resolvePipelineRepoDid(repo *tangled.Pipeline_TriggerRepo) (syntax.DID, error) {
888 if repo.RepoDid == nil || *repo.RepoDid == "" {
889 return "", fmt.Errorf("pipeline trigger missing repoDid")
890 }
891 repoDid, err := syntax.ParseDID(*repo.RepoDid)
892 if err != nil {
893 return "", fmt.Errorf("parse repoDid %s: %w", *repo.RepoDid, err)
894 }
895 if _, err := s.db.GetRepoByDid(repoDid); err != nil {
896 return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err)
897 }
898 return repoDid, nil
899}
900
901func (s *Spindle) configureOwner() error {
902 cfgOwner := s.cfg.Server.Owner
903
904 existing, err := s.e.GetSpindleUsersByRole("server:owner", rbacDomain)
905 if err != nil {
906 return err
907 }
908
909 switch len(existing) {
910 case 0:
911 // no owner configured, continue
912 case 1:
913 // find existing owner
914 existingOwner := existing[0]
915
916 // no ownership change, this is okay
917 if existingOwner == s.cfg.Server.Owner {
918 break
919 }
920
921 // remove existing owner
922 err = s.e.RemoveSpindleOwner(rbacDomain, existingOwner)
923 if err != nil {
924 return nil
925 }
926 default:
927 return fmt.Errorf("more than one owner in DB, try deleting %q and starting over", s.cfg.Server.DBPath)
928 }
929
930 return s.e.AddSpindleOwner(rbacDomain, cfgOwner)
931}