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