@cryptotaxi247 / kubo / commits / 950957240

address comments from PR

Jeromy committed Oct 29, 2014 at 18:34 UTC 950957240a503b7a06fdb52272358aa1e8c5ecb3
7 files changed +176 -90
Godeps/_workspace/src/github.com/jbenet/go-msgio/chan.go
+17 -9
@@ -30,19 +30,27 @@ func NewChanWithPool(chanSize int, pool *sync.Pool) *Chan {
30 }
31 }
32
33 -func (s *Chan) ReadFrom(r io.Reader, maxMsgLen int) {
34 - // new buffer per message
35 - // if bottleneck, cycle around a set of buffers
36 - mr := NewReader(r, s.BufPool)
33 +func (s *Chan) getBuffer(size int) []byte {
34 if s.BufPool == nil {
38 - s.BufPool = new(sync.Pool)
39 - s.BufPool.New = func() interface{} {
40 - return make([]byte, maxMsgLen)
35 + return make([]byte, size)
36 + } else {
37 + bufi := s.BufPool.Get()
38 + buf, ok := bufi.([]byte)
39 + if !ok {
40 + panic("Got invalid type from sync pool!")
41 }
42 + return buf
43 }
44 +}
45 +
46 +func (s *Chan) ReadFrom(r io.Reader, maxMsgLen int) {
47 + // new buffer per message
48 + // if bottleneck, cycle around a set of buffers
49 + mr := NewReader(r)
50 Loop:
51 for {
45 - buf, err := mr.ReadMsg()
52 + buf := s.getBuffer(maxMsgLen)
53 + l, err := mr.ReadMsg(buf)
54 if err != nil {
55 if err == io.EOF {
56 break Loop // done
@@ -56,7 +64,7 @@ Loop:
64 select {
65 case <-s.CloseChan:
66 break Loop // told we're done
59 - case s.MsgChan <- buf:
67 + case s.MsgChan <- buf[:l]:
68 // ok seems fine. send it away
69 }
70 }
Godeps/_workspace/src/github.com/jbenet/go-msgio/msgio.go
+10 -19
@@ -3,7 +3,6 @@ package msgio
3 import (
4 "encoding/binary"
5 "io"
6 - "sync"
6 )
7
8 var NBO = binary.BigEndian
@@ -18,7 +17,7 @@ type WriteCloser interface {
17 }
18
19 type Reader interface {
21 - ReadMsg() ([]byte, error)
20 + ReadMsg([]byte) (int, error)
21 }
22
23 type ReadCloser interface {
@@ -64,30 +63,22 @@ func (s *Writer_) Close() error {
63 type Reader_ struct {
64 R io.Reader
65 lbuf []byte
67 - bp *sync.Pool
66 }
67
70 -func NewReader(r io.Reader, bufpool *sync.Pool) ReadCloser {
71 - return &Reader_{R: r, lbuf: make([]byte, 4), bp: bufpool}
68 +func NewReader(r io.Reader) ReadCloser {
69 + return &Reader_{r, make([]byte, 4)}
70 }
71
74 -func (s *Reader_) ReadMsg() ([]byte, error) {
72 +func (s *Reader_) ReadMsg(msg []byte) (int, error) {
73 if _, err := io.ReadFull(s.R, s.lbuf); err != nil {
76 - return nil, err
74 + return 0, err
75 }
78 -
79 - bufi := s.bp.Get()
80 - buf, ok := bufi.([]byte)
81 - if !ok {
82 - panic("invalid type in pool!")
83 - }
84 -
76 length := int(NBO.Uint32(s.lbuf))
86 - if length < 0 || length > len(buf) {
87 - return nil, io.ErrShortBuffer
77 + if length < 0 || length > len(msg) {
78 + return 0, io.ErrShortBuffer
79 }
89 - _, err := io.ReadFull(s.R, buf[:length])
90 - return buf[:length], err
80 + _, err := io.ReadFull(s.R, msg[:length])
81 + return length, err
82 }
83
84 func (s *Reader_) Close() error {
@@ -104,7 +95,7 @@ type ReadWriter_ struct {
95
96 func NewReadWriter(rw io.ReadWriter) ReadWriter {
97 return &ReadWriter_{
107 - Reader: NewReader(rw, nil),
98 + Reader: NewReader(rw),
99 Writer: NewWriter(rw),
100 }
101 }
blocks/blocks.go
+11
@@ -1,6 +1,7 @@
1 package blocks
2
3 import (
4 + "errors"
5 "fmt"
6
7 mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
@@ -18,6 +19,16 @@ func NewBlock(data []byte) *Block {
19 return &Block{Data: data, Multihash: u.Hash(data)}
20 }
21
22 +func NewBlockWithHash(data []byte, h mh.Multihash) (*Block, error) {
23 + if u.Debug {
24 + chk := u.Hash(data)
25 + if string(chk) != string(h) {
26 + return nil, errors.New("Data did not match given hash!")
27 + }
28 + }
29 + return &Block{Data: data, Multihash: h}, nil
30 +}
31 +
32 // Key returns the block's Multihash as a Key value.
33 func (b *Block) Key() u.Key {
34 return u.Key(b.Multihash)
crypto/spipe/handshake.go
+4 -38
@@ -185,44 +185,10 @@ func (s *SecurePipe) handshake() error {
185
186 cmp := bytes.Compare(myPubKey, proposeResp.GetPubkey())
187
188 - if true {
189 - mIV, tIV, mCKey, tCKey, mMKey, tMKey := ci.KeyStretcher(cmp, cipherType, hashType, secret)
190 -
191 - go s.handleSecureIn(hashType, cipherType, tIV, tCKey, tMKey)
192 - go s.handleSecureOut(hashType, cipherType, mIV, mCKey, mMKey)
193 -
194 - } else {
195 - log.Critical("Secure Channel Disabled! PLEASE ENSURE YOU KNOW WHAT YOU ARE DOING")
196 - // Disable Secure Channel
197 - go func(sp *SecurePipe) {
198 - for {
199 - select {
200 - case <-sp.ctx.Done():
201 - return
202 - case m, ok := <-sp.insecure.In:
203 - if !ok {
204 - sp.cancel()
205 - return
206 - }
207 - sp.In <- m
208 - }
209 - }
210 - }(s)
211 - go func(sp *SecurePipe) {
212 - for {
213 - select {
214 - case <-sp.ctx.Done():
215 - return
216 - case m, ok := <-sp.Out:
217 - if !ok {
218 - sp.cancel()
219 - return
220 - }
221 - sp.insecure.Out <- m
222 - }
223 - }
224 - }(s)
225 - }
188 + mIV, tIV, mCKey, tCKey, mMKey, tMKey := ci.KeyStretcher(cmp, cipherType, hashType, secret)
189 +
190 + go s.handleSecureIn(hashType, cipherType, tIV, tCKey, tMKey)
191 + go s.handleSecureOut(hashType, cipherType, mIV, mCKey, mMKey)
192
193 finished := []byte("Finished")
194
crypto/spipe/internal/pb/spipe.pb.go
+13 -5
@@ -1,4 +1,4 @@
1 -// Code generated by protoc-gen-go.
1 +// Code generated by protoc-gen-gogo.
2 // source: spipe.proto
3 // DO NOT EDIT!
4
@@ -15,7 +15,7 @@ It has these top-level messages:
15 */
16 package spipe_pb
17
18 -import proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/goprotobuf/proto"
18 +import proto "code.google.com/p/gogoprotobuf/proto"
19 import math "math"
20
21 // Reference imports to suppress errors if they are not otherwise used.
@@ -95,9 +95,10 @@ func (m *Exchange) GetSignature() []byte {
95 }
96
97 type DataSig struct {
98 - Data []byte `protobuf:"bytes,1,opt,name=data" json:"data,omitempty"`
99 - Sig []byte `protobuf:"bytes,2,opt,name=sig" json:"sig,omitempty"`
100 - XXX_unrecognized []byte `json:"-"`
98 + Data []byte `protobuf:"bytes,1,opt,name=data" json:"data,omitempty"`
99 + Sig []byte `protobuf:"bytes,2,opt,name=sig" json:"sig,omitempty"`
100 + Id *uint64 `protobuf:"varint,3,opt,name=id" json:"id,omitempty"`
101 + XXX_unrecognized []byte `json:"-"`
102 }
103
104 func (m *DataSig) Reset() { *m = DataSig{} }
@@ -118,5 +119,12 @@ func (m *DataSig) GetSig() []byte {
119 return nil
120 }
121
122 +func (m *DataSig) GetId() uint64 {
123 + if m != nil && m.Id != nil {
124 + return *m.Id
125 + }
126 + return 0
127 +}
128 +
129 func init() {
130 }
crypto/spipe/internal/pb/spipe.proto
+1
@@ -16,4 +16,5 @@ message Exchange {
16 message DataSig {
17 optional bytes data = 1;
18 optional bytes sig = 2;
19 + optional uint64 id = 3;
20 }
crypto/spipe/signedpipe.go
+120 -19
@@ -24,18 +24,23 @@ type SignedPipe struct {
24
25 ctx context.Context
26 cancel context.CancelFunc
27 +
28 + mesid uint64
29 + theirmesid uint64
30 }
31
32 +// secureChallengeSize is a constant that determines the initial challenge, and every subsequent
33 +// sequence number. It should be large enough to be unguessable by adversaries (128+ bits).
34 +// (SECURITY WARNING)
35 +const secureChallengeSize = (256 / 32)
36 +
37 func NewSignedPipe(parctx context.Context, bufsize int, local peer.Peer,
38 peers peer.Peerstore, insecure pipes.Duplex) (*SignedPipe, error) {
39
40 ctx, cancel := context.WithCancel(parctx)
41
42 sp := &SignedPipe{
35 - Duplex: pipes.Duplex{
36 - In: make(chan []byte, bufsize),
37 - Out: make(chan []byte, bufsize),
38 - },
43 + Duplex: pipes.NewDuplex(bufsize),
44 local: local,
45 peers: peers,
46 insecure: insecure,
@@ -51,6 +56,36 @@ func NewSignedPipe(parctx context.Context, bufsize int, local peer.Peer,
56 return sp, nil
57 }
58
59 +func (sp *SignedPipe) trySend(b []byte) bool {
60 + select {
61 + case <-sp.ctx.Done():
62 + return false
63 + case sp.insecure.Out <- b:
64 + return true
65 + }
66 +}
67 +
68 +func (sp *SignedPipe) tryRecv() ([]byte, bool) {
69 + select {
70 + case <-sp.ctx.Done():
71 + return nil, false
72 + case data, ok := <-sp.insecure.In:
73 + if !ok {
74 + return nil, false
75 + }
76 + return data, true
77 + }
78 +}
79 +
80 +func reduceChallenge(cha []byte) uint64 {
81 + var out uint64
82 + for _, b := range cha {
83 + out ^= uint64(b)
84 + out = out << 1
85 + }
86 + return out
87 +}
88 +
89 func (sp *SignedPipe) handshake() error {
90 // Send them our public key
91 pubk := sp.local.PubKey()
@@ -59,7 +94,10 @@ func (sp *SignedPipe) handshake() error {
94 return err
95 }
96
62 - sp.insecure.Out <- pkb
97 + // Exchange public keys with remote peer
98 + if !sp.trySend(pkb) {
99 + return context.Canceled
100 + }
101 theirPkb := <-sp.insecure.In
102
103 theirPubKey, err := ci.UnmarshalPublicKey(theirPkb)
@@ -67,7 +105,7 @@ func (sp *SignedPipe) handshake() error {
105 return err
106 }
107
70 - challenge := make([]byte, 32)
108 + challenge := make([]byte, secureChallengeSize)
109 rand.Read(challenge)
110
111 enc, err := theirPubKey.Encrypt(challenge)
@@ -75,24 +113,62 @@ func (sp *SignedPipe) handshake() error {
113 return err
114 }
115
78 - sp.insecure.Out <- enc
79 - theirEnc := <-sp.insecure.In
116 + chsig, err := sp.local.PrivKey().Sign(challenge)
117 + if err != nil {
118 + return err
119 + }
120 +
121 + if !sp.trySend(enc) {
122 + return context.Canceled
123 + }
124 + if !sp.trySend(chsig) {
125 + return context.Canceled
126 + }
127 +
128 + theirEnc, ok := sp.tryRecv()
129 + if !ok {
130 + return context.Canceled
131 + }
132 + theirChSig, ok := sp.tryRecv()
133 + if !ok {
134 + return context.Canceled
135 + }
136
137 + // Unencrypt and verify their challenge
138 unenc, err := sp.local.PrivKey().Unencrypt(theirEnc)
139 if err != nil {
140 return err
141 }
142 + ok, err = theirPubKey.Verify(unenc, theirChSig)
143 + if err != nil {
144 + return err
145 + }
146 + if !ok {
147 + return errors.New("Invalid signature!")
148 + }
149
150 + // Sign the unencrypted challenge, and send it back
151 sig, err := sp.local.PrivKey().Sign(unenc)
152 if err != nil {
153 return err
154 }
155
91 - sp.insecure.Out <- unenc
92 - theirUnenc := <-sp.insecure.In
93 - sp.insecure.Out <- sig
94 - theirSig := <-sp.insecure.In
156 + if !sp.trySend(unenc) {
157 + return context.Canceled
158 + }
159 + if !sp.trySend(sig) {
160 + return context.Canceled
161 + }
162 + theirUnenc, ok := sp.tryRecv()
163 + if !ok {
164 + return context.Canceled
165 + }
166 + theirSig, ok := sp.tryRecv()
167 + if !ok {
168 + return context.Canceled
169 + }
170
171 + // Verify that they correctly unecrypted the challenge
172 if !bytes.Equal(theirUnenc, challenge) {
173 return errors.New("received bad challenge response")
174 }
@@ -106,12 +182,29 @@ func (sp *SignedPipe) handshake() error {
182 return errors.New("Incorrect signature on challenge")
183 }
184
185 + sp.theirmesid = reduceChallenge(challenge)
186 + sp.mesid = reduceChallenge(unenc)
187 +
188 go sp.handleIn(theirPubKey)
110 - go sp.handleOut()
189 + go sp.handleOut(sp.local.PrivKey())
190
191 finished := []byte("finished")
113 - sp.Out <- finished
114 - resp := <-sp.In
192 +
193 + select {
194 + case <-sp.ctx.Done():
195 + return context.Canceled
196 + case sp.Out <- finished:
197 + }
198 +
199 + var resp []byte
200 + select {
201 + case <-sp.ctx.Done():
202 + return context.Canceled
203 + case resp, ok = <-sp.In:
204 + if !ok {
205 + return errors.New("Channel closed before handshake finished.")
206 + }
207 + }
208 if !bytes.Equal(resp, finished) {
209 return errors.New("Handshake failed!")
210 }
@@ -119,7 +212,7 @@ func (sp *SignedPipe) handshake() error {
212 return nil
213 }
214
122 -func (sp *SignedPipe) handleOut() {
215 +func (sp *SignedPipe) handleOut(pk ci.PrivKey) {
216 for {
217 var data []byte
218 var ok bool
@@ -135,19 +228,21 @@ func (sp *SignedPipe) handleOut() {
228
229 sdata := new(pb.DataSig)
230
138 - sig, err := sp.local.PrivKey().Sign(data)
231 + sig, err := pk.Sign(data)
232 if err != nil {
233 log.Error("Error signing outgoing data: %s", err)
141 - continue
234 + return
235 }
236
237 sdata.Data = data
238 sdata.Sig = sig
239 + sdata.Id = proto.Uint64(sp.mesid)
240 b, err := proto.Marshal(sdata)
241 if err != nil {
242 log.Error("Error marshaling signed data object: %s", err)
149 - continue
243 + return
244 }
245 + sp.mesid++
246
247 select {
248 case sp.insecure.Out <- b:
@@ -188,6 +283,12 @@ func (sp *SignedPipe) handleIn(theirPubkey ci.PubKey) {
283 continue
284 }
285
286 + if sdata.GetId() != sp.theirmesid {
287 + log.Critical("Out of order message id!")
288 + return
289 + }
290 + sp.theirmesid++
291 +
292 select {
293 case <-sp.ctx.Done():
294 return