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 867 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 if s.cfg.Server.AdminPassword != "" { 362 mux.Mount("/admin", s.adminRouter()) 363 } else { 364 s.l.Warn("admin api disabled: SPINDLE_SERVER_ADMIN_PASSWORD is unset") 365 } 366 return mux 367} 368 369func (s *Spindle) XrpcRouter() http.Handler { 370 serviceAuth := serviceauth.NewServiceAuth(s.l, s.res.Directory(), s.cfg.Server.Did().String()) 371 372 l := log.SubLogger(s.l, "xrpc") 373 374 x := xrpc.Xrpc{ 375 Logger: l, 376 Db: s.db, 377 Enforcer: s.e, 378 Engines: s.engs, 379 Config: s.cfg, 380 Resolver: s.res, 381 Vault: s.vault, 382 Notifier: s.Notifier(), 383 ServiceAuth: serviceAuth, 384 Trigger: s, 385 } 386 387 return x.Router() 388} 389 390func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) error { 391 l := log.FromContext(ctx).With("handler", "processKnotStream") 392 l = l.With("src", src.Key(), "msg.Nsid", msg.Nsid, "msg.Rkey", msg.Rkey) 393 switch msg.Nsid { 394 case knotdb.RepoCollaboratorUpdateNSID: 395 return s.ingestKnotCollaborator(ctx, l, src, msg) 396 397 case tangled.GitRefUpdateNSID: 398 event := tangled.GitRefUpdate{} 399 if err := json.Unmarshal(msg.EventJson, &event); err != nil { 400 l.Error("error unmarshalling", "err", err) 401 return err 402 } 403 l = l.With("repo", event.Repo, "ref", event.Ref, "newSha", event.NewSha) 404 l.Debug("debug") 405 406 repoDid := syntax.DID(event.Repo) 407 repo, err := s.db.GetRepoByDid(repoDid) 408 if err != nil { 409 return fmt.Errorf("unknown repoDid %s: %w", repoDid, err) 410 } 411 412 if src.Host != repo.Knot { 413 return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, repo.Knot) 414 } 415 416 if kgit.HasSkipCIPushOption(event.PushOptions) { 417 l.Info("push event requested ci skip, skipping the event") 418 return nil 419 } 420 421 // NOTE: we are blindly trusting the knot that it will return only repos it own 422 repoCloneUri := s.newRepoCloneUrl(src.Host, repoDid) 423 repoPath := s.newRepoPath(repoDid) 424 425 triggerRepo, err := s.buildTriggerRepo(ctx, repo) 426 if err != nil { 427 return fmt.Errorf("building trigger repo: %w", err) 428 } 429 430 trigger := tangled.Pipeline_TriggerMetadata{ 431 Kind: string(workflow.TriggerKindPush), 432 Push: &tangled.Pipeline_PushTriggerData{ 433 Ref: event.Ref, 434 OldSha: event.OldSha, 435 NewSha: event.NewSha, 436 }, 437 Repo: triggerRepo, 438 } 439 440 pipelineId, err := s.runPipeline(ctx, repoDid, trigger, event.ChangedFiles, repoCloneUri, repoPath, event.NewSha, nil, triggerRepo) 441 if err != nil { 442 return err 443 } 444 if pipelineId.Rkey == "" { 445 l.Info("no workflow matched 'push' trigger, skipping the event") 446 return nil 447 } 448 l.Info("pipeline triggered", "pipeline", pipelineId.AtUri()) 449 } 450 451 return nil 452} 453 454func (s *Spindle) ingestKnotCollaborator(_ context.Context, l *slog.Logger, src eventconsumer.Source, msg eventstream.Event) error { 455 var rec knotdb.RepoCollaboratorUpdate 456 if err := json.Unmarshal(msg.EventJson, &rec); err != nil { 457 l.Error("error unmarshalling collaboratorUpdate", "err", err) 458 return err 459 } 460 461 subject, err := syntax.ParseDID(rec.Subject) 462 if err != nil { 463 l.Info("skipping collaboratorUpdate with malformed subject", "subject", rec.Subject, "err", err) 464 return nil 465 } 466 repoDid, err := syntax.ParseDID(rec.Repo) 467 if err != nil { 468 l.Info("skipping collaboratorUpdate with malformed repo", "repo", rec.Repo, "err", err) 469 return nil 470 } 471 472 repo, err := s.db.GetRepoByDid(repoDid) 473 if errors.Is(err, sql.ErrNoRows) { 474 l.Info("skipping collaboratorUpdate for unknown repo", "repo", repoDid) 475 return nil 476 } 477 if err != nil { 478 return fmt.Errorf("lookup repo %s: %w", repoDid, err) 479 } 480 if src.Host != repo.Knot { 481 l.Warn("dropping collaboratorUpdate from non-owning knot", "src", src.Host, "repoKnot", repo.Knot) 482 return nil 483 } 484 485 switch rec.Op { 486 case knotdb.AclOpAdd: 487 if err := s.grantCollaborator(subject, repoDid); err != nil { 488 return fmt.Errorf("add collaborator policy: %w", err) 489 } 490 l.Info("added knot-managed collaborator", "subject", subject, "repo", repoDid) 491 case knotdb.AclOpRemove: 492 // ponytail: no jc.RemoveDid here - the subject may still be a member or a 493 // collaborator elsewhere, and a stale filter entry is harmless (processRepo still 494 // rejects non-members). Add refcounting across members + acl_2 if the filter grows. 495 if err := s.e.RemoveRepoCollaborator(subject, repoDid); err != nil { 496 return fmt.Errorf("remove collaborator policy: %w", err) 497 } 498 l.Info("removed knot-managed collaborator", "subject", subject, "repo", repoDid) 499 default: 500 return fmt.Errorf("collaboratorUpdate unknown op %q", rec.Op) 501 } 502 return nil 503} 504 505// buildTriggerRepo gathers trigger metadata, resolving default branch from the knot 506func (s *Spindle) buildTriggerRepo(ctx context.Context, repo *db.Repo) (*tangled.Pipeline_TriggerRepo, error) { 507 rkey := string(repo.Rkey) 508 repoDid := repo.RepoDid.String() 509 return s.buildTriggerRepoFrom(ctx, repo.Knot, repo.Owner.String(), rkey, repoDid), nil 510} 511 512func (s *Spindle) buildTriggerRepoFrom(ctx context.Context, knot, did, rkey, repoDid string) *tangled.Pipeline_TriggerRepo { 513 scheme := "https" 514 if s.cfg.Server.Dev { 515 scheme = "http" 516 } 517 client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, knot)} 518 519 // this should maybe (?) be in the refUpdate event itself to save a roundtrip 520 defaultBranch := "" 521 if out, err := tangled.RepoGetDefaultBranch(ctx, client, repoDid); err == nil { 522 defaultBranch = out.Name 523 } 524 525 var rkeyPtr *string 526 if rkey != "" { 527 rkeyPtr = &rkey 528 } 529 return &tangled.Pipeline_TriggerRepo{ 530 Did: did, 531 Knot: knot, 532 Repo: rkeyPtr, 533 RepoDid: &repoDid, 534 DefaultBranch: defaultBranch, 535 } 536} 537 538func (s *Spindle) resolvePipelineSourceRepo(ctx context.Context, trigger *tangled.Pipeline_TriggerMetadata) (*tangled.Pipeline_TriggerRepo, error) { 539 if trigger == nil { 540 return nil, nil 541 } 542 if trigger.SourceRepo == nil || *trigger.SourceRepo == "" { 543 return trigger.Repo, nil 544 } 545 repoDid, err := syntax.ParseDID(*trigger.SourceRepo) 546 if err != nil { 547 return nil, fmt.Errorf("parse sourceRepo %s: %w", *trigger.SourceRepo, err) 548 } 549 return s.resolveSourceRepoInfo(ctx, repoDid) 550} 551 552// resolveSourceRepoInfo resolves trigger-repo metadata for a source repo DID. 553func (s *Spindle) resolveSourceRepoInfo(ctx context.Context, repoDid syntax.DID) (*tangled.Pipeline_TriggerRepo, error) { 554 repo, err := s.db.GetRepoByDid(repoDid) 555 if err == nil { 556 return s.buildTriggerRepo(ctx, repo) 557 } 558 559 // verify repo, we don't want git sync to point to arbitrary endpoints 560 res, err := s.verify(ctx, repoident.RepoDid(repoDid)) 561 if err != nil { 562 return nil, fmt.Errorf("verify sourceRepo %s: %w", repoDid, err) 563 } 564 return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host, res.OwnerDid.String(), res.Rkey, repoDid.String()), nil 565} 566 567// runPipeline compiles and enqueues the pipeline for the given revision. 568// sourceRepo is the resolved repo the code was checked out from, forwarded to 569// processPipeline for env vars. 570func (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) { 571 l := log.FromContext(ctx) 572 573 compiler := workflow.Compiler{ 574 ChangedFiles: changedFiles, 575 Trigger: trigger, 576 } 577 578 rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev) 579 if err != nil { 580 return models.PipelineId{}, fmt.Errorf("loading pipeline: %w", err) 581 } 582 if len(rawPipeline) == 0 { 583 return models.PipelineId{}, nil 584 } 585 586 tpl := compiler.Compile(compiler.Parse(rawPipeline)) 587 // todo(dawn): pass compile error to workflow log 588 for _, w := range compiler.Diagnostics.Errors { 589 l.Error(w.String()) 590 } 591 for _, w := range compiler.Diagnostics.Warnings { 592 l.Warn(w.String()) 593 } 594 595 if len(only) > 0 { 596 tpl.Workflows = filterWorkflows(tpl.Workflows, only) 597 } 598 if len(tpl.Workflows) == 0 { 599 return models.PipelineId{}, nil 600 } 601 602 pipelineId := models.PipelineId{ 603 Knot: trigger.Repo.Knot, 604 Rkey: tid.TID(), 605 } 606 if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { 607 return models.PipelineId{}, fmt.Errorf("creating pipeline event: %w", err) 608 } 609 err = s.processPipeline(repoDid, tpl, pipelineId, sourceRepo) 610 return pipelineId, err 611} 612 613// filterWorkflows filters workflows to the requested names 614func filterWorkflows(workflows []*tangled.Pipeline_Workflow, only []string) []*tangled.Pipeline_Workflow { 615 allowed := make(map[string]struct{}, len(only)) 616 for _, n := range only { 617 allowed[n] = struct{}{} 618 } 619 var filtered []*tangled.Pipeline_Workflow 620 for _, w := range workflows { 621 if w == nil { 622 continue 623 } 624 if _, ok := allowed[w.Name]; ok { 625 filtered = append(filtered, w) 626 } 627 } 628 return filtered 629} 630 631// TriggerManual dispatches a pipeline at sha, authorized against and recorded 632// under repoDid. sourceRepo, pull, and inputs are optional trigger payload. 633func (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) { 634 repo, err := s.db.GetRepoByDid(repoDid) 635 if err != nil { 636 return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) 637 } 638 639 triggerRepo, err := s.buildTriggerRepo(ctx, repo) 640 if err != nil { 641 return "", fmt.Errorf("building trigger repo: %w", err) 642 } 643 644 trigger := tangled.Pipeline_TriggerMetadata{Repo: triggerRepo} 645 if pull.IsPullRequest { 646 var pullAt *string 647 if pull.Pull != "" { 648 pullAtStr := pull.Pull.String() 649 pullAt = &pullAtStr 650 } 651 trigger.Kind = string(workflow.TriggerKindPullRequest) 652 trigger.PullRequest = &tangled.Pipeline_PullRequestTriggerData{ 653 SourceBranch: pull.SourceBranch, 654 TargetBranch: pull.TargetBranch, 655 SourceSha: sha, 656 Pull: pullAt, 657 } 658 } else { 659 var refPtr *string 660 if ref != "" { 661 refPtr = &ref 662 } 663 trigger.Kind = string(workflow.TriggerKindManual) 664 trigger.Manual = &tangled.Pipeline_ManualTriggerData{ 665 Sha: sha, 666 Ref: refPtr, 667 Inputs: inputs, 668 } 669 } 670 671 repoCloneUri := s.newRepoCloneUrl(repo.Knot, repoDid) 672 repoPath := s.newRepoPath(repoDid) 673 sourceInfo := triggerRepo // default: code comes from the repo itself 674 if sourceRepo != "" && sourceRepo != repoDid { 675 sourceInfo, err = s.resolveSourceRepoInfo(ctx, sourceRepo) 676 if err != nil { 677 return "", err 678 } 679 sourceRepoStr := sourceRepo.String() 680 trigger.SourceRepo = &sourceRepoStr 681 repoCloneUri = models.BuildRepoURL(sourceInfo) 682 repoPath = s.newRepoPath(sourceRepo) 683 } 684 685 pipelineId, err := s.runPipeline(ctx, repoDid, trigger, nil, repoCloneUri, repoPath, sha, workflows, sourceInfo) 686 if err != nil { 687 return "", err 688 } 689 if pipelineId.Rkey == "" { 690 return "", xrpc.ErrNoMatchingWorkflows 691 } 692 return pipelineId.AtUri(), nil 693} 694 695func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, error) { 696 if err := git.SparseSyncGitRepo(ctx, repoUri, repoPath, rev); err != nil { 697 return nil, fmt.Errorf("syncing git repo: %w", err) 698 } 699 gr, err := kgit.Open(repoPath, rev) 700 if err != nil { 701 return nil, fmt.Errorf("opening git repo: %w", err) 702 } 703 704 workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) 705 if errors.Is(err, object.ErrDirectoryNotFound) { 706 // return empty RawPipeline when directory doesn't exist 707 return nil, nil 708 } else if err != nil { 709 return nil, fmt.Errorf("loading file tree: %w", err) 710 } 711 712 var rawPipeline workflow.RawPipeline 713 for _, e := range workflowDir { 714 if !e.IsFile() { 715 continue 716 } 717 718 fpath := filepath.Join(workflow.WorkflowDir, e.Name) 719 contents, err := gr.RawContent(fpath) 720 if err != nil { 721 return nil, fmt.Errorf("reading raw content of '%s': %w", fpath, err) 722 } 723 724 rawPipeline = append(rawPipeline, workflow.RawWorkflow{ 725 Name: e.Name, 726 Contents: contents, 727 }) 728 } 729 730 return rawPipeline, nil 731} 732 733func (s *Spindle) StartJobWorkers(ctx context.Context) { 734 for range s.cfg.Server.MaxJobCount { 735 go func() { 736 for { 737 job, err := s.db.DequeueJob(ctx) 738 if err != nil { 739 s.l.Error("failed to dequeue job", "error", err) 740 } 741 if job == nil { 742 // sleep until a new job wakes us 743 select { 744 case <-ctx.Done(): 745 return 746 case <-s.jobWake: 747 } 748 continue 749 } 750 s.runJob(ctx, job) 751 } 752 }() 753 } 754} 755 756func (s *Spindle) runJob(_ context.Context, job *db.JobRow) { 757 pipelineId := models.PipelineId{ 758 Knot: job.PipelineIdKnot, 759 Rkey: job.PipelineIdRkey, 760 } 761 762 pipelineEnv := models.PipelineEnvVarsForSource(job.Tpl.TriggerMetadata, pipelineId, job.SourceRepo) 763 trustedSource := true 764 if tm := job.Tpl.TriggerMetadata; tm != nil && tm.SourceRepo != nil && 765 *tm.SourceRepo != "" && *tm.SourceRepo != job.RepoDid { 766 trustedSource = false 767 } 768 769 initTpl := job.Tpl 770 if job.SourceRepo != nil && job.Tpl.TriggerMetadata != nil { 771 tm := *job.Tpl.TriggerMetadata 772 tm.Repo = job.SourceRepo 773 initTpl.TriggerMetadata = &tm 774 } 775 776 workflows := make(map[models.Engine][]models.Workflow) 777 for _, w := range job.Tpl.Workflows { 778 if w == nil { 779 continue 780 } 781 eng, ok := s.engs[w.Engine] 782 if !ok { 783 _ = s.db.StatusFailed(models.WorkflowId{ 784 PipelineId: pipelineId, 785 Name: w.Name, 786 }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) 787 continue 788 } 789 790 ewf, err := eng.InitWorkflow(*w, initTpl) 791 if err != nil { 792 _ = s.db.StatusFailed(models.WorkflowId{ 793 PipelineId: pipelineId, 794 Name: w.Name, 795 }, fmt.Sprintf("init workflow: %s", err), -1, s.n) 796 continue 797 } 798 799 if ewf.Environment == nil { 800 ewf.Environment = make(map[string]string) 801 } 802 maps.Copy(ewf.Environment, pipelineEnv) 803 workflows[eng] = append(workflows[eng], *ewf) 804 } 805 806 engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.rootCtx, &models.Pipeline{ 807 RepoDid: syntax.DID(job.RepoDid), 808 Workflows: workflows, 809 TrustedSource: trustedSource, 810 }, pipelineId) 811} 812 813// enqueues the workflows in tpl. 814func (s *Spindle) processPipeline(repoDid syntax.DID, tpl tangled.Pipeline, pipelineId models.PipelineId, sourceRepo *tangled.Pipeline_TriggerRepo) error { 815 err := s.db.EnqueueJob(s.rootCtx, repoDid.String(), pipelineId, sourceRepo, tpl) 816 if err != nil { 817 return fmt.Errorf("failed to enqueue durable job: %w", err) 818 } 819 s.l.Info("pipeline enqueued successfully to db", "id", pipelineId) 820 821 // wake up an idle worker to pick up more jobs if any 822 select { 823 case s.jobWake <- struct{}{}: 824 default: 825 } 826 827 // pipelines visible from now on, they are sitting in queue 828 for _, w := range tpl.Workflows { 829 if w == nil { 830 continue 831 } 832 if err := s.db.StatusPending(models.WorkflowId{ 833 PipelineId: pipelineId, 834 Name: w.Name, 835 }, s.n); err != nil { 836 return fmt.Errorf("db.StatusPending: %w", err) 837 } 838 } 839 return nil 840} 841 842// newRepoPath creates a path to store repository by its did and rkey. 843// The path format would be: `/data/repos/did:plc:foo/sh.tangled.repo/repo-rkey 844func (s *Spindle) newRepoPath(repo syntax.DID) string { 845 return filepath.Join(s.cfg.Server.RepoDir, repo.String()) 846} 847 848func (s *Spindle) newRepoCloneUrl(knot string, did syntax.DID) string { 849 scheme := "https://" 850 if s.cfg.Server.Dev { 851 scheme = "http://" 852 } 853 return fmt.Sprintf("%s%s/%s", scheme, knot, did) 854} 855 856const RequiredVersion = "2.49.0" 857 858func ensureGitVersion() error { 859 v, err := git.Version() 860 if err != nil { 861 return fmt.Errorf("fetching git version: %w", err) 862 } 863 if v.LessThan(version.Must(version.NewVersion(RequiredVersion))) { 864 return fmt.Errorf("installed git version %q is not supported, Spindle requires git version >= %q", v, RequiredVersion) 865 } 866 return nil 867}