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