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