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