SubscribeMsg is a message subscription operation binding the P2PMessage
(chanBuffer int, peersFeed ...interface{})
| 516 | |
| 517 | // SubscribeMsg is a message subscription operation binding the P2PMessage |
| 518 | func (n *server) SubscribeMsg(chanBuffer int, peersFeed ...interface{}) (outch chan P2PMessage, err error) { |
| 519 | |
| 520 | var eventList []chan P2PMessage |
| 521 | for _, m := range peersFeed { |
| 522 | if chanBuffer > 0 { |
| 523 | outch = make(chan P2PMessage, chanBuffer) |
| 524 | } else { |
| 525 | outch = make(chan P2PMessage) |
| 526 | } |
| 527 | eventList = append(eventList, outch) |
| 528 | select { |
| 529 | case <-n.ctx.Done(): |
| 530 | case n.subscribeMsg <- &subscription{msgType: reflect.TypeOf(m).String(), msgCh: outch}: |
| 531 | } |
| 532 | } |
| 533 | outch = merge(n.ctx, eventList...) |
| 534 | return |
| 535 | } |
| 536 | |
| 537 | // UnSubscribeEvent is a un-subscription operation |
| 538 | func (n *server) UnSubscribeMsg(peersFeed ...interface{}) { |