@cryptotaxi247 / kubo / commits / 3183b1cb8

coreapi: Untangle from core.IpfsNode

License: MIT Signed-off-by: Łukasz Magiera <magik6k@gmail.com>

Łukasz Magiera committed Dec 6, 2018 at 21:00 UTC 3183b1cb8e44a3016776d9ba52e7c0d569442f0c
14 files changed +211 -125
core/commands/pin.go
+3 -3
@@ -88,7 +88,7 @@ var addPinCmd = &cmds.Command{
88 }
89
90 if !showProgress {
91 - added, err := corerepo.Pin(n, api, req.Context, req.Arguments, recursive)
91 + added, err := corerepo.Pin(n.Pinning, api, req.Context, req.Arguments, recursive)
92 if err != nil {
93 return err
94 }
@@ -105,7 +105,7 @@ var addPinCmd = &cmds.Command{
105
106 ch := make(chan pinResult, 1)
107 go func() {
108 - added, err := corerepo.Pin(n, api, ctx, req.Arguments, recursive)
108 + added, err := corerepo.Pin(n.Pinning, api, ctx, req.Arguments, recursive)
109 ch <- pinResult{pins: added, err: err}
110 }()
111
@@ -215,7 +215,7 @@ collected if needed. (By default, recursively. Use -r=false for direct pins.)
215 return err
216 }
217
218 - removed, err := corerepo.Unpin(n, api, req.Context, req.Arguments, recursive)
218 + removed, err := corerepo.Unpin(n.Pinning, api, req.Context, req.Arguments, recursive)
219 if err != nil {
220 return err
221 }
core/coreapi/block.go
+4 -4
@@ -43,7 +43,7 @@ func (api *BlockAPI) Put(ctx context.Context, src io.Reader, opts ...caopts.Bloc
43 return nil, err
44 }
45
46 - err = api.node.Blocks.AddBlock(b)
46 + err = api.blocks.AddBlock(b)
47 if err != nil {
48 return nil, err
49 }
@@ -57,7 +57,7 @@ func (api *BlockAPI) Get(ctx context.Context, p coreiface.Path) (io.Reader, erro
57 return nil, err
58 }
59
60 - b, err := api.node.Blocks.GetBlock(ctx, rp.Cid())
60 + b, err := api.blocks.GetBlock(ctx, rp.Cid())
61 if err != nil {
62 return nil, err
63 }
@@ -78,7 +78,7 @@ func (api *BlockAPI) Rm(ctx context.Context, p coreiface.Path, opts ...caopts.Bl
78 cids := []cid.Cid{rp.Cid()}
79 o := util.RmBlocksOpts{Force: settings.Force}
80
81 - out, err := util.RmBlocks(api.node.Blockstore, api.node.Pinning, cids, o)
81 + out, err := util.RmBlocks(api.blockstore, api.pinning, cids, o)
82 if err != nil {
83 return err
84 }
@@ -109,7 +109,7 @@ func (api *BlockAPI) Stat(ctx context.Context, p coreiface.Path) (coreiface.Bloc
109 return nil, err
110 }
111
112 - b, err := api.node.Blocks.GetBlock(ctx, rp.Cid())
112 + b, err := api.blocks.GetBlock(ctx, rp.Cid())
113 if err != nil {
114 return nil, err
115 }
core/coreapi/coreapi.go
+89 -7
@@ -15,25 +15,103 @@ package coreapi
15
16 import (
17 "context"
18 + "errors"
19
20 core "github.com/ipfs/go-ipfs/core"
21 coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
21 -
22 + options "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
23 + namesys "github.com/ipfs/go-ipfs/namesys"
24 + pin "github.com/ipfs/go-ipfs/pin"
25 + repo "github.com/ipfs/go-ipfs/repo"
26 +
27 + ci "gx/ipfs/QmNiJiXwWE3kRhZrC5ej3kSjWHm337pYfhjLGSCDNKJP2s/go-libp2p-crypto"
28 + exchange "gx/ipfs/QmP2g3VxmC7g7fyRJDj1VJ72KHZbJ9UW24YjSWEj1XTb4H/go-ipfs-exchange-interface"
29 + bserv "gx/ipfs/QmPoh3SrQzFBWtdGK6qmHDV4EanKR6kYPj4DD3J2NLoEmZ/go-blockservice"
30 + routing "gx/ipfs/QmRASJXJUFygM5qU4YrH7k7jD6S4Hg8nJmgqJ4bYJvLatd/go-libp2p-routing"
31 + blockstore "gx/ipfs/QmS2aqUZLJp8kF1ihE5rvDGE5LvmKDPnx32w9Z1BW9xLV5/go-ipfs-blockstore"
32 + peer "gx/ipfs/QmY5Grm8pJdiSSVsYxx4uNRgweY72EmYwuSDbRnbFok3iY/go-libp2p-peer"
33 + pstore "gx/ipfs/QmZ9zH2FnLcxv1xyzFeUpDUeo55xEhZQHgveZijcxr7TLj/go-libp2p-peerstore"
34 + pubsub "gx/ipfs/QmaqGyUhWLsJbVo1QAujSu13mxNjFJ98Kt2VWGSnShGE1Q/go-libp2p-pubsub"
35 ipld "gx/ipfs/QmcKKBwfz6FyQdHR2jsXrrF6XeSBXYL86anmWNewpFpoF5/go-ipld-format"
36 logging "gx/ipfs/QmcuXC5cxs79ro2cUuHs4HQ2bkDLJUYokwL8aivcX6HW3C/go-log"
37 dag "gx/ipfs/QmdV35UHnL1FM52baPkeUo6u7Fxm2CRUkPTLRPxeF8a4Ap/go-merkledag"
38 + record "gx/ipfs/QmfARXVCzpwFXQdepAJZuqyNDgV9doEsMnVCo1ssmuSe1U/go-libp2p-record"
39 + p2phost "gx/ipfs/QmfD51tKgJiTMnW9JEiDiPwsCY4mqUoxkhKhBfyW12spTC/go-libp2p-host"
40 )
41
42 var log = logging.Logger("core/coreapi")
43
44 type CoreAPI struct {
30 - node *core.IpfsNode
31 - dag ipld.DAGService
45 + nctx context.Context
46 +
47 + identity peer.ID //TODO: check mutable structs
48 + privateKey ci.PrivKey
49 +
50 + repo repo.Repo
51 + blockstore blockstore.GCBlockstore
52 + baseBlocks blockstore.Blockstore
53 + blocks bserv.BlockService
54 + dag ipld.DAGService
55 + pinning pin.Pinner
56 +
57 + peerstore pstore.Peerstore
58 + peerHost p2phost.Host
59 + namesys namesys.NameSystem
60 + recordValidator record.Validator
61 + exchange exchange.Interface
62 +
63 + routing routing.IpfsRouting
64 + pubSub *pubsub.PubSub
65 +
66 + checkRouting func(bool) error
67 +
68 + // TODO: this can be generalized to all functions when we implement some
69 + // api based security mechanism
70 + isPublishAllowed func() error
71 }
72
73 // NewCoreAPI creates new instance of IPFS CoreAPI backed by go-ipfs Node.
35 -func NewCoreAPI(n *core.IpfsNode) coreiface.CoreAPI {
36 - api := &CoreAPI{n, n.DAG}
74 +func NewCoreAPI(n *core.IpfsNode, opts ...options.ApiOption) coreiface.CoreAPI {
75 + api := &CoreAPI{
76 + nctx: n.Context(),
77 +
78 + identity: n.Identity,
79 + privateKey: n.PrivateKey,
80 +
81 + repo: n.Repo,
82 + blockstore: n.Blockstore,
83 + baseBlocks: n.BaseBlocks,
84 + blocks: n.Blocks,
85 + dag: n.DAG,
86 + pinning: n.Pinning,
87 +
88 + peerstore: n.Peerstore,
89 + peerHost: n.PeerHost,
90 + namesys: n.Namesys,
91 + recordValidator: n.RecordValidator,
92 + exchange: n.Exchange,
93 +
94 + routing: n.Routing,
95 + pubSub: n.PubSub,
96 +
97 + checkRouting: func(allowOffline bool) error {
98 + if !n.OnlineMode() {
99 + if !allowOffline {
100 + return coreiface.ErrOffline
101 + }
102 + return n.SetupOfflineRouting()
103 + }
104 + return nil
105 + },
106 +
107 + isPublishAllowed: func() error {
108 + if n.Mounts.Ipns != nil && n.Mounts.Ipns.IsActive() {
109 + return errors.New("cannot manually publish while IPNS is mounted")
110 + }
111 + return nil
112 + },
113 + }
114 +
115 return api
116 }
117
@@ -89,6 +167,10 @@ func (api *CoreAPI) PubSub() coreiface.PubSubAPI {
167
168 // getSession returns new api backed by the same node with a read-only session DAG
169 func (api *CoreAPI) getSession(ctx context.Context) *CoreAPI {
92 - ng := dag.NewReadOnlyDagService(dag.NewSession(ctx, api.dag))
93 - return &CoreAPI{api.node, ng}
170 + sesApi := *api
171 +
172 + //TODO: we may want to apply this to other things too
173 + sesApi.dag = dag.NewReadOnlyDagService(dag.NewSession(ctx, api.dag))
174 +
175 + return &sesApi
176 }
core/coreapi/dht.go
+6 -6
@@ -22,7 +22,7 @@ import (
22 type DhtAPI CoreAPI
23
24 func (api *DhtAPI) FindPeer(ctx context.Context, p peer.ID) (pstore.PeerInfo, error) {
25 - pi, err := api.node.Routing.FindPeer(ctx, peer.ID(p))
25 + pi, err := api.routing.FindPeer(ctx, peer.ID(p))
26 if err != nil {
27 return pstore.PeerInfo{}, err
28 }
@@ -46,7 +46,7 @@ func (api *DhtAPI) FindProviders(ctx context.Context, p coreiface.Path, opts ...
46 return nil, fmt.Errorf("number of providers must be greater than 0")
47 }
48
49 - pchan := api.node.Routing.FindProvidersAsync(ctx, rp.Cid(), numProviders)
49 + pchan := api.routing.FindProvidersAsync(ctx, rp.Cid(), numProviders)
50 return pchan, nil
51 }
52
@@ -56,7 +56,7 @@ func (api *DhtAPI) Provide(ctx context.Context, path coreiface.Path, opts ...cao
56 return err
57 }
58
59 - if api.node.Routing == nil {
59 + if api.routing == nil {
60 return errors.New("cannot provide in offline mode")
61 }
62
@@ -67,7 +67,7 @@ func (api *DhtAPI) Provide(ctx context.Context, path coreiface.Path, opts ...cao
67
68 c := rp.Cid()
69
70 - has, err := api.node.Blockstore.Has(c)
70 + has, err := api.blockstore.Has(c)
71 if err != nil {
72 return err
73 }
@@ -77,9 +77,9 @@ func (api *DhtAPI) Provide(ctx context.Context, path coreiface.Path, opts ...cao
77 }
78
79 if settings.Recursive {
80 - err = provideKeysRec(ctx, api.node.Routing, api.node.Blockstore, []cid.Cid{c})
80 + err = provideKeysRec(ctx, api.routing, api.blockstore, []cid.Cid{c})
81 } else {
82 - err = provideKeys(ctx, api.node.Routing, []cid.Cid{c})
82 + err = provideKeys(ctx, api.routing, []cid.Cid{c})
83 }
84 if err != nil {
85 return err
core/coreapi/key.go
+9 -9
@@ -54,7 +54,7 @@ func (api *KeyAPI) Generate(ctx context.Context, name string, opts ...caopts.Key
54 return nil, fmt.Errorf("cannot create key with name 'self'")
55 }
56
57 - _, err = api.node.Repo.Keystore().Get(name)
57 + _, err = api.repo.Keystore().Get(name)
58 if err == nil {
59 return nil, fmt.Errorf("key with name '%s' already exists", name)
60 }
@@ -87,7 +87,7 @@ func (api *KeyAPI) Generate(ctx context.Context, name string, opts ...caopts.Key
87 return nil, fmt.Errorf("unrecognized key type: %s", options.Algorithm)
88 }
89
90 - err = api.node.Repo.Keystore().Put(name, sk)
90 + err = api.repo.Keystore().Put(name, sk)
91 if err != nil {
92 return nil, err
93 }
@@ -102,7 +102,7 @@ func (api *KeyAPI) Generate(ctx context.Context, name string, opts ...caopts.Key
102
103 // List returns a list keys stored in keystore.
104 func (api *KeyAPI) List(ctx context.Context) ([]coreiface.Key, error) {
105 - keys, err := api.node.Repo.Keystore().List()
105 + keys, err := api.repo.Keystore().List()
106 if err != nil {
107 return nil, err
108 }
@@ -110,10 +110,10 @@ func (api *KeyAPI) List(ctx context.Context) ([]coreiface.Key, error) {
110 sort.Strings(keys)
111
112 out := make([]coreiface.Key, len(keys)+1)
113 - out[0] = &key{"self", api.node.Identity}
113 + out[0] = &key{"self", api.identity}
114
115 for n, k := range keys {
116 - privKey, err := api.node.Repo.Keystore().Get(k)
116 + privKey, err := api.repo.Keystore().Get(k)
117 if err != nil {
118 return nil, err
119 }
@@ -138,7 +138,7 @@ func (api *KeyAPI) Rename(ctx context.Context, oldName string, newName string, o
138 return nil, false, err
139 }
140
141 - ks := api.node.Repo.Keystore()
141 + ks := api.repo.Keystore()
142
143 if oldName == "self" {
144 return nil, false, fmt.Errorf("cannot rename key with name 'self'")
@@ -192,7 +192,7 @@ func (api *KeyAPI) Rename(ctx context.Context, oldName string, newName string, o
192
193 // Remove removes keys from keystore. Returns ipns path of the removed key.
194 func (api *KeyAPI) Remove(ctx context.Context, name string) (coreiface.Key, error) {
195 - ks := api.node.Repo.Keystore()
195 + ks := api.repo.Keystore()
196
197 if name == "self" {
198 return nil, fmt.Errorf("cannot remove key with name 'self'")
@@ -219,9 +219,9 @@ func (api *KeyAPI) Remove(ctx context.Context, name string) (coreiface.Key, erro
219 }
220
221 func (api *KeyAPI) Self(ctx context.Context) (coreiface.Key, error) {
222 - if api.node.Identity == "" {
222 + if api.identity == "" {
223 return nil, errors.New("identity not loaded")
224 }
225
226 - return &key{"self", api.node.Identity}, nil
226 + return &key{"self", api.identity}, nil
227 }
core/coreapi/name.go
+24 -21
@@ -7,13 +7,13 @@ import (
7 "strings"
8 "time"
9
10 - "github.com/ipfs/go-ipfs/core"
10 coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
11 caopts "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
12 "github.com/ipfs/go-ipfs/keystore"
13 "github.com/ipfs/go-ipfs/namesys"
14
15 "gx/ipfs/QmNiJiXwWE3kRhZrC5ej3kSjWHm337pYfhjLGSCDNKJP2s/go-libp2p-crypto"
16 + ci "gx/ipfs/QmNiJiXwWE3kRhZrC5ej3kSjWHm337pYfhjLGSCDNKJP2s/go-libp2p-crypto"
17 "gx/ipfs/QmY5Grm8pJdiSSVsYxx4uNRgweY72EmYwuSDbRnbFok3iY/go-libp2p-peer"
18 ipath "gx/ipfs/QmZErC2Ay6WuGi96CPg316PwitdwgLo6RxZRqVjJjRj2MR/go-path"
19 "gx/ipfs/QmdmWkx54g7VfVyxeG8ic84uf4G6Eq1GohuyKA3XDuJ8oC/go-ipfs-routing/offline"
@@ -38,20 +38,17 @@ func (e *ipnsEntry) Value() coreiface.Path {
38
39 // Publish announces new IPNS name and returns the new IPNS entry.
40 func (api *NameAPI) Publish(ctx context.Context, p coreiface.Path, opts ...caopts.NamePublishOption) (coreiface.IpnsEntry, error) {
41 - options, err := caopts.NamePublishOptions(opts...)
42 - if err != nil {
41 + if err := api.isPublishAllowed(); err != nil {
42 return nil, err
43 }
45 - n := api.node
44
47 - if !n.OnlineMode() {
48 - if !options.AllowOffline {
49 - return nil, coreiface.ErrOffline
50 - }
45 + options, err := caopts.NamePublishOptions(opts...)
46 + if err != nil {
47 + return nil, err
48 }
49
53 - if n.Mounts.Ipns != nil && n.Mounts.Ipns.IsActive() {
54 - return nil, errors.New("cannot manually publish while IPNS is mounted")
50 + if err := api.checkRouting(options.AllowOffline); err != nil {
51 + return nil, err
52 }
53
54 pth, err := ipath.ParsePath(p.String())
@@ -59,7 +56,7 @@ func (api *NameAPI) Publish(ctx context.Context, p coreiface.Path, opts ...caopt
56 return nil, err
57 }
58
62 - k, err := keylookup(n, options.Key)
59 + k, err := keylookup(api.privateKey, api.repo.Keystore(), options.Key)
60 if err != nil {
61 return nil, err
62 }
@@ -69,7 +66,7 @@ func (api *NameAPI) Publish(ctx context.Context, p coreiface.Path, opts ...caopt
66 }
67
68 eol := time.Now().Add(options.ValidTime)
72 - err = n.Namesys.PublishWithEOL(ctx, k, pth, eol)
69 + err = api.namesys.PublishWithEOL(ctx, k, pth, eol)
70 if err != nil {
71 return nil, err
72 }
@@ -91,21 +88,23 @@ func (api *NameAPI) Search(ctx context.Context, name string, opts ...caopts.Name
88 return nil, err
89 }
90
94 - n := api.node
91 + if err := api.checkRouting(true); err != nil {
92 + return nil, err
93 + }
94
96 - var resolver namesys.Resolver = n.Namesys
95 + var resolver namesys.Resolver = api.namesys
96
98 - if options.Local && !options.Cache {
97 + if options.Local && !options.Cache { //TODO: rm before offline/local global opt merge
98 return nil, errors.New("cannot specify both local and nocache")
99 }
100
101 if options.Local {
103 - offroute := offline.NewOfflineRouter(n.Repo.Datastore(), n.RecordValidator)
102 + offroute := offline.NewOfflineRouter(api.repo.Datastore(), api.recordValidator)
103 resolver = namesys.NewIpnsResolver(offroute)
104 }
105
106 if !options.Cache {
108 - resolver = namesys.NewNameSystem(n.Routing, n.Repo.Datastore(), 0)
107 + resolver = namesys.NewNameSystem(api.routing, api.repo.Datastore(), 0)
108 }
109
110 if !strings.HasPrefix(name, "/ipns/") {
@@ -150,8 +149,12 @@ func (api *NameAPI) Resolve(ctx context.Context, name string, opts ...caopts.Nam
149 return p, err
150 }
151
153 -func keylookup(n *core.IpfsNode, k string) (crypto.PrivKey, error) {
154 - res, err := n.GetKey(k)
152 +func keylookup(self ci.PrivKey, kstore keystore.Keystore, k string) (crypto.PrivKey, error) {
153 + if k == "self" {
154 + return self, nil
155 + }
156 +
157 + res, err := kstore.Get(k)
158 if res != nil {
159 return res, nil
160 }
@@ -160,13 +163,13 @@ func keylookup(n *core.IpfsNode, k string) (crypto.PrivKey, error) {
163 return nil, err
164 }
165
163 - keys, err := n.Repo.Keystore().List()
166 + keys, err := kstore.List()
167 if err != nil {
168 return nil, err
169 }
170
171 for _, key := range keys {
169 - privKey, err := n.Repo.Keystore().Get(key)
172 + privKey, err := kstore.Get(key)
173 if err != nil {
174 return nil, err
175 }
core/coreapi/object.go
+3 -3
@@ -118,7 +118,7 @@ func (api *ObjectAPI) Put(ctx context.Context, src io.Reader, opts ...caopts.Obj
118 }
119
120 if options.Pin {
121 - defer api.node.Blockstore.PinLock().Unlock()
121 + defer api.blockstore.PinLock().Unlock()
122 }
123
124 err = api.dag.Add(ctx, dagnode)
@@ -127,8 +127,8 @@ func (api *ObjectAPI) Put(ctx context.Context, src io.Reader, opts ...caopts.Obj
127 }
128
129 if options.Pin {
130 - api.node.Pinning.PinWithMode(dagnode.Cid(), pin.Recursive)
131 - err = api.node.Pinning.Flush()
130 + api.pinning.PinWithMode(dagnode.Cid(), pin.Recursive)
131 + err = api.pinning.Flush()
132 if err != nil {
133 return nil, err
134 }
core/coreapi/path.go
+1 -1
@@ -38,7 +38,7 @@ func (api *CoreAPI) ResolvePath(ctx context.Context, p coreiface.Path) (coreifac
38 }
39
40 ipath := ipfspath.Path(p.String())
41 - ipath, err := core.ResolveIPNS(ctx, api.node.Namesys, ipath)
41 + ipath, err := core.ResolveIPNS(ctx, api.namesys, ipath)
42 if err == core.ErrNoNamesys {
43 return nil, coreiface.ErrOffline
44 } else if err != nil {
core/coreapi/pin.go
+13 -13
@@ -27,14 +27,14 @@ func (api *PinAPI) Add(ctx context.Context, p coreiface.Path, opts ...caopts.Pin
27 return err
28 }
29
30 - defer api.node.Blockstore.PinLock().Unlock()
30 + defer api.blockstore.PinLock().Unlock()
31
32 - _, err = corerepo.Pin(api.node, api.core(), ctx, []string{rp.Cid().String()}, settings.Recursive)
32 + _, err = corerepo.Pin(api.pinning, api.core(), ctx, []string{rp.Cid().String()}, settings.Recursive)
33 if err != nil {
34 return err
35 }
36
37 - return api.node.Pinning.Flush()
37 + return api.pinning.Flush()
38 }
39
40 func (api *PinAPI) Ls(ctx context.Context, opts ...caopts.PinLsOption) ([]coreiface.Pin, error) {
@@ -53,12 +53,12 @@ func (api *PinAPI) Ls(ctx context.Context, opts ...caopts.PinLsOption) ([]coreif
53 }
54
55 func (api *PinAPI) Rm(ctx context.Context, p coreiface.Path) error {
56 - _, err := corerepo.Unpin(api.node, api.core(), ctx, []string{p.String()}, true)
56 + _, err := corerepo.Unpin(api.pinning, api.core(), ctx, []string{p.String()}, true)
57 if err != nil {
58 return err
59 }
60
61 - return api.node.Pinning.Flush()
61 + return api.pinning.Flush()
62 }
63
64 func (api *PinAPI) Update(ctx context.Context, from coreiface.Path, to coreiface.Path, opts ...caopts.PinUpdateOption) error {
@@ -77,14 +77,14 @@ func (api *PinAPI) Update(ctx context.Context, from coreiface.Path, to coreiface
77 return err
78 }
79
80 - defer api.node.Blockstore.PinLock().Unlock()
80 + defer api.blockstore.PinLock().Unlock()
81
82 - err = api.node.Pinning.Update(ctx, fp.Cid(), tp.Cid(), settings.Unpin)
82 + err = api.pinning.Update(ctx, fp.Cid(), tp.Cid(), settings.Unpin)
83 if err != nil {
84 return err
85 }
86
87 - return api.node.Pinning.Flush()
87 + return api.pinning.Flush()
88 }
89
90 type pinStatus struct {
@@ -117,10 +117,10 @@ func (n *badNode) Err() error {
117
118 func (api *PinAPI) Verify(ctx context.Context) (<-chan coreiface.PinStatus, error) {
119 visited := make(map[cid.Cid]*pinStatus)
120 - bs := api.node.Blocks.Blockstore()
120 + bs := api.blockstore
121 DAG := merkledag.NewDAGService(bserv.New(bs, offline.Exchange(bs)))
122 getLinks := merkledag.GetLinksWithDAG(DAG)
123 - recPins := api.node.Pinning.RecursiveKeys()
123 + recPins := api.pinning.RecursiveKeys()
124
125 var checkPin func(root cid.Cid) *pinStatus
126 checkPin = func(root cid.Cid) *pinStatus {
@@ -187,11 +187,11 @@ func (api *PinAPI) pinLsAll(typeStr string, ctx context.Context) ([]coreiface.Pi
187 }
188
189 if typeStr == "direct" || typeStr == "all" {
190 - AddToResultKeys(api.node.Pinning.DirectKeys(), "direct")
190 + AddToResultKeys(api.pinning.DirectKeys(), "direct")
191 }
192 if typeStr == "indirect" || typeStr == "all" {
193 set := cid.NewSet()
194 - for _, k := range api.node.Pinning.RecursiveKeys() {
194 + for _, k := range api.pinning.RecursiveKeys() {
195 err := merkledag.EnumerateChildren(ctx, merkledag.GetLinksWithDAG(api.dag), k, set.Visit)
196 if err != nil {
197 return nil, err
@@ -200,7 +200,7 @@ func (api *PinAPI) pinLsAll(typeStr string, ctx context.Context) ([]coreiface.Pi
200 AddToResultKeys(set.Keys(), "indirect")
201 }
202 if typeStr == "recursive" || typeStr == "all" {
203 - AddToResultKeys(api.node.Pinning.RecursiveKeys(), "recursive")
203 + AddToResultKeys(api.pinning.RecursiveKeys(), "recursive")
204 }
205
206 out := make([]coreiface.Pin, 0, len(keys))
core/coreapi/pubsub.go
+14 -13
@@ -7,14 +7,15 @@ import (
7 "sync"
8 "time"
9
10 - core "github.com/ipfs/go-ipfs/core"
10 coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
11 caopts "github.com/ipfs/go-ipfs/core/coreapi/interface/options"
12
13 cid "gx/ipfs/QmR8BauakNcBa3RbE4nbQu76PDiJgoQgz8AJdhJuiU4TAw/go-cid"
14 + routing "gx/ipfs/QmRASJXJUFygM5qU4YrH7k7jD6S4Hg8nJmgqJ4bYJvLatd/go-libp2p-routing"
15 peer "gx/ipfs/QmY5Grm8pJdiSSVsYxx4uNRgweY72EmYwuSDbRnbFok3iY/go-libp2p-peer"
16 pstore "gx/ipfs/QmZ9zH2FnLcxv1xyzFeUpDUeo55xEhZQHgveZijcxr7TLj/go-libp2p-peerstore"
17 pubsub "gx/ipfs/QmaqGyUhWLsJbVo1QAujSu13mxNjFJ98Kt2VWGSnShGE1Q/go-libp2p-pubsub"
18 + p2phost "gx/ipfs/QmfD51tKgJiTMnW9JEiDiPwsCY4mqUoxkhKhBfyW12spTC/go-libp2p-host"
19 )
20
21 type PubSubAPI CoreAPI
@@ -33,7 +34,7 @@ func (api *PubSubAPI) Ls(ctx context.Context) ([]string, error) {
34 return nil, err
35 }
36
36 - return api.node.PubSub.GetTopics(), nil
37 + return api.pubSub.GetTopics(), nil
38 }
39
40 func (api *PubSubAPI) Peers(ctx context.Context, opts ...caopts.PubSubPeersOption) ([]peer.ID, error) {
@@ -46,7 +47,7 @@ func (api *PubSubAPI) Peers(ctx context.Context, opts ...caopts.PubSubPeersOptio
47 return nil, err
48 }
49
49 - peers := api.node.PubSub.ListPeers(settings.Topic)
50 + peers := api.pubSub.ListPeers(settings.Topic)
51 out := make([]peer.ID, len(peers))
52
53 for i, peer := range peers {
@@ -61,7 +62,7 @@ func (api *PubSubAPI) Publish(ctx context.Context, topic string, data []byte) er
62 return err
63 }
64
64 - return api.node.PubSub.Publish(topic, data)
65 + return api.pubSub.Publish(topic, data)
66 }
67
68 func (api *PubSubAPI) Subscribe(ctx context.Context, topic string, opts ...caopts.PubSubSubscribeOption) (coreiface.PubSubSubscription, error) {
@@ -71,12 +72,12 @@ func (api *PubSubAPI) Subscribe(ctx context.Context, topic string, opts ...caopt
72 return nil, err
73 }
74
74 - sub, err := api.node.PubSub.Subscribe(topic)
75 + sub, err := api.pubSub.Subscribe(topic)
76 if err != nil {
77 return nil, err
78 }
79
79 - pubctx, cancel := context.WithCancel(api.node.Context())
80 + pubctx, cancel := context.WithCancel(api.nctx)
81
82 if options.Discover {
83 go func() {
@@ -86,18 +87,18 @@ func (api *PubSubAPI) Subscribe(ctx context.Context, topic string, opts ...caopt
87 return
88 }
89
89 - connectToPubSubPeers(pubctx, api.node, blk.Path().Cid())
90 + connectToPubSubPeers(pubctx, api.routing, api.peerHost, blk.Path().Cid())
91 }()
92 }
93
94 return &pubSubSubscription{cancel, sub}, nil
95 }
96
96 -func connectToPubSubPeers(ctx context.Context, n *core.IpfsNode, cid cid.Cid) {
97 +func connectToPubSubPeers(ctx context.Context, r routing.IpfsRouting, ph p2phost.Host, cid cid.Cid) {
98 ctx, cancel := context.WithCancel(ctx)
99 defer cancel()
100
100 - provs := n.Routing.FindProvidersAsync(ctx, cid, 10)
101 + provs := r.FindProvidersAsync(ctx, cid, 10)
102 var wg sync.WaitGroup
103 for p := range provs {
104 wg.Add(1)
@@ -105,7 +106,7 @@ func connectToPubSubPeers(ctx context.Context, n *core.IpfsNode, cid cid.Cid) {
106 defer wg.Done()
107 ctx, cancel := context.WithTimeout(ctx, time.Second*10)
108 defer cancel()
108 - err := n.PeerHost.Connect(ctx, pi)
109 + err := ph.Connect(ctx, pi)
110 if err != nil {
111 log.Info("pubsub discover: ", err)
112 return
@@ -118,11 +119,11 @@ func connectToPubSubPeers(ctx context.Context, n *core.IpfsNode, cid cid.Cid) {
119 }
120
121 func (api *PubSubAPI) checkNode() error {
121 - if !api.node.OnlineMode() {
122 - return coreiface.ErrOffline
122 + if err := api.checkRouting(false); err != nil {
123 + return err
124 }
125
125 - if api.node.PubSub == nil {
126 + if api.pubSub == nil {
127 return errors.New("experimental pubsub feature not enabled. Run daemon with --enable-pubsub-experiment to use.")
128 }
129
core/coreapi/swarm.go
+20 -21
@@ -5,7 +5,6 @@ import (
5 "sort"
6 "time"
7
8 - core "github.com/ipfs/go-ipfs/core"
8 coreiface "github.com/ipfs/go-ipfs/core/coreapi/interface"
9
10 inet "gx/ipfs/QmPtFaR7BWHLAjSwLh9kXcyrgTzDpuhcWLkx8ioa9RMYnx/go-libp2p-net"
@@ -21,9 +20,9 @@ import (
20 type SwarmAPI CoreAPI
21
22 type connInfo struct {
24 - node *core.IpfsNode
25 - conn net.Conn
26 - dir net.Direction
23 + peerstore pstore.Peerstore
24 + conn net.Conn
25 + dir net.Direction
26
27 addr ma.Multiaddr
28 peer peer.ID
@@ -31,19 +30,19 @@ type connInfo struct {
30 }
31
32 func (api *SwarmAPI) Connect(ctx context.Context, pi pstore.PeerInfo) error {
34 - if api.node.PeerHost == nil {
33 + if api.peerHost == nil {
34 return coreiface.ErrOffline
35 }
36
38 - if swrm, ok := api.node.PeerHost.Network().(*swarm.Swarm); ok {
37 + if swrm, ok := api.peerHost.Network().(*swarm.Swarm); ok {
38 swrm.Backoff().Clear(pi.ID)
39 }
40
42 - return api.node.PeerHost.Connect(ctx, pi)
41 + return api.peerHost.Connect(ctx, pi)
42 }
43
44 func (api *SwarmAPI) Disconnect(ctx context.Context, addr ma.Multiaddr) error {
46 - if api.node.PeerHost == nil {
45 + if api.peerHost == nil {
46 return coreiface.ErrOffline
47 }
48
@@ -54,7 +53,7 @@ func (api *SwarmAPI) Disconnect(ctx context.Context, addr ma.Multiaddr) error {
53
54 taddr := ia.Transport()
55 id := ia.ID()
57 - net := api.node.PeerHost.Network()
56 + net := api.peerHost.Network()
57
58 if taddr == nil {
59 if net.Connectedness(id) != inet.Connected {
@@ -78,12 +77,12 @@ func (api *SwarmAPI) Disconnect(ctx context.Context, addr ma.Multiaddr) error {
77 }
78
79 func (api *SwarmAPI) KnownAddrs(context.Context) (map[peer.ID][]ma.Multiaddr, error) {
81 - if api.node.PeerHost == nil {
80 + if api.peerHost == nil {
81 return nil, coreiface.ErrOffline
82 }
83
84 addrs := make(map[peer.ID][]ma.Multiaddr)
86 - ps := api.node.PeerHost.Network().Peerstore()
85 + ps := api.peerHost.Network().Peerstore()
86 for _, p := range ps.Peers() {
87 for _, a := range ps.Addrs(p) {
88 addrs[p] = append(addrs[p], a)
@@ -97,27 +96,27 @@ func (api *SwarmAPI) KnownAddrs(context.Context) (map[peer.ID][]ma.Multiaddr, er
96 }
97
98 func (api *SwarmAPI) LocalAddrs(context.Context) ([]ma.Multiaddr, error) {
100 - if api.node.PeerHost == nil {
99 + if api.peerHost == nil {
100 return nil, coreiface.ErrOffline
101 }
102
104 - return api.node.PeerHost.Addrs(), nil
103 + return api.peerHost.Addrs(), nil
104 }
105
106 func (api *SwarmAPI) ListenAddrs(context.Context) ([]ma.Multiaddr, error) {
108 - if api.node.PeerHost == nil {
107 + if api.peerHost == nil {
108 return nil, coreiface.ErrOffline
109 }
110
112 - return api.node.PeerHost.Network().InterfaceListenAddresses()
111 + return api.peerHost.Network().InterfaceListenAddresses()
112 }
113
114 func (api *SwarmAPI) Peers(context.Context) ([]coreiface.ConnectionInfo, error) {
116 - if api.node.PeerHost == nil {
115 + if api.peerHost == nil {
116 return nil, coreiface.ErrOffline
117 }
118
120 - conns := api.node.PeerHost.Network().Conns()
119 + conns := api.peerHost.Network().Conns()
120
121 var out []coreiface.ConnectionInfo
122 for _, c := range conns {
@@ -125,9 +124,9 @@ func (api *SwarmAPI) Peers(context.Context) ([]coreiface.ConnectionInfo, error)
124 addr := c.RemoteMultiaddr()
125
126 ci := &connInfo{
128 - node: api.node,
129 - conn: c,
130 - dir: c.Stat().Direction,
127 + peerstore: api.peerstore,
128 + conn: c,
129 + dir: c.Stat().Direction,
130
131 addr: addr,
132 peer: pid,
@@ -160,7 +159,7 @@ func (ci *connInfo) Direction() net.Direction {
159 }
160
161 func (ci *connInfo) Latency() (time.Duration, error) {
163 - return ci.node.Peerstore.LatencyEWMA(peer.ID(ci.ID())), nil
162 + return ci.peerstore.LatencyEWMA(peer.ID(ci.ID())), nil
163 }
164
165 func (ci *connInfo) Streams() ([]protocol.ID, error) {
core/coreapi/unixfs.go
+12 -11
@@ -34,9 +34,7 @@ func (api *UnixfsAPI) Add(ctx context.Context, files files.Node, opts ...options
34 return nil, err
35 }
36
37 - n := api.node
38 -
39 - cfg, err := n.Repo.Config()
37 + cfg, err := api.repo.Config()
38 if err != nil {
39 return nil, err
40 }
@@ -53,6 +51,13 @@ func (api *UnixfsAPI) Add(ctx context.Context, files files.Node, opts ...options
51 return nil, filestore.ErrFilestoreNotEnabled
52 }
53
54 + addblockstore := api.blockstore
55 + if !(settings.FsCache || settings.NoCopy) {
56 + addblockstore = bstore.NewGCBlockstore(api.baseBlocks, api.blockstore)
57 + }
58 + exch := api.exchange
59 + pinning := api.pinning
60 +
61 if settings.OnlyHash {
62 nilnode, err := core.NewNode(ctx, &core.BuildCfg{
63 //TODO: need this to be true or all files
@@ -62,15 +67,11 @@ func (api *UnixfsAPI) Add(ctx context.Context, files files.Node, opts ...options
67 if err != nil {
68 return nil, err
69 }
65 - n = nilnode
66 - }
67 -
68 - addblockstore := n.Blockstore
69 - if !(settings.FsCache || settings.NoCopy) {
70 - addblockstore = bstore.NewGCBlockstore(n.BaseBlocks, n.GCLocker)
70 + addblockstore = nilnode.Blockstore
71 + exch = nilnode.Exchange
72 + pinning = nilnode.Pinning
73 }
74
73 - exch := n.Exchange
75 if settings.Local {
76 exch = offline.Exchange(addblockstore)
77 }
@@ -78,7 +79,7 @@ func (api *UnixfsAPI) Add(ctx context.Context, files files.Node, opts ...options
79 bserv := blockservice.New(addblockstore, exch) // hash security 001
80 dserv := dag.NewDAGService(bserv)
81
81 - fileAdder, err := coreunix.NewAdder(ctx, n.Pinning, n.Blockstore, dserv)
82 + fileAdder, err := coreunix.NewAdder(ctx, pinning, addblockstore, dserv)
83 if err != nil {
84 return nil, err
85 }
core/corerepo/pinning.go
+7 -7
@@ -16,14 +16,14 @@ package corerepo
16 import (
17 "context"
18 "fmt"
19 + "github.com/ipfs/go-ipfs/pin"
20
20 - "github.com/ipfs/go-ipfs/core"
21 "github.com/ipfs/go-ipfs/core/coreapi/interface"
22
23 "gx/ipfs/QmR8BauakNcBa3RbE4nbQu76PDiJgoQgz8AJdhJuiU4TAw/go-cid"
24 )
25
26 -func Pin(n *core.IpfsNode, api iface.CoreAPI, ctx context.Context, paths []string, recursive bool) ([]cid.Cid, error) {
26 +func Pin(pinning pin.Pinner, api iface.CoreAPI, ctx context.Context, paths []string, recursive bool) ([]cid.Cid, error) {
27 out := make([]cid.Cid, len(paths))
28
29 for i, fpath := range paths {
@@ -36,14 +36,14 @@ func Pin(n *core.IpfsNode, api iface.CoreAPI, ctx context.Context, paths []strin
36 if err != nil {
37 return nil, fmt.Errorf("pin: %s", err)
38 }
39 - err = n.Pinning.Pin(ctx, dagnode, recursive)
39 + err = pinning.Pin(ctx, dagnode, recursive)
40 if err != nil {
41 return nil, fmt.Errorf("pin: %s", err)
42 }
43 out[i] = dagnode.Cid()
44 }
45
46 - err := n.Pinning.Flush()
46 + err := pinning.Flush()
47 if err != nil {
48 return nil, err
49 }
@@ -51,7 +51,7 @@ func Pin(n *core.IpfsNode, api iface.CoreAPI, ctx context.Context, paths []strin
51 return out, nil
52 }
53
54 -func Unpin(n *core.IpfsNode, api iface.CoreAPI, ctx context.Context, paths []string, recursive bool) ([]cid.Cid, error) {
54 +func Unpin(pinning pin.Pinner, api iface.CoreAPI, ctx context.Context, paths []string, recursive bool) ([]cid.Cid, error) {
55 unpinned := make([]cid.Cid, len(paths))
56
57 for i, p := range paths {
@@ -65,14 +65,14 @@ func Unpin(n *core.IpfsNode, api iface.CoreAPI, ctx context.Context, paths []str
65 return nil, err
66 }
67
68 - err = n.Pinning.Unpin(ctx, k.Cid(), recursive)
68 + err = pinning.Unpin(ctx, k.Cid(), recursive)
69 if err != nil {
70 return nil, err
71 }
72 unpinned[i] = k.Cid()
73 }
74
75 - err := n.Pinning.Flush()
75 + err := pinning.Flush()
76 if err != nil {
77 return nil, err
78 }
core/coreunix/add.go
+6 -6
@@ -48,13 +48,13 @@ type Object struct {
48 }
49
50 // NewAdder Returns a new Adder used for a file add operation.
51 -func NewAdder(ctx context.Context, p pin.Pinner, bs bstore.GCBlockstore, ds ipld.DAGService) (*Adder, error) {
51 +func NewAdder(ctx context.Context, p pin.Pinner, bs bstore.GCLocker, ds ipld.DAGService) (*Adder, error) {
52 bufferedDS := ipld.NewBufferedDAG(ctx, ds)
53
54 return &Adder{
55 ctx: ctx,
56 pinning: p,
57 - blockstore: bs,
57 + gcLocker: bs,
58 dagService: ds,
59 bufferedDS: bufferedDS,
60 Progress: false,
@@ -70,7 +70,7 @@ func NewAdder(ctx context.Context, p pin.Pinner, bs bstore.GCBlockstore, ds ipld
70 type Adder struct {
71 ctx context.Context
72 pinning pin.Pinner
73 - blockstore bstore.GCBlockstore
73 + gcLocker bstore.GCLocker
74 dagService ipld.DAGService
75 bufferedDS *ipld.BufferedDAG
76 Out chan<- interface{}
@@ -401,7 +401,7 @@ func (adder *Adder) addNode(node ipld.Node, path string) error {
401 // AddAllAndPin adds the given request's files and pin them.
402 func (adder *Adder) AddAllAndPin(file files.Node) (ipld.Node, error) {
403 if adder.Pin {
404 - adder.unlocker = adder.blockstore.PinLock()
404 + adder.unlocker = adder.gcLocker.PinLock()
405 }
406 defer func() {
407 if adder.unlocker != nil {
@@ -556,14 +556,14 @@ func (adder *Adder) addDir(path string, dir files.Directory) error {
556 }
557
558 func (adder *Adder) maybePauseForGC() error {
559 - if adder.unlocker != nil && adder.blockstore.GCRequested() {
559 + if adder.unlocker != nil && adder.gcLocker.GCRequested() {
560 err := adder.PinRoot()
561 if err != nil {
562 return err
563 }
564
565 adder.unlocker.Unlock()
566 - adder.unlocker = adder.blockstore.PinLock()
566 + adder.unlocker = adder.gcLocker.PinLock()
567 }
568 return nil
569 }