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