package spindle import ( "context" "crypto/sha256" _ "embed" "encoding/binary" "errors" "fmt" "io" "log/slog" "maps" "net/http" "path/filepath" "sort" "sync" "time" "github.com/bluesky-social/indigo/atproto/syntax" indigoxrpc "github.com/bluesky-social/indigo/xrpc" "github.com/go-chi/chi/v5" "github.com/go-git/go-git/v5/plumbing/object" "github.com/hashicorp/go-version" "go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "tangled.org/core/api/tangled" "tangled.org/core/gitutil" "tangled.org/core/idresolver" "tangled.org/core/jetstream" "tangled.org/core/knotfeed" kgit "tangled.org/core/knotserver/git" "tangled.org/core/log" "tangled.org/core/notifier" "tangled.org/core/rbac" "tangled.org/core/repoident" "tangled.org/core/repoverify" "tangled.org/core/spindle/artifactstore" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/engines/dummy" "tangled.org/core/spindle/engines/nixery" "tangled.org/core/spindle/feed" "tangled.org/core/spindle/mill" "tangled.org/core/spindle/mill/executor" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability" pipelinecodec "tangled.org/core/spindle/pipeline" "tangled.org/core/spindle/quota" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/storage" "tangled.org/core/spindle/webhook" "tangled.org/core/spindle/xrpc" "tangled.org/core/tid" "tangled.org/core/workflow" "tangled.org/core/xrpc/serviceauth" ) //go:embed motd var defaultMotd []byte const ( rbacDomain = "thisserver" executorShutdownTimeout = 8 * time.Minute serverShutdownTimeout = 10 * time.Second // sparse-checkout wants a leading slash to anchor the pattern at the root sparseWorkflowDir = "/" + workflow.WorkflowDir ) type executorClient interface { Connect(context.Context) Drain(context.Context) error RegisterMetrics(*observability.Metrics) SetQuotaClient(*executor.QuotaClient) } type Spindle struct { jc *jetstream.JetstreamClient tap *Tap embedTap *embeddedTap db *db.DB e *rbac.Enforcer l *slog.Logger n *notifier.Notifier wh *webhook.Service metrics *observability.Metrics engs map[string]models.Engine jobWake chan struct{} jobWorkers sync.WaitGroup scheduler *pipelineScheduler cfg *config.Config feed *feed.Feed res *idresolver.Resolver verify repoverify.Verifier vault secrets.Manager cache storage.Storage motd []byte motdMu sync.RWMutex rootCtx context.Context rootCancel context.CancelFunc store artifactstore.Store stores *artifactstore.Stores reader artifactstore.Reader qm *quota.Manager // set only when this spindle hosts the mill or joins one as an executor mill *mill.Mill exec executorClient listCollaborators func(ctx context.Context, knot string, repoDid syntax.DID) ([]syntax.DID, error) listRefRecords func(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) } func newCacheStore(ctx context.Context, cfg *config.Config) (storage.Storage, error) { switch cfg.Cache.Backend { case "": return nil, nil case "disk": dir := cfg.Cache.DiskDir if dir == "" { dir = filepath.Join(filepath.Dir(cfg.Server.DBPath), "cache") } return storage.NewDisk(dir) case "s3": return storage.NewUnversionedS3(ctx, cfg.Cache.S3Bucket, cfg.Cache.S3Prefix) default: return nil, fmt.Errorf("storage: unknown backend %q", cfg.Cache.Backend) } } // New creates a new Spindle server with the provided configuration and engines. func New(ctx context.Context, cfg *config.Config, d *db.DB, engines map[string]models.Engine, qm *quota.Manager) (*Spindle, error) { logger := log.FromContext(ctx) metrics := observability.GetMetrics(ctx) if metrics == nil { metrics = observability.NewMetrics() ctx = observability.WithMetrics(ctx, metrics) } n := notifier.New() lifecycleCtx, cancelLifecycle := context.WithCancel(context.WithoutCancel(ctx)) keepLifecycle := false defer func() { if !keepLifecycle { cancelLifecycle() } }() if cfg.Role == config.RoleExecutor { if err := cleanupOrphanRepos(ctx, d, logger); err != nil { return nil, fmt.Errorf("failed to run startup cleanup: %w", err) } } else if err := runStartupMigrations(ctx, d, cfg.Server.Tap.Embed, cfg.Server.Tap.DBPath, logger); err != nil { return nil, fmt.Errorf("failed to run startup migrations: %w", err) } cacheStore, err := newCacheStore(ctx, cfg) if err != nil { return nil, fmt.Errorf("failed to setup cache storage: %w", err) } if cacheStore != nil { logger.Info("cache storage enabled", "backend", cfg.Cache.Backend, "storeID", cfg.Cache.StoreID) } spindle := &Spindle{ db: d, l: logger, n: &n, wh: webhook.New(d, cfg.Server.Dev), metrics: metrics, engs: engines, cfg: cfg, cache: cacheStore, motd: defaultMotd, rootCtx: lifecycleCtx, rootCancel: cancelLifecycle, jobWake: make(chan struct{}, 1), qm: qm, } metrics.AttachDB(ctx, d) diskFallback := "" if cfg.Role == config.RoleStandalone { diskFallback = cfg.Server.LogDir if cfg.ArtifactStores.Disk.Dir == "" { logger.Warn("using SPINDLE_SERVER_LOG_DIR as the implicit disk artifact store; configure SPINDLE_ARTIFACT_STORES_DISK_DIR explicitly") } } stores, err := artifactstore.NewStores(cfg.ArtifactStores, diskFallback, cfg.LegacyS3.LogBucket) if err != nil { return nil, fmt.Errorf("failed to setup artifact stores: %w", err) } spindle.stores = stores if cfg.LegacyS3.LogBucket != "" { logger.Warn("SPINDLE_S3_LOG_BUCKET is deprecated; use SPINDLE_ARTIFACT_STORES_S3_BUCKET") } if cfg.Role != config.RoleExecutor { engine.StartCachePruner(ctx, logger, d, cacheStore, cfg.Cache.Retention, cfg.Cache.PruneInterval) } if cfg.Role == config.RoleStandalone { spindle.reader = stores } else { name := cfg.Mill.ArtifactStore if name == "" { names := stores.Names() if len(names) != 1 { return nil, fmt.Errorf("%s requires SPINDLE_MILL_ARTIFACT_STORE when %d artifact stores are configured", cfg.Role, len(names)) } name = names[0] logger.Warn("SPINDLE_MILL_ARTIFACT_STORE is not set; inferred the only configured store", "store", name) } store, ok := stores.Store(name) if !ok { return nil, fmt.Errorf("SPINDLE_MILL_ARTIFACT_STORE=%q is not configured", name) } spindle.store = store spindle.reader = store } if cfg.Role == config.RoleExecutor { keepLifecycle = true return spindle, nil } e, err := rbac.NewEnforcer(cfg.Server.DBPath) if err != nil { return nil, fmt.Errorf("failed to setup rbac enforcer: %w", err) } e.E.EnableAutoSave(true) spindle.e = e switch cfg.Server.Secrets.Provider { case "openbao": if cfg.Server.Secrets.OpenBao.ProxyAddr == "" { return nil, fmt.Errorf("openbao proxy address is required when using openbao secrets provider") } spindle.vault, err = secrets.NewOpenBaoManager( cfg.Server.Secrets.OpenBao.ProxyAddr, logger, secrets.WithMountPath(cfg.Server.Secrets.OpenBao.Mount), ) if err != nil { return nil, fmt.Errorf("failed to setup openbao secrets provider: %w", err) } logger.Info("using openbao secrets provider", "proxy_address", cfg.Server.Secrets.OpenBao.ProxyAddr, "mount", cfg.Server.Secrets.OpenBao.Mount) case "sqlite", "": spindle.vault, err = secrets.NewSQLiteManager(cfg.Server.DBPath, secrets.WithTableName("secrets")) if err != nil { return nil, fmt.Errorf("failed to setup sqlite secrets provider: %w", err) } logger.Info("using sqlite secrets provider", "path", cfg.Server.DBPath) default: return nil, fmt.Errorf("unknown secrets provider: %s", cfg.Server.Secrets.Provider) } collections := []string{ tangled.SpindleMemberNSID, tangled.RepoNSID, tangled.RepoCollaboratorNSID, tangled.RepoPullNSID, tangled.RepoPullStatusNSID, } jc, err := jetstream.NewJetstreamClient(cfg.Server.JetstreamEndpoint, "spindle", collections, nil, log.SubLogger(logger, "jetstream"), d, true, true) if err != nil { return nil, fmt.Errorf("failed to setup jetstream client: %w", err) } spindle.jc = jc jc.AddDid(cfg.Server.Owner) // pull (status) records are created by arbitrary users too, same hack as in tap jc.ExemptCollection(tangled.RepoPullNSID) jc.ExemptCollection(tangled.RepoPullStatusNSID) // Check if the spindle knows about any Dids; dids, err := d.GetAllDids() if err != nil { return nil, fmt.Errorf("failed to get all dids: %w", err) } for _, d := range dids { jc.AddDid(d) } knownRepos, err := d.AllRepos() if err != nil { return nil, fmt.Errorf("failed to get known repos: %w", err) } for _, r := range knownRepos { if r.Owner != "" { jc.AddDid(r.Owner.String()) } } spindle.res = idresolver.DefaultResolver(cfg.Server.PlcUrl) spindle.verify = repoverify.New(spindle.res, cfg.Server.Dev) spindle.scheduler = newPipelineScheduler(spindle) err = e.AddSpindle(rbacDomain) if err != nil { return nil, fmt.Errorf("failed to set rbac domain: %w", err) } err = spindle.configureOwner() if err != nil { return nil, err } logger.Info("owner set", "did", cfg.Server.Owner) err = jc.StartJetstream(lifecycleCtx, spindle.ingest()) if err != nil { return nil, fmt.Errorf("failed to start jetstream consumer: %w", err) } spindle.feed = feed.New(log.SubLogger(logger, "knotfeed"), feed.Hooks{ LoadCursor: func(ctx context.Context, knot string) (knotfeed.Cursor, error) { cursor, err := spindle.db.LoadFeedCursor(knot) if errors.Is(err, knotfeed.ErrUnrecognizedFeed) { logger.Warn("stored feed didn't parse, so resuming live", "knot", knot, "err", err) return knotfeed.Cursor{}, nil } return cursor, err }, StoreCursor: func(ctx context.Context, knot string, cursor knotfeed.Cursor) error { return spindle.db.StoreFeedCursor(knot, cursor) }, Handle: spindle.handleKnotFeed, OutdatedReplay: spindle.feedOutdatedReplay, OnConnectError: func(knot string, err error) { logger.Warn("cannot reach knot firehose", "knot", knot, "err", err) }, }) if cfg.Server.Tap.Embed { pw, err := randomAdminPassword() if err != nil { return nil, err } cfg.Server.Tap.AdminPassword = pw logger.Info("embedded tap: using random admin password") } spindle.tap = NewTapClient(spindle) keepLifecycle = true return spindle, nil } func (s *Spindle) DB() *db.DB { return s.db } func (s *Spindle) Engines() map[string]models.Engine { return s.engs } func (s *Spindle) Vault() secrets.Manager { return s.vault } func (s *Spindle) Notifier() *notifier.Notifier { return s.n } func (s *Spindle) Enforcer() *rbac.Enforcer { return s.e } func (s *Spindle) SetMotdContent(content []byte) { s.motdMu.Lock() defer s.motdMu.Unlock() s.motd = content } func (s *Spindle) GetMotdContent() []byte { s.motdMu.RLock() defer s.motdMu.RUnlock() return s.motd } // runs the server. blocks func (s *Spindle) Start(ctx context.Context) error { runCtx := s.rootCtx cancelRun := s.rootCancel if runCtx == nil || cancelRun == nil { runCtx, cancelRun = context.WithCancel(context.WithoutCancel(ctx)) s.rootCtx = runCtx s.rootCancel = cancelRun } defer cancelRun() _, err := observability.StartMetricsServer(runCtx, s.cfg.Server.MetricsListenAddr, s.l, s.metrics.Registry()) if err != nil { return fmt.Errorf("starting metrics listener: %w", err) } if s.scheduler != nil { s.scheduler.Start(runCtx) } // an executor dials out to its mill and takes work from it var execDone chan struct{} if s.exec != nil { execDone = make(chan struct{}) go func() { defer close(execDone) s.exec.Connect(runCtx) }() } waitForExecutor := func(waitCtx context.Context) error { if execDone == nil { return nil } select { case <-execDone: return nil case <-waitCtx.Done(): return waitCtx.Err() } } if s.mill != nil && s.cfg.Mill.JumpListenAddr != "" { go s.mill.ServeJump(runCtx, s.cfg.Mill.JumpListenAddr, s.cfg.Mill.JumpHostKeyPath, s.cfg.Mill.DebugExecutorPort, s.cfg.Mill.MaxJumpConnections) } if stopper, ok := s.vault.(secrets.Stopper); ok { defer stopper.Stop() } if s.cfg.Role != config.RoleExecutor { tapCtx, tapCancel := context.WithCancel(runCtx) if s.cfg.Server.Tap.Embed { emb, err := startEmbeddedTap(tapCtx, s.cfg, log.SubLogger(s.l, "embedtap")) if err != nil { tapCancel() return fmt.Errorf("starting embedded tap: %w", err) } s.embedTap = emb defer func() { tapCancel() s.embedTap.Shutdown() }() go s.watchTapDrain(tapCtx, tapCancel) } else { defer tapCancel() } // tap must be ready before owner wipes start s.resumeWipes(ctx) go func() { s.l.Info("starting knot firehose feed") knots, err := s.db.Knots() if err != nil { s.l.Error("listing known knots for the feed", "err", err) return } s.feed.Start(runCtx, knots) }() go s.reconcileCollaboratorsLoop(ctx) s.l.Info("starting tap client", "url", s.cfg.Server.Tap.Url) s.tap.Start(tapCtx) } // start workers after resuming wipes so banned jobs cannot run workersCtx, cancelWorkers := context.WithCancel(runCtx) s.StartJobWorkers(workersCtx) defer func() { cancelWorkers() done := make(chan struct{}) go func() { s.jobWorkers.Wait() close(done) }() select { case <-done: case <-time.After(5 * time.Second): s.l.Warn("timed out waiting for job workers to stop") } }() server := &http.Server{ Addr: s.cfg.Server.ListenAddr, Handler: s.Router(), } serverErr := make(chan error, 1) go func() { s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr) serverErr <- server.ListenAndServe() }() var listenErr error select { case listenErr = <-serverErr: case <-ctx.Done(): } if s.exec != nil { drainCtx, cancelDrain := context.WithTimeout(context.Background(), s.cfg.Mill.DrainTimeout) if err := s.exec.Drain(drainCtx); err != nil { s.l.Warn("executor drain did not complete", "err", err) } cancelDrain() } if s.cfg.Role != config.RoleExecutor { webhookCtx, cancelWebhooks := context.WithTimeout(context.Background(), serverShutdownTimeout) s.wh.Close(webhookCtx) cancelWebhooks() } cancelRun() waitCtx, cancelWait := context.WithTimeout(context.Background(), executorShutdownTimeout) if err := waitForExecutor(waitCtx); err != nil { s.l.Warn("executor shutdown did not complete", "err", err) } cancelWait() shutdownCtx, cancelShutdown := context.WithTimeout(context.Background(), serverShutdownTimeout) shutdownErr := server.Shutdown(shutdownCtx) cancelShutdown() if listenErr == nil { listenErr = <-serverErr } if errors.Is(listenErr, http.ErrServerClosed) { listenErr = nil } if shutdownErr != nil { shutdownErr = fmt.Errorf("shutting down spindle server: %w", shutdownErr) } return errors.Join(listenErr, shutdownErr) } func (s *Spindle) declareTapInterest(ctx context.Context) { repos, reposErr := s.db.AllRepos() if reposErr != nil { s.l.Warn("tap declare: failed to load known repos", "err", reposErr) } members, membersErr := s.db.GetAllDids() if membersErr != nil { s.l.Warn("tap declare: failed to load known members", "err", membersErr) } seen := make(map[syntax.DID]struct{}, len(repos)+len(members)+1) dids := make([]syntax.DID, 0, len(repos)+len(members)+1) add := func(did syntax.DID) { if did == "" { return } if _, ok := seen[did]; ok { return } seen[did] = struct{}{} dids = append(dids, did) } owner, err := syntax.ParseDID(s.cfg.Server.Owner) if err != nil { s.l.Warn("tap declare: invalid configured owner", "owner", s.cfg.Server.Owner, "err", err) } else { add(owner) } for _, member := range members { did, err := syntax.ParseDID(member) if err != nil { s.l.Warn("tap declare: invalid known member", "did", member, "err", err) continue } add(did) } for _, repo := range repos { add(repo.Owner) } sort.Slice(dids, func(i, j int) bool { return dids[i] < dids[j] }) if err := s.tap.AddOwnerDIDs(ctx, dids); err != nil { s.l.Warn("tap declare: AddRepos rejected", "count", len(dids), "err", err) return } s.l.Info("tap declare: owner DIDs registered", "count", len(dids)) } func Run(ctx context.Context) error { cfg, err := config.Load(ctx) if err != nil { return fmt.Errorf("failed to load config: %w", err) } shutdownTracing, err := observability.InitTracing(ctx, cfg.Tracing) if err != nil { return fmt.Errorf("failed to initialize tracing: %w", err) } defer func() { shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if err := shutdownTracing(shutdownCtx); err != nil { log.FromContext(ctx).Error("failed to shut down tracing", "err", err) } }() otelHandler, shutdownLogging, err := observability.InitLogging(ctx, cfg.Logging) if err != nil { return fmt.Errorf("failed to initialize logging: %w", err) } defer func() { shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if err := shutdownLogging(shutdownCtx); err != nil { log.FromContext(ctx).Error("failed to shut down logging", "err", err) } }() currentLogger := log.FromContext(ctx) combinedHandler := log.NewFanoutHandler(currentLogger.Handler(), otelHandler) combinedLogger := slog.New(combinedHandler) slog.SetDefault(combinedLogger) ctx = log.IntoContext(ctx, combinedLogger) metrics := observability.NewMetrics() ctx = observability.WithMetrics(ctx, metrics) if err := ensureGitVersion(); err != nil { return fmt.Errorf("ensuring git version: %w", err) } d, err := db.Make(ctx, cfg.Server.DBPath) if err != nil { return fmt.Errorf("failed to setup db: %w", err) } if err := d.MigratePipelineLogFiles(cfg.Server.LogDir); err != nil { return fmt.Errorf("failed to migrate pipeline log files: %w", err) } var qs *db.QuotaStore var qm *quota.Manager if cfg.Role != config.RoleExecutor { defaults, err := cfg.ToDefaults() if err != nil { return fmt.Errorf("invalid quota configuration: %w", err) } for scope, resources := range defaults { for resource, limit := range resources { metrics.SetQuotaDefaultLimit(string(scope), string(resource), limit) } } qs = db.NewQuotaStore(d, defaults) qm = quota.NewManager(qs, 50*time.Millisecond, metrics.QuotaObserver()) defer qm.Close() if cfg.Role != config.RoleMill { if err := qs.Recover(ctx, nil); err != nil { return fmt.Errorf("failed to recover quota reservations: %w", err) } } } logger := log.FromContext(ctx) var engines map[string]models.Engine var m *mill.Mill var quotaClient *executor.QuotaClient if cfg.Role == config.RoleMill { // on a mill host, engines place jobs on executors instead of running // them. all names share one Mill m = mill.New(log.SubLogger(logger, "mill"), mill.Config{ LogDir: cfg.Server.LogDir, MaxPending: cfg.Mill.MaxPending, ReconnectGrace: cfg.Mill.ReconnectGrace, CancelAckTimeout: cfg.Mill.CancelAckTimeout, CancelTeardownTimeout: cfg.Mill.CancelTeardownTimeout, CacheStoreID: cfg.Cache.StoreID, CacheMaxBytesPerOwner: cfg.Cache.MaxBytesPerOwner, }) engines = map[string]models.Engine{ "nixery": mill.NewEngine("nixery", m), "microvm": mill.NewEngine("microvm", m), "dummy": mill.NewEngine("dummy", m), } } else { // standalone and executor both run real engines locally nixeryEng, err := nixery.New(ctx, cfg) if err != nil { return err } var engineQStore quota.ReservationStore if cfg.Role == config.RoleExecutor { quotaClient = executor.NewQuotaClient() engineQStore = nil } else { engineQStore = qs } microvmEng, err := newMicrovmEngine(ctx, cfg, d, engineQStore) if err != nil { return err } engines = map[string]models.Engine{ "nixery": nixeryEng, "microvm": microvmEng, "dummy": dummy.New(logger), } } s, err := New(ctx, cfg, d, engines, qm) if err != nil { return err } if cfg.Role == config.RoleMill || cfg.Role == config.RoleStandalone { // flattens the store snapshot into bounded label tuples, the collector // drops anything outside the bounded scope, resource and status sets metrics.SetQuotaLoader(func() (observability.QuotaSnapshot, error) { snap, err := qs.MetricsSnapshot(context.Background()) if err != nil { return observability.QuotaSnapshot{}, err } var out observability.QuotaSnapshot for scope, resources := range snap.Usage { for resource, used := range resources { out.Usage = append(out.Usage, observability.QuotaUsage{ Scope: string(scope), Resource: resource, Used: used, }) } } for scope, resources := range snap.Subjects { for resource, statuses := range resources { for status, count := range statuses { out.Subjects = append(out.Subjects, observability.QuotaSubjectCount{ Scope: string(scope), Resource: string(resource), Status: status, Count: count, }) } } } return out, nil }) } if m != nil { // the engines built above hold the mill, but the mill's db and // notifier only exist after New, so attach them here m.Attach(s.DB(), s.Notifier(), qm) m.AttachCache(s.cache) s.mill = m s.mill.RegisterMetrics(s.metrics) if err := m.RestoreState(); err != nil { return fmt.Errorf("restoring mill state: %w", err) } liveIDs := m.LiveQuotaReservationIDs() for i := range liveIDs { liveIDs[i] = quota.StorageReservationID(liveIDs[i]) } if err := qs.Recover(ctx, liveIDs); err != nil { return fmt.Errorf("failed to recover mill quota reservations: %w", err) } } if cfg.Role == config.RoleExecutor { s.exec, err = executor.New(cfg, engines, s.DB(), s.Notifier(), log.SubLogger(logger, "executor"), s.store, s.cache) if err != nil { return err } s.exec.SetQuotaClient(quotaClient) s.exec.RegisterMetrics(s.metrics) } return s.Start(ctx) } func (s *Spindle) Router() http.Handler { mux := chi.NewRouter() mux.Use(observability.HTTPMiddleware(s.metrics)) if s.cfg.Tracing.Endpoint != "" { mux.Use(observability.OTelRouteMiddleware) } mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { w.Write(s.GetMotdContent()) }) mux.Get("/_health", func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") _, _ = io.WriteString(w, "{\"status\":\"ok\"}\n") }) if s.cfg.Role != config.RoleExecutor { if s.mill != nil { mux.HandleFunc("/mill", s.mill.HandleExecutorConn) } mux.Mount("/xrpc", s.XrpcRouter()) } if s.cfg.Tracing.Endpoint == "" { return mux } return otelhttp.NewHandler( mux, "HTTP", otelhttp.WithSpanNameFormatter(func(string, *http.Request) string { return "HTTP" }), ) } func (s *Spindle) XrpcRouter() http.Handler { serviceAuth := serviceauth.NewServiceAuth(s.l, s.res.Directory(), s.cfg.Server.Did().String()) l := log.SubLogger(s.l, "xrpc") x := xrpc.Xrpc{ Logger: l, Db: s.db, Enforcer: s.e, Engines: s.engs, Config: s.cfg, ArtifactReader: s.reader, Resolver: s.res, Vault: s.vault, Notifier: s.Notifier(), Webhooks: s.wh, ServiceAuth: serviceAuth, Trigger: s, QuotaStore: db.NewQuotaStore(s.db, quota.Defaults{}), Wiper: s, } return x.Router() } // buildTriggerRepo gathers trigger metadata, resolving default branch from the knot func (s *Spindle) buildTriggerRepo(ctx context.Context, repo *db.Repo) (*tangled.Pipeline_TriggerRepo, error) { rkey := string(repo.Rkey) repoDid := repo.RepoDid.String() return s.buildTriggerRepoFrom(ctx, repo.Knot, repo.Owner.String(), rkey, repoDid), nil } func (s *Spindle) buildTriggerRepoFrom(ctx context.Context, knot, did, rkey, repoDid string) *tangled.Pipeline_TriggerRepo { scheme := "https" if s.cfg != nil && s.cfg.Server.Dev { scheme = "http" } client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, knot)} // this should maybe (?) be in the refUpdate event itself to save a roundtrip defaultBranch := "" if out, err := tangled.RepoGetDefaultBranch(ctx, client, repoDid); err == nil { defaultBranch = out.Name } var rkeyPtr *string if rkey != "" { rkeyPtr = &rkey } return &tangled.Pipeline_TriggerRepo{ Did: did, Knot: knot, Repo: rkeyPtr, RepoDid: &repoDid, DefaultBranch: defaultBranch, } } func (s *Spindle) resolvePipelineSourceRepo(ctx context.Context, trigger *tangled.Pipeline_TriggerMetadata) (*tangled.Pipeline_TriggerRepo, error) { if trigger == nil { return nil, nil } if trigger.SourceRepo == nil || *trigger.SourceRepo == "" { return trigger.Repo, nil } repoDid, err := syntax.ParseDID(*trigger.SourceRepo) if err != nil { return nil, fmt.Errorf("parse sourceRepo %s: %w", *trigger.SourceRepo, err) } return s.resolveSourceRepoInfo(ctx, repoDid) } func (s *Spindle) resolveSourceRepoInfo(ctx context.Context, repoDid syntax.DID) (*tangled.Pipeline_TriggerRepo, error) { repo, err := s.db.GetRepoByDid(repoDid) if err == nil { return s.buildTriggerRepo(ctx, repo) } // verify repo, we don't want git sync to point to arbitrary endpoints res, err := s.verify(ctx, repoident.RepoDid(repoDid)) if err != nil { return nil, fmt.Errorf("verify sourceRepo %s: %w", repoDid, err) } ownership, ok := res.Ownership() if !ok { return nil, fmt.Errorf("verify sourceRepo %s: knot %s answered %s", repoDid, res.KnotURL, res.Answer()) } return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host(), ownership.OwnerDid.String(), ownership.Rkey.String(), repoDid.String()), nil } func (s *Spindle) createPipeline(id models.PipelineId, raw tangled.Pipeline) error { record, err := pipelinecodec.FromTangled(id, time.Now(), raw) if err != nil { return err } return s.db.CreatePipeline(record) } func (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) { l := log.FromContext(ctx) compiler := workflow.Compiler{ ChangedFiles: changedFiles, Trigger: trigger, } rawPipeline, _, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev) if err != nil { return "", fmt.Errorf("loading pipeline: %w", err) } if len(rawPipeline) == 0 { return "", nil } tpl := compiler.Compile(compiler.Parse(rawPipeline)) // todo(dawn): pass compile error to workflow log for _, w := range compiler.Diagnostics.Errors { l.Error(w.String()) } for _, w := range compiler.Diagnostics.Warnings { l.Warn(w.String()) } if len(only) > 0 { tpl.Workflows = filterWorkflows(tpl.Workflows, only) } if len(tpl.Workflows) == 0 { return "", nil } pipelineId := models.PipelineId(tid.TID()) if err := s.createPipeline(pipelineId, tpl); err != nil { return "", fmt.Errorf("creating pipeline: %w", err) } err = s.processPipeline(ctx, repoDid, tpl, pipelineId, sourceRepo) return pipelineId, err } func filterWorkflows(workflows []*tangled.Pipeline_Workflow, only []string) []*tangled.Pipeline_Workflow { allowed := make(map[string]struct{}, len(only)) for _, n := range only { allowed[n] = struct{}{} } var filtered []*tangled.Pipeline_Workflow for _, w := range workflows { if w == nil { continue } if _, ok := allowed[w.Name]; ok { filtered = append(filtered, w) } } return filtered } // TriggerManual dispatches a pipeline at sha, authorized against and recorded // under repoDid. sourceRepo, pull, and inputs are optional trigger payload. func (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) { repo, err := s.db.GetRepoByDid(repoDid) if err != nil { return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) } if ban, err := s.db.IsBanned(repoDid, repo.Owner); err != nil { return "", fmt.Errorf("checking bans: %w", err) } else if ban != nil { return "", fmt.Errorf("repo %s is banned", repoDid) } triggerRepo, err := s.buildTriggerRepo(ctx, repo) if err != nil { return "", fmt.Errorf("building trigger repo: %w", err) } trigger := tangled.Pipeline_TriggerMetadata{Repo: triggerRepo} if pull.IsPullRequest { var pullAt *string if pull.Pull != "" { pullAtStr := pull.Pull.String() pullAt = &pullAtStr } trigger.Kind = string(workflow.TriggerKindPullRequest) trigger.PullRequest = &tangled.Pipeline_PullRequestTriggerData{ SourceBranch: pull.SourceBranch, TargetBranch: pull.TargetBranch, SourceSha: sha, Pull: pullAt, } } else { var refPtr *string if ref != "" { refPtr = &ref } trigger.Kind = string(workflow.TriggerKindManual) trigger.Manual = &tangled.Pipeline_ManualTriggerData{ Sha: sha, Ref: refPtr, Inputs: inputs, } } repoCloneUri, repoPath, sourceInfo, err := s.resolveCheckout(ctx, repoDid, sourceRepo) if err != nil { return "", err } if sourceInfo == nil { sourceInfo = triggerRepo } else { // reject a banned source before running its workflow sourceOwner, _ := syntax.ParseDID(sourceInfo.Did) if ban, err := s.db.IsBanned(sourceRepo, sourceOwner); err != nil { return "", fmt.Errorf("checking bans: %w", err) } else if ban != nil { return "", fmt.Errorf("source repo %s is banned", sourceRepo) } sourceRepoStr := sourceRepo.String() trigger.SourceRepo = &sourceRepoStr } pipelineId, err := s.runPipeline(ctx, repoDid, trigger, nil, repoCloneUri, repoPath, sha, workflows, sourceInfo) if err != nil { return "", err } if pipelineId == "" { return "", xrpc.ErrNoMatchingWorkflows } return syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", serviceauth.DidWeb(repo.Knot), tangled.PipelineNSID, pipelineId)), nil } // sourceInfo is nil when the checkout comes from the target repo. func (s *Spindle) resolveCheckout(ctx context.Context, repoDid syntax.DID, sourceRepo syntax.DID) (cloneUri, repoPath string, sourceInfo *tangled.Pipeline_TriggerRepo, err error) { repo, err := s.db.GetRepoByDid(repoDid) if err != nil { return "", "", nil, fmt.Errorf("unknown repoDid %s: %w", repoDid, err) } cloneUri = s.newRepoCloneUrl(repo.Knot, repoDid) repoPath = s.newRepoPath(repoDid) if sourceRepo != "" && sourceRepo != repoDid { sourceInfo, err = s.resolveSourceRepoInfo(ctx, sourceRepo) if err != nil { return "", "", nil, err } cloneUri = models.BuildRepoURL(sourceInfo) repoPath = s.newRepoPath(sourceRepo) } return cloneUri, repoPath, sourceInfo, nil } // resolves the workflow definition at sha without executing it // returns a deterministic fingerprint over the resolved files. func (s *Spindle) ListWorkflowDefinitions(ctx context.Context, repoDid syntax.DID, ref string) ([]*models.WorkflowDefinition, error) { if ref == "" { ref = "HEAD" } repoCloneURI, repoPath, _, err := s.resolveCheckout(ctx, repoDid, "") if err != nil { return nil, err } rawPipeline, commit, err := s.loadPipeline(ctx, repoCloneURI, repoPath, ref) if err != nil { return nil, fmt.Errorf("loading pipeline: %w", err) } out := make([]*models.WorkflowDefinition, 0, len(rawPipeline)) for _, raw := range rawPipeline { definition, err := pipelinecodec.FileDefinition(raw.Name, raw.Contents, repoDid.String(), commit) if err != nil { return nil, fmt.Errorf("parsing workflow %s: %w", raw.Name, err) } out = append(out, definition) } return out, nil } func (s *Spindle) DescribeWorkflowDefinition(ctx context.Context, repoDid syntax.DID, sha string, sourceRepo syntax.DID) (*tangled.CiDescribeWorkflowDefinition_Output, error) { repoCloneUri, repoPath, _, err := s.resolveCheckout(ctx, repoDid, sourceRepo) if err != nil { return nil, err } rawPipeline, _, err := s.loadPipeline(ctx, repoCloneUri, repoPath, sha) if err != nil { return nil, fmt.Errorf("loading pipeline: %w", err) } hash := fingerprintWorkflowDefinition(rawPipeline) workflows := make([]string, 0, len(rawPipeline)) for _, w := range rawPipeline { workflows = append(workflows, w.Name) } return &tangled.CiDescribeWorkflowDefinition_Output{ Derived: true, Hash: &hash, Workflows: workflows, }, nil } func fingerprintWorkflowDefinition(rawPipeline workflow.RawPipeline) string { sorted := make([]workflow.RawWorkflow, len(rawPipeline)) copy(sorted, rawPipeline) sort.Slice(sorted, func(i, j int) bool { return sorted[i].Name < sorted[j].Name }) h := sha256.New() var lenBuf [8]byte for _, w := range sorted { binary.LittleEndian.PutUint64(lenBuf[:], uint64(len(w.Contents))) h.Write([]byte(w.Name)) // terminate name to avoid ["fo", "o"] == ["f", "oo"] h.Write([]byte{0}) h.Write(lenBuf[:]) h.Write(w.Contents) } return fmt.Sprintf("sha256:%x", h.Sum(nil)) } func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, string, error) { if err := gitutil.SparseSync(ctx, repoUri, repoPath, rev, sparseWorkflowDir); err != nil { return nil, "", fmt.Errorf("syncing git repo: %w", err) } gr, err := kgit.Open(repoPath, rev) if err != nil { return nil, "", fmt.Errorf("opening git repo: %w", err) } workflowDir, err := gr.FileTree(ctx, workflow.WorkflowDir) if errors.Is(err, object.ErrDirectoryNotFound) { return nil, gr.Hash().String(), nil } else if err != nil { return nil, "", fmt.Errorf("loading file tree: %w", err) } var rawPipeline workflow.RawPipeline for _, e := range workflowDir { if !e.IsFile() { continue } fpath := filepath.Join(workflow.WorkflowDir, e.Name) contents, err := gr.RawContent(fpath) if err != nil { return nil, "", fmt.Errorf("reading raw content of '%s': %w", fpath, err) } rawPipeline = append(rawPipeline, workflow.RawWorkflow{ Name: e.Name, Contents: contents, }) } return rawPipeline, gr.Hash().String(), nil } func (s *Spindle) newRepoPath(repo syntax.DID) string { return filepath.Join(s.cfg.Server.RepoDir, repo.String()) } func (s *Spindle) newRepoCloneUrl(knot string, did syntax.DID) string { scheme := "https://" if s.cfg != nil && s.cfg.Server.Dev { scheme = "http://" } return fmt.Sprintf("%s%s/%s", scheme, knot, did) } const RequiredVersion = "2.49.0" func ensureGitVersion() error { v, err := gitutil.Version() if err != nil { return fmt.Errorf("fetching git version: %w", err) } if v.LessThan(version.Must(version.NewVersion(RequiredVersion))) { return fmt.Errorf("installed git version %q is not supported, Spindle requires git version >= %q", v, RequiredVersion) } return nil } func (s *Spindle) configureOwner() error { cfgOwner := s.cfg.Server.Owner existing, err := s.e.GetSpindleUsersByRole("server:owner", rbacDomain) if err != nil { return err } switch len(existing) { case 0: // no owner configured, continue case 1: // find existing owner existingOwner := existing[0] // no ownership change, this is okay if existingOwner == s.cfg.Server.Owner { break } // remove existing owner err = s.e.RemoveSpindleOwner(rbacDomain, existingOwner) if err != nil { return nil } default: return fmt.Errorf("more than one owner in DB, try deleting %q and starting over", s.cfg.Server.DBPath) } return s.e.AddSpindleOwner(rbacDomain, cfgOwner) } func (s *Spindle) StartJobWorkers(ctx context.Context) { for range s.cfg.Server.MaxJobCount { s.jobWorkers.Add(1) go func() { defer s.jobWorkers.Done() for { job, err := s.db.DequeueJob(ctx) if err != nil { s.l.Error("failed to dequeue job", "error", err) } if job == nil { // sleep until a new job wakes us select { case <-ctx.Done(): return case <-s.jobWake: } continue } s.metrics.RecordJobQueueActivity("dequeue") s.runJob(ctx, job) } }() } } func (s *Spindle) runJob(ctx context.Context, job *db.JobRow) { pipelineId := job.PipelineId repoDID := job.RepoDid if job.SourceRepo != nil && job.SourceRepo.RepoDid != nil && *job.SourceRepo.RepoDid != "" { repoDID = *job.SourceRepo.RepoDid } jobCtx := observability.ExtractFromTraceparentAndTracestate(ctx, job.Traceparent, job.Tracestate) if job.CreatedAtNs > 0 { queueDelay := time.Since(time.Unix(0, job.CreatedAtNs)) if queueDelay < 0 { queueDelay = 0 } s.metrics.RecordJobDequeueLatency(jobCtx, queueDelay) } jobCtx, span := observability.Tracer().Start(jobCtx, "job.run") if span.IsRecording() { attrs := []attribute.KeyValue{ attribute.Int64(observability.JobIDKey, job.Id), attribute.String(observability.PipelineIDKey, pipelineId.String()), } if repoDID != "" { attrs = append(attrs, attribute.String(observability.RepoDIDKey, repoDID)) } if job.RepoDid != "" && job.RepoDid != repoDID { attrs = append(attrs, attribute.String(observability.TargetRepoDIDKey, job.RepoDid)) } span.SetAttributes(attrs...) } failed := false defer func() { if failed { span.SetStatus(codes.Error, "job failed") } else { span.SetStatus(codes.Ok, "success") } span.End() }() l := log.SubLogger(log.SubLogger(s.l, "job"), "engine").With( "job_id", job.Id, ) if job.RepoDid != "" && job.RepoDid != repoDID { l = l.With("target_repo_did", job.RepoDid) } pipelineEnv := models.PipelineEnvVarsForSource(job.Tpl.TriggerMetadata, pipelineId, job.SourceRepo) initTpl := job.Tpl if job.SourceRepo != nil && job.Tpl.TriggerMetadata != nil { tm := *job.Tpl.TriggerMetadata tm.Repo = job.SourceRepo initTpl.TriggerMetadata = &tm } trustedSource := models.TrustedPipelineSource(initTpl.TriggerMetadata, job.RepoDid) workflows := make(map[models.Engine][]models.Workflow) for _, w := range job.Tpl.Workflows { if w == nil { continue } eng, ok := s.engs[w.Engine] if !ok { failed = true _ = s.db.StatusFailed(models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) s.metrics.RecordWorkflowTerminal( fmt.Sprint(w.Engine), "failure", string(engine.FailureClassUser), string(engine.FailureReasonWorkflowInvalid), ) continue } ewf, err := eng.InitWorkflow(*w, initTpl) if err != nil { failed = true _ = s.db.StatusFailed(models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }, fmt.Sprintf("init workflow: %s", err), -1, s.n) failureClass, failureReason := engine.FailureAttribution("failure", err) s.metrics.RecordWorkflowTerminal( fmt.Sprint(w.Engine), "failure", failureClass, failureReason, ) continue } ewf.RunID = fmt.Sprintf("%d", job.Id) if ewf.Environment == nil { ewf.Environment = make(map[string]string) } maps.Copy(ewf.Environment, pipelineEnv) ewf.Engine = w.Engine ewf.RepoDID = job.RepoDid if !trustedSource { ewf.Caches = nil } workflows[eng] = append(workflows[eng], *ewf) } engine.StartWorkflows(l, s.vault, s.cfg, s.qm, s.stores, s.db, s.n, s.cache, nil, jobCtx, &models.Pipeline{ RepoDid: syntax.DID(job.RepoDid), Workflows: workflows, TrustedSource: trustedSource, TriggerMetadata: job.Tpl.TriggerMetadata, }, pipelineId) } func (s *Spindle) processPipeline(ctx context.Context, repoDid syntax.DID, tpl tangled.Pipeline, pipelineId models.PipelineId, sourceRepo *tangled.Pipeline_TriggerRepo) error { traceparent, tracestate := observability.InjectToTraceparentAndTracestate(ctx) if err := s.db.EnqueueJobWithPending( s.rootCtx, repoDid.String(), pipelineId, sourceRepo, tpl, traceparent, tracestate, s.n, ); err != nil { return fmt.Errorf("failed to enqueue durable job: %w", err) } s.l.Info("pipeline enqueued successfully to db", "id", pipelineId) select { case s.jobWake <- struct{}{}: default: } return nil }