@cryptotaxi247 / kubo / commits / 7ad2f527a

Send empty struct to pubsub cmd output to flush

So the HTTP headers get sent License: MIT Signed-off-by: Jan Winkelmann <j-winkelmann@tuhh.de>

Jan Winkelmann committed Nov 23, 2016 at 19:32 UTC 7ad2f527a5e137fc5d5146dc1836911106b32466
1 file changed +21 -3
core/commands/pubsub.go
+21 -3
@@ -6,6 +6,7 @@ import (
6 "encoding/binary"
7 "fmt"
8 "io"
9 + "strings"
10 "sync"
11 "time"
12
@@ -89,10 +90,13 @@ To use, the daemon must be run with '--enable-pubsub-experiment'.
90 go func() {
91 defer sub.Cancel()
92 defer close(out)
93 +
94 + out <- floodsub.Message{}
95 +
96 for {
97 msg, err := sub.Next(req.Context())
98 if err == io.EOF || err == context.Canceled {
95 - break
99 + return
100 } else if err != nil {
101 res.SetError(err, cmds.ErrNormal)
102 return
@@ -118,16 +122,30 @@ To use, the daemon must be run with '--enable-pubsub-experiment'.
122 },
123 Marshalers: cmds.MarshalerMap{
124 cmds.Text: getPsMsgMarshaler(func(m *floodsub.Message) (io.Reader, error) {
125 + if m.Message == nil {
126 + return strings.NewReader(""), nil
127 + }
128 +
129 return bytes.NewReader(m.Data), nil
130 }),
131 "ndpayload": getPsMsgMarshaler(func(m *floodsub.Message) (io.Reader, error) {
132 + if m.Message == nil {
133 + return strings.NewReader("\n"), nil
134 + }
135 +
136 m.Data = append(m.Data, '\n')
137 return bytes.NewReader(m.Data), nil
138 }),
139 "lenpayload": getPsMsgMarshaler(func(m *floodsub.Message) (io.Reader, error) {
140 buf := make([]byte, 8)
129 - n := binary.PutUvarint(buf, uint64(len(m.Data)))
130 - return io.MultiReader(bytes.NewReader(buf[:n]), bytes.NewReader(m.Data)), nil
141 +
142 + var data []byte
143 + if m.Message != nil {
144 + data = m.Data
145 + }
146 +
147 + n := binary.PutUvarint(buf, uint64(len(data)))
148 + return io.MultiReader(bytes.NewReader(buf[:n]), bytes.NewReader(data)), nil
149 }),
150 },
151 Type: floodsub.Message{},