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