This repository has no description
0

Configure Feed

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

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