@cryptotaxi247 / kubo / commits / d0304def6

refactor(blockstore, blockservice) use Blockstore and offline.Exchange

License: MIT Signed-off-by: Brian Tiger Chow <brian@perfmode.com>

Brian Tiger Chow committed Nov 20, 2014 at 18:08 UTC d0304def6beec57e6e5bc2c8f5325f8f2a4ef71b
8 files changed +65 -44
blocks/blockstore/blockstore.go
+10
@@ -15,6 +15,8 @@ import (
15 var ValueTypeMismatch = errors.New("The retrieved value is not a Block")
16
17 type Blockstore interface {
18 + DeleteBlock(u.Key) error
19 + Has(u.Key) (bool, error)
20 Get(u.Key) (*blocks.Block, error)
21 Put(*blocks.Block) error
22 }
@@ -45,3 +47,11 @@ func (bs *blockstore) Get(k u.Key) (*blocks.Block, error) {
47 func (bs *blockstore) Put(block *blocks.Block) error {
48 return bs.datastore.Put(block.Key().DsKey(), block.Data)
49 }
50 +
51 +func (bs *blockstore) Has(k u.Key) (bool, error) {
52 + return bs.datastore.Has(k.DsKey())
53 +}
54 +
55 +func (s *blockstore) DeleteBlock(k u.Key) error {
56 + return s.datastore.Delete(k.DsKey())
57 +}
blockservice/blocks_test.go
+5 -2
@@ -6,15 +6,18 @@ import (
6 "time"
7
8 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
9 -
9 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
10 + dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
11 blocks "github.com/jbenet/go-ipfs/blocks"
12 + blockstore "github.com/jbenet/go-ipfs/blocks/blockstore"
13 + offline "github.com/jbenet/go-ipfs/exchange/offline"
14 u "github.com/jbenet/go-ipfs/util"
15 )
16
17 func TestBlocks(t *testing.T) {
18 d := ds.NewMapDatastore()
17 - bs, err := NewBlockService(d, nil)
19 + tsds := dssync.MutexWrap(d)
20 + bs, err := New(blockstore.NewBlockstore(tsds), offline.Exchange())
21 if err != nil {
22 t.Error("failed to construct block service", err)
23 return
blockservice/blockservice.go
+25 -34
@@ -9,9 +9,9 @@ import (
9
10 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
11 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
12 - mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
12
13 blocks "github.com/jbenet/go-ipfs/blocks"
14 + "github.com/jbenet/go-ipfs/blocks/blockstore"
15 exchange "github.com/jbenet/go-ipfs/exchange"
16 u "github.com/jbenet/go-ipfs/util"
17 )
@@ -19,25 +19,28 @@ import (
19 var log = u.Logger("blockservice")
20 var ErrNotFound = errors.New("blockservice: key not found")
21
22 -// BlockService is a block datastore.
22 +// BlockService is a hybrid block datastore. It stores data in a local
23 +// datastore and may retrieve data from a remote Exchange.
24 // It uses an internal `datastore.Datastore` instance to store values.
25 type BlockService struct {
25 - Datastore ds.Datastore
26 - Remote exchange.Interface
26 + // TODO don't expose underlying impl details
27 + Blockstore blockstore.Blockstore
28 + Remote exchange.Interface
29 }
30
31 // NewBlockService creates a BlockService with given datastore instance.
30 -func NewBlockService(d ds.Datastore, rem exchange.Interface) (*BlockService, error) {
31 - if d == nil {
32 - return nil, fmt.Errorf("BlockService requires valid datastore")
32 +func New(bs blockstore.Blockstore, rem exchange.Interface) (*BlockService, error) {
33 + if bs == nil {
34 + return nil, fmt.Errorf("BlockService requires valid blockstore")
35 }
36 if rem == nil {
37 log.Warning("blockservice running in local (offline) mode.")
38 }
37 - return &BlockService{Datastore: d, Remote: rem}, nil
39 + return &BlockService{Blockstore: bs, Remote: rem}, nil
40 }
41
42 // AddBlock adds a particular block to the service, Putting it into the datastore.
43 +// TODO pass a context into this if the remote.HasBlock is going to remain here.
44 func (s *BlockService) AddBlock(b *blocks.Block) (u.Key, error) {
45 k := b.Key()
46 log.Debugf("blockservice: storing [%s] in datastore", k)
@@ -47,7 +50,7 @@ func (s *BlockService) AddBlock(b *blocks.Block) (u.Key, error) {
50 // check if we have it before adding. this is an extra read, but large writes
51 // are more expensive.
52 // TODO(jbenet) cheaper has. https://github.com/jbenet/go-datastore/issues/6
50 - has, err := s.Datastore.Has(k.DsKey())
53 + has, err := s.Blockstore.Has(k)
54 if err != nil {
55 return k, err
56 }
@@ -55,12 +58,14 @@ func (s *BlockService) AddBlock(b *blocks.Block) (u.Key, error) {
58 log.Debugf("blockservice: storing [%s] in datastore (already stored)", k)
59 } else {
60 log.Debugf("blockservice: storing [%s] in datastore", k)
58 - err := s.Datastore.Put(k.DsKey(), b.Data)
61 + err := s.Blockstore.Put(b)
62 if err != nil {
63 return k, err
64 }
65 }
66
67 + // TODO this operation rate-limits blockservice operations, we should
68 + // consider moving this to an sync process.
69 if s.Remote != nil {
70 ctx := context.TODO()
71 err = s.Remote.HasBlock(ctx, *b)
@@ -72,17 +77,11 @@ func (s *BlockService) AddBlock(b *blocks.Block) (u.Key, error) {
77 // Getting it from the datastore using the key (hash).
78 func (s *BlockService) GetBlock(ctx context.Context, k u.Key) (*blocks.Block, error) {
79 log.Debugf("BlockService GetBlock: '%s'", k)
75 - datai, err := s.Datastore.Get(k.DsKey())
80 + block, err := s.Blockstore.Get(k)
81 if err == nil {
77 - log.Debug("Blockservice: Got data in datastore.")
78 - bdata, ok := datai.([]byte)
79 - if !ok {
80 - return nil, fmt.Errorf("data associated with %s is not a []byte", k)
81 - }
82 - return &blocks.Block{
83 - Multihash: mh.Multihash(k),
84 - Data: bdata,
85 - }, nil
82 + return block, nil
83 + // TODO be careful checking ErrNotFound. If the underlying
84 + // implementation changes, this will break.
85 } else if err == ds.ErrNotFound && s.Remote != nil {
86 log.Debug("Blockservice: Searching bitswap.")
87 blk, err := s.Remote.GetBlock(ctx, k)
@@ -101,21 +100,13 @@ func (s *BlockService) GetBlocks(ctx context.Context, ks []u.Key) <-chan *blocks
100 go func() {
101 var toFetch []u.Key
102 for _, k := range ks {
104 - datai, err := s.Datastore.Get(k.DsKey())
105 - if err == nil {
106 - log.Debug("Blockservice: Got data in datastore.")
107 - bdata, ok := datai.([]byte)
108 - if !ok {
109 - log.Criticalf("data associated with %s is not a []byte", k)
110 - continue
111 - }
112 - out <- &blocks.Block{
113 - Multihash: mh.Multihash(k),
114 - Data: bdata,
115 - }
116 - } else {
103 + block, err := s.Blockstore.Get(k)
104 + if err != nil {
105 toFetch = append(toFetch, k)
106 + continue
107 }
108 + log.Debug("Blockservice: Got data in datastore.")
109 + out <- block
110 }
111 }()
112 return out
@@ -123,5 +114,5 @@ func (s *BlockService) GetBlocks(ctx context.Context, ks []u.Key) <-chan *blocks
114
115 // DeleteBlock deletes a block in the blockservice from the datastore
116 func (s *BlockService) DeleteBlock(k u.Key) error {
126 - return s.Datastore.Delete(k.DsKey())
117 + return s.Blockstore.DeleteBlock(k)
118 }
core/core.go
+4 -1
@@ -15,6 +15,7 @@ import (
15 exchange "github.com/jbenet/go-ipfs/exchange"
16 bitswap "github.com/jbenet/go-ipfs/exchange/bitswap"
17 bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
18 + "github.com/jbenet/go-ipfs/exchange/offline"
19 mount "github.com/jbenet/go-ipfs/fuse/mount"
20 merkledag "github.com/jbenet/go-ipfs/merkledag"
21 namesys "github.com/jbenet/go-ipfs/namesys"
@@ -127,6 +128,8 @@ func NewIpfsNode(cfg *config.Config, online bool) (n *IpfsNode, err error) {
128 return nil, debugerror.Wrap(err)
129 }
130
131 + n.Exchange = offline.Exchange()
132 +
133 // setup online services
134 if online {
135
@@ -178,7 +181,7 @@ func NewIpfsNode(cfg *config.Config, online bool) (n *IpfsNode, err error) {
181
182 // TODO(brian): when offline instantiate the BlockService with a bitswap
183 // session that simply doesn't return blocks
181 - n.Blocks, err = bserv.NewBlockService(n.Datastore, n.Exchange)
184 + n.Blocks, err = bserv.New(blockstore.NewBlockstore(n.Datastore), n.Exchange)
185 if err != nil {
186 return nil, debugerror.Wrap(err)
187 }
core/mock.go
+4 -4
@@ -3,8 +3,10 @@ package core
3 import (
4 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
5 syncds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
6 - bs "github.com/jbenet/go-ipfs/blockservice"
6 + "github.com/jbenet/go-ipfs/blocks/blockstore"
7 + blockservice "github.com/jbenet/go-ipfs/blockservice"
8 ci "github.com/jbenet/go-ipfs/crypto"
9 + "github.com/jbenet/go-ipfs/exchange/offline"
10 mdag "github.com/jbenet/go-ipfs/merkledag"
11 nsys "github.com/jbenet/go-ipfs/namesys"
12 path "github.com/jbenet/go-ipfs/path"
@@ -43,9 +45,7 @@ func NewMockNode() (*IpfsNode, error) {
45 nd.Routing = dht
46
47 // Bitswap
46 - //??
47 -
48 - bserv, err := bs.NewBlockService(nd.Datastore, nil)
48 + bserv, err := blockservice.New(blockstore.NewBlockstore(nd.Datastore), offline.Exchange())
49 if err != nil {
50 return nil, err
51 }
pin/pin_test.go
+6 -1
@@ -4,7 +4,10 @@ import (
4 "testing"
5
6 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
7 + dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
8 + "github.com/jbenet/go-ipfs/blocks/blockstore"
9 bs "github.com/jbenet/go-ipfs/blockservice"
10 + "github.com/jbenet/go-ipfs/exchange/offline"
11 mdag "github.com/jbenet/go-ipfs/merkledag"
12 "github.com/jbenet/go-ipfs/util"
13 )
@@ -19,13 +22,15 @@ func randNode() (*mdag.Node, util.Key) {
22
23 func TestPinnerBasic(t *testing.T) {
24 dstore := ds.NewMapDatastore()
22 - bserv, err := bs.NewBlockService(dstore, nil)
25 + bstore := blockstore.NewBlockstore(dssync.MutexWrap(dstore))
26 + bserv, err := bs.New(bstore, offline.Exchange())
27 if err != nil {
28 t.Fatal(err)
29 }
30
31 dserv := mdag.NewDAGService(bserv)
32
33 + // TODO does pinner need to share datastore with blockservice?
34 p := NewPinner(dstore, dserv)
35
36 a, ak := randNode()
unixfs/io/dagmodifier_test.go
+6 -1
@@ -6,7 +6,10 @@ import (
6 "io/ioutil"
7 "testing"
8
9 + "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
10 + "github.com/jbenet/go-ipfs/blocks/blockstore"
11 bs "github.com/jbenet/go-ipfs/blockservice"
12 + "github.com/jbenet/go-ipfs/exchange/offline"
13 imp "github.com/jbenet/go-ipfs/importer"
14 "github.com/jbenet/go-ipfs/importer/chunk"
15 mdag "github.com/jbenet/go-ipfs/merkledag"
@@ -19,7 +22,9 @@ import (
22
23 func getMockDagServ(t *testing.T) mdag.DAGService {
24 dstore := ds.NewMapDatastore()
22 - bserv, err := bs.NewBlockService(dstore, nil)
25 + tsds := sync.MutexWrap(dstore)
26 + bstore := blockstore.NewBlockstore(tsds)
27 + bserv, err := bs.New(bstore, offline.Exchange())
28 if err != nil {
29 t.Fatal(err)
30 }
util/testutil/gen.go
+5 -1
@@ -4,9 +4,12 @@ import (
4 crand "crypto/rand"
5 "testing"
6
7 + dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
8 + "github.com/jbenet/go-ipfs/exchange/offline"
9 "github.com/jbenet/go-ipfs/peer"
10
11 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
12 + "github.com/jbenet/go-ipfs/blocks/blockstore"
13 bsrv "github.com/jbenet/go-ipfs/blockservice"
14 dag "github.com/jbenet/go-ipfs/merkledag"
15 u "github.com/jbenet/go-ipfs/util"
@@ -14,7 +17,8 @@ import (
17
18 func GetDAGServ(t testing.TB) dag.DAGService {
19 dstore := ds.NewMapDatastore()
17 - bserv, err := bsrv.NewBlockService(dstore, nil)
20 + tsds := dssync.MutexWrap(dstore)
21 + bserv, err := bsrv.New(blockstore.NewBlockstore(tsds), offline.Exchange())
22 if err != nil {
23 t.Fatal(err)
24 }