diff --git a/bzz/netstore.go b/bzz/netstore.go index 6fddd75691..1a0fbb23b7 100644 --- a/bzz/netstore.go +++ b/bzz/netstore.go @@ -130,7 +130,7 @@ func (self *NetStore) addStoreRequest(req *storeRequestMsgData) { } else { return } - chunk.source = &req.peer + chunk.source = req.peer self.put(chunk) } @@ -170,12 +170,18 @@ func (self *NetStore) get(key Key) (chunk *Chunk) { } if chunk.req == nil { - chunk.req = new(requestStatus) - chunk.req.C = make(chan bool) + chunk.req = newRequestStatus() } return } +func newRequestStatus() *requestStatus { + return &requestStatus{ + requesters: make(map[int64][]*retrieveRequestMsgData), + C: make(chan bool), + } +} + func (self *NetStore) addRetrieveRequest(req *retrieveRequestMsgData) { self.lock.Lock() @@ -214,6 +220,7 @@ func (self *NetStore) startSearch(chunk *Chunk, id int64, timeout *time.Time) { } for _, peer := range peers { dpaLogger.Debugf("NetStore.startSearch: sending retrieveRequests to peer [%064x]", req.Key) + dpaLogger.Debugf("req.requesters: %v", chunk.req.requesters) var requester bool OUT: for _, recipients := range chunk.req.requesters { @@ -238,7 +245,7 @@ func generateId() int64 { /* adds a new peer to an existing open request only add if less than requesterCount peers forwarded the same request id so far -note this is done irrespective of status (searching or found/timedOut) +note this is done irrespective of status (searching or found) */ func (self *NetStore) addRequester(rs *requestStatus, req *retrieveRequestMsgData) { dpaLogger.Debugf("NetStore.addRequester: key %064x - add peer [%v] to req.Id %064x", req.Key, req.peer, req.Id) @@ -260,7 +267,7 @@ this is the most simplistic implementation: */ func (self *NetStore) strategyUpdateRequest(rs *requestStatus, req *retrieveRequestMsgData) (timeout *time.Time) { dpaLogger.Debugf("NetStore.strategyUpdateRequest: key %064x", req.Key) - + self.addRequester(rs, req) if rs.status == reqSearching { timeout = self.searchTimeout(rs, req) } diff --git a/bzz/protocol.go b/bzz/protocol.go index 40951e9c35..74d6325a6d 100644 --- a/bzz/protocol.go +++ b/bzz/protocol.go @@ -119,11 +119,17 @@ type storeRequestMsgData struct { storageTimeout *time.Time // expiry of content Metadata metaData // // - peer peer + peer *peer } func (self storeRequestMsgData) String() string { - return fmt.Sprintf("From: %v, Key: %x; ID: %v, requestTimeout: %v, storageTimeout: %v, SData %x", self.peer.Addr(), self.Key[:4], self.Id, self.requestTimeout, self.storageTimeout, self.SData[:10]) + var from string + if self.peer == nil { + from = "self" + } else { + from = self.peer.Addr().String() + } + return fmt.Sprintf("From: %v, Key: %x; ID: %v, requestTimeout: %v, storageTimeout: %v, SData %x", from, self.Key[:4], self.Id, self.requestTimeout, self.storageTimeout, self.SData[:10]) } /* @@ -142,11 +148,17 @@ type retrieveRequestMsgData struct { timeout *time.Time // //Metadata metaData // // - peer peer // protocol registers the requester + peer *peer // protocol registers the requester } func (self retrieveRequestMsgData) String() string { - return fmt.Sprintf("From: %v, Key: %x; ID: %v, MaxSize: %v, MaxPeers: %v", self.peer.Addr(), self.Key[:4], self.Id, self.MaxSize, self.MaxPeers) + var from string + if self.peer == nil { + from = "self" + } else { + from = self.peer.Addr().String() + } + return fmt.Sprintf("From: %v, Key: %x; ID: %v, MaxSize: %v, MaxPeers: %v", from, self.Key[:4], self.Id, self.MaxSize, self.MaxPeers) } type peerAddr struct { @@ -187,7 +199,7 @@ type peersMsgData struct { Key Key // if a response to a retrieval request Id uint64 // if a response to a retrieval request // - peer peer + peer *peer } /* @@ -280,7 +292,7 @@ func (self *bzzProtocol) handle() error { if err := msg.Decode(&req); err != nil { return self.protoError(ErrDecode, "msg %v: %v", msg, err) } - req.peer = peer{bzzProtocol: self} + req.peer = &peer{bzzProtocol: self} self.netStore.addStoreRequest(&req) case retrieveRequestMsg: @@ -292,7 +304,7 @@ func (self *bzzProtocol) handle() error { if req.Key == nil { return self.protoError(ErrDecode, "protocol handler: req.Key == nil || req.Timeout == nil") } - req.peer = peer{bzzProtocol: self} + req.peer = &peer{bzzProtocol: self} self.netStore.addRetrieveRequest(&req) case peersMsg: @@ -300,7 +312,7 @@ func (self *bzzProtocol) handle() error { if err := msg.Decode(&req); err != nil { return self.protoError(ErrDecode, "->msg %v: %v", msg, err) } - req.peer = peer{bzzProtocol: self} + req.peer = &peer{bzzProtocol: self} self.netStore.hive.addPeerEntries(&req) default: diff --git a/common/kademlia/kademlia.go b/common/kademlia/kademlia.go index cd501b177b..4e679876e0 100644 --- a/common/kademlia/kademlia.go +++ b/common/kademlia/kademlia.go @@ -52,6 +52,10 @@ type Kademlia struct { type Address common.Hash +func (a Address) String() string { + return fmt.Sprintf("%x", a[:]) +} + type Node interface { Addr() Address Url() string