This repository has no description
0

Configure Feed

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

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