This repository has no description
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}