| 200 | } |
| 201 | |
| 202 | func (n *server) callHandler() (err error) { |
| 203 | addrToid := make(map[string][]byte) |
| 204 | clients := make(map[string]*client) |
| 205 | watchDog := time.NewTicker(5 * time.Second) |
| 206 | for { |
| 207 | select { |
| 208 | case <-n.ctx.Done(): |
| 209 | for _, client := range clients { |
| 210 | if err := client.close(); err != nil { |
| 211 | err = &P2PError{err: errors.Errorf("conn close failed: %w", err), t: time.Now()} |
| 212 | n.logger.Error(err) |
| 213 | } |
| 214 | } |
| 215 | err = n.ctx.Err() |
| 216 | return |
| 217 | case <-watchDog.C: |
| 218 | if !n.members.IsAlive() { |
| 219 | err = errors.New("p2p cluster status is not alive") |
| 220 | n.logger.Error(err) |
| 221 | return |
| 222 | } |
| 223 | case id, ok := <-n.removeCallingC: |
| 224 | if ok { |
| 225 | c := clients[string(id)] |
| 226 | if c != nil { |
| 227 | delete(addrToid, c.conn.RemoteAddr().String()) |
| 228 | delete(clients, string(id)) |
| 229 | } |
| 230 | n.callingNum = len(clients) |
| 231 | } |
| 232 | case req, ok := <-n.calling: |
| 233 | if ok { |
| 234 | var c *client |
| 235 | if c = clients[string(req.id)]; c == nil { |
| 236 | //TODO : performance bottleneck |
| 237 | req.addr = n.members.Lookup(req.id) |
| 238 | if c = n.handleCallReq(req); c == nil { |
| 239 | continue |
| 240 | } |
| 241 | clients[string(req.id)] = c |
| 242 | go n.runClient(c, false) |
| 243 | } |
| 244 | go c.send(req) |
| 245 | } |
| 246 | } |
| 247 | } |
| 248 | return |
| 249 | } |
| 250 | |
| 251 | func (n *server) runClient(c *client, inBound bool) { |
| 252 | if err := c.run(); err != nil { |