Monorepo for Tangled forked from tangled.org/core
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453package spindle
import ( "context" _ "embed" "encoding/json" "fmt" "log/slog" "maps" "net/http" "sync"
"github.com/go-chi/chi/v5" "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" "tangled.org/core/eventconsumer/cursor" "tangled.org/core/idresolver" "tangled.org/core/jetstream" "tangled.org/core/log" "tangled.org/core/notifier" "tangled.org/core/rbac" "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/models" "tangled.org/core/spindle/queue" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/xrpc" "tangled.org/core/xrpc/serviceauth")
//go:embed motdvar defaultMotd []byte
const ( rbacDomain = "thisserver")
type Spindle struct { jc *jetstream.JetstreamClient db *db.DB e *rbac.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 workflowSem chan struct{}}
// 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)
d, err := db.Make(cfg.Server.DBPath) if err != nil { return nil, fmt.Errorf("failed to setup db: %w", err) }
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) }
jq := queue.NewQueue(cfg.Server.QueueSize, cfg.Server.MaxJobCount) logger.Info("initialized queue", "queueSize", cfg.Server.QueueSize, "numWorkers", cfg.Server.MaxJobCount)
workflowSem := make(chan struct{}, cfg.Server.MaxConcurrentWorkflows) logger.Info("initialized workflow semaphore", "maxConcurrentWorkflows", cfg.Server.MaxConcurrentWorkflows)
collections := []string{ tangled.SpindleMemberNSID, tangled.RepoNSID, tangled.RepoCollaboratorNSID, } 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)
// 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) }
resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl)
spindle := &Spindle{ jc: jc, e: e, db: d, l: logger, n: &n, engs: engines, jq: jq, cfg: cfg, res: resolver, vault: vault, motd: defaultMotd, workflowSem: workflowSem, }
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) }
// for each incoming sh.tangled.pipeline, we execute // spindle.processPipeline, which in turn enqueues the pipeline // job in the above registered queue. ccfg := eventconsumer.NewConsumerConfig() ccfg.Logger = log.SubLogger(logger, "eventconsumer") ccfg.Dev = cfg.Server.Dev ccfg.ProcessFunc = spindle.processPipeline 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() *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.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) }()
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) processPipeline(ctx context.Context, src eventconsumer.Source, msg eventconsumer.Message) error { if msg.Nsid == tangled.PipelineNSID { tpl := tangled.Pipeline{} err := json.Unmarshal(msg.EventJson, &tpl) if err != nil { s.l.Error("failed to unmarshal pipeline event", "err", err) return err }
if tpl.TriggerMetadata == nil { return fmt.Errorf("no trigger metadata found") }
if tpl.TriggerMetadata.Repo == nil { return fmt.Errorf("no repo data found") }
if src.Key() != tpl.TriggerMetadata.Repo.Knot { return fmt.Errorf("repo knot does not match event source: %s != %s", src.Key(), tpl.TriggerMetadata.Repo.Knot) }
// filter by repos repoName := "" if tpl.TriggerMetadata.Repo.Repo != nil { repoName = *tpl.TriggerMetadata.Repo.Repo }
_, err = s.db.GetRepo( tpl.TriggerMetadata.Repo.Knot, tpl.TriggerMetadata.Repo.Did, repoName, ) if err != nil { return fmt.Errorf("failed to get repo: %w", err) }
pipelineId := models.PipelineId{ Knot: src.Key(), Rkey: msg.Rkey, }
workflows := make(map[models.Engine][]models.Workflow)
// Build pipeline environment variables once for all workflows pipelineEnv := models.PipelineEnvVars(tpl.TriggerMetadata, pipelineId, s.cfg.Server.Dev)
for _, w := range tpl.Workflows { if w != nil { 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 { err = s.db.StatusFailed(models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }, fmt.Sprintf("init workflow: %s", err), -1, s.n) if err != nil { return fmt.Errorf("db.StatusFailed: %w", err) }
continue }
// 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)
err = s.db.StatusPending(models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }, s.n) if err != nil { return fmt.Errorf("db.StatusPending: %w", err) } } }
ok := s.jq.Enqueue(queue.Job{ Run: func() error { engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.workflowSem, ctx, &models.Pipeline{ RepoOwner: tpl.TriggerMetadata.Repo.Did, RepoName: repoName, Workflows: workflows, }, pipelineId) return nil }, OnFail: func(jobError error) { s.l.Error("pipeline run failed", "error", jobError) }, }) if ok { s.l.Info("pipeline enqueued successfully", "id", msg.Rkey) } else { s.l.Error("failed to enqueue pipeline: queue is full") } }
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)}