update go-msgio dependency
License: MIT Signed-off-by: Jeromy <jeromyj@gmail.com>
Jeromy committed
Oct 8, 2015 at 15:49 UTC
8e79770b4a78bfe11a62ceec07855b4c0a4ec2b8
8 files changed
+234
-11
Godeps/Godeps.json
+1
-1
@@ -162,7 +162,7 @@
162
},
163
{
164
"ImportPath": "github.com/jbenet/go-msgio",
165
- "Rev": "b4f3f1e1c7ec0cbf2fe35d8a45d1c253d224dc72"
165
+ "Rev": "9399b44f6bf265b30bedaf2af8c0604bbc8d5275"
166
},
167
{
168
"ImportPath": "github.com/jbenet/go-multiaddr",
Godeps/_workspace/src/github.com/jbenet/go-msgio/limit.go
new
+45
@@ -0,0 +1,45 @@
1
+package msgio
2
+
3
+import (
4
+ "bytes"
5
+ "io"
6
+ "sync"
7
+)
8
+
9
+// LimitedReader wraps an io.Reader with a msgio framed reader. The LimitedReader
10
+// will return a reader which will io.EOF when the msg length is done.
11
+func LimitedReader(r io.Reader) (io.Reader, error) {
12
+ l, err := ReadLen(r, nil)
13
+ return io.LimitReader(r, int64(l)), err
14
+}
15
+
16
+// LimitedWriter wraps an io.Writer with a msgio framed writer. It is the inverse
17
+// of LimitedReader: it will buffer all writes until "Flush" is called. When Flush
18
+// is called, it will write the size of the buffer first, flush the buffer, reset
19
+// the buffer, and begin accept more incoming writes.
20
+func NewLimitedWriter(w io.Writer) *LimitedWriter {
21
+ return &LimitedWriter{W: w}
22
+}
23
+
24
+type LimitedWriter struct {
25
+ W io.Writer
26
+ B bytes.Buffer
27
+ M sync.Mutex
28
+}
29
+
30
+func (w *LimitedWriter) Write(buf []byte) (n int, err error) {
31
+ w.M.Lock()
32
+ n, err = w.B.Write(buf)
33
+ w.M.Unlock()
34
+ return n, err
35
+}
36
+
37
+func (w *LimitedWriter) Flush() error {
38
+ w.M.Lock()
39
+ defer w.M.Unlock()
40
+ if err := WriteLen(w.W, w.B.Len()); err != nil {
41
+ return err
42
+ }
43
+ _, err := w.B.WriteTo(w.W)
44
+ return err
45
+}
Godeps/_workspace/src/github.com/jbenet/go-msgio/msgio.go
+7
-10
@@ -1,7 +1,6 @@
1
package msgio
2
3
import (
4
- "encoding/binary"
4
"errors"
5
"io"
6
"sync"
@@ -9,9 +8,6 @@ import (
8
mpool "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-msgio/mpool"
9
)
10
12
-// NBO is NetworkByteOrder
13
-var NBO = binary.BigEndian
14
-
11
// ErrMsgTooLarge is returned when the message length is exessive
12
var ErrMsgTooLarge = errors.New("message too large")
13
@@ -101,9 +97,7 @@ func (s *writer) Write(msg []byte) (int, error) {
97
func (s *writer) WriteMsg(msg []byte) (err error) {
98
s.lock.Lock()
99
defer s.lock.Unlock()
104
-
105
- length := uint32(len(msg))
106
- if err := binary.Write(s.W, NBO, &length); err != nil {
100
+ if err := WriteLen(s.W, len(msg)); err != nil {
101
return err
102
}
103
_, err = s.W.Write(msg)
@@ -166,10 +160,12 @@ func (s *reader) NextMsgLen() (int, error) {
160
161
func (s *reader) nextMsgLen() (int, error) {
162
if s.next == -1 {
169
- if _, err := io.ReadFull(s.R, s.lbuf); err != nil {
163
+ n, err := ReadLen(s.R, s.lbuf)
164
+ if err != nil {
165
return 0, err
166
}
172
- s.next = int(NBO.Uint32(s.lbuf))
167
+
168
+ s.next = n
169
}
170
return s.next, nil
171
}
@@ -186,6 +182,7 @@ func (s *reader) Read(msg []byte) (int, error) {
182
if length > len(msg) {
183
return 0, io.ErrShortBuffer
184
}
185
+
186
_, err = io.ReadFull(s.R, msg[:length])
187
s.next = -1 // signal we've consumed this msg
188
return length, err
@@ -200,7 +197,7 @@ func (s *reader) ReadMsg() ([]byte, error) {
197
return nil, err
198
}
199
203
- if length > s.max {
200
+ if length > s.max || length < 0 {
201
return nil, ErrMsgTooLarge
202
}
203
Godeps/_workspace/src/github.com/jbenet/go-msgio/msgio/.gitignore
new
+1
@@ -0,0 +1 @@
1
+msgio
Godeps/_workspace/src/github.com/jbenet/go-msgio/msgio/README.md
new
+24
@@ -0,0 +1,24 @@
1
+# msgio headers tool
2
+
3
+Conveniently output msgio headers.
4
+
5
+## Install
6
+
7
+```
8
+go get github.com/jbenet/go-msgio/msgio
9
+```
10
+
11
+## Usage
12
+
13
+```
14
+> msgio -h
15
+msgio - tool to wrap messages with msgio header
16
+
17
+Usage
18
+ msgio header 1020 >header
19
+ cat file | msgio wrap >wrapped
20
+
21
+Commands
22
+ header <size> output a msgio header of given size
23
+ wrap wrap incoming stream with msgio
24
+```
Godeps/_workspace/src/github.com/jbenet/go-msgio/msgio/msgio.go
new
+108
@@ -0,0 +1,108 @@
1
+package main
2
+
3
+import (
4
+ "flag"
5
+ "fmt"
6
+ "io"
7
+ "io/ioutil"
8
+ "os"
9
+ "strconv"
10
+ "strings"
11
+
12
+ msgio "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-msgio"
13
+)
14
+
15
+var Args ArgType
16
+
17
+type ArgType struct {
18
+ Command string
19
+ Args []string
20
+}
21
+
22
+func (a *ArgType) Arg(i int) string {
23
+ n := i + 1
24
+ if len(a.Args) < n {
25
+ die(fmt.Sprintf("expected %d argument(s)", n))
26
+ }
27
+ return a.Args[i]
28
+}
29
+
30
+var usageStr = `
31
+msgio - tool to wrap messages with msgio header
32
+
33
+Usage
34
+ msgio header 1020 >header
35
+ cat file | msgio wrap >wrapped
36
+
37
+Commands
38
+ header <size> output a msgio header of given size
39
+ wrap wrap incoming stream with msgio
40
+`
41
+
42
+func usage() {
43
+ fmt.Println(strings.TrimSpace(usageStr))
44
+ os.Exit(0)
45
+}
46
+
47
+func die(err string) {
48
+ fmt.Fprintf(os.Stderr, "error: %s\n", err)
49
+ os.Exit(-1)
50
+}
51
+
52
+func main() {
53
+ if err := run(); err != nil {
54
+ die(err.Error())
55
+ }
56
+}
57
+
58
+func argParse() {
59
+ flag.Usage = usage
60
+ flag.Parse()
61
+
62
+ args := flag.Args()
63
+ if l := len(args); l < 1 || l > 2 {
64
+ usage()
65
+ }
66
+
67
+ Args.Command = flag.Args()[0]
68
+ Args.Args = flag.Args()[1:]
69
+}
70
+
71
+func run() error {
72
+ argParse()
73
+
74
+ w := os.Stdout
75
+ r := os.Stdin
76
+
77
+ switch Args.Command {
78
+ case "header":
79
+ size, err := strconv.Atoi(Args.Arg(0))
80
+ if err != nil {
81
+ return err
82
+ }
83
+ return header(w, size)
84
+ case "wrap":
85
+ return wrap(w, r)
86
+ default:
87
+ usage()
88
+ return nil
89
+ }
90
+}
91
+
92
+func header(w io.Writer, size int) error {
93
+ return msgio.WriteLen(w, size)
94
+}
95
+
96
+func wrap(w io.Writer, r io.Reader) error {
97
+ buf, err := ioutil.ReadAll(r)
98
+ if err != nil {
99
+ return err
100
+ }
101
+
102
+ if err := msgio.WriteLen(w, len(buf)); err != nil {
103
+ return err
104
+ }
105
+
106
+ _, err = w.Write(buf)
107
+ return err
108
+}
Godeps/_workspace/src/github.com/jbenet/go-msgio/msgio_test.go
+15
@@ -180,3 +180,18 @@ func SubtestReadWriteMsgSync(t *testing.T, writer WriteCloser, reader ReadCloser
180
t.Error(e)
181
}
182
}
183
+
184
+func TestBadSizes(t *testing.T) {
185
+ data := make([]byte, 4)
186
+
187
+ // on a 64 bit system, this will fail because its too large
188
+ // on a 32 bit system, this will fail because its too small
189
+ NBO.PutUint32(data, 4000000000)
190
+ buf := bytes.NewReader(data)
191
+ read := NewReader(buf)
192
+ msg, err := read.ReadMsg()
193
+ if err == nil {
194
+ t.Fatal(err)
195
+ }
196
+ _ = msg
197
+}
Godeps/_workspace/src/github.com/jbenet/go-msgio/num.go
new
+33
@@ -0,0 +1,33 @@
1
+package msgio
2
+
3
+import (
4
+ "encoding/binary"
5
+ "io"
6
+)
7
+
8
+// NBO is NetworkByteOrder
9
+var NBO = binary.BigEndian
10
+
11
+// WriteLen writes a length to the given writer.
12
+func WriteLen(w io.Writer, l int) error {
13
+ ul := uint32(l)
14
+ return binary.Write(w, NBO, &ul)
15
+}
16
+
17
+// ReadLen reads a length from the given reader.
18
+// if buf is non-nil, it reuses the buffer. Ex:
19
+// l, err := ReadLen(r, nil)
20
+// _, err := ReadLen(r, buf)
21
+func ReadLen(r io.Reader, buf []byte) (int, error) {
22
+ if len(buf) < 4 {
23
+ buf = make([]byte, 4)
24
+ }
25
+ buf = buf[:4]
26
+
27
+ if _, err := io.ReadFull(r, buf); err != nil {
28
+ return 0, err
29
+ }
30
+
31
+ n := int(NBO.Uint32(buf))
32
+ return n, nil
33
+}