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