@cryptotaxi247 / kubo / commits / 87407a99b

add context to blockservice Get

Jeromy committed Oct 25, 2014 at 12:38 UTC 87407a99b9904cd2ad30caf1dac5344a5db7a0c2
5 files changed +23 -13
blockservice/blocks_test.go
+5 -1
@@ -3,6 +3,9 @@ package blockservice
3 import (
4 "bytes"
5 "testing"
6 + "time"
7 +
8 + "code.google.com/p/go.net/context"
9
10 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
11 blocks "github.com/jbenet/go-ipfs/blocks"
@@ -37,7 +40,8 @@ func TestBlocks(t *testing.T) {
40 t.Error("returned key is not equal to block key", err)
41 }
42
40 - b2, err := bs.GetBlock(b.Key())
43 + ctx, _ := context.WithTimeout(context.TODO(), time.Second*5)
44 + b2, err := bs.GetBlock(ctx, b.Key())
45 if err != nil {
46 t.Error("failed to retrieve block from BlockService", err)
47 return
blockservice/blockservice.go
+1 -3
@@ -2,7 +2,6 @@ package blockservice
2
3 import (
4 "fmt"
5 - "time"
5
6 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
@@ -52,7 +51,7 @@ func (s *BlockService) AddBlock(b *blocks.Block) (u.Key, error) {
51
52 // GetBlock retrieves a particular block from the service,
53 // Getting it from the datastore using the key (hash).
55 -func (s *BlockService) GetBlock(k u.Key) (*blocks.Block, error) {
54 +func (s *BlockService) GetBlock(ctx context.Context, k u.Key) (*blocks.Block, error) {
55 log.Debug("BlockService GetBlock: '%s'", k)
56 datai, err := s.Datastore.Get(k.DsKey())
57 if err == nil {
@@ -67,7 +66,6 @@ func (s *BlockService) GetBlock(k u.Key) (*blocks.Block, error) {
66 }, nil
67 } else if err == ds.ErrNotFound && s.Remote != nil {
68 log.Debug("Blockservice: Searching bitswap.")
70 - ctx, _ := context.WithTimeout(context.TODO(), 5*time.Second)
69 blk, err := s.Remote.Block(ctx, k)
70 if err != nil {
71 return nil, err
core/commands/block.go
+5 -1
@@ -5,6 +5,9 @@ import (
5 "io"
6 "io/ioutil"
7 "os"
8 + "time"
9 +
10 + "code.google.com/p/go.net/context"
11
12 mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
13 "github.com/jbenet/go-ipfs/blocks"
@@ -26,7 +29,8 @@ func BlockGet(n *core.IpfsNode, args []string, opts map[string]interface{}, out
29
30 k := u.Key(h)
31 log.Debug("BlockGet key: '%q'", k)
29 - b, err := n.Blocks.GetBlock(k)
32 + ctx, _ := context.WithTimeout(context.TODO(), time.Second*5)
33 + b, err := n.Blocks.GetBlock(ctx, k)
34 if err != nil {
35 return fmt.Errorf("block get: %v", err)
36 }
merkledag/merkledag.go
+5 -1
@@ -2,6 +2,9 @@ package merkledag
2
3 import (
4 "fmt"
5 + "time"
6 +
7 + "code.google.com/p/go.net/context"
8
9 mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
10 blocks "github.com/jbenet/go-ipfs/blocks"
@@ -204,7 +207,8 @@ func (n *DAGService) Get(k u.Key) (*Node, error) {
207 return nil, fmt.Errorf("DAGService is nil")
208 }
209
207 - b, err := n.Blocks.GetBlock(k)
210 + ctx, _ := context.WithTimeout(context.TODO(), time.Second*5)
211 + b, err := n.Blocks.GetBlock(ctx, k)
212 if err != nil {
213 return nil, err
214 }
net/conn/multiconn.go
+7 -7
@@ -67,7 +67,7 @@ func (c *MultiConn) Add(conns ...Conn) {
67 defer c.Unlock()
68
69 for _, c2 := range conns {
70 - log.Info("MultiConn: adding %s", c2)
70 + log.Infof("MultiConn: adding %s", c2)
71 if c.LocalPeer() != c2.LocalPeer() || c.RemotePeer() != c2.RemotePeer() {
72 log.Error(c2)
73 c.Unlock() // ok to unlock (to log). panicing.
@@ -82,7 +82,7 @@ func (c *MultiConn) Add(conns ...Conn) {
82
83 c.conns[c2.ID()] = c2
84 go c.fanInSingle(c2)
85 - log.Info("MultiConn: added %s", c2)
85 + log.Infof("MultiConn: added %s", c2)
86 }
87 }
88
@@ -146,7 +146,7 @@ func (c *MultiConn) fanOut() {
146 // send data out through our "best connection"
147 case m, more := <-c.duplex.Out:
148 if !more {
149 - log.Info("%s out channel closed", c)
149 + log.Infof("%s out channel closed", c)
150 return
151 }
152 sc := c.BestConn()
@@ -156,7 +156,7 @@ func (c *MultiConn) fanOut() {
156 }
157
158 i++
159 - log.Info("%s sending (%d)", sc, i)
159 + log.Infof("%s sending (%d)", sc, i)
160 sc.Out() <- m
161 }
162 }
@@ -170,7 +170,7 @@ func (c *MultiConn) fanInSingle(child Conn) {
170
171 // cleanup all data associated with this child Connection.
172 defer func() {
173 - log.Info("closing: %s", child)
173 + log.Infof("closing: %s", child)
174
175 // in case it still is in the map, remove it.
176 c.Lock()
@@ -197,11 +197,11 @@ func (c *MultiConn) fanInSingle(child Conn) {
197
198 case m, more := <-child.In(): // receiving data
199 if !more {
200 - log.Info("%s in channel closed", child)
200 + log.Infof("%s in channel closed", child)
201 return // closed
202 }
203 i++
204 - log.Info("%s received (%d)", child, i)
204 + log.Infof("%s received (%d)", child, i)
205 c.duplex.In <- m
206 }
207 }