Monorepo for Tangled
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932package spindle
import ( "context" "database/sql" _ "embed" "encoding/json" "errors" "fmt" "log/slog" "maps" "net/http" "path/filepath" "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" "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" "tangled.org/core/eventconsumer/cursor" "tangled.org/core/eventstream" "tangled.org/core/idresolver" "tangled.org/core/jetstream" knotdb "tangled.org/core/knotserver/db" 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/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/engines/dummy" "tangled.org/core/spindle/engines/microvm" "tangled.org/core/spindle/engines/nixery" "tangled.org/core/spindle/git" "tangled.org/core/spindle/models" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/xrpc" "tangled.org/core/tid" "tangled.org/core/workflow" "tangled.org/core/xrpc/serviceauth")
//go:embed motdvar defaultMotd []byte
const ( rbacDomain = "thisserver")
type Spindle struct { jc *jetstream.JetstreamClient tap *Tap embedTap *embeddedTap db *db.DB e *rbac.Enforcer l *slog.Logger n *notifier.Notifier engs map[string]models.Engine cfg *config.Config ks *eventconsumer.Consumer res *idresolver.Resolver verify repoverify.Verifier vault secrets.Manager motd []byte motdMu sync.RWMutex rootCtx context.Context jobWake chan struct{}}
// 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) (*Spindle, error) { logger := log.FromContext(ctx)
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)
n := notifier.New()
var vault secrets.Manager 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") } 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", "": 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) }
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) }
collections := []string{ tangled.SpindleMemberNSID, tangled.RepoNSID, tangled.RepoCollaboratorNSID, tangled.RepoPullNSID, } 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) } jc.AddDid(cfg.Server.Owner) // pull records are created by arbitrary users too, same hack as in tap jc.ExemptCollection(tangled.RepoPullNSID)
// 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()) } }
resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl)
spindle := &Spindle{ jc: jc, e: e, db: d, l: logger, n: &n, engs: engines, cfg: cfg, res: resolver, verify: repoverify.New(resolver, cfg.Server.Dev), vault: vault, motd: defaultMotd, rootCtx: ctx, jobWake: make(chan struct{}, 1), }
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)
cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath) if err != nil { return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err) }
err = jc.StartJetstream(ctx, spindle.ingest()) if err != nil { return nil, fmt.Errorf("failed to start jetstream consumer: %w", err) }
// spindle listen to knot stream for sh.tangled.git.refUpdate // which will sync the local workflow files in spindle and enqueues the // pipeline job for on-push workflows ccfg := eventconsumer.NewConsumerConfig() ccfg.Logger = log.SubLogger(logger, "eventconsumer") ccfg.ProcessFunc = spindle.processKnotStream ccfg.CursorStore = cursorStore ccfg.WorkerCount = 16 ccfg.QueueSize = 200 if cfg.Server.Dev { ccfg.RetryInterval = 5 * time.Second ccfg.MaxRetryInterval = 10 * time.Second } else { ccfg.RetryInterval = 1 * time.Minute ccfg.MaxRetryInterval = 10 * time.Minute } knownKnots, err := d.Knots() if err != nil { return nil, err } for _, knot := range knownKnots { logger.Info("adding source start", "knot", knot) src := eventconsumer.NewKnotSource(knot) eventconsumer.MigrateLegacyCursor(cursorStore, src) ccfg.Sources[src] = struct{}{} } spindle.ks = eventconsumer.NewConsumer(*ccfg)
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)
return spindle, nil}
// DB returns the database instance.func (s *Spindle) DB() *db.DB { return s.db}
// Engines returns the map of available engines.func (s *Spindle) Engines() map[string]models.Engine { return s.engs}
// Vault returns the secrets manager instance.func (s *Spindle) Vault() secrets.Manager { return s.vault}
// Notifier returns the notifier instance.func (s *Spindle) Notifier() *notifier.Notifier { return s.n}
// Enforcer returns the RBAC enforcer instance.func (s *Spindle) Enforcer() *rbac.Enforcer { return s.e}
// SetMotdContent sets custom MOTD content, replacing the embedded default.func (s *Spindle) SetMotdContent(content []byte) { s.motdMu.Lock() defer s.motdMu.Unlock() s.motd = content}
// GetMotdContent returns the current MOTD content.func (s *Spindle) GetMotdContent() []byte { s.motdMu.RLock() defer s.motdMu.RUnlock() return s.motd}
// Start starts the Spindle server (blocking).func (s *Spindle) Start(ctx context.Context) error { // starts a job queue runner in the background s.StartJobWorkers(ctx)
// Stop vault token renewal if it implements Stopper if stopper, ok := s.vault.(secrets.Stopper); ok { defer stopper.Stop() }
tapCtx, tapCancel := context.WithCancel(ctx)
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() }
go func() { s.l.Info("starting knot event consumer") s.ks.Start(ctx) }()
s.l.Info("starting tap client", "url", s.cfg.Server.Tap.Url) s.tap.Start(tapCtx)
s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr) return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router())}
func (s *Spindle) declareTapInterest(ctx context.Context) { repos, err := s.db.AllRepos() if err != nil { s.l.Warn("tap declare: failed to load known repos", "err", err) return } seen := make(map[syntax.DID]struct{}, len(repos)) dids := make([]syntax.DID, 0, len(repos)) for _, r := range repos { if r.Owner == "" { continue } if _, ok := seen[r.Owner]; ok { continue } seen[r.Owner] = struct{}{} dids = append(dids, r.Owner) } 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: known 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) }
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) }
nixeryEng, err := nixery.New(ctx, cfg) if err != nil { return err }
microvmEng, err := microvm.New(ctx, cfg, d) if err != nil { return err }
s, err := New(ctx, cfg, d, map[string]models.Engine{ "nixery": nixeryEng, "microvm": microvmEng, "dummy": dummy.New(log.FromContext(ctx)), }) if err != nil { return err }
return s.Start(ctx)}
func (s *Spindle) Router() http.Handler { mux := chi.NewRouter()
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { w.Write(s.GetMotdContent()) }) mux.HandleFunc("/events", s.Events) mux.HandleFunc("/logs/{knot}/{rkey}/{name}", s.Logs)
mux.Mount("/xrpc", s.XrpcRouter()) return mux}
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, Resolver: s.res, Vault: s.vault, Notifier: s.Notifier(), ServiceAuth: serviceAuth, Trigger: s, }
return x.Router()}
func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) error { l := log.FromContext(ctx).With("handler", "processKnotStream") l = l.With("src", src.Key(), "msg.Nsid", msg.Nsid, "msg.Rkey", msg.Rkey) if msg.Nsid == knotdb.RepoCollaboratorUpdateNSID { return s.ingestKnotCollaborator(ctx, l, src, msg) } if msg.Nsid == tangled.GitRefUpdateNSID { event := tangled.GitRefUpdate{} if err := json.Unmarshal(msg.EventJson, &event); err != nil { l.Error("error unmarshalling", "err", err) return err } l = l.With("repo", event.Repo, "ref", event.Ref, "newSha", event.NewSha) l.Debug("debug")
repoDid := syntax.DID(event.Repo) repo, err := s.db.GetRepoByDid(repoDid) if err != nil { return fmt.Errorf("unknown repoDid %s: %w", repoDid, err) }
if src.Host != repo.Knot { return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, repo.Knot) }
if kgit.HasSkipCIPushOption(event.PushOptions) { l.Info("push event requested ci skip, skipping the event") return nil }
// NOTE: we are blindly trusting the knot that it will return only repos it own repoCloneUri := s.newRepoCloneUrl(src.Host, repoDid) repoPath := s.newRepoPath(repoDid)
triggerRepo, err := s.buildTriggerRepo(ctx, repo) if err != nil { return fmt.Errorf("building trigger repo: %w", err) }
trigger := tangled.Pipeline_TriggerMetadata{ Kind: string(workflow.TriggerKindPush), Push: &tangled.Pipeline_PushTriggerData{ Ref: event.Ref, OldSha: event.OldSha, NewSha: event.NewSha, }, Repo: triggerRepo, }
pipelineId, err := s.runPipeline(ctx, repoDid, trigger, event.ChangedFiles, repoCloneUri, repoPath, event.NewSha, nil, triggerRepo) if err != nil { return err } if pipelineId.Rkey == "" { l.Info("no workflow matched 'push' trigger, skipping the event") return nil } l.Info("pipeline triggered", "pipeline", pipelineId.AtUri()) }
return nil}
func (s *Spindle) ingestKnotCollaborator(ctx context.Context, l *slog.Logger, src eventconsumer.Source, msg eventstream.Event) error { var rec knotdb.RepoCollaboratorUpdate if err := json.Unmarshal(msg.EventJson, &rec); err != nil { l.Error("error unmarshalling collaboratorUpdate", "err", err) return err }
subject, err := syntax.ParseDID(rec.Subject) if err != nil { l.Info("skipping collaboratorUpdate with malformed subject", "subject", rec.Subject, "err", err) return nil } repoDid, err := syntax.ParseDID(rec.Repo) if err != nil { l.Info("skipping collaboratorUpdate with malformed repo", "repo", rec.Repo, "err", err) return nil }
repo, err := s.db.GetRepoByDid(repoDid) if errors.Is(err, sql.ErrNoRows) { l.Info("skipping collaboratorUpdate for unknown repo", "repo", repoDid) return nil } if err != nil { return fmt.Errorf("lookup repo %s: %w", repoDid, err) } if src.Host != repo.Knot { l.Warn("dropping collaboratorUpdate from non-owning knot", "src", src.Host, "repoKnot", repo.Knot) return nil }
switch rec.Op { case knotdb.AclOpAdd: if err := s.e.AddCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { return fmt.Errorf("add collaborator policy: %w", err) } if err := s.db.AddKnotCollaborator(repoDid, subject); err != nil { return fmt.Errorf("track collaborator: %w", err) } l.Info("added knot-managed collaborator", "subject", subject, "repo", repoDid) case knotdb.AclOpRemove: if err := s.e.RemoveCollaborator(subject.String(), rbac.ThisServer, repoDid.String()); err != nil { return fmt.Errorf("remove collaborator policy: %w", err) } if err := s.db.DeleteRepoCollaboratorBySubjectRepo(subject, repoDid); err != nil { return fmt.Errorf("delete collaborator row: %w", err) } l.Info("removed knot-managed collaborator", "subject", subject, "repo", repoDid) default: return fmt.Errorf("collaboratorUpdate unknown op %q", rec.Op) } return nil}
// buildTriggerRepo gathers trigger metadata, resolving default branch from the knotfunc (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.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)}
// resolveSourceRepoInfo resolves trigger-repo metadata for a source repo DID.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) } return s.buildTriggerRepoFrom(ctx, res.KnotURL.Host(), res.OwnerDid.String(), res.Rkey.String(), repoDid.String()), nil}
// runPipeline compiles and enqueues the pipeline for the given revision.// sourceRepo is the resolved repo the code was checked out from, forwarded to// processPipeline for env vars.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 models.PipelineId{}, fmt.Errorf("loading pipeline: %w", err) } if len(rawPipeline) == 0 { return models.PipelineId{}, 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 models.PipelineId{}, nil }
pipelineId := models.PipelineId{ Knot: trigger.Repo.Knot, Rkey: tid.TID(), } if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { return models.PipelineId{}, fmt.Errorf("creating pipeline event: %w", err) } err = s.processPipeline(repoDid, tpl, pipelineId, sourceRepo) return pipelineId, err}
// filterWorkflows filters workflows to the requested namesfunc 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) }
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 := s.newRepoCloneUrl(repo.Knot, repoDid) repoPath := s.newRepoPath(repoDid) sourceInfo := triggerRepo // default: code comes from the repo itself if sourceRepo != "" && sourceRepo != repoDid { sourceInfo, err = s.resolveSourceRepoInfo(ctx, sourceRepo) if err != nil { return "", err } sourceRepoStr := sourceRepo.String() trigger.SourceRepo = &sourceRepoStr repoCloneUri = models.BuildRepoURL(sourceInfo) repoPath = s.newRepoPath(sourceRepo) }
pipelineId, err := s.runPipeline(ctx, repoDid, trigger, nil, repoCloneUri, repoPath, sha, workflows, sourceInfo) if err != nil { return "", err } if pipelineId.Rkey == "" { return "", xrpc.ErrNoMatchingWorkflows } return pipelineId.AtUri(), nil}
func (s *Spindle) loadPipeline(ctx context.Context, repoUri, repoPath, rev string) (workflow.RawPipeline, error) { if err := git.SparseSyncGitRepo(ctx, repoUri, repoPath, rev); 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 empty RawPipeline when directory doesn't exist return nil, 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, nil}
func (s *Spindle) StartJobWorkers(ctx context.Context) { for range s.cfg.Server.MaxJobCount { go func() { 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.runJob(ctx, job) } }() }}
func (s *Spindle) runJob(ctx context.Context, job *db.JobRow) { pipelineId := models.PipelineId{ Knot: job.PipelineIdKnot, Rkey: job.PipelineIdRkey, }
pipelineEnv := models.PipelineEnvVarsForSource(job.Tpl.TriggerMetadata, pipelineId, job.SourceRepo) trustedSource := true if tm := job.Tpl.TriggerMetadata; tm != nil && tm.SourceRepo != nil && *tm.SourceRepo != "" && *tm.SourceRepo != job.RepoDid { trustedSource = false }
initTpl := job.Tpl if job.SourceRepo != nil && job.Tpl.TriggerMetadata != nil { tm := *job.Tpl.TriggerMetadata tm.Repo = job.SourceRepo initTpl.TriggerMetadata = &tm }
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 { _ = s.db.StatusFailed(models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) continue }
ewf, err := eng.InitWorkflow(*w, initTpl) if err != nil { _ = s.db.StatusFailed(models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }, fmt.Sprintf("init workflow: %s", err), -1, s.n) continue }
if ewf.Environment == nil { ewf.Environment = make(map[string]string) } maps.Copy(ewf.Environment, pipelineEnv) workflows[eng] = append(workflows[eng], *ewf) }
engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.rootCtx, &models.Pipeline{ RepoDid: syntax.DID(job.RepoDid), Workflows: workflows, TrustedSource: trustedSource, }, pipelineId)}
// enqueues the workflows in tpl.func (s *Spindle) processPipeline(repoDid syntax.DID, tpl tangled.Pipeline, pipelineId models.PipelineId, sourceRepo *tangled.Pipeline_TriggerRepo) error { err := s.db.EnqueueJob(s.rootCtx, repoDid.String(), pipelineId, sourceRepo, tpl) if err != nil { return fmt.Errorf("failed to enqueue durable job: %w", err) } s.l.Info("pipeline enqueued successfully to db", "id", pipelineId)
// wake up an idle worker to pick up more jobs if any select { case s.jobWake <- struct{}{}: default: }
// pipelines visible from now on, they are sitting in queue for _, w := range tpl.Workflows { if w == nil { continue } if err := s.db.StatusPending(models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }, s.n); err != nil { return fmt.Errorf("db.StatusPending: %w", err) } } return nil}
// newRepoPath creates a path to store repository by its did and rkey.// The path format would be: `/data/repos/did:plc:foo/sh.tangled.repo/repo-rkeyfunc (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.Server.Dev { scheme = "http://" } return fmt.Sprintf("%s%s/%s", scheme, knot, did)}
const RequiredVersion = "2.49.0"
func ensureGitVersion() error { v, err := git.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) resolvePipelineRepoDid(repo *tangled.Pipeline_TriggerRepo) (syntax.DID, error) { if repo.RepoDid == nil || *repo.RepoDid == "" { return "", fmt.Errorf("pipeline trigger missing repoDid") } repoDid, err := syntax.ParseDID(*repo.RepoDid) if err != nil { return "", fmt.Errorf("parse repoDid %s: %w", *repo.RepoDid, err) } if _, err := s.db.GetRepoByDid(repoDid); err != nil { return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) } return repoDid, 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)}