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