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