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