Something went wrong. Try again.
Monorepo for Tangled forked from tangled.org/core
Something went wrong. Try again.
Go
at sl/rbac2test
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507package spindle
import ( "context" _ "embed" "encoding/json" "errors" "fmt" "log/slog" "maps" "net/http" "path/filepath" "sync"
"github.com/bluesky-social/indigo/atproto/syntax" "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/idresolver" kgit "tangled.org/core/knotserver/git" "tangled.org/core/log" "tangled.org/core/notifier" "tangled.org/core/rbac2" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/engines/nixery" "tangled.org/core/spindle/git" "tangled.org/core/spindle/models" "tangled.org/core/spindle/queue" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/xrpc" "tangled.org/core/tap" "tangled.org/core/tid" "tangled.org/core/workflow" "tangled.org/core/xrpc/serviceauth")
//go:embed motdvar defaultMotd []byte
type Spindle struct { tap *tap.Client db *db.DB e *rbac2.Enforcer l *slog.Logger n *notifier.Notifier engs map[string]models.Engine jq *queue.Queue cfg *config.Config ks *eventconsumer.Consumer res *idresolver.Resolver vault secrets.Manager motd []byte motdMu sync.RWMutex}
// New creates a new Spindle server with the provided configuration and engines.func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engine) (*Spindle, error) { logger := log.FromContext(ctx)
if err := ensureGitVersion(); err != nil { return nil, fmt.Errorf("ensuring git version: %w", err) }
d, err := db.Make(ctx, cfg.Server.DBPath()) if err != nil { return nil, fmt.Errorf("failed to setup db: %w", err) }
e, err := rbac2.NewEnforcer(cfg.Server.DBPath()) if err != nil { return nil, fmt.Errorf("failed to setup rbac enforcer: %w", err) }
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) }
jq := queue.NewQueue(cfg.Server.QueueSize, cfg.Server.MaxJobCount) logger.Info("initialized queue", "queueSize", cfg.Server.QueueSize, "numWorkers", cfg.Server.MaxJobCount)
tap := tap.NewClient(cfg.Server.TapUrl, "")
resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl)
spindle := &Spindle{ tap: &tap, e: e, db: d, l: logger, n: &n, engs: engines, jq: jq, cfg: cfg, res: resolver, vault: vault, motd: defaultMotd, }
err = e.SetSpindleOwner(spindle.cfg.Server.Owner, spindle.cfg.Server.Did()) 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) }
// 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.Dev = cfg.Server.Dev ccfg.ProcessFunc = spindle.processKnotStream ccfg.CursorStore = cursorStore knownKnots, err := d.Knots() if err != nil { return nil, err } for _, knot := range knownKnots { logger.Info("adding source start", "knot", knot) ccfg.Sources[eventconsumer.NewKnotSource(knot)] = struct{}{} } spindle.ks = eventconsumer.NewConsumer(*ccfg)
return spindle, nil}
// DB returns the database instance.func (s *Spindle) DB() *db.DB { return s.db}
// Queue returns the job queue instance.func (s *Spindle) Queue() *queue.Queue { return s.jq}
// 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() *rbac2.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.jq.Start() defer s.jq.Stop()
// Stop vault token renewal if it implements Stopper if stopper, ok := s.vault.(secrets.Stopper); ok { defer stopper.Stop() }
go func() { s.l.Info("starting knot event consumer") s.ks.Start(ctx) }()
// ensure server owner is tracked if err := s.tap.AddRepos(ctx, []syntax.DID{s.cfg.Server.Owner}); err != nil { return err }
go func() { s.l.Info("starting tap stream consumer") s.tap.Connect(ctx, &tap.SimpleIndexer{ EventHandler: s.processEvent, }) }()
s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr) return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router())}
func Run(ctx context.Context) error { cfg, err := config.Load(ctx) if err != nil { return fmt.Errorf("failed to load config: %w", err) }
nixeryEng, err := nixery.New(ctx, cfg) if err != nil { return err }
s, err := New(ctx, cfg, map[string]models.Engine{ "nixery": nixeryEng, }) 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, 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, }
return x.Router()}
func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Source, msg eventconsumer.Message) 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 == 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("repoDid", event.RepoDid, "repoName", event.RepoName)
// resolve repo name to rkey // TODO: git.refUpdate should respond with rkey instead of repo name repo, err := s.db.GetRepoWithName(syntax.DID(event.RepoDid), event.RepoName) if err != nil { return fmt.Errorf("get repo with did and name (%s/%s): %w", event.RepoDid, event.RepoName, err) }
// NOTE: we are blindly trusting the knot that it will return only repos it own repoCloneUri := s.newRepoCloneUrl(src.Key(), event.RepoDid, event.RepoName) repoPath := s.newRepoPath(repo.Did, repo.Rkey) if err := git.SparseSyncGitRepo(ctx, repoCloneUri, repoPath, event.NewSha); err != nil { return fmt.Errorf("sync git repo: %w", err) } l.Info("synced git repo")
compiler := workflow.Compiler{ Trigger: tangled.Pipeline_TriggerMetadata{ Kind: string(workflow.TriggerKindPush), Push: &tangled.Pipeline_PushTriggerData{ Ref: event.Ref, OldSha: event.OldSha, NewSha: event.NewSha, }, Repo: &tangled.Pipeline_TriggerRepo{ Did: repo.Did.String(), Knot: repo.Knot, Repo: repo.Name, }, }, }
// load workflow definitions from rev (without spindle context) rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, event.NewSha) if err != nil { return fmt.Errorf("loading pipeline: %w", err) } if len(rawPipeline) == 0 { l.Info("no workflow definition find for the repo. skipping the event") return nil } tpl := compiler.Compile(compiler.Parse(rawPipeline)) // TODO: 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()) }
pipelineId := models.PipelineId{ Knot: tpl.TriggerMetadata.Repo.Knot, Rkey: tid.TID(), } if err := s.db.CreatePipelineEvent(pipelineId.Rkey, tpl, s.n); err != nil { l.Error("failed to create pipeline event", "err", err) return nil } err = s.processPipeline(ctx, tpl, pipelineId) if err != nil { return err } }
return 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) processPipeline(ctx context.Context, tpl tangled.Pipeline, pipelineId models.PipelineId) error { // Build pipeline environment variables once for all workflows pipelineEnv := models.PipelineEnvVars(tpl.TriggerMetadata, pipelineId, s.cfg.Server.Dev)
// filter & init workflows workflows := make(map[models.Engine][]models.Workflow) for _, w := range tpl.Workflows { if w == nil { continue } if _, ok := s.engs[w.Engine]; !ok { err := s.db.StatusFailed(models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) if err != nil { return fmt.Errorf("db.StatusFailed: %w", err) }
continue }
eng := s.engs[w.Engine]
if _, ok := workflows[eng]; !ok { workflows[eng] = []models.Workflow{} }
ewf, err := s.engs[w.Engine].InitWorkflow(*w, tpl) if err != nil { return fmt.Errorf("init workflow: %w", err) }
// inject TANGLED_* env vars after InitWorkflow // This prevents user-defined env vars from overriding them if ewf.Environment == nil { ewf.Environment = make(map[string]string) } maps.Copy(ewf.Environment, pipelineEnv)
workflows[eng] = append(workflows[eng], *ewf) }
// enqueue pipeline ok := s.jq.Enqueue(queue.Job{ Run: func() error { engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, ctx, &models.Pipeline{ RepoOwner: tpl.TriggerMetadata.Repo.Did, RepoName: tpl.TriggerMetadata.Repo.Repo, Workflows: workflows, }, pipelineId) return nil }, OnFail: func(jobError error) { s.l.Error("pipeline run failed", "error", jobError) }, }) if !ok { return fmt.Errorf("failed to enqueue pipeline: queue is full") } s.l.Info("pipeline enqueued successfully", "id", pipelineId)
// emit StatusPending for all workflows here (after successful enqueue) for _, ewfs := range workflows { for _, ewf := range ewfs { err := s.db.StatusPending(models.WorkflowId{ PipelineId: pipelineId, Name: ewf.Name, }, s.n) if 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(did syntax.DID, rkey syntax.RecordKey) string { return filepath.Join(s.cfg.Server.RepoDir(), did.String(), tangled.RepoNSID, rkey.String())}
func (s *Spindle) newRepoCloneUrl(knot, did, name string) string { scheme := "https://" if s.cfg.Server.Dev { scheme = "http://" } return fmt.Sprintf("%s%s/%s/%s", scheme, knot, did, name)}
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}