@cryptotaxi247 / kubo / commits / 3711d5409

getValueSingle using SendRequest

Juan Batiz-Benet committed Sep 16, 2014 at 02:43 UTC 3711d540988b5edba085a6fa99fd193148db0bb6
1 file changed +38 -33
routing/dht/dht.go
+38 -33
@@ -159,8 +159,37 @@ func (dht *IpfsDHT) HandleMessage(ctx context.Context, mes msg.NetMessage) (msg.
159 return rmes, nil
160 }
161
162 -func (dht *IpfsDHT) getValueOrPeers(p *peer.Peer, key u.Key, timeout time.Duration, level int) ([]byte, []*peer.Peer, error) {
163 - pmes, err := dht.getValueSingle(p, key, timeout, level)
162 +// sendRequest sends out a request using dht.sender, but also makes sure to
163 +// measure the RTT for latency measurements.
164 +func (dht *IpfsDHT) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message) (*Message, error) {
165 +
166 + mes, err := msg.FromObject(p, pmes)
167 + if err != nil {
168 + return nil, err
169 + }
170 +
171 + start := time.Now()
172 +
173 + rmes, err := dht.sender.SendRequest(ctx, mes)
174 + if err != nil {
175 + return nil, err
176 + }
177 +
178 + rtt := time.Since(start)
179 + rmes.Peer().SetLatency(rtt)
180 +
181 + rpmes := new(Message)
182 + if err := proto.Unmarshal(rmes.Data(), rpmes); err != nil {
183 + return nil, err
184 + }
185 +
186 + return rpmes, nil
187 +}
188 +
189 +func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
190 + key u.Key, level int) ([]byte, []*peer.Peer, error) {
191 +
192 + pmes, err := dht.getValueSingle(ctx, p, key, level)
193 if err != nil {
194 return nil, nil, err
195 }
@@ -202,39 +231,15 @@ func (dht *IpfsDHT) getValueOrPeers(p *peer.Peer, key u.Key, timeout time.Durati
231 }
232
233 // getValueSingle simply performs the get value RPC with the given parameters
205 -func (dht *IpfsDHT) getValueSingle(p *peer.Peer, key u.Key, timeout time.Duration, level int) (*Message, error) {
206 - pmes := Message{
207 - Type: Message_GET_VALUE,
208 - Key: string(key),
209 - Value: []byte{byte(level)},
210 - ID: swarm.GenerateMessageID(),
211 - }
212 - responseChan := dht.listener.Listen(pmes.ID, 1, time.Minute)
234 +func (dht *IpfsDHT) getValueSingle(ctx context.Context, p *peer.Peer,
235 + key u.Key, level int) (*Message, error) {
236
214 - mes := swarm.NewMessage(p, pmes.ToProtobuf())
215 - t := time.Now()
216 - dht.netChan.Outgoing <- mes
237 + typ := Message_GET_VALUE
238 + skey := string(key)
239 + pmes := &Message{Type: &typ, Key: &skey}
240 + pmes.SetClusterLevel(int32(level))
241
218 - // Wait for either the response or a timeout
219 - timeup := time.After(timeout)
220 - select {
221 - case <-timeup:
222 - dht.listener.Unlisten(pmes.ID)
223 - return nil, u.ErrTimeout
224 - case resp, ok := <-responseChan:
225 - if !ok {
226 - u.PErr("response channel closed before timeout, please investigate.\n")
227 - return nil, u.ErrTimeout
228 - }
229 - roundtrip := time.Since(t)
230 - resp.Peer.SetLatency(roundtrip)
231 - pmesOut := new(Message)
232 - err := proto.Unmarshal(resp.Data, pmesOut)
233 - if err != nil {
234 - return nil, err
235 - }
236 - return pmesOut, nil
237 - }
242 + return dht.sendRequest(ctx, p, pmes)
243 }
244
245 // TODO: Im not certain on this implementation, we get a list of peers/providers