2018-09-05 01:07:44 +08:00
|
|
|
package pstoremem
|
2015-10-01 06:42:55 +08:00
|
|
|
|
|
|
|
import (
|
2016-10-05 07:19:16 +08:00
|
|
|
"context"
|
2016-04-20 06:18:24 +08:00
|
|
|
"sort"
|
2015-10-01 06:42:55 +08:00
|
|
|
"sync"
|
2019-05-09 07:27:58 +08:00
|
|
|
"sync/atomic"
|
2015-10-01 06:42:55 +08:00
|
|
|
"time"
|
|
|
|
|
2018-09-02 19:10:55 +08:00
|
|
|
logging "github.com/ipfs/go-log"
|
2018-09-08 01:46:23 +08:00
|
|
|
peer "github.com/libp2p/go-libp2p-peer"
|
2018-09-02 19:10:55 +08:00
|
|
|
ma "github.com/multiformats/go-multiaddr"
|
2018-08-30 23:24:09 +08:00
|
|
|
|
|
|
|
pstore "github.com/libp2p/go-libp2p-peerstore"
|
2018-09-08 01:46:23 +08:00
|
|
|
addr "github.com/libp2p/go-libp2p-peerstore/addr"
|
2015-10-01 06:42:55 +08:00
|
|
|
)
|
|
|
|
|
2018-08-29 22:12:41 +08:00
|
|
|
var log = logging.Logger("peerstore")
|
|
|
|
|
2015-10-01 06:42:55 +08:00
|
|
|
type expiringAddr struct {
|
2018-01-05 05:00:37 +08:00
|
|
|
Addr ma.Multiaddr
|
|
|
|
TTL time.Duration
|
|
|
|
Expires time.Time
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
func (e *expiringAddr) ExpiredBy(t time.Time) bool {
|
2018-01-05 05:00:37 +08:00
|
|
|
return t.After(e.Expires)
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
|
2019-05-17 17:33:33 +08:00
|
|
|
type addrSegments [256]*addrSegment
|
2019-05-09 07:27:58 +08:00
|
|
|
|
2019-05-17 17:33:33 +08:00
|
|
|
type addrSegment struct {
|
2019-05-17 17:37:46 +08:00
|
|
|
sync.RWMutex
|
2019-05-09 07:27:58 +08:00
|
|
|
|
2019-05-17 17:37:46 +08:00
|
|
|
// Use pointers to save memory. Maps always leave some fraction of their
|
|
|
|
// space unused. storing the *values* directly in the map will
|
|
|
|
// drastically increase the space waste. In our case, by 6x.
|
2019-05-09 07:27:58 +08:00
|
|
|
addrs map[peer.ID]map[string]*expiringAddr
|
2019-05-17 17:37:46 +08:00
|
|
|
size uint32
|
2019-05-09 07:27:58 +08:00
|
|
|
}
|
|
|
|
|
2019-05-17 18:14:39 +08:00
|
|
|
func (s *addrSegments) get(p peer.ID) *addrSegment {
|
|
|
|
b := []byte(p)
|
2019-05-17 17:33:33 +08:00
|
|
|
return s[b[len(b)-1]]
|
2019-05-09 07:27:58 +08:00
|
|
|
}
|
2018-08-29 22:12:41 +08:00
|
|
|
|
2018-08-30 23:43:40 +08:00
|
|
|
// memoryAddrBook manages addresses.
|
|
|
|
type memoryAddrBook struct {
|
2019-05-17 17:33:33 +08:00
|
|
|
segments addrSegments
|
2016-04-16 00:30:59 +08:00
|
|
|
|
2019-04-28 17:17:48 +08:00
|
|
|
ctx context.Context
|
|
|
|
cancel func()
|
2018-10-03 01:49:55 +08:00
|
|
|
|
2018-06-13 05:37:26 +08:00
|
|
|
subManager *AddrSubManager
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
|
2019-05-09 07:27:58 +08:00
|
|
|
var _ pstore.AddrBook = (*memoryAddrBook)(nil)
|
|
|
|
|
2018-08-30 23:43:40 +08:00
|
|
|
func NewAddrBook() pstore.AddrBook {
|
2019-04-28 17:17:48 +08:00
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
|
|
|
ab := &memoryAddrBook{
|
2019-05-17 17:33:33 +08:00
|
|
|
segments: func() (ret addrSegments) {
|
2019-05-09 07:27:58 +08:00
|
|
|
for i, _ := range ret {
|
2019-05-17 17:33:33 +08:00
|
|
|
ret[i] = &addrSegment{addrs: make(map[peer.ID]map[string]*expiringAddr)}
|
2019-05-09 07:27:58 +08:00
|
|
|
}
|
|
|
|
return ret
|
|
|
|
}(),
|
2018-08-30 23:43:40 +08:00
|
|
|
subManager: NewAddrSubManager(),
|
2019-04-28 17:17:48 +08:00
|
|
|
ctx: ctx,
|
|
|
|
cancel: cancel,
|
|
|
|
}
|
|
|
|
|
|
|
|
go ab.background()
|
|
|
|
return ab
|
|
|
|
}
|
|
|
|
|
|
|
|
// background periodically schedules a gc
|
|
|
|
func (mab *memoryAddrBook) background() {
|
|
|
|
ticker := time.NewTicker(1 * time.Hour)
|
|
|
|
defer ticker.Stop()
|
|
|
|
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-ticker.C:
|
|
|
|
mab.gc()
|
|
|
|
|
|
|
|
case <-mab.ctx.Done():
|
|
|
|
return
|
|
|
|
}
|
2016-04-16 00:30:59 +08:00
|
|
|
}
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
|
2019-04-28 17:19:26 +08:00
|
|
|
func (mab *memoryAddrBook) Close() error {
|
|
|
|
mab.cancel()
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2019-04-28 21:21:07 +08:00
|
|
|
// gc garbage collects the in-memory address book.
|
2018-10-03 01:49:55 +08:00
|
|
|
func (mab *memoryAddrBook) gc() {
|
|
|
|
now := time.Now()
|
2019-05-09 07:27:58 +08:00
|
|
|
for _, s := range mab.segments {
|
2019-05-17 17:37:46 +08:00
|
|
|
s.Lock()
|
2019-05-09 07:27:58 +08:00
|
|
|
for p, amap := range s.addrs {
|
|
|
|
for k, addr := range amap {
|
|
|
|
if addr.ExpiredBy(now) {
|
|
|
|
delete(amap, k)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
if len(amap) == 0 {
|
|
|
|
delete(s.addrs, p)
|
2018-10-03 01:49:55 +08:00
|
|
|
}
|
|
|
|
}
|
2019-05-09 16:20:12 +08:00
|
|
|
atomic.StoreUint32(&s.size, uint32(len(s.addrs)))
|
2019-05-17 17:37:46 +08:00
|
|
|
s.Unlock()
|
2018-10-03 01:49:55 +08:00
|
|
|
}
|
2019-05-09 07:27:58 +08:00
|
|
|
|
2018-10-03 01:49:55 +08:00
|
|
|
}
|
|
|
|
|
2018-08-31 19:59:46 +08:00
|
|
|
func (mab *memoryAddrBook) PeersWithAddrs() peer.IDSlice {
|
2019-05-09 07:27:58 +08:00
|
|
|
var length uint32
|
|
|
|
for _, s := range mab.segments {
|
|
|
|
length += atomic.LoadUint32(&s.size)
|
|
|
|
}
|
2015-10-01 06:42:55 +08:00
|
|
|
|
2019-05-09 07:27:58 +08:00
|
|
|
pids := make(peer.IDSlice, 0, length)
|
|
|
|
for _, s := range mab.segments {
|
2019-05-17 17:37:46 +08:00
|
|
|
s.RLock()
|
2019-05-09 07:27:58 +08:00
|
|
|
for pid, _ := range s.addrs {
|
|
|
|
pids = append(pids, pid)
|
|
|
|
}
|
2019-05-17 17:37:46 +08:00
|
|
|
s.RUnlock()
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
return pids
|
|
|
|
}
|
|
|
|
|
|
|
|
// AddAddr calls AddAddrs(p, []ma.Multiaddr{addr}, ttl)
|
2018-08-30 23:43:40 +08:00
|
|
|
func (mab *memoryAddrBook) AddAddr(p peer.ID, addr ma.Multiaddr, ttl time.Duration) {
|
|
|
|
mab.AddAddrs(p, []ma.Multiaddr{addr}, ttl)
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
|
2018-08-30 23:43:40 +08:00
|
|
|
// AddAddrs gives memoryAddrBook addresses to use, with a given ttl
|
2015-10-01 06:42:55 +08:00
|
|
|
// (time-to-live), after which the address is no longer valid.
|
|
|
|
// If the manager has a longer TTL, the operation is a no-op for that address
|
2018-08-30 23:43:40 +08:00
|
|
|
func (mab *memoryAddrBook) AddAddrs(p peer.ID, addrs []ma.Multiaddr, ttl time.Duration) {
|
2015-10-01 06:42:55 +08:00
|
|
|
// if ttl is zero, exit. nothing to do.
|
|
|
|
if ttl <= 0 {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2019-05-09 07:27:58 +08:00
|
|
|
s := mab.segments.get(p)
|
2019-05-17 17:37:46 +08:00
|
|
|
s.Lock()
|
|
|
|
defer s.Unlock()
|
2019-05-09 07:27:58 +08:00
|
|
|
|
|
|
|
// update the segment size
|
|
|
|
defer atomic.StoreUint32(&s.size, uint32(len(s.addrs)))
|
|
|
|
|
|
|
|
amap := s.addrs[p]
|
2018-10-02 13:16:08 +08:00
|
|
|
if amap == nil {
|
2018-10-14 22:11:35 +08:00
|
|
|
amap = make(map[string]*expiringAddr, len(addrs))
|
2019-05-09 07:27:58 +08:00
|
|
|
s.addrs[p] = amap
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
exp := time.Now().Add(ttl)
|
|
|
|
for _, addr := range addrs {
|
2016-05-17 08:04:08 +08:00
|
|
|
if addr == nil {
|
|
|
|
log.Warningf("was passed nil multiaddr for %s", p)
|
|
|
|
continue
|
|
|
|
}
|
2016-05-17 12:23:25 +08:00
|
|
|
addrstr := string(addr.Bytes())
|
2015-10-01 06:42:55 +08:00
|
|
|
a, found := amap[addrstr]
|
2018-01-05 05:00:37 +08:00
|
|
|
if !found || exp.After(a.Expires) {
|
2018-10-14 22:11:35 +08:00
|
|
|
amap[addrstr] = &expiringAddr{Addr: addr, Expires: exp, TTL: ttl}
|
2016-04-16 00:30:59 +08:00
|
|
|
|
2018-08-30 23:43:40 +08:00
|
|
|
mab.subManager.BroadcastAddr(p, addr)
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// SetAddr calls mgr.SetAddrs(p, addr, ttl)
|
2018-08-30 23:43:40 +08:00
|
|
|
func (mab *memoryAddrBook) SetAddr(p peer.ID, addr ma.Multiaddr, ttl time.Duration) {
|
|
|
|
mab.SetAddrs(p, []ma.Multiaddr{addr}, ttl)
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
// SetAddrs sets the ttl on addresses. This clears any TTL there previously.
|
|
|
|
// This is used when we receive the best estimate of the validity of an address.
|
2018-08-30 23:43:40 +08:00
|
|
|
func (mab *memoryAddrBook) SetAddrs(p peer.ID, addrs []ma.Multiaddr, ttl time.Duration) {
|
2019-05-09 07:27:58 +08:00
|
|
|
s := mab.segments.get(p)
|
2019-05-17 17:37:46 +08:00
|
|
|
s.Lock()
|
|
|
|
defer s.Unlock()
|
2019-05-09 07:27:58 +08:00
|
|
|
|
|
|
|
// update the segment size
|
|
|
|
defer atomic.StoreUint32(&s.size, uint32(len(s.addrs)))
|
2015-10-01 06:42:55 +08:00
|
|
|
|
2019-05-09 07:27:58 +08:00
|
|
|
amap := s.addrs[p]
|
2018-10-02 13:16:08 +08:00
|
|
|
if amap == nil {
|
2018-10-14 22:11:35 +08:00
|
|
|
amap = make(map[string]*expiringAddr, len(addrs))
|
2019-05-09 07:27:58 +08:00
|
|
|
s.addrs[p] = amap
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
exp := time.Now().Add(ttl)
|
|
|
|
for _, addr := range addrs {
|
2016-05-17 08:04:08 +08:00
|
|
|
if addr == nil {
|
|
|
|
log.Warningf("was passed nil multiaddr for %s", p)
|
|
|
|
continue
|
|
|
|
}
|
2015-10-01 06:42:55 +08:00
|
|
|
// re-set all of them for new ttl.
|
2018-10-02 13:16:08 +08:00
|
|
|
addrstr := string(addr.Bytes())
|
2015-10-01 06:42:55 +08:00
|
|
|
|
|
|
|
if ttl > 0 {
|
2018-10-14 22:11:35 +08:00
|
|
|
amap[addrstr] = &expiringAddr{Addr: addr, Expires: exp, TTL: ttl}
|
2018-08-30 23:43:40 +08:00
|
|
|
mab.subManager.BroadcastAddr(p, addr)
|
2015-10-01 06:42:55 +08:00
|
|
|
} else {
|
2018-10-02 13:16:08 +08:00
|
|
|
delete(amap, addrstr)
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-01-05 05:00:37 +08:00
|
|
|
// UpdateAddrs updates the addresses associated with the given peer that have
|
|
|
|
// the given oldTTL to have the given newTTL.
|
2018-08-30 23:43:40 +08:00
|
|
|
func (mab *memoryAddrBook) UpdateAddrs(p peer.ID, oldTTL time.Duration, newTTL time.Duration) {
|
2019-05-09 07:27:58 +08:00
|
|
|
s := mab.segments.get(p)
|
2019-05-17 17:37:46 +08:00
|
|
|
s.Lock()
|
|
|
|
defer s.Unlock()
|
2019-05-09 07:27:58 +08:00
|
|
|
|
|
|
|
// update the segment size
|
|
|
|
defer atomic.StoreUint32(&s.size, uint32(len(s.addrs)))
|
2018-01-05 05:00:37 +08:00
|
|
|
|
2019-05-09 07:27:58 +08:00
|
|
|
amap, found := s.addrs[p]
|
2018-01-05 05:00:37 +08:00
|
|
|
if !found {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
exp := time.Now().Add(newTTL)
|
2018-10-02 13:16:08 +08:00
|
|
|
for k, addr := range amap {
|
|
|
|
if oldTTL == addr.TTL {
|
|
|
|
addr.TTL = newTTL
|
|
|
|
addr.Expires = exp
|
|
|
|
amap[k] = addr
|
2018-01-05 05:00:37 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2015-10-01 06:42:55 +08:00
|
|
|
// Addresses returns all known (and valid) addresses for a given
|
2018-08-30 23:43:40 +08:00
|
|
|
func (mab *memoryAddrBook) Addrs(p peer.ID) []ma.Multiaddr {
|
2019-05-09 07:27:58 +08:00
|
|
|
s := mab.segments.get(p)
|
2019-05-17 17:37:46 +08:00
|
|
|
s.RLock()
|
|
|
|
defer s.RUnlock()
|
2015-10-01 06:42:55 +08:00
|
|
|
|
2019-05-09 07:27:58 +08:00
|
|
|
amap, found := s.addrs[p]
|
2015-10-01 06:42:55 +08:00
|
|
|
if !found {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
now := time.Now()
|
2018-10-02 13:16:08 +08:00
|
|
|
good := make([]ma.Multiaddr, 0, len(amap))
|
2019-04-28 16:35:23 +08:00
|
|
|
for _, m := range amap {
|
2018-01-17 07:36:29 +08:00
|
|
|
if !m.ExpiredBy(now) {
|
2015-10-01 06:42:55 +08:00
|
|
|
good = append(good, m.Addr)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return good
|
|
|
|
}
|
|
|
|
|
2018-06-16 01:54:04 +08:00
|
|
|
// ClearAddrs removes all previously stored addresses
|
2018-08-30 23:43:40 +08:00
|
|
|
func (mab *memoryAddrBook) ClearAddrs(p peer.ID) {
|
2019-05-09 07:27:58 +08:00
|
|
|
s := mab.segments.get(p)
|
2019-05-17 17:37:46 +08:00
|
|
|
s.Lock()
|
|
|
|
defer s.Unlock()
|
2019-05-09 07:27:58 +08:00
|
|
|
|
|
|
|
// update the segment size
|
|
|
|
defer atomic.StoreUint32(&s.size, uint32(len(s.addrs)))
|
2015-10-01 06:42:55 +08:00
|
|
|
|
2019-05-09 07:27:58 +08:00
|
|
|
delete(s.addrs, p)
|
2015-10-01 06:42:55 +08:00
|
|
|
}
|
2016-04-16 00:30:59 +08:00
|
|
|
|
2018-06-16 01:54:04 +08:00
|
|
|
// AddrStream returns a channel on which all new addresses discovered for a
|
|
|
|
// given peer ID will be published.
|
2018-08-30 23:43:40 +08:00
|
|
|
func (mab *memoryAddrBook) AddrStream(ctx context.Context, p peer.ID) <-chan ma.Multiaddr {
|
2019-05-09 07:27:58 +08:00
|
|
|
s := mab.segments.get(p)
|
2019-05-17 17:37:46 +08:00
|
|
|
s.RLock()
|
|
|
|
defer s.RUnlock()
|
2018-06-15 04:17:50 +08:00
|
|
|
|
2019-05-09 07:27:58 +08:00
|
|
|
baseaddrslice := s.addrs[p]
|
2018-06-13 05:37:26 +08:00
|
|
|
initial := make([]ma.Multiaddr, 0, len(baseaddrslice))
|
|
|
|
for _, a := range baseaddrslice {
|
|
|
|
initial = append(initial, a.Addr)
|
|
|
|
}
|
|
|
|
|
2018-08-30 23:43:40 +08:00
|
|
|
return mab.subManager.AddrStream(ctx, p, initial)
|
|
|
|
}
|
|
|
|
|
|
|
|
type addrSub struct {
|
|
|
|
pubch chan ma.Multiaddr
|
|
|
|
lk sync.Mutex
|
|
|
|
buffer []ma.Multiaddr
|
|
|
|
ctx context.Context
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *addrSub) pubAddr(a ma.Multiaddr) {
|
|
|
|
select {
|
|
|
|
case s.pubch <- a:
|
|
|
|
case <-s.ctx.Done():
|
|
|
|
}
|
2018-06-13 05:37:26 +08:00
|
|
|
}
|
|
|
|
|
2018-06-16 01:54:04 +08:00
|
|
|
// An abstracted, pub-sub manager for address streams. Extracted from
|
2018-08-30 23:43:40 +08:00
|
|
|
// memoryAddrBook in order to support additional implementations.
|
2018-06-13 05:37:26 +08:00
|
|
|
type AddrSubManager struct {
|
|
|
|
mu sync.RWMutex
|
|
|
|
subs map[peer.ID][]*addrSub
|
|
|
|
}
|
|
|
|
|
2018-06-16 01:54:04 +08:00
|
|
|
// NewAddrSubManager initializes an AddrSubManager.
|
2018-06-13 05:37:26 +08:00
|
|
|
func NewAddrSubManager() *AddrSubManager {
|
|
|
|
return &AddrSubManager{
|
|
|
|
subs: make(map[peer.ID][]*addrSub),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-06-16 01:54:04 +08:00
|
|
|
// Used internally by the address stream coroutine to remove a subscription
|
|
|
|
// from the manager.
|
2018-06-13 05:37:26 +08:00
|
|
|
func (mgr *AddrSubManager) removeSub(p peer.ID, s *addrSub) {
|
|
|
|
mgr.mu.Lock()
|
|
|
|
defer mgr.mu.Unlock()
|
2018-08-31 20:23:02 +08:00
|
|
|
|
2018-06-13 05:37:26 +08:00
|
|
|
subs := mgr.subs[p]
|
2017-11-22 10:08:58 +08:00
|
|
|
if len(subs) == 1 {
|
|
|
|
if subs[0] != s {
|
|
|
|
return
|
|
|
|
}
|
2018-06-13 05:37:26 +08:00
|
|
|
delete(mgr.subs, p)
|
2017-11-22 10:08:58 +08:00
|
|
|
return
|
|
|
|
}
|
2018-08-30 23:43:40 +08:00
|
|
|
|
2017-11-22 10:08:58 +08:00
|
|
|
for i, v := range subs {
|
|
|
|
if v == s {
|
|
|
|
subs[i] = subs[len(subs)-1]
|
|
|
|
subs[len(subs)-1] = nil
|
2018-06-13 05:37:26 +08:00
|
|
|
mgr.subs[p] = subs[:len(subs)-1]
|
2017-11-22 10:08:58 +08:00
|
|
|
return
|
2016-04-16 00:30:59 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-06-16 01:54:04 +08:00
|
|
|
// BroadcastAddr broadcasts a new address to all subscribed streams.
|
2018-06-13 05:37:26 +08:00
|
|
|
func (mgr *AddrSubManager) BroadcastAddr(p peer.ID, addr ma.Multiaddr) {
|
2018-06-16 01:54:04 +08:00
|
|
|
mgr.mu.RLock()
|
|
|
|
defer mgr.mu.RUnlock()
|
2018-06-13 05:37:26 +08:00
|
|
|
|
2018-06-16 01:54:04 +08:00
|
|
|
if subs, ok := mgr.subs[p]; ok {
|
|
|
|
for _, sub := range subs {
|
|
|
|
sub.pubAddr(addr)
|
|
|
|
}
|
2016-04-16 00:30:59 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-06-16 01:54:04 +08:00
|
|
|
// AddrStream creates a new subscription for a given peer ID, pre-populating the
|
|
|
|
// channel with any addresses we might already have on file.
|
2018-06-13 05:37:26 +08:00
|
|
|
func (mgr *AddrSubManager) AddrStream(ctx context.Context, p peer.ID, initial []ma.Multiaddr) <-chan ma.Multiaddr {
|
2016-04-16 00:30:59 +08:00
|
|
|
sub := &addrSub{pubch: make(chan ma.Multiaddr), ctx: ctx}
|
|
|
|
out := make(chan ma.Multiaddr)
|
|
|
|
|
2018-06-16 05:33:18 +08:00
|
|
|
mgr.mu.Lock()
|
2018-06-15 04:17:50 +08:00
|
|
|
if _, ok := mgr.subs[p]; ok {
|
|
|
|
mgr.subs[p] = append(mgr.subs[p], sub)
|
|
|
|
} else {
|
|
|
|
mgr.subs[p] = []*addrSub{sub}
|
|
|
|
}
|
2018-06-16 05:33:18 +08:00
|
|
|
mgr.mu.Unlock()
|
2016-04-16 00:30:59 +08:00
|
|
|
|
2016-04-20 06:18:24 +08:00
|
|
|
sort.Sort(addr.AddrList(initial))
|
2016-04-16 00:30:59 +08:00
|
|
|
|
|
|
|
go func(buffer []ma.Multiaddr) {
|
|
|
|
defer close(out)
|
|
|
|
|
2017-12-06 09:37:20 +08:00
|
|
|
sent := make(map[string]bool, len(buffer))
|
2016-04-21 02:25:03 +08:00
|
|
|
var outch chan ma.Multiaddr
|
2016-04-16 00:30:59 +08:00
|
|
|
|
|
|
|
for _, a := range buffer {
|
2017-12-06 09:38:43 +08:00
|
|
|
sent[string(a.Bytes())] = true
|
2016-04-16 00:30:59 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
var next ma.Multiaddr
|
|
|
|
if len(buffer) > 0 {
|
|
|
|
next = buffer[0]
|
|
|
|
buffer = buffer[1:]
|
2016-04-21 02:25:03 +08:00
|
|
|
outch = out
|
2016-04-16 00:30:59 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case outch <- next:
|
|
|
|
if len(buffer) > 0 {
|
|
|
|
next = buffer[0]
|
|
|
|
buffer = buffer[1:]
|
|
|
|
} else {
|
|
|
|
outch = nil
|
|
|
|
next = nil
|
|
|
|
}
|
|
|
|
case naddr := <-sub.pubch:
|
2017-12-06 09:38:43 +08:00
|
|
|
if sent[string(naddr.Bytes())] {
|
2016-04-16 00:30:59 +08:00
|
|
|
continue
|
|
|
|
}
|
|
|
|
|
2017-12-06 09:38:43 +08:00
|
|
|
sent[string(naddr.Bytes())] = true
|
2016-04-16 00:30:59 +08:00
|
|
|
if next == nil {
|
|
|
|
next = naddr
|
|
|
|
outch = out
|
|
|
|
} else {
|
|
|
|
buffer = append(buffer, naddr)
|
|
|
|
}
|
|
|
|
case <-ctx.Done():
|
|
|
|
mgr.removeSub(p, sub)
|
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
}(initial)
|
|
|
|
|
|
|
|
return out
|
|
|
|
}
|