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