server_lifecycle_test.mx raw
1 package server
2
3 import (
4 "testing"
5 "time"
6
7 "git.smesh.lol/smesh/pkg/relay/ratelimit"
8 )
9
10 // --- constructor wiring ---
11
12 var stReadyCalls int32
13
14 func stReady() { stReadyCalls++ }
15
16 func TestNewWiresPools(t *testing.T) {
17 c := stCfg()
18 c.RelayURL = "wss://relay.example.com"
19 c.IngestWorkers = 1
20 c.MediaProxyWorkers = 1
21 c.BlossomWorkers = 1
22 c.FreeWriteLimit = 5
23 c.FreeWriteWindow = 60
24 c.MaxConnPerIP = 3
25 db := stFakeDB()
26
27 s := New(db, c)
28 if s.selfHost != "relay.example.com/" {
29 t.Fatalf("selfHost = %s", s.selfHost)
30 return
31 }
32 if s.writeLimiter == nil {
33 t.Fatal("a positive free-write limit must create a limiter")
34 return
35 }
36 if s.workers.Len() != 1 || s.proxyPool.Len() != 1 || s.blossomPool.Len() != 1 {
37 t.Fatalf("pools: ingest=%d proxy=%d blossom=%d", s.workers.Len(), s.proxyPool.Len(), s.blossomPool.Len())
38 return
39 }
40 if s.bcast == nil {
41 t.Fatal("the broadcast worker must be started")
42 return
43 }
44 if len(s.proxyBusyTime) != 1 {
45 t.Fatalf("proxyBusyTime = %d, want 1", len(s.proxyBusyTime))
46 return
47 }
48 s.Close()
49
50 // Zero worker counts skip the pools entirely.
51 c2 := stCfg()
52 c2.IngestWorkers = 0
53 c2.MediaProxyWorkers = 0
54 c2.BlossomWorkers = 0
55 s2 := New(stFakeDB(), c2)
56 if s2.workers.Len() != 0 || s2.proxyPool.Len() != 0 || s2.blossomPool.Len() != 0 {
57 t.Fatal("zero worker counts must not start pools")
58 return
59 }
60 s2.Close()
61 }
62
63 func TestCloseWithStartedPools(t *testing.T) {
64 s := stServer(stCfg())
65 s.db = stFakeDB()
66 s.startIngestWorkers(1)
67 s.startProxyWorkers(1)
68 s.startBlossomWorkers(1)
69 s.startBroadcastWorker()
70
71 s.Close()
72 s.Close()
73 if !s.closed {
74 t.Fatal("Close must set the closed flag")
75 return
76 }
77 }
78
79 // --- connection-open path ---
80
81 func TestOnWSConnected(t *testing.T) {
82 c := stCfg()
83 c.RelayURL = "wss://relay.example.com"
84 s, b := stSubSrv(c)
85 defer b.Close()
86 s.OnWSConnected(7)
87 cs := s.conns[7]
88 if cs == nil {
89 t.Fatal("OnWSConnected must create the connection state")
90 return
91 }
92 if cs.subs == nil {
93 t.Fatal("OnWSConnected must initialise the subscription map")
94 return
95 }
96 if len(cs.challenge) == 0 {
97 t.Fatal("OnWSConnected must issue a NIP-42 challenge when RelayURL is set")
98 return
99 }
100
101 // No relay URL means no challenge, but the state still exists.
102 c2 := stCfg()
103 s2, b2 := stSubSrv(c2)
104 defer b2.Close()
105 s2.OnWSConnected(8)
106 if s2.conns[8] == nil {
107 t.Fatal("OnWSConnected must create state without a relay URL")
108 return
109 }
110 if s2.conns[8].challenge != nil {
111 t.Fatal("no challenge must be issued when RelayURL is empty")
112 return
113 }
114 }
115
116 func TestWireOnReadyAndOnWSMessage(t *testing.T) {
117 s := stServer(stCfg())
118 stReadyCalls = 0
119 s.OnReady = stReady
120 s.wireOnReady()
121 s.t.OnReady()
122 if stReadyCalls != 1 {
123 t.Fatalf("OnReady calls = %d, want 1", stReadyCalls)
124 return
125 }
126 // No connection for the fd: the message is dropped by dispatch.
127 s.OnWSMessage(1, []byte("[\"NOTICE\",\"x\"]"))
128 }
129
130 // --- broadcast respawn ---
131
132 func TestRespawnBroadcastWorker(t *testing.T) {
133 s, old := stSubSrv(stCfg())
134 fd := int32(4)
135 cs := &cstate{subs: map[string]*sub{}}
136 cs.authedPubkey = []byte("pubkey")
137 cs.subs["s1"] = &sub{id: "s1", rawReq: []byte("[\"REQ\",\"s1\",{}]")}
138 s.conns[fd] = cs
139
140 s.respawnBroadcastWorker()
141 old.Close()
142 if s.bcast == nil {
143 t.Fatal("respawnBroadcastWorker must install a new broadcaster")
144 return
145 }
146 s.bcast.Close()
147 }
148
149 // --- poll loop ---
150
151 func TestOnPollRunsAllDrains(t *testing.T) {
152 s, b := stSubSrv(stCfg())
153 defer b.Close()
154 s.OnPoll()
155 }
156
157 func TestOnTickReaperAndLimiterCleanup(t *testing.T) {
158 s := stServer(stCfg())
159 s.writeLimiter = ratelimit.New(1.0, 1)
160 // tickCount is a process-global; enough calls guarantee one multiple of 60.
161 for i := 0; i < 61; i++ {
162 s.OnTick()
163 }
164 if _, ok := s.asyncPending[1]; ok {
165 t.Fatal("OnTick must reap stale async entries")
166 return
167 }
168 }
169
170 func TestProxyReapStuck(t *testing.T) {
171 s := stProxySrv(1)
172 old := s.proxyIn[0]
173 s.proxyPool.Busy[0] = true
174 s.proxyBusyTime[0] = time.Now().UnixNano() - 20_000_000_000
175 s.proxyReapStuck()
176 if s.proxyPool.Busy[0] {
177 t.Fatal("a stuck worker must be respawned idle")
178 return
179 }
180 if s.proxyBusyTime[0] != 0 {
181 t.Fatal("a reaped worker must clear its busy timestamp")
182 return
183 }
184 close(old)
185 close(s.proxyIn[0])
186
187 // A freshly busy worker is not reaped.
188 s2 := stProxySrv(1)
189 s2.proxyPool.Busy[0] = true
190 s2.proxyBusyTime[0] = time.Now().UnixNano()
191 s2.proxyReapStuck()
192 if !s2.proxyPool.Busy[0] {
193 t.Fatal("a fresh worker must not be reaped")
194 return
195 }
196 }
197
198 // --- database-unavailable REQ ---
199
200 func TestHandleReqDBUnavailable(t *testing.T) {
201 s, b := stSubSrv(stCfg())
202 defer b.Close()
203 s.db = nil
204 fd := int32(4)
205 s.conns[fd] = &cstate{subs: map[string]*sub{}}
206 s.handleReq(fd, []byte("[\"REQ\",\"s1\",{}]"))
207 if len(s.conns[fd].subs) != 1 {
208 t.Fatal("the subscription must register even when the database is gone")
209 return
210 }
211 }
212