This repository has no description
0

Configure Feed

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

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