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 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}