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