| 176 | } |
| 177 | |
| 178 | func (n *Node) runTransport() error { |
| 179 | logger := logger.Get() |
| 180 | idString := fmt.Sprintf("%d", n.config.ID) |
| 181 | transport := &rafthttp.Transport{ |
| 182 | ID: types.ID(n.config.ID), |
| 183 | Logger: logger, |
| 184 | ClusterID: 0x6666, |
| 185 | Raft: n, |
| 186 | LeaderStats: stats.NewLeaderStats(logger, idString), |
| 187 | ServerStats: stats.NewServerStats("raft", idString), |
| 188 | ErrorC: make(chan error), |
| 189 | } |
| 190 | if err := transport.Start(); err != nil { |
| 191 | return fmt.Errorf("unable to start transport: %w", err) |
| 192 | } |
| 193 | for i, peer := range n.config.Peers { |
| 194 | // Don't add self to transport |
| 195 | if uint64(i+1) != n.config.ID { |
| 196 | transport.AddPeer(types.ID(i+1), []string{peer}) |
| 197 | } |
| 198 | n.peers.Store(uint64(i+1), peer) |
| 199 | } |
| 200 | |
| 201 | n.addr = n.config.Peers[n.config.ID-1] |
| 202 | url, err := url.Parse(n.addr) |
| 203 | if err != nil { |
| 204 | return err |
| 205 | } |
| 206 | httpServer := &http.Server{ |
| 207 | Addr: url.Host, |
| 208 | Handler: transport.Handler(), |
| 209 | } |
| 210 | |
| 211 | n.wg.Add(1) |
| 212 | go func() { |
| 213 | defer n.wg.Done() |
| 214 | if err := httpServer.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { |
| 215 | n.logger.Fatal("Unable to start http server", zap.Error(err)) |
| 216 | os.Exit(1) |
| 217 | } |
| 218 | }() |
| 219 | |
| 220 | n.transport = transport |
| 221 | n.httpServer = httpServer |
| 222 | return nil |
| 223 | } |
| 224 | |
| 225 | func (n *Node) watchLeaderChange() { |
| 226 | n.wg.Add(1) |