working on upper level dht implementations, protbuf, etc
Jeromy committed
Jul 28, 2014 at 22:14 UTC
8bc80124a47849d3ad818782c939e27a71e27340
4 files changed
+158
-2
routing/dht/dht.go
+12
@@ -1,5 +1,9 @@
1
package dht
2
3
+import (
4
+ swarm "github.com/jbenet/go-ipfs/swarm"
5
+)
6
+
7
// TODO. SEE https://github.com/jbenet/node-ipfs/blob/master/submodules/ipfs-dht/index.js
8
9
@@ -7,4 +11,12 @@ package dht
11
// It is used to implement the base IpfsRouting module.
12
type IpfsDHT struct {
13
routes RoutingTable
14
+
15
+ network *swarm.Swarm
16
+}
17
+
18
+func (dht *IpfsDHT) handleMessages() {
19
+ for mes := range dht.network.Chan.Incoming {
20
+
21
+ }
22
}
routing/dht/messages.pb.go
new
+96
@@ -0,0 +1,96 @@
1
+// Code generated by protoc-gen-go.
2
+// source: messages.proto
3
+// DO NOT EDIT!
4
+
5
+/*
6
+Package dht is a generated protocol buffer package.
7
+
8
+It is generated from these files:
9
+ messages.proto
10
+
11
+It has these top-level messages:
12
+ DHTMessage
13
+*/
14
+package dht
15
+
16
+import proto "code.google.com/p/goprotobuf/proto"
17
+import math "math"
18
+
19
+// Reference imports to suppress errors if they are not otherwise used.
20
+var _ = proto.Marshal
21
+var _ = math.Inf
22
+
23
+type DHTMessage_MessageType int32
24
+
25
+const (
26
+ DHTMessage_PUT_VALUE DHTMessage_MessageType = 0
27
+ DHTMessage_GET_VALUE DHTMessage_MessageType = 1
28
+ DHTMessage_PING DHTMessage_MessageType = 2
29
+ DHTMessage_FIND_NODE DHTMessage_MessageType = 3
30
+)
31
+
32
+var DHTMessage_MessageType_name = map[int32]string{
33
+ 0: "PUT_VALUE",
34
+ 1: "GET_VALUE",
35
+ 2: "PING",
36
+ 3: "FIND_NODE",
37
+}
38
+var DHTMessage_MessageType_value = map[string]int32{
39
+ "PUT_VALUE": 0,
40
+ "GET_VALUE": 1,
41
+ "PING": 2,
42
+ "FIND_NODE": 3,
43
+}
44
+
45
+func (x DHTMessage_MessageType) Enum() *DHTMessage_MessageType {
46
+ p := new(DHTMessage_MessageType)
47
+ *p = x
48
+ return p
49
+}
50
+func (x DHTMessage_MessageType) String() string {
51
+ return proto.EnumName(DHTMessage_MessageType_name, int32(x))
52
+}
53
+func (x *DHTMessage_MessageType) UnmarshalJSON(data []byte) error {
54
+ value, err := proto.UnmarshalJSONEnum(DHTMessage_MessageType_value, data, "DHTMessage_MessageType")
55
+ if err != nil {
56
+ return err
57
+ }
58
+ *x = DHTMessage_MessageType(value)
59
+ return nil
60
+}
61
+
62
+type DHTMessage struct {
63
+ Type *DHTMessage_MessageType `protobuf:"varint,1,req,name=type,enum=dht.DHTMessage_MessageType" json:"type,omitempty"`
64
+ Key *string `protobuf:"bytes,2,opt,name=key" json:"key,omitempty"`
65
+ Value []byte `protobuf:"bytes,3,opt,name=value" json:"value,omitempty"`
66
+ XXX_unrecognized []byte `json:"-"`
67
+}
68
+
69
+func (m *DHTMessage) Reset() { *m = DHTMessage{} }
70
+func (m *DHTMessage) String() string { return proto.CompactTextString(m) }
71
+func (*DHTMessage) ProtoMessage() {}
72
+
73
+func (m *DHTMessage) GetType() DHTMessage_MessageType {
74
+ if m != nil && m.Type != nil {
75
+ return *m.Type
76
+ }
77
+ return DHTMessage_PUT_VALUE
78
+}
79
+
80
+func (m *DHTMessage) GetKey() string {
81
+ if m != nil && m.Key != nil {
82
+ return *m.Key
83
+ }
84
+ return ""
85
+}
86
+
87
+func (m *DHTMessage) GetValue() []byte {
88
+ if m != nil {
89
+ return m.Value
90
+ }
91
+ return nil
92
+}
93
+
94
+func init() {
95
+ proto.RegisterEnum("dht.DHTMessage_MessageType", DHTMessage_MessageType_name, DHTMessage_MessageType_value)
96
+}
routing/dht/messages.proto
new
+16
@@ -0,0 +1,16 @@
1
+package dht;
2
+
3
+//run `protoc --go_out=. *.proto` to generate
4
+
5
+message DHTMessage {
6
+ enum MessageType {
7
+ PUT_VALUE = 0;
8
+ GET_VALUE = 1;
9
+ PING = 2;
10
+ FIND_NODE = 3;
11
+ }
12
+
13
+ required MessageType type = 1;
14
+ optional string key = 2;
15
+ optional bytes value = 3;
16
+}
routing/dht/routing.go
+34
-2
@@ -4,6 +4,7 @@ import (
4
"time"
5
peer "github.com/jbenet/go-ipfs/peer"
6
u "github.com/jbenet/go-ipfs/util"
7
+ swarm "github.com/jbenet/go-ipfs/swarm"
8
)
9
10
@@ -13,12 +14,43 @@ import (
14
15
// PutValue adds value corresponding to given Key.
16
func (s *IpfsDHT) PutValue(key u.Key, value []byte) (error) {
16
- return u.ErrNotImplemented
17
+ var p *peer.Peer
18
+ p = s.routes.NearestNode(key)
19
+
20
+ pmes := new(PutValue)
21
+ pmes.Key = &key
22
+ pmes.Value = value
23
+
24
+ mes := new(swarm.Message)
25
+ mes.Data = []byte(pmes.String())
26
+ mes.Peer = p
27
+
28
+ s.network.Chan.Outgoing <- mes
29
+ return nil
30
}
31
32
// GetValue searches for the value corresponding to given Key.
33
func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
21
- return nil, u.ErrNotImplemented
34
+ var p *peer.Peer
35
+ p = s.routes.NearestNode(key)
36
+
37
+ // protobuf structure
38
+ pmes := new(GetValue)
39
+ pmes.Key = &key
40
+ pmes.Id = GenerateMessageID()
41
+
42
+ mes := new(swarm.Message)
43
+ mes.Data = []byte(pmes.String())
44
+ mes.Peer = p
45
+
46
+ response_chan := s.network.ListenFor(pmes.Id)
47
+
48
+ timeup := time.After(timeout)
49
+ select {
50
+ case <-timeup:
51
+ return nil, timeoutError
52
+ case resp := <-response_chan:
53
+ }
54
}
55
56