comment comment comment comment
Jeromy committed
Nov 3, 2014 at 03:02 UTC
88bf39118c5f4368f2888ef377c5541ba307e46e
7 files changed
+55
-8
blockservice/blockservice.go
+1
@@ -93,6 +93,7 @@ func (s *BlockService) GetBlock(ctx context.Context, k u.Key) (*blocks.Block, er
93
}
94
}
95
96
+// DeleteBlock deletes a block in the blockservice from the datastore
97
func (s *BlockService) DeleteBlock(k u.Key) error {
98
return s.Datastore.Delete(k.DsKey())
99
}
crypto/key.go
+9
@@ -28,6 +28,7 @@ const (
28
RSA = iota
29
)
30
31
+// Key represents a crypto key that can be compared to another key
32
type Key interface {
33
// Bytes returns a serialized, storeable representation of this key
34
Bytes() ([]byte, error)
@@ -39,6 +40,8 @@ type Key interface {
40
Equals(Key) bool
41
}
42
43
+// PrivKey represents a private key that can be used to generate a public key,
44
+// sign data, and decrypt data that was encrypted with a public key
45
type PrivKey interface {
46
Key
47
@@ -60,12 +63,14 @@ type PubKey interface {
63
// Verify that 'sig' is the signed hash of 'data'
64
Verify(data []byte, sig []byte) (bool, error)
65
66
+ // Encrypt data in a way that can be decrypted by a paired private key
67
Encrypt(data []byte) ([]byte, error)
68
}
69
70
// Given a public key, generates the shared key.
71
type GenSharedKey func([]byte) ([]byte, error)
72
73
+// Generates a keypair of the given type and bitsize
74
func GenerateKeyPair(typ, bits int) (PrivKey, PubKey, error) {
75
switch typ {
76
case RSA:
@@ -217,6 +222,8 @@ func KeyStretcher(cmp int, cipherType string, hashType string, secret []byte) ([
222
return myIV, theirIV, myCKey, theirCKey, myMKey, theirMKey
223
}
224
225
+// UnmarshalPublicKey converts a protobuf serialized public key into its
226
+// representative object
227
func UnmarshalPublicKey(data []byte) (PubKey, error) {
228
pmes := new(pb.PublicKey)
229
err := proto.Unmarshal(data, pmes)
@@ -232,6 +239,8 @@ func UnmarshalPublicKey(data []byte) (PubKey, error) {
239
}
240
}
241
242
+// UnmarshalPrivateKey converts a protobuf serialized private key into its
243
+// representative object
244
func UnmarshalPrivateKey(data []byte) (PrivKey, error) {
245
pmes := new(pb.PrivateKey)
246
err := proto.Unmarshal(data, pmes)
diagnostics/diag.go
+28
-7
@@ -24,6 +24,8 @@ var log = util.Logger("diagnostics")
24
25
const ResponseTimeout = time.Second * 10
26
27
+// Diagnostics is a net service that manages requesting and responding to diagnostic
28
+// requests
29
type Diagnostics struct {
30
network net.Network
31
sender net.Sender
@@ -34,6 +36,7 @@ type Diagnostics struct {
36
birth time.Time
37
}
38
39
+// NewDiagnostics instantiates a new diagnostics service running on the given network
40
func NewDiagnostics(self peer.Peer, inet net.Network, sender net.Sender) *Diagnostics {
41
return &Diagnostics{
42
network: inet,
@@ -50,15 +53,30 @@ type connDiagInfo struct {
53
}
54
55
type DiagInfo struct {
53
- ID string
56
+ // This nodes ID
57
+ ID string
58
+
59
+ // A list of peers this node currently has open connections to
60
Connections []connDiagInfo
55
- Keys []string
56
- LifeSpan time.Duration
57
- BwIn uint64
58
- BwOut uint64
61
+
62
+ // A list of keys provided by this node
63
+ // (currently not filled)
64
+ Keys []string
65
+
66
+ // How long this node has been running for
67
+ LifeSpan time.Duration
68
+
69
+ // Incoming Bandwidth Usage
70
+ BwIn uint64
71
+
72
+ // Outgoing Bandwidth Usage
73
+ BwOut uint64
74
+
75
+ // Information about the version of code this node is running
76
CodeVersion string
77
}
78
79
+// Marshal to json
80
func (di *DiagInfo) Marshal() []byte {
81
b, err := json.Marshal(di)
82
if err != nil {
@@ -93,6 +111,7 @@ func newID() string {
111
return string(id)
112
}
113
114
+// GetDiagnostic runs a diagnostics request across the entire network
115
func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*DiagInfo, error) {
116
log.Debug("Getting diagnostic.")
117
ctx, _ := context.WithTimeout(context.TODO(), timeout)
@@ -134,12 +153,12 @@ func (d *Diagnostics) GetDiagnostic(timeout time.Duration) ([]*DiagInfo, error)
153
if data == nil {
154
continue
155
}
137
- out = AppendDiagnostics(data, out)
156
+ out = appendDiagnostics(data, out)
157
}
158
return out, nil
159
}
160
142
-func AppendDiagnostics(data []byte, cur []*DiagInfo) []*DiagInfo {
161
+func appendDiagnostics(data []byte, cur []*DiagInfo) []*DiagInfo {
162
buf := bytes.NewBuffer(data)
163
dec := json.NewDecoder(buf)
164
for {
@@ -202,6 +221,8 @@ func (d *Diagnostics) sendRequest(ctx context.Context, p peer.Peer, pmes *pb.Mes
221
func (d *Diagnostics) handleDiagnostic(p peer.Peer, pmes *pb.Message) (*pb.Message, error) {
222
log.Debugf("HandleDiagnostic from %s for id = %s", p, pmes.GetDiagID())
223
resp := newMessage(pmes.GetDiagID())
224
+
225
+ // Make sure we havent already handled this request to prevent loops
226
d.diagLock.Lock()
227
_, found := d.diagMap[pmes.GetDiagID()]
228
if found {
importer/importer.go
+5
@@ -29,6 +29,7 @@ func NewDagFromReader(r io.Reader) (*dag.Node, error) {
29
return NewDagFromReaderWithSplitter(r, chunk.DefaultSplitter)
30
}
31
32
+// Creates an in memory DAG from data in the given reader
33
func NewDagFromReaderWithSplitter(r io.Reader, spl chunk.BlockSplitter) (*dag.Node, error) {
34
blkChan := spl.Split(r)
35
first := <-blkChan
@@ -75,6 +76,8 @@ func NewDagFromFile(fpath string) (*dag.Node, error) {
76
return NewDagFromReader(f)
77
}
78
79
+// Builds a DAG from the given file, writing created blocks to disk as they are
80
+// created
81
func BuildDagFromFile(fpath string, ds dag.DAGService, mp pin.ManualPinner) (*dag.Node, error) {
82
stat, err := os.Stat(fpath)
83
if err != nil {
@@ -94,6 +97,8 @@ func BuildDagFromFile(fpath string, ds dag.DAGService, mp pin.ManualPinner) (*da
97
return BuildDagFromReader(f, ds, mp, chunk.DefaultSplitter)
98
}
99
100
+// Builds a DAG from the data in the given reader, writing created blocks to disk
101
+// as they are created
102
func BuildDagFromReader(r io.Reader, ds dag.DAGService, mp pin.ManualPinner, spl chunk.BlockSplitter) (*dag.Node, error) {
103
blkChan := spl.Split(r)
104
net/conn/conn.go
+3
@@ -28,6 +28,7 @@ const (
28
HandshakeTimeout = time.Second * 5
29
)
30
31
+// global static buffer pool for byte arrays of size MaxMessageSize
32
var BufferPool *sync.Pool
33
34
func init() {
@@ -38,6 +39,8 @@ func init() {
39
}
40
}
41
42
+// ReleaseBuffer puts the given byte array back into the buffer pool,
43
+// first verifying that it is the correct size
44
func ReleaseBuffer(b []byte) {
45
log.Warningf("Releasing buffer! (cap,size = %d, %d)", cap(b), len(b))
46
if cap(b) != MaxMessageSize {
routing/dht/dht.go
+6
@@ -228,6 +228,8 @@ func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p peer.Peer,
228
return nil
229
}
230
231
+// putProvider sends a message to peer 'p' saying that the local node
232
+// can provide the value of 'key'
233
func (dht *IpfsDHT) putProvider(ctx context.Context, p peer.Peer, key string) error {
234
235
pmes := pb.NewMessage(pb.Message_ADD_PROVIDER, string(key), 0)
@@ -384,6 +386,7 @@ func (dht *IpfsDHT) FindLocal(id peer.ID) (peer.Peer, *kb.RoutingTable) {
386
return nil, nil
387
}
388
389
+// findPeerSingle asks peer 'p' if they know where the peer with id 'id' is
390
func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p peer.Peer, id peer.ID, level int) (*pb.Message, error) {
391
pmes := pb.NewMessage(pb.Message_FIND_NODE, string(id), level)
392
return dht.sendRequest(ctx, p, pmes)
@@ -457,6 +460,7 @@ func (dht *IpfsDHT) betterPeersToQuery(pmes *pb.Message, count int) []peer.Peer
460
return filtered
461
}
462
463
+// getPeer searches the peerstore for a peer with the given peer ID
464
func (dht *IpfsDHT) getPeer(id peer.ID) (peer.Peer, error) {
465
p, err := dht.peerstore.Get(id)
466
if err != nil {
@@ -467,6 +471,8 @@ func (dht *IpfsDHT) getPeer(id peer.ID) (peer.Peer, error) {
471
return p, nil
472
}
473
474
+// peerFromInfo returns a peer using info in the protobuf peer struct
475
+// to lookup or create a peer
476
func (dht *IpfsDHT) peerFromInfo(pbp *pb.Message_Peer) (peer.Peer, error) {
477
478
id := peer.ID(pbp.GetId())
unixfs/io/dagreader.go
+3
-1
@@ -48,7 +48,7 @@ func NewDagReader(n *mdag.Node, serv mdag.DAGService) (io.Reader, error) {
48
}
49
}
50
51
-// Follows the next link in line and loads it from the DAGService,
51
+// precalcNextBuf follows the next link in line and loads it from the DAGService,
52
// setting the next buffer to read from
53
func (dr *DagReader) precalcNextBuf() error {
54
if dr.position >= len(dr.node.Links) {
@@ -67,6 +67,7 @@ func (dr *DagReader) precalcNextBuf() error {
67
68
switch pb.GetType() {
69
case ftpb.Data_Directory:
70
+ // A directory should not exist within a file
71
return ft.ErrInvalidDirLocation
72
case ftpb.Data_File:
73
//TODO: this *should* work, needs testing first
@@ -85,6 +86,7 @@ func (dr *DagReader) precalcNextBuf() error {
86
}
87
}
88
89
+// Read reads data from the DAG structured file
90
func (dr *DagReader) Read(b []byte) (int, error) {
91
// If no cached buffer, load one
92
if dr.buf == nil {