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