| 159 | } |
| 160 | |
| 161 | func (n *server) receiveHandler() { |
| 162 | clients := make(map[string]*client) |
| 163 | for { |
| 164 | select { |
| 165 | case <-n.ctx.Done(): |
| 166 | for _, client := range clients { |
| 167 | if err := client.close(); err != nil { |
| 168 | err = &P2PError{err: errors.Errorf("conn close failed: %w", err), t: time.Now()} |
| 169 | n.logger.Error(err) |
| 170 | } |
| 171 | } |
| 172 | return |
| 173 | case c := <-n.addIncomingC: |
| 174 | if clients[string(c.remoteID)] != nil { |
| 175 | if err := c.close(); err != nil { |
| 176 | err = &P2PError{err: errors.Errorf("conn close failed: %w", err), t: time.Now()} |
| 177 | n.logger.Error(err) |
| 178 | } |
| 179 | continue |
| 180 | } |
| 181 | clients[string(c.remoteID)] = c |
| 182 | n.incomingNum = len(clients) |
| 183 | go n.runClient(c, true) |
| 184 | case id := <-n.removeIncomingC: |
| 185 | if c := clients[string(id)]; c != nil { |
| 186 | delete(clients, string(id)) |
| 187 | } |
| 188 | n.incomingNum = len(clients) |
| 189 | case req := <-n.replying: |
| 190 | client := clients[string(req.id)] |
| 191 | if client == nil { |
| 192 | err := &P2PError{err: errors.Errorf("reply failed: %w", ErrCanNotFindClient), t: time.Now()} |
| 193 | n.logger.Error(err) |
| 194 | req.replyResult(nil, err) |
| 195 | continue |
| 196 | } |
| 197 | go client.send(req) |
| 198 | } |
| 199 | } |
| 200 | } |
| 201 | |
| 202 | func (n *server) callHandler() (err error) { |
| 203 | addrToid := make(map[string][]byte) |