This repository has no description
0

Configure Feed

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

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