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