Monorepo for Tangled tangled.org
1

Configure Feed

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

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