Forked monorepo for Tangled
0

Configure Feed

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

core / spindle / server.go
25 kB 868 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 if kgit.HasSkipCIPushOption(event.PushOptions) { 431 l.Info("push event requested ci skip, skipping the event") 432 return nil 433 } 434 435 // NOTE: we are blindly trusting the knot that it will return only repos it own 436 repoCloneUri := s.newRepoCloneUrl(src.Host, repoDid) 437 repoPath := s.newRepoPath(repoDid) 438 if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, event.NewSha); err != nil { 439 return fmt.Errorf("sync git repo: %w", err) 440 } 441 l.Info("synced git repo") 442 443 triggerRepo, err := s.buildTriggerRepo(ctx, repo) 444 if err != nil { 445 return fmt.Errorf("building trigger repo: %w", err) 446 } 447 448 trigger := tangled.Pipeline_TriggerMetadata{ 449 Kind: string(workflow.TriggerKindPush), 450 Push: &tangled.Pipeline_PushTriggerData{ 451 Ref: event.Ref, 452 OldSha: event.OldSha, 453 NewSha: event.NewSha, 454 }, 455 Repo: triggerRepo, 456 } 457 458 pipelineId, err := s.runPipeline(ctx, repoDid, trigger, event.ChangedFiles, repoCloneUri, repoPath, event.NewSha, nil, triggerRepo) 459 if err != nil { 460 return err 461 } 462 if pipelineId.Rkey == "" { 463 l.Info("no workflow matched 'push' trigger, skipping the event") 464 return nil 465 } 466 l.Info("pipeline triggered", "pipeline", pipelineId.AtUri()) 467 } 468 469 return nil 470} 471 472// buildTriggerRepo gathers trigger metadata, resolving default branch from the knot 473func (s *Spindle) buildTriggerRepo(ctx context.Context, repo *db.Repo) (*tangled.Pipeline_TriggerRepo, error) { 474 rkey := string(repo.Rkey) 475 repoDid := repo.RepoDid.String() 476 return s.buildTriggerRepoFrom(ctx, repo.Knot, repo.Owner.String(), rkey, repoDid), nil 477} 478 479func (s *Spindle) buildTriggerRepoFrom(ctx context.Context, knot, did, rkey, repoDid string) *tangled.Pipeline_TriggerRepo { 480 scheme := "https" 481 if s.cfg.Server.Dev { 482 scheme = "http" 483 } 484 client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, knot)} 485 486 // this should maybe (?) be in the refUpdate event itself to save a roundtrip 487 defaultBranch := "" 488 if out, err := tangled.RepoGetDefaultBranch(ctx, client, repoDid); err == nil { 489 defaultBranch = out.Name 490 } 491 492 var rkeyPtr *string 493 if rkey != "" { 494 rkeyPtr = &rkey 495 } 496 return &tangled.Pipeline_TriggerRepo{ 497 Did: did, 498 Knot: knot, 499 Repo: rkeyPtr, 500 RepoDid: &repoDid, 501 DefaultBranch: defaultBranch, 502 } 503} 504 505func (s *Spindle) resolvePipelineSourceRepo(ctx context.Context, trigger *tangled.Pipeline_TriggerMetadata) (*tangled.Pipeline_TriggerRepo, error) { 506 if trigger == nil { 507 return nil, nil 508 } 509 if trigger.SourceRepo == nil || *trigger.SourceRepo == "" { 510 return trigger.Repo, nil 511 } 512 repoDid, err := syntax.ParseDID(*trigger.SourceRepo) 513 if err != nil { 514 return nil, fmt.Errorf("parse sourceRepo %s: %w", *trigger.SourceRepo, err) 515 } 516 return s.resolveSourceRepoInfo(ctx, repoDid) 517} 518 519// resolveSourceRepoInfo resolves trigger-repo metadata for a source repo DID. 520func (s *Spindle) resolveSourceRepoInfo(ctx context.Context, repoDid syntax.DID) (*tangled.Pipeline_TriggerRepo, error) { 521 repo, err := s.db.GetRepoByDid(repoDid) 522 if err == nil { 523 return s.buildTriggerRepo(ctx, repo) 524 } 525 526 // verify repo, we don't want git sync to point to arbitrary endpoints 527 res, err := s.verify(ctx, repoverify.RepoDid(repoDid)) 528 if err != nil { 529 return nil, fmt.Errorf("verify sourceRepo %s: %w", repoDid, err) 530 } 531 return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host, res.OwnerDid.String(), res.Rkey, repoDid.String()), nil 532} 533 534// runPipeline compiles and enqueues the pipeline for the given revision. 535// sourceRepo is the resolved repo the code was checked out from, forwarded to 536// processPipeline for env vars. 537func (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) { 538 l := log.FromContext(ctx) 539 540 compiler := workflow.Compiler{ 541 ChangedFiles: changedFiles, 542 Trigger: trigger, 543 } 544 545 rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev) 546 if err != nil { 547 return models.PipelineId{}, fmt.Errorf("loading pipeline: %w", err) 548 } 549 if len(rawPipeline) == 0 { 550 return models.PipelineId{}, nil 551 } 552 553 tpl := compiler.Compile(compiler.Parse(rawPipeline)) 554 // todo(dawn): pass compile error to workflow log 555 for _, w := range compiler.Diagnostics.Errors { 556 l.Error(w.String()) 557 } 558 for _, w := range compiler.Diagnostics.Warnings { 559 l.Warn(w.String()) 560 } 561 562 if len(only) > 0 { 563 tpl.Workflows = filterWorkflows(tpl.Workflows, only) 564 } 565 if len(tpl.Workflows) == 0 { 566 return models.PipelineId{}, nil 567 } 568 569 pipelineId := models.PipelineId{ 570 Knot: trigger.Repo.Knot, 571 Rkey: tid.TID(), 572 } 573 if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { 574 return models.PipelineId{}, fmt.Errorf("creating pipeline event: %w", err) 575 } 576 err = s.processPipeline(repoDid, tpl, pipelineId, sourceRepo) 577 return pipelineId, err 578} 579 580// filterWorkflows filters workflows to the requested names 581func filterWorkflows(workflows []*tangled.Pipeline_Workflow, only []string) []*tangled.Pipeline_Workflow { 582 allowed := make(map[string]struct{}, len(only)) 583 for _, n := range only { 584 allowed[n] = struct{}{} 585 } 586 var filtered []*tangled.Pipeline_Workflow 587 for _, w := range workflows { 588 if w == nil { 589 continue 590 } 591 if _, ok := allowed[w.Name]; ok { 592 filtered = append(filtered, w) 593 } 594 } 595 return filtered 596} 597 598// TriggerManual dispatches a pipeline at sha, authorized against and recorded 599// under repoDid. sourceRepo, pull, and inputs are optional trigger payload. 600func (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) { 601 repo, err := s.db.GetRepoByDid(repoDid) 602 if err != nil { 603 return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) 604 } 605 606 triggerRepo, err := s.buildTriggerRepo(ctx, repo) 607 if err != nil { 608 return "", fmt.Errorf("building trigger repo: %w", err) 609 } 610 611 trigger := tangled.Pipeline_TriggerMetadata{Repo: triggerRepo} 612 if pull.IsPullRequest { 613 var pullAt *string 614 if pull.Pull != "" { 615 pullAtStr := pull.Pull.String() 616 pullAt = &pullAtStr 617 } 618 trigger.Kind = string(workflow.TriggerKindPullRequest) 619 trigger.PullRequest = &tangled.Pipeline_PullRequestTriggerData{ 620 SourceBranch: pull.SourceBranch, 621 TargetBranch: pull.TargetBranch, 622 SourceSha: sha, 623 Pull: pullAt, 624 } 625 } else { 626 var refPtr *string 627 if ref != "" { 628 refPtr = &ref 629 } 630 trigger.Kind = string(workflow.TriggerKindManual) 631 trigger.Manual = &tangled.Pipeline_ManualTriggerData{ 632 Sha: sha, 633 Ref: refPtr, 634 Inputs: inputs, 635 } 636 } 637 638 repoCloneUri := s.newRepoCloneUrl(repo.Knot, repoDid) 639 repoPath := s.newRepoPath(repoDid) 640 sourceInfo := triggerRepo // default: code comes from the repo itself 641 if sourceRepo != "" && sourceRepo != repoDid { 642 sourceInfo, err = s.resolveSourceRepoInfo(ctx, sourceRepo) 643 if err != nil { 644 return "", err 645 } 646 sourceRepoStr := sourceRepo.String() 647 trigger.SourceRepo = &sourceRepoStr 648 repoCloneUri = models.BuildRepoURL(sourceInfo) 649 repoPath = s.newRepoPath(sourceRepo) 650 } 651 652 pipelineId, err := s.runPipeline(ctx, repoDid, trigger, nil, repoCloneUri, repoPath, sha, workflows, sourceInfo) 653 if err != nil { 654 return "", err 655 } 656 if pipelineId.Rkey == "" { 657 return "", xrpc.ErrNoMatchingWorkflows 658 } 659 return pipelineId.AtUri(), nil 660} 661 662func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, error) { 663 if err := git.SparseSyncGitRepo(ctx, repoUri, repoPath, rev); err != nil { 664 return nil, fmt.Errorf("syncing git repo: %w", err) 665 } 666 gr, err := kgit.Open(repoPath, rev) 667 if err != nil { 668 return nil, fmt.Errorf("opening git repo: %w", err) 669 } 670 671 workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) 672 if errors.Is(err, object.ErrDirectoryNotFound) { 673 // return empty RawPipeline when directory doesn't exist 674 return nil, nil 675 } else if err != nil { 676 return nil, fmt.Errorf("loading file tree: %w", err) 677 } 678 679 var rawPipeline workflow.RawPipeline 680 for _, e := range workflowDir { 681 if !e.IsFile() { 682 continue 683 } 684 685 fpath := filepath.Join(workflow.WorkflowDir, e.Name) 686 contents, err := gr.RawContent(fpath) 687 if err != nil { 688 return nil, fmt.Errorf("reading raw content of '%s': %w", fpath, err) 689 } 690 691 rawPipeline = append(rawPipeline, workflow.RawWorkflow{ 692 Name: e.Name, 693 Contents: contents, 694 }) 695 } 696 697 return rawPipeline, nil 698} 699 700// processPipeline enqueues the workflows in tpl. 701func (s *Spindle) processPipeline(repoDid syntax.DID, tpl tangled.Pipeline, pipelineId models.PipelineId, sourceRepo *tangled.Pipeline_TriggerRepo) error { 702 // derive security-relevant things like whether this run is trusted and can be passed 703 // secrets to from the original metadata. 704 pipelineEnv := models.PipelineEnvVarsForSource(tpl.TriggerMetadata, pipelineId, sourceRepo) 705 trustedSource := true 706 if tm := tpl.TriggerMetadata; tm != nil && tm.SourceRepo != nil && 707 *tm.SourceRepo != "" && *tm.SourceRepo != repoDid.String() { 708 trustedSource = false 709 } 710 711 // swap the repo with our sourceRepo if we are running a pipeline on a fork. 712 // the metadata stays the same. we check whether the repo is trusted above, 713 // so this only affects the clone URL. 714 initTpl := tpl 715 if sourceRepo != nil && tpl.TriggerMetadata != nil { 716 tm := *tpl.TriggerMetadata 717 tm.Repo = sourceRepo 718 initTpl.TriggerMetadata = &tm 719 } 720 721 // filter & init workflows 722 workflows := make(map[models.Engine][]models.Workflow) 723 for _, w := range tpl.Workflows { 724 if w == nil { 725 continue 726 } 727 eng, ok := s.engs[w.Engine] 728 if !ok { 729 err := s.db.StatusFailed(models.WorkflowId{ 730 PipelineId: pipelineId, 731 Name: w.Name, 732 }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) 733 if err != nil { 734 return fmt.Errorf("db.StatusFailed: %w", err) 735 } 736 737 continue 738 } 739 740 ewf, err := eng.InitWorkflow(*w, initTpl) 741 if err != nil { 742 err = s.db.StatusFailed(models.WorkflowId{ 743 PipelineId: pipelineId, 744 Name: w.Name, 745 }, fmt.Sprintf("init workflow: %s", err), -1, s.n) 746 if err != nil { 747 return fmt.Errorf("db.StatusFailed: %w", err) 748 } 749 750 continue 751 } 752 753 // inject TANGLED_* env vars after InitWorkflow 754 // This prevents user-defined env vars from overriding them 755 if ewf.Environment == nil { 756 ewf.Environment = make(map[string]string) 757 } 758 maps.Copy(ewf.Environment, pipelineEnv) 759 760 workflows[eng] = append(workflows[eng], *ewf) 761 } 762 763 // enqueue pipeline 764 ok := s.jq.Enqueue(repoDid, queue.Job{ 765 Run: func() error { 766 engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.rootCtx, &models.Pipeline{ 767 RepoDid: repoDid, 768 Workflows: workflows, 769 TrustedSource: trustedSource, 770 }, pipelineId) 771 return nil 772 }, 773 OnFail: func(jobError error) { 774 s.l.Error("pipeline run failed", "error", jobError) 775 }, 776 }) 777 if !ok { 778 return fmt.Errorf("failed to enqueue pipeline: queue is full") 779 } 780 s.l.Info("pipeline enqueued successfully", "id", pipelineId) 781 782 // after successful enqueue, emit StatusPending for all workflows 783 for _, ewfs := range workflows { 784 for _, ewf := range ewfs { 785 err := s.db.StatusPending(models.WorkflowId{ 786 PipelineId: pipelineId, 787 Name: ewf.Name, 788 }, s.n) 789 if err != nil { 790 return fmt.Errorf("db.StatusPending: %w", err) 791 } 792 } 793 } 794 return nil 795} 796 797// newRepoPath creates a path to store repository by its did and rkey. 798// The path format would be: `/data/repos/did:plc:foo/sh.tangled.repo/repo-rkey 799func (s *Spindle) newRepoPath(repo syntax.DID) string { 800 return filepath.Join(s.cfg.Server.RepoDir, repo.String()) 801} 802 803func (s *Spindle) newRepoCloneUrl(knot string, did syntax.DID) string { 804 scheme := "https://" 805 if s.cfg.Server.Dev { 806 scheme = "http://" 807 } 808 return fmt.Sprintf("%s%s/%s", scheme, knot, did) 809} 810 811const RequiredVersion = "2.49.0" 812 813func ensureGitVersion() error { 814 v, err := git.Version() 815 if err != nil { 816 return fmt.Errorf("fetching git version: %w", err) 817 } 818 if v.LessThan(version.Must(version.NewVersion(RequiredVersion))) { 819 return fmt.Errorf("installed git version %q is not supported, Spindle requires git version >= %q", v, RequiredVersion) 820 } 821 return nil 822} 823 824func (s *Spindle) resolvePipelineRepoDid(repo *tangled.Pipeline_TriggerRepo) (syntax.DID, error) { 825 if repo.RepoDid == nil || *repo.RepoDid == "" { 826 return "", fmt.Errorf("pipeline trigger missing repoDid") 827 } 828 repoDid, err := syntax.ParseDID(*repo.RepoDid) 829 if err != nil { 830 return "", fmt.Errorf("parse repoDid %s: %w", *repo.RepoDid, err) 831 } 832 if _, err := s.db.GetRepoByDid(repoDid); err != nil { 833 return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) 834 } 835 return repoDid, nil 836} 837 838func (s *Spindle) configureOwner() error { 839 cfgOwner := s.cfg.Server.Owner 840 841 existing, err := s.e.GetSpindleUsersByRole("server:owner", rbacDomain) 842 if err != nil { 843 return err 844 } 845 846 switch len(existing) { 847 case 0: 848 // no owner configured, continue 849 case 1: 850 // find existing owner 851 existingOwner := existing[0] 852 853 // no ownership change, this is okay 854 if existingOwner == s.cfg.Server.Owner { 855 break 856 } 857 858 // remove existing owner 859 err = s.e.RemoveSpindleOwner(rbacDomain, existingOwner) 860 if err != nil { 861 return nil 862 } 863 default: 864 return fmt.Errorf("more than one owner in DB, try deleting %q and starting over", s.cfg.Server.DBPath) 865 } 866 867 return s.e.AddSpindleOwner(rbacDomain, cfgOwner) 868}