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