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 136 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 FullNetworkMode: !cfg.Server.InviteOnly, 67 CollectionFilters: []string{tangled.RepoNSID}, 68 AdminPassword: cfg.Server.Tap.AdminPassword, 69 RetryTimeout: 60 * time.Second, 70 } 71} 72 73func startEmbeddedTap(ctx context.Context, cfg *config.Config, logger *slog.Logger) (*embeddedTap, error) { 74 if err := assertLoopbackBind(cfg.Server.Tap.Bind); err != nil { 75 return nil, err 76 } 77 78 tcfg := newEmbeddedTapConfig(cfg) 79 80 t, err := tap.New(tcfg) 81 if err != nil { 82 return nil, fmt.Errorf("tap.New: %w", err) 83 } 84 85 go func() { 86 if err := t.Firehose.Run(ctx); err != nil && !errors.Is(err, context.Canceled) { 87 logger.Error("firehose terminated", "err", err) 88 } 89 }() 90 t.Run(ctx) 91 go func() { 92 logger.Info("tap http server listening", "bind", cfg.Server.Tap.Bind) 93 if err := t.Server.Start(cfg.Server.Tap.Bind); err != nil && !errors.Is(err, http.ErrServerClosed) { 94 logger.Error("tap http server terminated", "err", err) 95 } 96 }() 97 98 if err := waitForListener(ctx, cfg.Server.Tap.Bind, time.Now().Add(10*time.Second)); err != nil { 99 logger.Warn("tap http server unreachable before timeout", "bind", cfg.Server.Tap.Bind, "err", err) 100 } 101 102 return &embeddedTap{tap: t, logger: logger}, nil 103} 104 105func waitForListener(ctx context.Context, addr string, deadline time.Time) error { 106 if ctx.Err() != nil { 107 return ctx.Err() 108 } 109 if time.Now().After(deadline) { 110 return fmt.Errorf("timed out waiting for %s", addr) 111 } 112 c, err := net.DialTimeout("tcp", addr, 100*time.Millisecond) 113 if err == nil { 114 c.Close() 115 return nil 116 } 117 time.Sleep(50 * time.Millisecond) 118 return waitForListener(ctx, addr, deadline) 119} 120 121func (e *embeddedTap) Shutdown() { 122 if e == nil || e.tap == nil { 123 return 124 } 125 if e.closed.Swap(true) { 126 return 127 } 128 shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) 129 defer cancel() 130 if err := e.tap.Server.Shutdown(shutdownCtx); err != nil { 131 e.logger.Error("tap server shutdown failed", "err", err) 132 } 133 if err := e.tap.CloseDb(shutdownCtx); err != nil { 134 e.logger.Error("tap db close failed", "err", err) 135 } 136}