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