Monorepo for Tangled tangled.org
1

Configure Feed

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

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