|
| 1 | +// SPDX-License-Identifier: AGPL-3.0-or-later |
| 2 | + |
| 3 | +package beacon |
| 4 | + |
| 5 | +import ( |
| 6 | + "encoding/binary" |
| 7 | + "fmt" |
| 8 | + "net" |
| 9 | + "net/http" |
| 10 | + "testing" |
| 11 | + "time" |
| 12 | + |
| 13 | + "github.com/TeoSlayer/pilotprotocol/pkg/protocol" |
| 14 | +) |
| 15 | + |
| 16 | +// helper: send a discover message to register a node with a beacon |
| 17 | +func registerNode(t *testing.T, beaconAddr *net.UDPAddr, nodeID uint32) *net.UDPConn { |
| 18 | + t.Helper() |
| 19 | + conn, err := net.DialUDP("udp", nil, beaconAddr) |
| 20 | + if err != nil { |
| 21 | + t.Fatalf("dial beacon: %v", err) |
| 22 | + } |
| 23 | + |
| 24 | + msg := make([]byte, 5) |
| 25 | + msg[0] = protocol.BeaconMsgDiscover |
| 26 | + binary.BigEndian.PutUint32(msg[1:5], nodeID) |
| 27 | + if _, err := conn.Write(msg); err != nil { |
| 28 | + t.Fatalf("send discover: %v", err) |
| 29 | + } |
| 30 | + |
| 31 | + // Read discover reply |
| 32 | + buf := make([]byte, 64) |
| 33 | + conn.SetReadDeadline(time.Now().Add(2 * time.Second)) |
| 34 | + n, err := conn.Read(buf) |
| 35 | + if err != nil { |
| 36 | + t.Fatalf("read discover reply: %v", err) |
| 37 | + } |
| 38 | + if n < 1 || buf[0] != protocol.BeaconMsgDiscoverReply { |
| 39 | + t.Fatalf("unexpected reply type: 0x%02x", buf[0]) |
| 40 | + } |
| 41 | + |
| 42 | + return conn |
| 43 | +} |
| 44 | + |
| 45 | +func beaconUDPAddr(t *testing.T, s *Server) *net.UDPAddr { |
| 46 | + t.Helper() |
| 47 | + addr, err := net.ResolveUDPAddr("udp", s.Addr().String()) |
| 48 | + if err != nil { |
| 49 | + t.Fatalf("resolve beacon addr: %v", err) |
| 50 | + } |
| 51 | + return addr |
| 52 | +} |
| 53 | + |
| 54 | +func TestGossip(t *testing.T) { |
| 55 | + t.Parallel() |
| 56 | + |
| 57 | + // Start two beacons — they'll be peers of each other |
| 58 | + b1 := NewWithPeers(1, nil) // peers set after both bind |
| 59 | + b2 := NewWithPeers(2, nil) |
| 60 | + |
| 61 | + go b1.ListenAndServe("127.0.0.1:0") |
| 62 | + go b2.ListenAndServe("127.0.0.1:0") |
| 63 | + <-b1.Ready() |
| 64 | + <-b2.Ready() |
| 65 | + defer b1.Close() |
| 66 | + defer b2.Close() |
| 67 | + |
| 68 | + b1Addr := beaconUDPAddr(t, b1) |
| 69 | + b2Addr := beaconUDPAddr(t, b2) |
| 70 | + |
| 71 | + // Set peers manually (after bind, so we know the ports) |
| 72 | + b1.peers = []*net.UDPAddr{b2Addr} |
| 73 | + b2.peers = []*net.UDPAddr{b1Addr} |
| 74 | + |
| 75 | + // Register node 100 on beacon 1 |
| 76 | + conn1 := registerNode(t, b1Addr, 100) |
| 77 | + defer conn1.Close() |
| 78 | + |
| 79 | + // Register node 200 on beacon 2 |
| 80 | + conn2 := registerNode(t, b2Addr, 200) |
| 81 | + defer conn2.Close() |
| 82 | + |
| 83 | + // Verify local counts |
| 84 | + if b1.LocalNodeCount() != 1 { |
| 85 | + t.Fatalf("b1 local nodes: got %d, want 1", b1.LocalNodeCount()) |
| 86 | + } |
| 87 | + if b2.LocalNodeCount() != 1 { |
| 88 | + t.Fatalf("b2 local nodes: got %d, want 1", b2.LocalNodeCount()) |
| 89 | + } |
| 90 | + |
| 91 | + // Trigger gossip manually |
| 92 | + b1.sendGossip() |
| 93 | + b2.sendGossip() |
| 94 | + |
| 95 | + // Give gossip time to propagate |
| 96 | + time.Sleep(200 * time.Millisecond) |
| 97 | + |
| 98 | + // Each beacon should know about the other's node via gossip |
| 99 | + if b1.PeerNodeCount() != 1 { |
| 100 | + t.Errorf("b1 peer nodes: got %d, want 1", b1.PeerNodeCount()) |
| 101 | + } |
| 102 | + if b2.PeerNodeCount() != 1 { |
| 103 | + t.Errorf("b2 peer nodes: got %d, want 1", b2.PeerNodeCount()) |
| 104 | + } |
| 105 | +} |
| 106 | + |
| 107 | +func TestCrossBeaconRelay(t *testing.T) { |
| 108 | + t.Parallel() |
| 109 | + |
| 110 | + b1 := NewWithPeers(1, nil) |
| 111 | + b2 := NewWithPeers(2, nil) |
| 112 | + |
| 113 | + go b1.ListenAndServe("127.0.0.1:0") |
| 114 | + go b2.ListenAndServe("127.0.0.1:0") |
| 115 | + <-b1.Ready() |
| 116 | + <-b2.Ready() |
| 117 | + defer b1.Close() |
| 118 | + defer b2.Close() |
| 119 | + |
| 120 | + b1Addr := beaconUDPAddr(t, b1) |
| 121 | + b2Addr := beaconUDPAddr(t, b2) |
| 122 | + |
| 123 | + b1.peers = []*net.UDPAddr{b2Addr} |
| 124 | + b2.peers = []*net.UDPAddr{b1Addr} |
| 125 | + |
| 126 | + // Register node 10 on beacon 1 |
| 127 | + conn1 := registerNode(t, b1Addr, 10) |
| 128 | + defer conn1.Close() |
| 129 | + |
| 130 | + // Register node 20 on beacon 2 |
| 131 | + conn2 := registerNode(t, b2Addr, 20) |
| 132 | + defer conn2.Close() |
| 133 | + |
| 134 | + // Gossip so b1 knows node 20 is on b2 |
| 135 | + b1.sendGossip() |
| 136 | + b2.sendGossip() |
| 137 | + time.Sleep(200 * time.Millisecond) |
| 138 | + |
| 139 | + // Node 10 sends relay to node 20 via beacon 1 |
| 140 | + // beacon 1 should forward to beacon 2, which delivers to node 20 |
| 141 | + payload := []byte("hello from node 10") |
| 142 | + relayMsg := make([]byte, 1+4+4+len(payload)) |
| 143 | + relayMsg[0] = protocol.BeaconMsgRelay |
| 144 | + binary.BigEndian.PutUint32(relayMsg[1:5], 10) // sender |
| 145 | + binary.BigEndian.PutUint32(relayMsg[5:9], 20) // dest |
| 146 | + copy(relayMsg[9:], payload) |
| 147 | + |
| 148 | + if _, err := conn1.Write(relayMsg); err != nil { |
| 149 | + t.Fatalf("send relay: %v", err) |
| 150 | + } |
| 151 | + |
| 152 | + // Node 20 should receive a RelayDeliver |
| 153 | + buf := make([]byte, 1500) |
| 154 | + conn2.SetReadDeadline(time.Now().Add(2 * time.Second)) |
| 155 | + n, err := conn2.Read(buf) |
| 156 | + if err != nil { |
| 157 | + t.Fatalf("read relay deliver: %v", err) |
| 158 | + } |
| 159 | + |
| 160 | + if buf[0] != protocol.BeaconMsgRelayDeliver { |
| 161 | + t.Fatalf("expected RelayDeliver (0x%02x), got 0x%02x", protocol.BeaconMsgRelayDeliver, buf[0]) |
| 162 | + } |
| 163 | + |
| 164 | + senderID := binary.BigEndian.Uint32(buf[1:5]) |
| 165 | + if senderID != 10 { |
| 166 | + t.Fatalf("sender ID: got %d, want 10", senderID) |
| 167 | + } |
| 168 | + |
| 169 | + received := string(buf[5:n]) |
| 170 | + if received != "hello from node 10" { |
| 171 | + t.Fatalf("payload: got %q, want %q", received, "hello from node 10") |
| 172 | + } |
| 173 | +} |
| 174 | + |
| 175 | +func TestHealthEndpoint(t *testing.T) { |
| 176 | + t.Parallel() |
| 177 | + |
| 178 | + s := New() |
| 179 | + go s.ListenAndServe("127.0.0.1:0") |
| 180 | + <-s.Ready() |
| 181 | + defer s.Close() |
| 182 | + |
| 183 | + // Find a free port for health |
| 184 | + ln, err := net.Listen("tcp", "127.0.0.1:0") |
| 185 | + if err != nil { |
| 186 | + t.Fatalf("find free port: %v", err) |
| 187 | + } |
| 188 | + healthAddr := ln.Addr().String() |
| 189 | + ln.Close() |
| 190 | + |
| 191 | + go s.ServeHealth(healthAddr) |
| 192 | + time.Sleep(100 * time.Millisecond) // let HTTP server start |
| 193 | + |
| 194 | + url := fmt.Sprintf("http://%s/healthz", healthAddr) |
| 195 | + |
| 196 | + // Should be healthy by default |
| 197 | + resp, err := http.Get(url) |
| 198 | + if err != nil { |
| 199 | + t.Fatalf("GET /healthz: %v", err) |
| 200 | + } |
| 201 | + if resp.StatusCode != 200 { |
| 202 | + t.Fatalf("expected 200, got %d", resp.StatusCode) |
| 203 | + } |
| 204 | + resp.Body.Close() |
| 205 | + |
| 206 | + // Set unhealthy |
| 207 | + s.SetHealthy(false) |
| 208 | + resp, err = http.Get(url) |
| 209 | + if err != nil { |
| 210 | + t.Fatalf("GET /healthz after unhealthy: %v", err) |
| 211 | + } |
| 212 | + if resp.StatusCode != 503 { |
| 213 | + t.Fatalf("expected 503, got %d", resp.StatusCode) |
| 214 | + } |
| 215 | + resp.Body.Close() |
| 216 | + |
| 217 | + // Set healthy again |
| 218 | + s.SetHealthy(true) |
| 219 | + resp, err = http.Get(url) |
| 220 | + if err != nil { |
| 221 | + t.Fatalf("GET /healthz after re-healthy: %v", err) |
| 222 | + } |
| 223 | + if resp.StatusCode != 200 { |
| 224 | + t.Fatalf("expected 200, got %d", resp.StatusCode) |
| 225 | + } |
| 226 | + resp.Body.Close() |
| 227 | +} |
| 228 | + |
| 229 | +func TestSyncMessageParsing(t *testing.T) { |
| 230 | + t.Parallel() |
| 231 | + |
| 232 | + s := NewWithPeers(1, nil) |
| 233 | + go s.ListenAndServe("127.0.0.1:0") |
| 234 | + <-s.Ready() |
| 235 | + defer s.Close() |
| 236 | + |
| 237 | + // Build a sync message with 3 nodes |
| 238 | + nodeIDs := []uint32{100, 200, 300} |
| 239 | + msg := make([]byte, 1+4+2+4*len(nodeIDs)) |
| 240 | + msg[0] = protocol.BeaconMsgSync |
| 241 | + binary.BigEndian.PutUint32(msg[1:5], 2) // peer beacon ID |
| 242 | + binary.BigEndian.PutUint16(msg[5:7], uint16(len(nodeIDs))) |
| 243 | + for i, id := range nodeIDs { |
| 244 | + binary.BigEndian.PutUint32(msg[7+4*i:7+4*i+4], id) |
| 245 | + } |
| 246 | + |
| 247 | + // Send the sync message to the beacon |
| 248 | + conn, err := net.DialUDP("udp", nil, beaconUDPAddr(t, s)) |
| 249 | + if err != nil { |
| 250 | + t.Fatalf("dial: %v", err) |
| 251 | + } |
| 252 | + defer conn.Close() |
| 253 | + |
| 254 | + if _, err := conn.Write(msg); err != nil { |
| 255 | + t.Fatalf("send sync: %v", err) |
| 256 | + } |
| 257 | + |
| 258 | + time.Sleep(100 * time.Millisecond) |
| 259 | + |
| 260 | + if s.PeerNodeCount() != 3 { |
| 261 | + t.Fatalf("peer nodes: got %d, want 3", s.PeerNodeCount()) |
| 262 | + } |
| 263 | +} |
0 commit comments