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