Monorepo for Tangled tangled.org
1

Configure Feed

Select the types of activity you want to include in your feed.

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