@cryptotaxi247 / kubo / commits / fd40702f7

commands: Changed Marshaler to return a io.Reader instead of a []byte

core/commands: Refactored command marshalers

Matt Bell committed Dec 16, 2014 at 04:31 UTC fd40702f7399f8c29b8e9aa7200cfa286b481978
19 files changed +135 -84
commands/command.go
+11 -6
@@ -3,6 +3,7 @@ package commands
3 import (
4 "errors"
5 "fmt"
6 + "io"
7 "reflect"
8 "strings"
9
@@ -15,9 +16,9 @@ var log = u.Logger("command")
16 // It reads from the Request, and writes results to the Response.
17 type Function func(Request) (interface{}, error)
18
18 -// Marshaler is a function that takes in a Response, and returns a marshalled []byte
19 +// Marshaler is a function that takes in a Response, and returns an io.Reader
20 // (or an error on failure)
20 -type Marshaler func(Response) ([]byte, error)
21 +type Marshaler func(Response) (io.Reader, error)
22
23 // MarshalerMap is a map of Marshaler functions, keyed by EncodingType
24 // (or an error on failure)
@@ -113,12 +114,16 @@ func (c *Command) Call(req Request) Response {
114 return res
115 }
116
117 + actualType := reflect.ValueOf(output).Type()
118 +
119 + // test if output is a channel
120 + isChan := actualType.Kind() == reflect.Chan
121 +
122 // If the command specified an output type, ensure the actual value returned is of that type
117 - if cmd.Type != nil {
118 - definedType := reflect.ValueOf(cmd.Type).Type()
119 - actualType := reflect.ValueOf(output).Type()
123 + if cmd.Type != nil && !isChan {
124 + expectedType := reflect.ValueOf(cmd.Type).Type()
125
121 - if definedType != actualType {
126 + if actualType != expectedType {
127 res.SetError(ErrIncorrectType, ErrNormal)
128 return res
129 }
commands/response.go
+35 -21
@@ -41,17 +41,33 @@ const (
41 )
42
43 var marshallers = map[EncodingType]Marshaler{
44 - JSON: func(res Response) ([]byte, error) {
44 + JSON: func(res Response) (io.Reader, error) {
45 + var value interface{}
46 if res.Error() != nil {
46 - return json.MarshalIndent(res.Error(), "", " ")
47 + value = res.Error()
48 + } else {
49 + value = res.Output()
50 }
48 - return json.MarshalIndent(res.Output(), "", " ")
51 +
52 + b, err := json.MarshalIndent(value, "", " ")
53 + if err != nil {
54 + return nil, err
55 + }
56 + return bytes.NewReader(b), nil
57 },
50 - XML: func(res Response) ([]byte, error) {
58 + XML: func(res Response) (io.Reader, error) {
59 + var value interface{}
60 if res.Error() != nil {
52 - return xml.Marshal(res.Error())
61 + value = res.Error()
62 + } else {
63 + value = res.Output()
64 + }
65 +
66 + b, err := xml.Marshal(value)
67 + if err != nil {
68 + return nil, err
69 }
54 - return xml.Marshal(res.Output())
70 + return bytes.NewReader(b), nil
71 },
72 }
73
@@ -70,7 +86,7 @@ type Response interface {
86
87 // Marshal marshals out the response into a buffer. It uses the EncodingType
88 // on the Request to chose a Marshaler (Codec).
73 - Marshal() ([]byte, error)
89 + Marshal() (io.Reader, error)
90
91 // Gets a io.Reader that reads the marshalled output
92 Reader() (io.Reader, error)
@@ -103,9 +119,9 @@ func (r *response) SetError(err error, code ErrorType) {
119 r.err = &Error{Message: err.Error(), Code: code}
120 }
121
106 -func (r *response) Marshal() ([]byte, error) {
122 +func (r *response) Marshal() (io.Reader, error) {
123 if r.err == nil && r.value == nil {
108 - return []byte{}, nil
124 + return bytes.NewReader([]byte{}), nil
125 }
126
127 enc, found, err := r.req.Option(EncShort).String()
@@ -119,7 +135,7 @@ func (r *response) Marshal() ([]byte, error) {
135
136 // Special case: if text encoding and an error, just print it out.
137 if encType == Text && r.Error() != nil {
122 - return []byte(r.Error().Error()), nil
138 + return strings.NewReader(r.Error().Error()), nil
139 }
140
141 var marshaller Marshaler
@@ -140,22 +156,20 @@ func (r *response) Marshal() ([]byte, error) {
156 // Reader returns an `io.Reader` representing marshalled output of this Response
157 // Note that multiple calls to this will return a reference to the same io.Reader
158 func (r *response) Reader() (io.Reader, error) {
143 - // if command set value to a io.Reader, use that as our reader
159 if r.out == nil {
160 if out, ok := r.value.(io.Reader); ok {
161 + // if command returned a io.Reader, use that as our reader
162 r.out = out
147 - }
148 - }
163
150 - if r.out == nil {
151 - // no reader set, so marshal the error or value
152 - marshalled, err := r.Marshal()
153 - if err != nil {
154 - return nil, err
155 - }
164 + } else {
165 + // otherwise, use the response marshaler output
166 + marshalled, err := r.Marshal()
167 + if err != nil {
168 + return nil, err
169 + }
170
157 - // create a Reader from the marshalled data
158 - r.out = bytes.NewReader(marshalled)
171 + r.out = marshalled
172 + }
173 }
174
175 return r.out, nil
core/commands/add.go
+2 -2
@@ -76,7 +76,7 @@ remains to be implemented.
76 return added, nil
77 },
78 Marshalers: cmds.MarshalerMap{
79 - cmds.Text: func(res cmds.Response) ([]byte, error) {
79 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
80 val, ok := res.Output().(*AddOutput)
81 if !ok {
82 return nil, u.ErrCast()
@@ -93,7 +93,7 @@ remains to be implemented.
93 buf.Write([]byte(fmt.Sprintf("added %s %s\n", obj.Hash, val.Names[i])))
94 }
95 }
96 - return buf.Bytes(), nil
96 + return &buf, nil
97 },
98 },
99 Type: &AddOutput{},
core/commands/block.go
+4 -2
@@ -2,7 +2,9 @@ package commands
2
3 import (
4 "bytes"
5 + "io"
6 "io/ioutil"
7 + "strings"
8 "time"
9
10 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -123,9 +125,9 @@ It reads from stdin, and <key> is a base58 encoded multihash.
125 },
126 Type: &Block{},
127 Marshalers: cmds.MarshalerMap{
126 - cmds.Text: func(res cmds.Response) ([]byte, error) {
128 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
129 block := res.Output().(*Block)
128 - return []byte(block.Key + "\n"), nil
130 + return strings.NewReader(block.Key + "\n"), nil
131 },
132 },
133 }
core/commands/bootstrap.go
+6 -6
@@ -114,7 +114,7 @@ in the bootstrap list).
114 },
115 Type: &BootstrapOutput{},
116 Marshalers: cmds.MarshalerMap{
117 - cmds.Text: func(res cmds.Response) ([]byte, error) {
117 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
118 v, ok := res.Output().(*BootstrapOutput)
119 if !ok {
120 return nil, u.ErrCast()
@@ -122,7 +122,7 @@ in the bootstrap list).
122
123 var buf bytes.Buffer
124 err := bootstrapWritePeers(&buf, "added ", v.Peers)
125 - return buf.Bytes(), err
125 + return &buf, err
126 },
127 },
128 }
@@ -175,7 +175,7 @@ var bootstrapRemoveCmd = &cmds.Command{
175 },
176 Type: &BootstrapOutput{},
177 Marshalers: cmds.MarshalerMap{
178 - cmds.Text: func(res cmds.Response) ([]byte, error) {
178 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
179 v, ok := res.Output().(*BootstrapOutput)
180 if !ok {
181 return nil, u.ErrCast()
@@ -183,7 +183,7 @@ var bootstrapRemoveCmd = &cmds.Command{
183
184 var buf bytes.Buffer
185 err := bootstrapWritePeers(&buf, "removed ", v.Peers)
186 - return buf.Bytes(), err
186 + return &buf, err
187 },
188 },
189 }
@@ -209,7 +209,7 @@ var bootstrapListCmd = &cmds.Command{
209 },
210 }
211
212 -func bootstrapMarshaler(res cmds.Response) ([]byte, error) {
212 +func bootstrapMarshaler(res cmds.Response) (io.Reader, error) {
213 v, ok := res.Output().(*BootstrapOutput)
214 if !ok {
215 return nil, u.ErrCast()
@@ -217,7 +217,7 @@ func bootstrapMarshaler(res cmds.Response) ([]byte, error) {
217
218 var buf bytes.Buffer
219 err := bootstrapWritePeers(&buf, "", v.Peers)
220 - return buf.Bytes(), err
220 + return &buf, err
221 }
222
223 func bootstrapWritePeers(w io.Writer, prefix string, peers []config.BootstrapPeer) error {
core/commands/commands.go
+3 -2
@@ -2,6 +2,7 @@ package commands
2
3 import (
4 "bytes"
5 + "io"
6 "sort"
7
8 cmds "github.com/jbenet/go-ipfs/commands"
@@ -26,13 +27,13 @@ func CommandsCmd(root *cmds.Command) *cmds.Command {
27 return &root, nil
28 },
29 Marshalers: cmds.MarshalerMap{
29 - cmds.Text: func(res cmds.Response) ([]byte, error) {
30 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
31 v := res.Output().(*Command)
32 var buf bytes.Buffer
33 for _, s := range cmdPathStrings(v) {
34 buf.Write([]byte(s + "\n"))
35 }
35 - return buf.Bytes(), nil
36 + return &buf, nil
37 },
38 },
39 Type: &Command{},
core/commands/config.go
+2 -2
@@ -72,7 +72,7 @@ Set the value of the 'datastore.path' key:
72 }
73 },
74 Marshalers: cmds.MarshalerMap{
75 - cmds.Text: func(res cmds.Response) ([]byte, error) {
75 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
76 if len(res.Request().Arguments()) == 2 {
77 return nil, nil // dont output anything
78 }
@@ -92,7 +92,7 @@ Set the value of the 'datastore.path' key:
92 return nil, err
93 }
94 buf = append(buf, byte('\n'))
95 - return buf, nil
95 + return bytes.NewReader(buf), nil
96 },
97 },
98 Type: &ConfigField{},
core/commands/diag.go
+2 -2
@@ -86,7 +86,7 @@ connected peers and latencies between them.
86 },
87 Type: &DiagnosticOutput{},
88 Marshalers: cmds.MarshalerMap{
89 - cmds.Text: func(r cmds.Response) ([]byte, error) {
89 + cmds.Text: func(r cmds.Response) (io.Reader, error) {
90 output, ok := r.Output().(*DiagnosticOutput)
91 if !ok {
92 return nil, util.ErrCast()
@@ -96,7 +96,7 @@ connected peers and latencies between them.
96 if err != nil {
97 return nil, err
98 }
99 - return buf.Bytes(), nil
99 + return &buf, nil
100 },
101 },
102 }
core/commands/id.go
+8 -2
@@ -1,9 +1,11 @@
1 package commands
2
3 import (
4 + "bytes"
5 "encoding/base64"
6 "encoding/json"
7 "errors"
8 + "io"
9 "time"
10
11 "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -76,13 +78,17 @@ if no peer is specified, prints out local peers info.
78 return printPeer(node.Peerstore, p.ID)
79 },
80 Marshalers: cmds.MarshalerMap{
79 - cmds.Text: func(res cmds.Response) ([]byte, error) {
81 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
82 val, ok := res.Output().(*IdOutput)
83 if !ok {
84 return nil, u.ErrCast()
85 }
86
85 - return json.MarshalIndent(val, "", "\t")
87 + marshaled, err := json.MarshalIndent(val, "", "\t")
88 + if err != nil {
89 + return nil, err
90 + }
91 + return bytes.NewReader(marshaled), nil
92 },
93 },
94 Type: &IdOutput{},
core/commands/ls.go
+4 -2
@@ -2,6 +2,8 @@ package commands
2
3 import (
4 "fmt"
5 + "io"
6 + "strings"
7
8 cmds "github.com/jbenet/go-ipfs/commands"
9 merkledag "github.com/jbenet/go-ipfs/merkledag"
@@ -70,7 +72,7 @@ it contains, with the following format:
72 return &LsOutput{output}, nil
73 },
74 Marshalers: cmds.MarshalerMap{
73 - cmds.Text: func(res cmds.Response) ([]byte, error) {
75 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
76 s := ""
77 output := res.Output().(*LsOutput).Objects
78
@@ -84,7 +86,7 @@ it contains, with the following format:
86 }
87 }
88
87 - return []byte(s), nil
89 + return strings.NewReader(s), nil
90 },
91 },
92 Type: &LsOutput{},
core/commands/mount_unix.go
+3 -2
@@ -4,6 +4,7 @@ package commands
4
5 import (
6 "fmt"
7 + "io"
8 "strings"
9 "time"
10
@@ -134,11 +135,11 @@ baz
135 },
136 Type: &config.Mounts{},
137 Marshalers: cmds.MarshalerMap{
137 - cmds.Text: func(res cmds.Response) ([]byte, error) {
138 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
139 v := res.Output().(*config.Mounts)
140 s := fmt.Sprintf("IPFS mounted at: %s\n", v.IPFS)
141 s += fmt.Sprintf("IPNS mounted at: %s\n", v.IPNS)
141 - return []byte(s), nil
142 + return strings.NewReader(s), nil
143 },
144 },
145 }
core/commands/object.go
+12 -6
@@ -6,6 +6,7 @@ import (
6 "errors"
7 "io"
8 "io/ioutil"
9 + "strings"
10
11 mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
12
@@ -101,10 +102,10 @@ multihash.
102 return objectLinks(n, key)
103 },
104 Marshalers: cmds.MarshalerMap{
104 - cmds.Text: func(res cmds.Response) ([]byte, error) {
105 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
106 object := res.Output().(*Object)
107 marshalled := marshalLinks(object.Links)
107 - return []byte(marshalled), nil
108 + return strings.NewReader(marshalled), nil
109 },
110 },
111 Type: &Object{},
@@ -163,13 +164,18 @@ This command outputs data in the following encodings:
164 },
165 Type: &Node{},
166 Marshalers: cmds.MarshalerMap{
166 - cmds.EncodingType("protobuf"): func(res cmds.Response) ([]byte, error) {
167 + cmds.EncodingType("protobuf"): func(res cmds.Response) (io.Reader, error) {
168 node := res.Output().(*Node)
169 object, err := deserializeNode(node)
170 if err != nil {
171 return nil, err
172 }
172 - return object.Marshal()
173 +
174 + marshaled, err := object.Marshal()
175 + if err != nil {
176 + return nil, err
177 + }
178 + return bytes.NewReader(marshaled), nil
179 },
180 },
181 }
@@ -221,9 +227,9 @@ Data should be in the format specified by <encoding>.
227 return output, nil
228 },
229 Marshalers: cmds.MarshalerMap{
224 - cmds.Text: func(res cmds.Response) ([]byte, error) {
230 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
231 object := res.Output().(*Object)
226 - return []byte("added " + object.Hash), nil
232 + return strings.NewReader("added " + object.Hash), nil
233 },
234 },
235 Type: &Object{},
core/commands/publish.go
+4 -2
@@ -3,6 +3,8 @@ package commands
3 import (
4 "errors"
5 "fmt"
6 + "io"
7 + "strings"
8
9 cmds "github.com/jbenet/go-ipfs/commands"
10 core "github.com/jbenet/go-ipfs/core"
@@ -78,10 +80,10 @@ Publish a <ref> to another public key:
80 return publish(n, n.PrivateKey, ref)
81 },
82 Marshalers: cmds.MarshalerMap{
81 - cmds.Text: func(res cmds.Response) ([]byte, error) {
83 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
84 v := res.Output().(*IpnsEntry)
85 s := fmt.Sprintf("Published name %s to %s\n", v.Name, v.Value)
84 - return []byte(s), nil
86 + return strings.NewReader(s), nil
87 },
88 },
89 Type: &IpnsEntry{},
core/commands/refs.go
+6 -4
@@ -1,7 +1,9 @@
1 package commands
2
3 import (
4 + "bytes"
5 "fmt"
6 + "io"
7
8 mh "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multihash"
9 cmds "github.com/jbenet/go-ipfs/commands"
@@ -16,13 +18,13 @@ type KeyList struct {
18 }
19
20 // KeyListTextMarshaler outputs a KeyList as plaintext, one key per line
19 -func KeyListTextMarshaler(res cmds.Response) ([]byte, error) {
21 +func KeyListTextMarshaler(res cmds.Response) (io.Reader, error) {
22 output := res.Output().(*KeyList)
21 - s := ""
23 + var buf bytes.Buffer
24 for _, key := range output.Keys {
23 - s += key.B58String() + "\n"
25 + buf.WriteString(key.B58String() + "\n")
26 }
25 - return []byte(s), nil
27 + return &buf, nil
28 }
29
30 var RefsCmd = &cmds.Command{
core/commands/resolve.go
+4 -2
@@ -2,6 +2,8 @@ package commands
2
3 import (
4 "errors"
5 + "io"
6 + "strings"
7
8 cmds "github.com/jbenet/go-ipfs/commands"
9 )
@@ -71,9 +73,9 @@ Resolve te value of another name:
73 return output, nil
74 },
75 Marshalers: cmds.MarshalerMap{
74 - cmds.Text: func(res cmds.Response) ([]byte, error) {
76 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
77 output := res.Output().(string)
76 - return []byte(output), nil
78 + return strings.NewReader(output), nil
79 },
80 },
81 }
core/commands/root.go
+5 -2
@@ -1,6 +1,9 @@
1 package commands
2
3 import (
4 + "io"
5 + "strings"
6 +
7 cmds "github.com/jbenet/go-ipfs/commands"
8 u "github.com/jbenet/go-ipfs/util"
9 )
@@ -90,6 +93,6 @@ type MessageOutput struct {
93 Message string
94 }
95
93 -func MessageTextMarshaler(res cmds.Response) ([]byte, error) {
94 - return []byte(res.Output().(*MessageOutput).Message), nil
96 +func MessageTextMarshaler(res cmds.Response) (io.Reader, error) {
97 + return strings.NewReader(res.Output().(*MessageOutput).Message), nil
98 }
core/commands/swarm.go
+5 -4
@@ -3,6 +3,7 @@ package commands
3 import (
4 "bytes"
5 "fmt"
6 + "io"
7 "path"
8
9 cmds "github.com/jbenet/go-ipfs/commands"
@@ -124,7 +125,7 @@ ipfs swarm connect /ip4/104.131.131.82/tcp/4001/QmaCpDMGvV2BGHeYERUEnRQAwe3N8Szb
125 Type: &stringList{},
126 }
127
127 -func stringListMarshaler(res cmds.Response) ([]byte, error) {
128 +func stringListMarshaler(res cmds.Response) (io.Reader, error) {
129 list, ok := res.Output().(*stringList)
130 if !ok {
131 return nil, errors.New("failed to cast []string")
@@ -132,10 +133,10 @@ func stringListMarshaler(res cmds.Response) ([]byte, error) {
133
134 var buf bytes.Buffer
135 for _, s := range list.Strings {
135 - buf.Write([]byte(s))
136 - buf.Write([]byte("\n"))
136 + buf.WriteString(s)
137 + buf.WriteString("\n")
138 }
138 - return buf.Bytes(), nil
139 + return &buf, nil
140 }
141
142 // splitAddresses is a function that takes in a slice of string peer addresses
core/commands/update.go
+14 -12
@@ -1,8 +1,10 @@
1 package commands
2
3 import (
4 + "bytes"
5 "errors"
6 "fmt"
7 + "io"
8
9 cmds "github.com/jbenet/go-ipfs/commands"
10 "github.com/jbenet/go-ipfs/core"
@@ -33,16 +35,16 @@ var UpdateCmd = &cmds.Command{
35 "log": UpdateLogCmd,
36 },
37 Marshalers: cmds.MarshalerMap{
36 - cmds.Text: func(res cmds.Response) ([]byte, error) {
38 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
39 v := res.Output().(*UpdateOutput)
38 - s := ""
40 + var buf bytes.Buffer
41 if v.NewVersion != v.OldVersion {
40 - s = fmt.Sprintf("Successfully updated to IPFS version '%s' (from '%s')\n",
41 - v.NewVersion, v.OldVersion)
42 + buf.WriteString(fmt.Sprintf("Successfully updated to IPFS version '%s' (from '%s')\n",
43 + v.NewVersion, v.OldVersion))
44 } else {
43 - s = fmt.Sprintf("Already updated to latest version ('%s')\n", v.NewVersion)
45 + buf.WriteString(fmt.Sprintf("Already updated to latest version ('%s')\n", v.NewVersion))
46 }
45 - return []byte(s), nil
47 + return &buf, nil
48 },
49 },
50 }
@@ -62,16 +64,16 @@ var UpdateCheckCmd = &cmds.Command{
64 },
65 Type: &UpdateOutput{},
66 Marshalers: cmds.MarshalerMap{
65 - cmds.Text: func(res cmds.Response) ([]byte, error) {
67 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
68 v := res.Output().(*UpdateOutput)
67 - s := ""
69 + var buf bytes.Buffer
70 if v.NewVersion != v.OldVersion {
69 - s = fmt.Sprintf("A new version of IPFS is available ('%s', currently running '%s')\n",
70 - v.NewVersion, v.OldVersion)
71 + buf.WriteString(fmt.Sprintf("A new version of IPFS is available ('%s', currently running '%s')\n",
72 + v.NewVersion, v.OldVersion))
73 } else {
72 - s = fmt.Sprintf("Already updated to latest version ('%s')\n", v.NewVersion)
74 + buf.WriteString(fmt.Sprintf("Already updated to latest version ('%s')\n", v.NewVersion))
75 }
74 - return []byte(s), nil
76 + return &buf, nil
77 },
78 },
79 }
core/commands/version.go
+5 -3
@@ -2,6 +2,8 @@ package commands
2
3 import (
4 "fmt"
5 + "io"
6 + "strings"
7
8 cmds "github.com/jbenet/go-ipfs/commands"
9 config "github.com/jbenet/go-ipfs/config"
@@ -26,7 +28,7 @@ var VersionCmd = &cmds.Command{
28 }, nil
29 },
30 Marshalers: cmds.MarshalerMap{
29 - cmds.Text: func(res cmds.Response) ([]byte, error) {
31 + cmds.Text: func(res cmds.Response) (io.Reader, error) {
32 v := res.Output().(*VersionOutput)
33
34 number, found, err := res.Request().Option("number").Bool()
@@ -34,9 +36,9 @@ var VersionCmd = &cmds.Command{
36 return nil, err
37 }
38 if found && number {
37 - return []byte(fmt.Sprintln(v.Version)), nil
39 + return strings.NewReader(fmt.Sprintln(v.Version)), nil
40 }
39 - return []byte(fmt.Sprintf("ipfs version %s\n", v.Version)), nil
41 + return strings.NewReader(fmt.Sprintf("ipfs version %s\n", v.Version)), nil
42 },
43 },
44 Type: &VersionOutput{},