This repository has no description
0

Configure Feed

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

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