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