This repository has no description
0

Configure Feed

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

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