Skip to content

Commit 38f3b51

Browse files
committed
Test coverage push: 72% → 76%
- Unskip 6 nameserver integration tests (server, client, persistence, concurrency) - Add 5 IPC round-trip tests: SetHostname, SetVisibility, Deregister, SetWebhook, Disconnect - Add 6 beacon/registry tests: register, list, validation, TTL filtering, stale reaping - Fix Driver.Disconnect deadlock: cmdCloseOK handler now dispatches to sendAndWait - Add SetClock/Reap to registry Server for deterministic time-based testing
1 parent effd506 commit 38f3b51

5 files changed

Lines changed: 482 additions & 10 deletions

File tree

pkg/driver/ipc.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -130,6 +130,16 @@ func (c *ipcClient) readLoop() {
130130
}
131131
c.recvMu.Unlock()
132132
}
133+
// Also dispatch to sendAndWait handlers (for Driver.Disconnect)
134+
c.mu.Lock()
135+
if chs, ok := c.handlers[cmd]; ok && len(chs) > 0 {
136+
ch := chs[0]
137+
c.handlers[cmd] = chs[1:]
138+
c.mu.Unlock()
139+
ch <- append([]byte(nil), payload...)
140+
} else {
141+
c.mu.Unlock()
142+
}
133143
case cmdRecvFrom:
134144
// Datagram: [6-byte src_addr][2-byte src_port][2-byte dst_port][data]
135145
if len(payload) >= protocol.AddrSize+4 {

pkg/registry/server.go

Lines changed: 21 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,9 @@ type Server struct {
103103
// Prometheus metrics
104104
metrics *registryMetrics
105105

106+
// Clock (overridable for testing)
107+
now func() time.Time
108+
106109
// Shutdown
107110
done chan struct{}
108111
}
@@ -329,6 +332,7 @@ func NewWithStore(beaconAddr, storePath string) *Server {
329332
done: make(chan struct{}),
330333
saveCh: make(chan struct{}, 1),
331334
saveDone: make(chan struct{}),
335+
now: time.Now,
332336
}
333337

334338
go s.saveLoop()
@@ -384,6 +388,19 @@ func (s *Server) SetAdminToken(token string) {
384388
s.mu.Unlock()
385389
}
386390

391+
// SetClock overrides the time source for testing.
392+
func (s *Server) SetClock(fn func() time.Time) {
393+
s.mu.Lock()
394+
s.now = fn
395+
s.mu.Unlock()
396+
}
397+
398+
// Reap triggers stale node and beacon cleanup (for testing).
399+
func (s *Server) Reap() {
400+
s.reapStaleNodes()
401+
s.reapStaleBeacons()
402+
}
403+
387404
// SetReplicationToken sets the token required for subscribe_replication (H4 fix).
388405
// If empty, replication subscription is disabled.
389406
func (s *Server) SetReplicationToken(token string) {
@@ -520,7 +537,7 @@ func (s *Server) reapLoop() {
520537
}
521538

522539
func (s *Server) reapStaleNodes() {
523-
threshold := time.Now().Add(-staleNodeThreshold)
540+
threshold := s.now().Add(-staleNodeThreshold)
524541
s.mu.Lock()
525542
defer s.mu.Unlock()
526543

@@ -554,7 +571,7 @@ func (s *Server) reapStaleNodes() {
554571
}
555572

556573
func (s *Server) reapStaleBeacons() {
557-
now := time.Now()
574+
now := s.now()
558575
s.mu.Lock()
559576
defer s.mu.Unlock()
560577
for id, b := range s.beacons {
@@ -1953,7 +1970,7 @@ func (s *Server) handleBeaconRegister(msg map[string]interface{}) (map[string]in
19531970
s.beacons[beaconID] = &beaconEntry{
19541971
ID: beaconID,
19551972
Addr: addr,
1956-
LastSeen: time.Now(),
1973+
LastSeen: s.now(),
19571974
}
19581975
s.mu.Unlock()
19591976

@@ -1970,7 +1987,7 @@ func (s *Server) handleBeaconList() (map[string]interface{}, error) {
19701987
s.mu.RLock()
19711988
defer s.mu.RUnlock()
19721989

1973-
now := time.Now()
1990+
now := s.now()
19741991
beacons := make([]map[string]interface{}, 0, len(s.beacons))
19751992
for _, b := range s.beacons {
19761993
if now.Sub(b.LastSeen) > beaconTTL {

tests/beacon_registry_test.go

Lines changed: 256 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,256 @@
1+
package tests
2+
3+
import (
4+
"testing"
5+
"time"
6+
7+
"github.com/TeoSlayer/pilotprotocol/internal/crypto"
8+
"github.com/TeoSlayer/pilotprotocol/pkg/registry"
9+
)
10+
11+
// TestBeaconRegisterAndList verifies beacon registration and listing via the registry.
12+
func TestBeaconRegisterAndList(t *testing.T) {
13+
t.Parallel()
14+
rc, _, cleanup := startTestRegistry(t)
15+
defer cleanup()
16+
17+
// Register a beacon
18+
resp, err := rc.Send(map[string]interface{}{
19+
"type": "beacon_register",
20+
"beacon_id": float64(42),
21+
"addr": "10.0.0.1:9001",
22+
})
23+
if err != nil {
24+
t.Fatalf("beacon_register: %v", err)
25+
}
26+
if resp["type"] != "beacon_register_ok" {
27+
t.Fatalf("expected beacon_register_ok, got %v", resp["type"])
28+
}
29+
30+
// Register a second beacon
31+
_, err = rc.Send(map[string]interface{}{
32+
"type": "beacon_register",
33+
"beacon_id": float64(99),
34+
"addr": "10.0.0.2:9001",
35+
})
36+
if err != nil {
37+
t.Fatalf("beacon_register second: %v", err)
38+
}
39+
40+
// List beacons
41+
resp, err = rc.Send(map[string]interface{}{"type": "beacon_list"})
42+
if err != nil {
43+
t.Fatalf("beacon_list: %v", err)
44+
}
45+
if resp["type"] != "beacon_list_ok" {
46+
t.Fatalf("expected beacon_list_ok, got %v", resp["type"])
47+
}
48+
49+
beacons, ok := resp["beacons"].([]interface{})
50+
if !ok || len(beacons) != 2 {
51+
t.Fatalf("expected 2 beacons, got %v", resp["beacons"])
52+
}
53+
t.Logf("beacons: %v", beacons)
54+
}
55+
56+
// TestBeaconRegisterValidation verifies that beacon_register rejects missing fields.
57+
func TestBeaconRegisterValidation(t *testing.T) {
58+
t.Parallel()
59+
rc, _, cleanup := startTestRegistry(t)
60+
defer cleanup()
61+
62+
// Missing beacon_id
63+
_, err := rc.Send(map[string]interface{}{
64+
"type": "beacon_register",
65+
"addr": "10.0.0.1:9001",
66+
})
67+
if err == nil {
68+
t.Fatal("expected error for missing beacon_id")
69+
}
70+
71+
// Missing addr
72+
_, err = rc.Send(map[string]interface{}{
73+
"type": "beacon_register",
74+
"beacon_id": float64(1),
75+
})
76+
if err == nil {
77+
t.Fatal("expected error for missing addr")
78+
}
79+
}
80+
81+
// TestBeaconListFiltersExpired verifies that expired beacons are excluded from list.
82+
func TestBeaconListFiltersExpired(t *testing.T) {
83+
t.Parallel()
84+
85+
clk := newTestClock()
86+
reg := registry.New("127.0.0.1:9001")
87+
reg.SetClock(clk.Now)
88+
go reg.ListenAndServe(":0")
89+
select {
90+
case <-reg.Ready():
91+
case <-time.After(5 * time.Second):
92+
t.Fatal("registry failed to start")
93+
}
94+
defer reg.Close()
95+
96+
rc, err := registry.Dial(reg.Addr().String())
97+
if err != nil {
98+
t.Fatalf("dial: %v", err)
99+
}
100+
defer rc.Close()
101+
102+
// Register beacon at current clock time
103+
_, err = rc.Send(map[string]interface{}{
104+
"type": "beacon_register",
105+
"beacon_id": float64(1),
106+
"addr": "10.0.0.1:9001",
107+
})
108+
if err != nil {
109+
t.Fatalf("register: %v", err)
110+
}
111+
112+
// List — should see 1 beacon
113+
resp, err := rc.Send(map[string]interface{}{"type": "beacon_list"})
114+
if err != nil {
115+
t.Fatalf("list before expire: %v", err)
116+
}
117+
beacons := resp["beacons"].([]interface{})
118+
if len(beacons) != 1 {
119+
t.Fatalf("expected 1 beacon before expiry, got %d", len(beacons))
120+
}
121+
122+
// Advance clock past beacon TTL (60s)
123+
clk.Advance(90 * time.Second)
124+
125+
// List — beacon should be filtered out
126+
resp, err = rc.Send(map[string]interface{}{"type": "beacon_list"})
127+
if err != nil {
128+
t.Fatalf("list after expire: %v", err)
129+
}
130+
beacons = resp["beacons"].([]interface{})
131+
if len(beacons) != 0 {
132+
t.Fatalf("expected 0 beacons after expiry, got %d", len(beacons))
133+
}
134+
}
135+
136+
// TestReapStaleNodes verifies that nodes without heartbeats are reaped.
137+
func TestReapStaleNodes(t *testing.T) {
138+
t.Parallel()
139+
140+
clk := newTestClock()
141+
reg := registry.New("127.0.0.1:9001")
142+
reg.SetClock(clk.Now)
143+
go reg.ListenAndServe(":0")
144+
select {
145+
case <-reg.Ready():
146+
case <-time.After(5 * time.Second):
147+
t.Fatal("registry failed to start")
148+
}
149+
defer reg.Close()
150+
151+
rc, err := registry.Dial(reg.Addr().String())
152+
if err != nil {
153+
t.Fatalf("dial: %v", err)
154+
}
155+
defer rc.Close()
156+
157+
// Register a node
158+
id, _ := crypto.GenerateIdentity()
159+
resp, err := rc.RegisterWithKey("127.0.0.1:4000", crypto.EncodePublicKey(id.PublicKey), "")
160+
if err != nil {
161+
t.Fatalf("register: %v", err)
162+
}
163+
nodeID := uint32(resp["node_id"].(float64))
164+
165+
// Verify node exists
166+
_, err = rc.Lookup(nodeID)
167+
if err != nil {
168+
t.Fatalf("lookup before reap: %v", err)
169+
}
170+
171+
// Advance clock past stale threshold (3 minutes)
172+
clk.Advance(4 * time.Minute)
173+
174+
// Trigger reap
175+
reg.Reap()
176+
177+
// Node should be gone
178+
_, err = rc.Lookup(nodeID)
179+
if err == nil {
180+
t.Fatal("expected lookup to fail after reap, node still exists")
181+
}
182+
t.Logf("correctly reaped: %v", err)
183+
}
184+
185+
// TestReapStaleBeacons verifies that stale beacons are removed by reap.
186+
func TestReapStaleBeacons(t *testing.T) {
187+
t.Parallel()
188+
189+
clk := newTestClock()
190+
reg := registry.New("127.0.0.1:9001")
191+
reg.SetClock(clk.Now)
192+
go reg.ListenAndServe(":0")
193+
select {
194+
case <-reg.Ready():
195+
case <-time.After(5 * time.Second):
196+
t.Fatal("registry failed to start")
197+
}
198+
defer reg.Close()
199+
200+
rc, err := registry.Dial(reg.Addr().String())
201+
if err != nil {
202+
t.Fatalf("dial: %v", err)
203+
}
204+
defer rc.Close()
205+
206+
// Register beacon
207+
_, err = rc.Send(map[string]interface{}{
208+
"type": "beacon_register",
209+
"beacon_id": float64(1),
210+
"addr": "10.0.0.1:9001",
211+
})
212+
if err != nil {
213+
t.Fatalf("register: %v", err)
214+
}
215+
216+
// Advance clock past beacon TTL
217+
clk.Advance(90 * time.Second)
218+
219+
// Trigger reap
220+
reg.Reap()
221+
222+
// Beacon list should be empty (both reap cleanup and list filter agree)
223+
resp, err := rc.Send(map[string]interface{}{"type": "beacon_list"})
224+
if err != nil {
225+
t.Fatalf("list after reap: %v", err)
226+
}
227+
beacons := resp["beacons"].([]interface{})
228+
if len(beacons) != 0 {
229+
t.Fatalf("expected 0 beacons after reap, got %d", len(beacons))
230+
}
231+
}
232+
233+
// TestRegistryPunch verifies the punch message handler.
234+
func TestRegistryPunch(t *testing.T) {
235+
t.Parallel()
236+
rc, _, cleanup := startTestRegistry(t)
237+
defer cleanup()
238+
239+
// Register two nodes with endpoints
240+
nodeA, idA := registerTestNode(t, rc)
241+
nodeB, _ := registerTestNode(t, rc)
242+
243+
// Punch requires signature from requester
244+
setClientSigner(rc, idA)
245+
resp, err := rc.Punch(nodeA, nodeA, nodeB)
246+
if err != nil {
247+
t.Fatalf("punch: %v", err)
248+
}
249+
if resp["type"] != "punch_ok" {
250+
t.Fatalf("expected punch_ok, got %v", resp["type"])
251+
}
252+
if resp["node_a_addr"] == nil || resp["node_b_addr"] == nil {
253+
t.Fatal("expected both node addresses in punch response")
254+
}
255+
t.Logf("punch: A=%v B=%v", resp["node_a_addr"], resp["node_b_addr"])
256+
}

0 commit comments

Comments
 (0)