This repository has no description
1package spindle
2
3import (
4 "context"
5 _ "embed"
6 "encoding/json"
7 "fmt"
8 "log/slog"
9 "maps"
10 "net/http"
11 "sync"
12
13 "github.com/bluesky-social/indigo/atproto/syntax"
14 "github.com/go-chi/chi/v5"
15 "tangled.org/core/api/tangled"
16 "tangled.org/core/eventconsumer"
17 "tangled.org/core/eventconsumer/cursor"
18 "tangled.org/core/eventstream"
19 "tangled.org/core/idresolver"
20 "tangled.org/core/jetstream"
21 "tangled.org/core/log"
22 "tangled.org/core/notifier"
23 "tangled.org/core/rbac"
24 "tangled.org/core/spindle/config"
25 "tangled.org/core/spindle/db"
26 "tangled.org/core/spindle/engine"
27 "tangled.org/core/spindle/engines/dummy"
28 "tangled.org/core/spindle/engines/nixery"
29 "tangled.org/core/spindle/models"
30 "tangled.org/core/spindle/queue"
31 "tangled.org/core/spindle/secrets"
32 "tangled.org/core/spindle/xrpc"
33 "tangled.org/core/xrpc/serviceauth"
34)
35
36//go:embed motd
37var defaultMotd []byte
38
39const (
40 rbacDomain = "thisserver"
41)
42
43type Spindle struct {
44 jc *jetstream.JetstreamClient
45 tap *Tap
46 embedTap *embeddedTap
47 db *db.DB
48 e *rbac.Enforcer
49 l *slog.Logger
50 n *notifier.Notifier
51 engs map[string]models.Engine
52 jq *queue.Queue
53 cfg *config.Config
54 ks *eventconsumer.Consumer
55 res *idresolver.Resolver
56 vault secrets.Manager
57 motd []byte
58 motdMu sync.RWMutex
59 workflowSem chan struct{}
60 rootCtx context.Context
61}
62
63// New creates a new Spindle server with the provided configuration and engines.
64func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engine) (*Spindle, error) {
65 logger := log.FromContext(ctx)
66
67 d, err := db.Make(ctx, cfg.Server.DBPath)
68 if err != nil {
69 return nil, fmt.Errorf("failed to setup db: %w", err)
70 }
71
72 e, err := rbac.NewEnforcer(cfg.Server.DBPath)
73 if err != nil {
74 return nil, fmt.Errorf("failed to setup rbac enforcer: %w", err)
75 }
76 e.E.EnableAutoSave(true)
77
78 n := notifier.New()
79
80 var vault secrets.Manager
81 switch cfg.Server.Secrets.Provider {
82 case "openbao":
83 if cfg.Server.Secrets.OpenBao.ProxyAddr == "" {
84 return nil, fmt.Errorf("openbao proxy address is required when using openbao secrets provider")
85 }
86 vault, err = secrets.NewOpenBaoManager(
87 cfg.Server.Secrets.OpenBao.ProxyAddr,
88 logger,
89 secrets.WithMountPath(cfg.Server.Secrets.OpenBao.Mount),
90 )
91 if err != nil {
92 return nil, fmt.Errorf("failed to setup openbao secrets provider: %w", err)
93 }
94 logger.Info("using openbao secrets provider", "proxy_address", cfg.Server.Secrets.OpenBao.ProxyAddr, "mount", cfg.Server.Secrets.OpenBao.Mount)
95 case "sqlite", "":
96 vault, err = secrets.NewSQLiteManager(cfg.Server.DBPath, secrets.WithTableName("secrets"))
97 if err != nil {
98 return nil, fmt.Errorf("failed to setup sqlite secrets provider: %w", err)
99 }
100 logger.Info("using sqlite secrets provider", "path", cfg.Server.DBPath)
101 default:
102 return nil, fmt.Errorf("unknown secrets provider: %s", cfg.Server.Secrets.Provider)
103 }
104
105 if err := runStartupMigrations(ctx, d, cfg.Server.Tap.Embed, cfg.Server.Tap.DBPath, logger); err != nil {
106 return nil, fmt.Errorf("failed to run startup migrations: %w", err)
107 }
108
109 jq := queue.NewQueue(cfg.Server.QueueSize, cfg.Server.MaxJobCount)
110 logger.Info("initialized queue", "queueSize", cfg.Server.QueueSize, "numWorkers", cfg.Server.MaxJobCount)
111
112 workflowSem := make(chan struct{}, cfg.Server.MaxConcurrentWorkflows)
113 logger.Info("initialized workflow semaphore", "maxConcurrentWorkflows", cfg.Server.MaxConcurrentWorkflows)
114
115 collections := []string{
116 tangled.SpindleMemberNSID,
117 tangled.RepoNSID,
118 tangled.RepoCollaboratorNSID,
119 }
120 jc, err := jetstream.NewJetstreamClient(cfg.Server.JetstreamEndpoint, "spindle", collections, nil, log.SubLogger(logger, "jetstream"), d, true, true)
121 if err != nil {
122 return nil, fmt.Errorf("failed to setup jetstream client: %w", err)
123 }
124 jc.AddDid(cfg.Server.Owner)
125
126 // Check if the spindle knows about any Dids;
127 dids, err := d.GetAllDids()
128 if err != nil {
129 return nil, fmt.Errorf("failed to get all dids: %w", err)
130 }
131 for _, d := range dids {
132 jc.AddDid(d)
133 }
134
135 knownRepos, err := d.AllRepos()
136 if err != nil {
137 return nil, fmt.Errorf("failed to get known repos: %w", err)
138 }
139 for _, r := range knownRepos {
140 if r.Owner != "" {
141 jc.AddDid(r.Owner.String())
142 }
143 }
144
145 resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl)
146
147 spindle := &Spindle{
148 jc: jc,
149 e: e,
150 db: d,
151 l: logger,
152 n: &n,
153 engs: engines,
154 jq: jq,
155 cfg: cfg,
156 res: resolver,
157 vault: vault,
158 motd: defaultMotd,
159 workflowSem: workflowSem,
160 rootCtx: ctx,
161 }
162
163 err = e.AddSpindle(rbacDomain)
164 if err != nil {
165 return nil, fmt.Errorf("failed to set rbac domain: %w", err)
166 }
167 err = spindle.configureOwner()
168 if err != nil {
169 return nil, err
170 }
171 logger.Info("owner set", "did", cfg.Server.Owner)
172
173 cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath)
174 if err != nil {
175 return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err)
176 }
177
178 err = jc.StartJetstream(ctx, spindle.ingest())
179 if err != nil {
180 return nil, fmt.Errorf("failed to start jetstream consumer: %w", err)
181 }
182
183 // for each incoming sh.tangled.pipeline, we execute
184 // spindle.processPipeline, which in turn enqueues the pipeline
185 // job in the above registered queue.
186 ccfg := eventconsumer.NewConsumerConfig()
187 ccfg.Logger = log.SubLogger(logger, "eventconsumer")
188 ccfg.URLFunc = eventconsumer.DefaultURL(cfg.Server.Dev)
189 ccfg.ProcessFunc = spindle.processPipeline
190 ccfg.CursorStore = cursorStore
191 knownKnots, err := d.Knots()
192 if err != nil {
193 return nil, err
194 }
195 for _, knot := range knownKnots {
196 logger.Info("adding source start", "knot", knot)
197 ccfg.Sources[eventconsumer.NewKnotSource(knot)] = struct{}{}
198 }
199 spindle.ks = eventconsumer.NewConsumer(*ccfg)
200
201 if cfg.Server.Tap.Embed {
202 pw, err := randomAdminPassword()
203 if err != nil {
204 return nil, err
205 }
206 cfg.Server.Tap.AdminPassword = pw
207 logger.Info("embedded tap: using random admin password")
208 }
209 spindle.tap = NewTapClient(spindle)
210
211 return spindle, nil
212}
213
214// DB returns the database instance.
215func (s *Spindle) DB() *db.DB {
216 return s.db
217}
218
219// Queue returns the job queue instance.
220func (s *Spindle) Queue() *queue.Queue {
221 return s.jq
222}
223
224// Engines returns the map of available engines.
225func (s *Spindle) Engines() map[string]models.Engine {
226 return s.engs
227}
228
229// Vault returns the secrets manager instance.
230func (s *Spindle) Vault() secrets.Manager {
231 return s.vault
232}
233
234// Notifier returns the notifier instance.
235func (s *Spindle) Notifier() *notifier.Notifier {
236 return s.n
237}
238
239// Enforcer returns the RBAC enforcer instance.
240func (s *Spindle) Enforcer() *rbac.Enforcer {
241 return s.e
242}
243
244// SetMotdContent sets custom MOTD content, replacing the embedded default.
245func (s *Spindle) SetMotdContent(content []byte) {
246 s.motdMu.Lock()
247 defer s.motdMu.Unlock()
248 s.motd = content
249}
250
251// GetMotdContent returns the current MOTD content.
252func (s *Spindle) GetMotdContent() []byte {
253 s.motdMu.RLock()
254 defer s.motdMu.RUnlock()
255 return s.motd
256}
257
258// Start starts the Spindle server (blocking).
259func (s *Spindle) Start(ctx context.Context) error {
260 // starts a job queue runner in the background
261 s.jq.Start()
262 defer s.jq.Stop()
263
264 // Stop vault token renewal if it implements Stopper
265 if stopper, ok := s.vault.(secrets.Stopper); ok {
266 defer stopper.Stop()
267 }
268
269 tapCtx, tapCancel := context.WithCancel(ctx)
270
271 if s.cfg.Server.Tap.Embed {
272 emb, err := startEmbeddedTap(tapCtx, s.cfg, log.SubLogger(s.l, "embedtap"))
273 if err != nil {
274 tapCancel()
275 return fmt.Errorf("starting embedded tap: %w", err)
276 }
277 s.embedTap = emb
278 defer func() {
279 tapCancel()
280 s.embedTap.Shutdown()
281 }()
282
283 go s.watchTapDrain(tapCtx, tapCancel)
284 } else {
285 defer tapCancel()
286 }
287
288 go func() {
289 s.l.Info("starting knot event consumer")
290 s.ks.Start(ctx)
291 }()
292
293 s.l.Info("starting tap client", "url", s.cfg.Server.Tap.Url)
294 s.tap.Start(tapCtx)
295
296 s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr)
297 return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router())
298}
299
300func (s *Spindle) declareTapInterest(ctx context.Context) {
301 repos, err := s.db.AllRepos()
302 if err != nil {
303 s.l.Warn("tap declare: failed to load known repos", "err", err)
304 return
305 }
306 seen := make(map[syntax.DID]struct{}, len(repos))
307 dids := make([]syntax.DID, 0, len(repos))
308 for _, r := range repos {
309 if r.Owner == "" {
310 continue
311 }
312 if _, ok := seen[r.Owner]; ok {
313 continue
314 }
315 seen[r.Owner] = struct{}{}
316 dids = append(dids, r.Owner)
317 }
318 if err := s.tap.AddOwnerDIDs(ctx, dids); err != nil {
319 s.l.Warn("tap declare: AddRepos rejected", "count", len(dids), "err", err)
320 return
321 }
322 s.l.Info("tap declare: known owner DIDs registered", "count", len(dids))
323}
324
325func Run(ctx context.Context) error {
326 cfg, err := config.Load(ctx)
327 if err != nil {
328 return fmt.Errorf("failed to load config: %w", err)
329 }
330
331 nixeryEng, err := nixery.New(ctx, cfg)
332 if err != nil {
333 return err
334 }
335
336 s, err := New(ctx, cfg, map[string]models.Engine{
337 "nixery": nixeryEng,
338 "dummy": dummy.New(log.FromContext(ctx)),
339 })
340 if err != nil {
341 return err
342 }
343
344 return s.Start(ctx)
345}
346
347func (s *Spindle) Router() http.Handler {
348 mux := chi.NewRouter()
349
350 mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
351 w.Write(s.GetMotdContent())
352 })
353 mux.HandleFunc("/events", s.Events)
354 mux.HandleFunc("/logs/{knot}/{rkey}/{name}", s.Logs)
355
356 mux.Mount("/xrpc", s.XrpcRouter())
357 return mux
358}
359
360func (s *Spindle) XrpcRouter() http.Handler {
361 serviceAuth := serviceauth.NewServiceAuth(s.l, s.res, s.cfg.Server.Did().String())
362
363 l := log.SubLogger(s.l, "xrpc")
364
365 x := xrpc.Xrpc{
366 Logger: l,
367 Db: s.db,
368 Enforcer: s.e,
369 Engines: s.engs,
370 Config: s.cfg,
371 Resolver: s.res,
372 Vault: s.vault,
373 Notifier: s.Notifier(),
374 ServiceAuth: serviceAuth,
375 }
376
377 return x.Router()
378}
379
380func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, msg eventstream.Event) error {
381 if msg.Nsid == tangled.PipelineNSID {
382 tpl := tangled.Pipeline{}
383 err := json.Unmarshal(msg.EventJson, &tpl)
384 if err != nil {
385 s.l.Error("failed to unmarshal pipeline event", "err", err)
386 return err
387 }
388
389 if tpl.TriggerMetadata == nil {
390 return fmt.Errorf("no trigger metadata found")
391 }
392
393 if tpl.TriggerMetadata.Repo == nil {
394 return fmt.Errorf("no repo data found")
395 }
396
397 if src.Host != tpl.TriggerMetadata.Repo.Knot {
398 return fmt.Errorf("repo knot does not match event source: %s != %s", src.Host, tpl.TriggerMetadata.Repo.Knot)
399 }
400
401 repoDid, err := s.resolvePipelineRepoDid(tpl.TriggerMetadata.Repo)
402 if err != nil {
403 return err
404 }
405
406 pipelineId := models.PipelineId{
407 Knot: src.Host,
408 Rkey: msg.Rkey,
409 }
410
411 workflows := make(map[models.Engine][]models.Workflow)
412
413 // Build pipeline environment variables once for all workflows
414 pipelineEnv := models.PipelineEnvVars(tpl.TriggerMetadata, pipelineId, s.cfg.Server.Dev)
415
416 for _, w := range tpl.Workflows {
417 if w != nil {
418 if _, ok := s.engs[w.Engine]; !ok {
419 err = s.db.StatusFailed(models.WorkflowId{
420 PipelineId: pipelineId,
421 Name: w.Name,
422 }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n)
423 if err != nil {
424 return fmt.Errorf("db.StatusFailed: %w", err)
425 }
426
427 continue
428 }
429
430 eng := s.engs[w.Engine]
431
432 if _, ok := workflows[eng]; !ok {
433 workflows[eng] = []models.Workflow{}
434 }
435
436 ewf, err := s.engs[w.Engine].InitWorkflow(*w, tpl)
437 if err != nil {
438 err = s.db.StatusFailed(models.WorkflowId{
439 PipelineId: pipelineId,
440 Name: w.Name,
441 }, fmt.Sprintf("init workflow: %s", err), -1, s.n)
442 if err != nil {
443 return fmt.Errorf("db.StatusFailed: %w", err)
444 }
445
446 continue
447 }
448
449 // inject TANGLED_* env vars after InitWorkflow
450 // This prevents user-defined env vars from overriding them
451 if ewf.Environment == nil {
452 ewf.Environment = make(map[string]string)
453 }
454 maps.Copy(ewf.Environment, pipelineEnv)
455
456 workflows[eng] = append(workflows[eng], *ewf)
457
458 err = s.db.StatusPending(models.WorkflowId{
459 PipelineId: pipelineId,
460 Name: w.Name,
461 }, s.n)
462 if err != nil {
463 return fmt.Errorf("db.StatusPending: %w", err)
464 }
465 }
466 }
467
468 ok := s.jq.Enqueue(queue.Job{
469 Run: func() error {
470 engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.workflowSem, ctx, &models.Pipeline{
471 RepoDid: repoDid,
472 Workflows: workflows,
473 }, pipelineId)
474 return nil
475 },
476 OnFail: func(jobError error) {
477 s.l.Error("pipeline run failed", "error", jobError)
478 },
479 })
480 if ok {
481 s.l.Info("pipeline enqueued successfully", "id", msg.Rkey)
482 } else {
483 s.l.Error("failed to enqueue pipeline: queue is full")
484 }
485 }
486
487 return nil
488}
489
490func (s *Spindle) resolvePipelineRepoDid(repo *tangled.Pipeline_TriggerRepo) (syntax.DID, error) {
491 if repo.RepoDid == nil || *repo.RepoDid == "" {
492 return "", fmt.Errorf("pipeline trigger missing repoDid")
493 }
494 repoDid, err := syntax.ParseDID(*repo.RepoDid)
495 if err != nil {
496 return "", fmt.Errorf("parse repoDid %s: %w", *repo.RepoDid, err)
497 }
498 if _, err := s.db.GetRepoByDid(repoDid); err != nil {
499 return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err)
500 }
501 return repoDid, nil
502}
503
504func (s *Spindle) configureOwner() error {
505 cfgOwner := s.cfg.Server.Owner
506
507 existing, err := s.e.GetSpindleUsersByRole("server:owner", rbacDomain)
508 if err != nil {
509 return err
510 }
511
512 switch len(existing) {
513 case 0:
514 // no owner configured, continue
515 case 1:
516 // find existing owner
517 existingOwner := existing[0]
518
519 // no ownership change, this is okay
520 if existingOwner == s.cfg.Server.Owner {
521 break
522 }
523
524 // remove existing owner
525 err = s.e.RemoveSpindleOwner(rbacDomain, existingOwner)
526 if err != nil {
527 return nil
528 }
529 default:
530 return fmt.Errorf("more than one owner in DB, try deleting %q and starting over", s.cfg.Server.DBPath)
531 }
532
533 return s.e.AddSpindleOwner(rbacDomain, cfgOwner)
534}