mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-21 04:06:44 +00:00
extract hive/peerPool into hive.go
This commit is contained in:
parent
279219eb40
commit
c5eee0fdd8
3 changed files with 36 additions and 39 deletions
30
bzz/hive.go
Normal file
30
bzz/hive.go
Normal file
|
|
@ -0,0 +1,30 @@
|
||||||
|
package bzz
|
||||||
|
|
||||||
|
type peer struct {
|
||||||
|
*bzzProtocol
|
||||||
|
pubkey []byte
|
||||||
|
}
|
||||||
|
|
||||||
|
// This is a mock implementation with a fixed peer pool with no distinction between peers
|
||||||
|
type hive struct {
|
||||||
|
pool map[string]peer
|
||||||
|
}
|
||||||
|
|
||||||
|
func (self *hive) addPeer(p peer) {
|
||||||
|
self.pool[string(p.pubkey)] = p
|
||||||
|
}
|
||||||
|
|
||||||
|
func (self *hive) removePeer(p peer) {
|
||||||
|
delete(self.pool, string(p.pubkey))
|
||||||
|
}
|
||||||
|
|
||||||
|
func (self *hive) getPeers(target Key) (peers []peer) {
|
||||||
|
for _, value := range self.pool {
|
||||||
|
peers = append(peers, value)
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
func (self *hive) addPeers(req *peersMsgData) (err error) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
@ -17,30 +17,10 @@ import (
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
// This is a mock implementation with a fixed peer pool with no distinction between peers
|
|
||||||
type peerPool struct {
|
|
||||||
pool map[string]peer
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *peerPool) addPeer(p peer) {
|
|
||||||
self.pool[string(p.pubkey)] = p
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *peerPool) removePeer(p peer) {
|
|
||||||
delete(self.pool, string(p.pubkey))
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *peerPool) getPeers(target Key) (peers []peer) {
|
|
||||||
for _, value := range self.pool {
|
|
||||||
peers = append(peers, value)
|
|
||||||
}
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
type netStore struct {
|
type netStore struct {
|
||||||
localStore *localStore
|
localStore *localStore
|
||||||
lock sync.Mutex
|
lock sync.Mutex
|
||||||
peerPool *peerPool
|
hive *hive
|
||||||
}
|
}
|
||||||
|
|
||||||
/*
|
/*
|
||||||
|
|
@ -66,11 +46,6 @@ var (
|
||||||
searchTimeout = 3 * time.Second
|
searchTimeout = 3 * time.Second
|
||||||
)
|
)
|
||||||
|
|
||||||
type peer struct {
|
|
||||||
*bzzProtocol
|
|
||||||
pubkey []byte
|
|
||||||
}
|
|
||||||
|
|
||||||
type requestStatus struct {
|
type requestStatus struct {
|
||||||
key Key
|
key Key
|
||||||
status int
|
status int
|
||||||
|
|
@ -257,7 +232,7 @@ func (self *netStore) store(chunk *Chunk) {
|
||||||
Id: r.Int63(),
|
Id: r.Int63(),
|
||||||
Size: chunk.Size,
|
Size: chunk.Size,
|
||||||
}
|
}
|
||||||
for _, peer := range self.peerPool.getPeers(chunk.Key) {
|
for _, peer := range self.hive.getPeers(chunk.Key) {
|
||||||
go peer.store(req)
|
go peer.store(req)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -280,12 +255,3 @@ func (self *netStore) searchTimeout(rs *requestStatus, req *retrieveRequestMsgDa
|
||||||
return t
|
return t
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// these should go to cademlia
|
|
||||||
func (self *netStore) addPeers(req *peersMsgData) (err error) {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *netStore) removePeer(p peer) {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
|
||||||
|
|
@ -33,6 +33,7 @@ const (
|
||||||
// instance is running on each peer
|
// instance is running on each peer
|
||||||
type bzzProtocol struct {
|
type bzzProtocol struct {
|
||||||
netStore *netStore
|
netStore *netStore
|
||||||
|
hive *hive
|
||||||
peer *p2p.Peer
|
peer *p2p.Peer
|
||||||
rw p2p.MsgReadWriter
|
rw p2p.MsgReadWriter
|
||||||
}
|
}
|
||||||
|
|
@ -164,7 +165,7 @@ func runBzzProtocol(netStore *netStore, p *p2p.Peer, rw p2p.MsgReadWriter) (err
|
||||||
for {
|
for {
|
||||||
err = self.handle()
|
err = self.handle()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
self.netStore.removePeer(peer{bzzProtocol: self})
|
self.hive.removePeer(peer{bzzProtocol: self})
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -214,7 +215,7 @@ func (self *bzzProtocol) handle() error {
|
||||||
return self.protoError(ErrDecode, "->msg %v: %v", msg, err)
|
return self.protoError(ErrDecode, "->msg %v: %v", msg, err)
|
||||||
}
|
}
|
||||||
req.peer = peer{bzzProtocol: self}
|
req.peer = peer{bzzProtocol: self}
|
||||||
self.netStore.addPeers(&req)
|
self.hive.addPeers(&req)
|
||||||
|
|
||||||
default:
|
default:
|
||||||
return self.protoError(ErrInvalidMsgCode, "%v", msg.Code)
|
return self.protoError(ErrInvalidMsgCode, "%v", msg.Code)
|
||||||
|
|
@ -274,7 +275,7 @@ func (self *bzzProtocol) handleStatus() error {
|
||||||
peer: peer{bzzProtocol: self, pubkey: status.NodeID},
|
peer: peer{bzzProtocol: self, pubkey: status.NodeID},
|
||||||
}
|
}
|
||||||
|
|
||||||
self.netStore.addPeers(req)
|
self.hive.addPeers(req)
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue