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