This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / server.go
14 kB 534 lines
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}