MCPcopy Create free account
hub / github.com/apache/kvrocks-controller / runTransport

Method runTransport

store/engine/raft/node.go:178–223  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

176}
177
178func (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
225func (n *Node) watchLeaderChange() {
226 n.wg.Add(1)

Callers 1

runMethod · 0.95

Implementers 1

ClusterNodestore/cluster_node.go

Calls 7

GetFunction · 0.92
ErrorfMethod · 0.80
AddPeerMethod · 0.80
FatalMethod · 0.80
ErrorMethod · 0.80
IDMethod · 0.65
StartMethod · 0.45

Tested by

no test coverage detected