This repository has no description
0

Configure Feed

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

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