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