This repository has no description
0

Configure Feed

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

core / spindle / embedtap.go
3.7 kB 135 lines
1package spindle 2 3import ( 4 "context" 5 "crypto/rand" 6 "encoding/hex" 7 "errors" 8 "fmt" 9 "log/slog" 10 "net" 11 "net/http" 12 "strings" 13 "sync/atomic" 14 "time" 15 16 "github.com/bluesky-social/indigo/service/tap" 17 "tangled.org/core/api/tangled" 18 "tangled.org/core/spindle/config" 19) 20 21func randomAdminPassword() (string, error) { 22 var b [32]byte 23 if _, err := rand.Read(b[:]); err != nil { 24 return "", fmt.Errorf("generate tap admin password: %w", err) 25 } 26 return hex.EncodeToString(b[:]), nil 27} 28 29func assertLoopbackBind(bind string) error { 30 host, _, err := net.SplitHostPort(bind) 31 if err != nil { 32 return fmt.Errorf("parse tap bind %q: %w", bind, err) 33 } 34 if host == "" { 35 return fmt.Errorf("embedded mode requires loopback host in tap bind %q", bind) 36 } 37 if strings.EqualFold(host, "localhost") { 38 return nil 39 } 40 ip := net.ParseIP(host) 41 if ip == nil || !ip.IsLoopback() { 42 return fmt.Errorf("embedded tap bind %q must be loopback like 127.0.0.1 or ::1", bind) 43 } 44 return nil 45} 46 47type embeddedTap struct { 48 tap *tap.Tap 49 logger *slog.Logger 50 closed atomic.Bool 51} 52 53func newEmbeddedTapConfig(cfg *config.Config) tap.Config { 54 return tap.Config{ 55 DatabaseURL: "sqlite://" + cfg.Server.Tap.DBPath, 56 DBMaxConns: 32, 57 PLCURL: cfg.Server.PlcUrl, 58 RelayUrl: cfg.Server.Tap.RelayUrl, 59 FirehoseParallelism: 4, 60 ResyncParallelism: 2, 61 OutboxParallelism: 1, 62 FirehoseCursorSaveInterval: time.Second, 63 RepoFetchTimeout: 5 * time.Minute, 64 IdentityCacheSize: 50_000, 65 EventCacheSize: 10_000, 66 CollectionFilters: []string{tangled.RepoNSID, tangled.RepoCollaboratorNSID}, 67 AdminPassword: cfg.Server.Tap.AdminPassword, 68 RetryTimeout: 60 * time.Second, 69 } 70} 71 72func startEmbeddedTap(ctx context.Context, cfg *config.Config, logger *slog.Logger) (*embeddedTap, error) { 73 if err := assertLoopbackBind(cfg.Server.Tap.Bind); err != nil { 74 return nil, err 75 } 76 77 tcfg := newEmbeddedTapConfig(cfg) 78 79 t, err := tap.New(tcfg) 80 if err != nil { 81 return nil, fmt.Errorf("tap.New: %w", err) 82 } 83 84 go func() { 85 if err := t.Firehose.Run(ctx); err != nil && !errors.Is(err, context.Canceled) { 86 logger.Error("firehose terminated", "err", err) 87 } 88 }() 89 t.Run(ctx) 90 go func() { 91 logger.Info("tap http server listening", "bind", cfg.Server.Tap.Bind) 92 if err := t.Server.Start(cfg.Server.Tap.Bind); err != nil && !errors.Is(err, http.ErrServerClosed) { 93 logger.Error("tap http server terminated", "err", err) 94 } 95 }() 96 97 if err := waitForListener(ctx, cfg.Server.Tap.Bind, time.Now().Add(10*time.Second)); err != nil { 98 logger.Warn("tap http server unreachable before timeout", "bind", cfg.Server.Tap.Bind, "err", err) 99 } 100 101 return &embeddedTap{tap: t, logger: logger}, nil 102} 103 104func waitForListener(ctx context.Context, addr string, deadline time.Time) error { 105 if ctx.Err() != nil { 106 return ctx.Err() 107 } 108 if time.Now().After(deadline) { 109 return fmt.Errorf("timed out waiting for %s", addr) 110 } 111 c, err := net.DialTimeout("tcp", addr, 100*time.Millisecond) 112 if err == nil { 113 c.Close() 114 return nil 115 } 116 time.Sleep(50 * time.Millisecond) 117 return waitForListener(ctx, addr, deadline) 118} 119 120func (e *embeddedTap) Shutdown() { 121 if e == nil || e.tap == nil { 122 return 123 } 124 if e.closed.Swap(true) { 125 return 126 } 127 shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) 128 defer cancel() 129 if err := e.tap.Server.Shutdown(shutdownCtx); err != nil { 130 e.logger.Error("tap server shutdown failed", "err", err) 131 } 132 if err := e.tap.CloseDb(shutdownCtx); err != nil { 133 e.logger.Error("tap db close failed", "err", err) 134 } 135}