fix(gcr/s) proto marshaling bugs
Brian Tiger Chow committed
Jan 30, 2015 at 08:22 UTC
0721a589341fdc575ed5ec42fb20b9a87192aefa
1 file changed
+103
-61
routing/grandcentral/server.go
+103
-61
@@ -1,6 +1,8 @@
1
package grandcentral
2
3
import (
4
+ "fmt"
5
+
6
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
7
proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
8
datastore "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
@@ -48,36 +50,17 @@ func (s *Server) handleMessage(
50
switch req.GetType() {
51
52
case dhtpb.Message_GET_VALUE:
51
- dskey := util.Key(req.GetKey()).DsKey()
52
- val, err := s.datastore.Get(dskey)
53
+ rawRecord, err := getRoutingRecord(s.datastore, util.Key(req.GetKey()))
54
if err != nil {
54
- log.Debug(errors.Wrap(err))
55
- return "", nil
56
- }
57
- rawRecord, ok := val.([]byte)
58
- if !ok {
59
- log.Debugf("datastore had non byte-slice value for %v", dskey)
60
- return "", nil
61
- }
62
- if err := proto.Unmarshal(rawRecord, response.Record); err != nil {
63
- log.Debug("failed to unmarshal dht record from datastore")
55
return "", nil
56
}
57
+ response.Record = rawRecord
58
// TODO before merging: if we know any providers for the requested value, return those.
59
return p, response
60
61
case dhtpb.Message_PUT_VALUE:
62
// TODO before merging: verifyRecord(req.GetRecord())
71
- data, err := proto.Marshal(req.GetRecord())
72
- if err != nil {
73
- log.Debug(err)
74
- return "", nil
75
- }
76
- dskey := util.Key(req.GetKey()).DsKey()
77
- if err := s.datastore.Put(dskey, data); err != nil {
78
- log.Debug(err)
79
- return "", nil
80
- }
63
+ putRoutingRecord(s.datastore, util.Key(req.GetKey()), req.GetRecord())
64
return p, req // TODO before merging: verify that we should return record
65
66
case dhtpb.Message_FIND_NODE:
@@ -92,51 +75,19 @@ func (s *Server) handleMessage(
75
return p.ID, response
76
77
case dhtpb.Message_ADD_PROVIDER:
95
- for _, provider := range req.GetProviderPeers() {
96
- providerID := peer.ID(provider.GetId())
97
- if providerID != p {
98
- log.Debugf("provider message came from third-party %s", p)
99
- continue
100
- }
101
- for _, maddr := range provider.Addresses() {
102
- // FIXME do we actually want to store to peerstore
103
- s.peerstore.AddAddr(p, maddr, peer.TempAddrTTL)
104
- }
105
- }
106
- var providers []dhtpb.Message_Peer
107
- pkey := datastore.KeyWithNamespaces([]string{"routing", "providers", req.GetKey()})
108
- if v, err := s.datastore.Get(pkey); err == nil {
109
- if protopeers, ok := v.([]dhtpb.Message_Peer); ok {
110
- providers = append(providers, protopeers...)
111
- }
112
- }
113
- if err := s.datastore.Put(pkey, providers); err != nil {
114
- log.Debug(err)
78
+ storeProvidersToPeerstore(s.peerstore, p, req.GetProviderPeers())
79
+
80
+ if err := putRoutingProviders(s.datastore, util.Key(req.GetKey()), req.GetProviderPeers()); err != nil {
81
return "", nil
82
}
83
return "", nil
84
85
case dhtpb.Message_GET_PROVIDERS:
120
- dskey := util.Key(req.GetKey()).DsKey()
121
- exists, err := s.datastore.Has(dskey)
122
- if err == nil && exists {
123
- pri := []dhtpb.PeerRoutingInfo{
124
- dhtpb.PeerRoutingInfo{
125
- // Connectedness: TODO how is connectedness defined for the local node
126
- PeerInfo: peer.PeerInfo{ID: s.local},
127
- },
128
- }
129
- response.ProviderPeers = append(response.ProviderPeers, dhtpb.PeerRoutingInfosToPBPeers(pri)...)
130
- }
131
- // FIXME(btc) is this how we want to persist this data?
132
- pkey := datastore.KeyWithNamespaces([]string{"routing", "providers", req.GetKey()})
133
- if v, err := s.datastore.Get(pkey); err == nil {
134
- if protopeers, ok := v.([]dhtpb.Message_Peer); ok {
135
- for _, p := range protopeers {
136
- response.ProviderPeers = append(response.ProviderPeers, &p)
137
- }
138
- }
86
+ providers, err := getRoutingProviders(s.local, s.datastore, util.Key(req.GetKey()))
87
+ if err != nil {
88
+ return "", nil
89
}
90
+ response.ProviderPeers = providers
91
return p, response
92
93
case dhtpb.Message_PING:
@@ -148,3 +99,94 @@ func (s *Server) handleMessage(
99
100
var _ proxy.RequestHandler = &Server{}
101
var _ proxy.Proxy = &Server{}
102
+
103
+func getRoutingRecord(ds datastore.Datastore, k util.Key) (*dhtpb.Record, error) {
104
+ dskey := k.DsKey()
105
+ val, err := ds.Get(dskey)
106
+ if err != nil {
107
+ return nil, errors.Wrap(err)
108
+ }
109
+ recordBytes, ok := val.([]byte)
110
+ if !ok {
111
+ return nil, fmt.Errorf("datastore had non byte-slice value for %v", dskey)
112
+ }
113
+ var record dhtpb.Record
114
+ if err := proto.Unmarshal(recordBytes, &record); err != nil {
115
+ return nil, errors.New("failed to unmarshal dht record from datastore")
116
+ }
117
+ return &record, nil
118
+}
119
+
120
+func putRoutingRecord(ds datastore.Datastore, k util.Key, value *dhtpb.Record) error {
121
+ data, err := proto.Marshal(value)
122
+ if err != nil {
123
+ return err
124
+ }
125
+ dskey := k.DsKey()
126
+ // TODO namespace
127
+ if err := ds.Put(dskey, data); err != nil {
128
+ return err
129
+ }
130
+ return nil
131
+}
132
+
133
+func putRoutingProviders(ds datastore.Datastore, k util.Key, providers []*dhtpb.Message_Peer) error {
134
+ pkey := datastore.KeyWithNamespaces([]string{"routing", "providers", k.String()})
135
+ if v, err := ds.Get(pkey); err == nil {
136
+ if msg, ok := v.([]byte); ok {
137
+ var protomsg dhtpb.Message
138
+ if err := proto.Unmarshal(msg, &protomsg); err != nil {
139
+ log.Error("failed to unmarshal routing provider record. programmer error")
140
+ } else {
141
+ providers = append(providers, protomsg.ProviderPeers...)
142
+ }
143
+ }
144
+ }
145
+ var protomsg dhtpb.Message
146
+ protomsg.ProviderPeers = providers
147
+ data, err := proto.Marshal(&protomsg)
148
+ if err != nil {
149
+ return err
150
+ }
151
+ return ds.Put(pkey, data)
152
+}
153
+
154
+func storeProvidersToPeerstore(ps peer.Peerstore, p peer.ID, providers []*dhtpb.Message_Peer) {
155
+ for _, provider := range providers {
156
+ providerID := peer.ID(provider.GetId())
157
+ if providerID != p {
158
+ log.Errorf("provider message came from third-party %s", p)
159
+ continue
160
+ }
161
+ for _, maddr := range provider.Addresses() {
162
+ // as a router, we want to store addresses for peers who have provided
163
+ ps.AddAddr(p, maddr, peer.AddressTTL)
164
+ }
165
+ }
166
+}
167
+
168
+func getRoutingProviders(local peer.ID, ds datastore.Datastore, k util.Key) ([]*dhtpb.Message_Peer, error) {
169
+ var providers []*dhtpb.Message_Peer
170
+ exists, err := ds.Has(k.DsKey()) // TODO store values in a local datastore?
171
+ if err == nil && exists {
172
+ pri := []dhtpb.PeerRoutingInfo{
173
+ dhtpb.PeerRoutingInfo{
174
+ // Connectedness: TODO how is connectedness defined for the local node
175
+ PeerInfo: peer.PeerInfo{ID: local},
176
+ },
177
+ }
178
+ providers = append(providers, dhtpb.PeerRoutingInfosToPBPeers(pri)...)
179
+ }
180
+
181
+ pkey := datastore.KeyWithNamespaces([]string{"routing", "providers", k.String()}) // TODO key fmt
182
+ if v, err := ds.Get(pkey); err == nil {
183
+ if data, ok := v.([]byte); ok {
184
+ var msg dhtpb.Message
185
+ if err := proto.Unmarshal(data, &msg); err != nil {
186
+ return nil, err
187
+ }
188
+ providers = append(providers, msg.GetProviderPeers()...)
189
+ }
190
+ }
191
+ return providers, nil
192
+}