feat(routing.grandcentral): skeleton
fixes breakages: - peer.Peer -> peer.ID - peer.X -> peer.PeerInfo - netmsg -> p2p streams
Brian Tiger Chow committed
Jan 26, 2015 at 20:48 UTC
577baaf6218df6f350588101515a80b8ba4e62c2
4 files changed
+387
routing/grandcentral/client.go
new
+121
@@ -0,0 +1,121 @@
1
+package grandcentral
2
+
3
+import (
4
+ "bytes"
5
+ "time"
6
+
7
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8
+ proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
9
+ inet "github.com/jbenet/go-ipfs/p2p/net"
10
+ peer "github.com/jbenet/go-ipfs/p2p/peer"
11
+ routing "github.com/jbenet/go-ipfs/routing"
12
+ pb "github.com/jbenet/go-ipfs/routing/dht/pb"
13
+ proxy "github.com/jbenet/go-ipfs/routing/grandcentral/proxy"
14
+ eventlog "github.com/jbenet/go-ipfs/thirdparty/eventlog"
15
+ u "github.com/jbenet/go-ipfs/util"
16
+ errors "github.com/jbenet/go-ipfs/util/debugerror"
17
+)
18
+
19
+var log = eventlog.Logger("grandcentral")
20
+
21
+var ErrTODO = errors.New("TODO")
22
+
23
+type Client struct {
24
+ peerstore peer.Peerstore
25
+ proxy proxy.Proxy
26
+ dialer inet.Network
27
+ local peer.ID
28
+}
29
+
30
+// TODO take in datastore/cache
31
+func NewClient(d inet.Network, px proxy.Proxy, ps peer.Peerstore, local peer.ID) (*Client, error) {
32
+ return &Client{
33
+ dialer: d,
34
+ proxy: px,
35
+ local: local,
36
+ peerstore: ps,
37
+ }, nil
38
+}
39
+
40
+func (c *Client) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan peer.PeerInfo {
41
+ ch := make(chan peer.PeerInfo)
42
+ go func() {
43
+ defer close(ch)
44
+ request := pb.NewMessage(pb.Message_GET_PROVIDERS, string(k), 0)
45
+ response, err := c.proxy.SendRequest(ctx, request)
46
+ if err != nil {
47
+ log.Error(errors.Wrap(err))
48
+ return
49
+ }
50
+ for _, p := range pb.PBPeersToPeerInfos(response.GetProviderPeers()) {
51
+ select {
52
+ case <-ctx.Done():
53
+ log.Error(errors.Wrap(ctx.Err()))
54
+ return
55
+ case ch <- p:
56
+ }
57
+ }
58
+ }()
59
+ return ch
60
+}
61
+
62
+func (c *Client) PutValue(ctx context.Context, k u.Key, v []byte) error {
63
+ r, err := makeRecord(c.peerstore, c.local, k, v)
64
+ if err != nil {
65
+ return err
66
+ }
67
+ pmes := pb.NewMessage(pb.Message_PUT_VALUE, string(k), 0)
68
+ pmes.Record = r
69
+ return c.proxy.SendMessage(ctx, pmes) // wrap to hide the remote
70
+}
71
+
72
+func (c *Client) GetValue(ctx context.Context, k u.Key) ([]byte, error) {
73
+ msg := pb.NewMessage(pb.Message_GET_VALUE, string(k), 0)
74
+ response, err := c.proxy.SendRequest(ctx, msg) // TODO wrap to hide the remote
75
+ if err != nil {
76
+ return nil, errors.Wrap(err)
77
+ }
78
+ return response.Record.GetValue(), nil
79
+}
80
+
81
+func (c *Client) Provide(ctx context.Context, k u.Key) error {
82
+ msg := pb.NewMessage(pb.Message_ADD_PROVIDER, string(k), 0)
83
+ // TODO wrap this to hide the dialer and the local/remote peers
84
+ msg.ProviderPeers = pb.PeerInfosToPBPeers(c.dialer, []peer.PeerInfo{peer.PeerInfo{ID: c.local}}) // FIXME how is connectedness defined for the local node
85
+ return c.proxy.SendMessage(ctx, msg) // TODO wrap to hide remote
86
+}
87
+
88
+func (c *Client) FindPeer(ctx context.Context, id peer.ID) (peer.PeerInfo, error) {
89
+ request := pb.NewMessage(pb.Message_FIND_NODE, string(id), 0)
90
+ response, err := c.proxy.SendRequest(ctx, request) // hide remote
91
+ if err != nil {
92
+ return peer.PeerInfo{}, errors.Wrap(err)
93
+ }
94
+ for _, p := range pb.PBPeersToPeerInfos(response.GetCloserPeers()) {
95
+ if p.ID == id {
96
+ return p, nil
97
+ }
98
+ }
99
+ return peer.PeerInfo{}, errors.New("could not find peer")
100
+}
101
+
102
+// creates and signs a record for the given key/value pair
103
+func makeRecord(ps peer.Peerstore, p peer.ID, k u.Key, v []byte) (*pb.Record, error) {
104
+ blob := bytes.Join([][]byte{[]byte(k), v, []byte(p)}, []byte{})
105
+ sig, err := ps.PrivKey(p).Sign(blob)
106
+ if err != nil {
107
+ return nil, err
108
+ }
109
+ return &pb.Record{
110
+ Key: proto.String(string(k)),
111
+ Value: v,
112
+ Author: proto.String(string(p)),
113
+ Signature: sig,
114
+ }, nil
115
+}
116
+
117
+func (c *Client) Ping(ctx context.Context, id peer.ID) (time.Duration, error) {
118
+ return time.Nanosecond, errors.New("grandcentral routing does not support the ping method")
119
+}
120
+
121
+var _ routing.IpfsRouting = &Client{}
routing/grandcentral/proxy/loopback.go
new
+53
@@ -0,0 +1,53 @@
1
+package proxy
2
+
3
+import (
4
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
+ inet "github.com/jbenet/go-ipfs/p2p/net"
6
+ peer "github.com/jbenet/go-ipfs/p2p/peer"
7
+ ggio "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/io"
8
+ dhtpb "github.com/jbenet/go-ipfs/routing/dht/pb"
9
+ errors "github.com/jbenet/go-ipfs/util/debugerror"
10
+)
11
+
12
+// RequestHandler handles routing requests locally
13
+type RequestHandler interface {
14
+ HandleRequest(ctx context.Context, p peer.ID, m *dhtpb.Message) *dhtpb.Message
15
+}
16
+
17
+// Loopback forwards requests to a local handler
18
+type Loopback struct {
19
+ Handler RequestHandler
20
+ Local peer.ID
21
+}
22
+
23
+// SendMessage intercepts local requests, forwarding them to a local handler
24
+func (lb *Loopback) SendMessage(ctx context.Context, m *dhtpb.Message) error {
25
+ response := lb.Handler.HandleRequest(ctx, lb.Local, m)
26
+ if response != nil {
27
+ log.Warning("loopback handler returned unexpected message")
28
+ }
29
+ return nil
30
+}
31
+
32
+// SendRequest intercepts local requests, forwarding them to a local handler
33
+func (lb *Loopback) SendRequest(ctx context.Context, m *dhtpb.Message) (*dhtpb.Message, error) {
34
+ return lb.Handler.HandleRequest(ctx, lb.Local, m), nil
35
+}
36
+
37
+func (lb *Loopback) handleNewStream(s inet.Stream) {
38
+ defer s.Close()
39
+ pbr := ggio.NewDelimitedReader(s, inet.MessageSizeMax)
40
+ var incoming dhtpb.Message
41
+ if err := pbr.ReadMsg(&incoming); err != nil {
42
+ log.Error(errors.Wrap(err))
43
+ return
44
+ }
45
+ ctx := context.TODO()
46
+ outgoing := lb.Handler.HandleRequest(ctx, s.Conn().RemotePeer(), &incoming)
47
+
48
+ pbw := ggio.NewDelimitedWriter(s)
49
+
50
+ if err := pbw.WriteMsg(outgoing); err != nil {
51
+ return // TODO logerr
52
+ }
53
+}
routing/grandcentral/proxy/standard.go
new
+73
@@ -0,0 +1,73 @@
1
+package proxy
2
+
3
+import (
4
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
+ host "github.com/jbenet/go-ipfs/p2p/host"
6
+ ggio "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/io"
7
+ inet "github.com/jbenet/go-ipfs/p2p/net"
8
+ peer "github.com/jbenet/go-ipfs/p2p/peer"
9
+ dhtpb "github.com/jbenet/go-ipfs/routing/dht/pb"
10
+ errors "github.com/jbenet/go-ipfs/util/debugerror"
11
+ eventlog "github.com/jbenet/go-ipfs/thirdparty/eventlog"
12
+)
13
+
14
+var log = eventlog.Logger("proxy")
15
+
16
+type Proxy interface {
17
+ SendMessage(ctx context.Context, m *dhtpb.Message) error
18
+ SendRequest(ctx context.Context, m *dhtpb.Message) (*dhtpb.Message, error)
19
+}
20
+
21
+type standard struct {
22
+ Host host.Host
23
+ Remote peer.ID
24
+}
25
+
26
+func Standard(h host.Host, remote peer.ID) Proxy {
27
+ return &standard{h, remote}
28
+}
29
+
30
+const ProtocolGCR = "/ipfs/grandcentral"
31
+
32
+func (px *standard) SendMessage(ctx context.Context, m *dhtpb.Message) error {
33
+ if err := px.Host.Connect(ctx, peer.PeerInfo{ID: px.Remote}); err != nil {
34
+ return err
35
+ }
36
+ s, err := px.Host.NewStream(ProtocolGCR, px.Remote)
37
+ if err != nil {
38
+ return err
39
+ }
40
+ defer s.Close()
41
+ pbw := ggio.NewDelimitedWriter(s)
42
+ if err := pbw.WriteMsg(m); err != nil {
43
+ return errors.Wrap(err)
44
+ }
45
+ return nil
46
+}
47
+
48
+func (px *standard) SendRequest(ctx context.Context, m *dhtpb.Message) (*dhtpb.Message, error) {
49
+ if err := px.Host.Connect(ctx, peer.PeerInfo{ID: px.Remote}); err != nil {
50
+ return nil, err
51
+ }
52
+ s, err := px.Host.NewStream(ProtocolGCR, px.Remote)
53
+ if err != nil {
54
+ return nil, err
55
+ }
56
+ defer s.Close()
57
+ r := ggio.NewDelimitedReader(s, inet.MessageSizeMax)
58
+ w := ggio.NewDelimitedWriter(s)
59
+ if err := w.WriteMsg(m); err != nil {
60
+ return nil, err
61
+ }
62
+
63
+ var reply dhtpb.Message
64
+ if err := r.ReadMsg(&reply); err != nil {
65
+ return nil, err
66
+ }
67
+ // need ctx expiration?
68
+ if &reply == nil {
69
+ return nil, errors.New("no response to request")
70
+ }
71
+ return &reply, nil
72
+}
73
+
routing/grandcentral/server.go
new
+140
@@ -0,0 +1,140 @@
1
+package grandcentral
2
+
3
+import (
4
+ context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
+ proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
6
+ datastore "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
7
+ inet "github.com/jbenet/go-ipfs/p2p/net"
8
+ peer "github.com/jbenet/go-ipfs/p2p/peer"
9
+ dhtpb "github.com/jbenet/go-ipfs/routing/dht/pb"
10
+ proxy "github.com/jbenet/go-ipfs/routing/grandcentral/proxy"
11
+ util "github.com/jbenet/go-ipfs/util"
12
+ errors "github.com/jbenet/go-ipfs/util/debugerror"
13
+)
14
+
15
+// Server handles routing queries using a database backend
16
+type Server struct {
17
+ local peer.ID
18
+ datastore datastore.ThreadSafeDatastore
19
+ dialer inet.Network
20
+ peerstore peer.Peerstore
21
+ *proxy.Loopback // so server can be injected into client
22
+}
23
+
24
+// NewServer creates a new GrandCentral routing Server
25
+func NewServer(ds datastore.ThreadSafeDatastore, d inet.Network, ps peer.Peerstore, local peer.ID) (*Server, error) {
26
+ s := &Server{local, ds, d, ps, nil}
27
+ s.Loopback = &proxy.Loopback{
28
+ Handler: s,
29
+ Local: local,
30
+ }
31
+ return s, nil
32
+}
33
+
34
+// HandleLocalRequest implements the proxy.RequestHandler interface. This is
35
+// where requests are received from the outside world.
36
+func (s *Server) HandleRequest(ctx context.Context, p peer.ID, req *dhtpb.Message) *dhtpb.Message {
37
+ _, response := s.handleMessage(ctx, p, req) // ignore response peer. it's local.
38
+ return response
39
+}
40
+
41
+// TODO extract backend. backend can be implemented with whatever database we desire
42
+func (s *Server) handleMessage(
43
+ ctx context.Context, p peer.ID, req *dhtpb.Message) (peer.ID, *dhtpb.Message) {
44
+
45
+ // FIXME threw everything into this switch statement to get things going.
46
+ // Once each operation is well-defined, extract pluggable backend so any
47
+ // database may be used.
48
+
49
+ var response = dhtpb.NewMessage(req.GetType(), req.GetKey(), req.GetClusterLevel())
50
+ switch req.GetType() {
51
+
52
+ case dhtpb.Message_GET_VALUE:
53
+ dskey := util.Key(req.GetKey()).DsKey()
54
+ val, err := s.datastore.Get(dskey)
55
+ if err != nil {
56
+ log.Error(errors.Wrap(err))
57
+ return "", nil
58
+ }
59
+ rawRecord, ok := val.([]byte)
60
+ if !ok {
61
+ log.Errorf("datastore had non byte-slice value for %v", dskey)
62
+ return "", nil
63
+ }
64
+ if err := proto.Unmarshal(rawRecord, response.Record); err != nil {
65
+ log.Error("failed to unmarshal dht record from datastore")
66
+ return "", nil
67
+ }
68
+ // TODO before merging: if we know any providers for the requested value, return those.
69
+ return p, response
70
+
71
+ case dhtpb.Message_PUT_VALUE:
72
+ // TODO before merging: verifyRecord(req.GetRecord())
73
+ data, err := proto.Marshal(req.GetRecord())
74
+ if err != nil {
75
+ log.Error(err)
76
+ return "", nil
77
+ }
78
+ dskey := util.Key(req.GetKey()).DsKey()
79
+ if err := s.datastore.Put(dskey, data); err != nil {
80
+ log.Error(err)
81
+ return "", nil
82
+ }
83
+ return p, req // TODO before merging: verify that we should return record
84
+
85
+ case dhtpb.Message_FIND_NODE:
86
+ p := s.peerstore.PeerInfo(peer.ID(req.GetKey()))
87
+ response.CloserPeers = dhtpb.PeerInfosToPBPeers(s.dialer, []peer.PeerInfo{p})
88
+ return p.ID, response
89
+
90
+ case dhtpb.Message_ADD_PROVIDER:
91
+ for _, provider := range req.GetProviderPeers() {
92
+ providerID := peer.ID(provider.GetId())
93
+ if providerID != p {
94
+ log.Errorf("provider message came from third-party %s", p)
95
+ continue
96
+ }
97
+ for _, maddr := range provider.Addresses() {
98
+ // FIXME do we actually want to store to peerstore
99
+ s.peerstore.AddAddress(p, maddr)
100
+ }
101
+ }
102
+ var providers []dhtpb.Message_Peer
103
+ pkey := datastore.KeyWithNamespaces([]string{"routing", "providers", req.GetKey()})
104
+ if v, err := s.datastore.Get(pkey); err == nil {
105
+ if protopeers, ok := v.([]dhtpb.Message_Peer); ok {
106
+ providers = append(providers, protopeers...)
107
+ }
108
+ }
109
+ if err := s.datastore.Put(pkey, providers); err != nil {
110
+ log.Error(err)
111
+ return "", nil
112
+ }
113
+ return "", nil
114
+
115
+ case dhtpb.Message_GET_PROVIDERS:
116
+ dskey := util.Key(req.GetKey()).DsKey()
117
+ exists, err := s.datastore.Has(dskey)
118
+ if err == nil && exists {
119
+ response.ProviderPeers = append(response.ProviderPeers, dhtpb.PeerInfosToPBPeers(s.dialer, []peer.PeerInfo{peer.PeerInfo{ID: s.local}})...)
120
+ }
121
+ // FIXME(btc) is this how we want to persist this data?
122
+ pkey := datastore.KeyWithNamespaces([]string{"routing", "providers", req.GetKey()})
123
+ if v, err := s.datastore.Get(pkey); err == nil {
124
+ if protopeers, ok := v.([]dhtpb.Message_Peer); ok {
125
+ for _, p := range protopeers {
126
+ response.ProviderPeers = append(response.ProviderPeers, &p)
127
+ }
128
+ }
129
+ }
130
+ return p, response
131
+
132
+ case dhtpb.Message_PING:
133
+ return p, req
134
+ default:
135
+ }
136
+ return "", nil
137
+}
138
+
139
+var _ proxy.RequestHandler = &Server{}
140
+var _ proxy.Proxy = &Server{}