commit 2f7e40d9e97ce934b511c4468f323f9fbaa27f29
parent d65c457cdb78bfe473db953cb4cd0f8243ec21be
Author: Oliver Lowe <o@olowe.co>
Date: Thu, 27 Apr 2023 13:19:35 +1000
Correctly return scan errors
Diffstat:
1 file changed, 3 insertions(+), 3 deletions(-)
diff --git a/stream.go b/stream.go
@@ -41,7 +41,6 @@ func (c *Client) Subscribe(typ, queue, filter string) (<-chan Event, error) {
if err := json.NewEncoder(buf).Encode(m); err != nil {
return nil, fmt.Errorf("encode stream parameters: %w", err)
}
- ch := make(chan Event)
resp, err := c.post("/events", buf)
if err != nil {
return nil, err
@@ -54,17 +53,18 @@ func (c *Client) Subscribe(typ, queue, filter string) (<-chan Event, error) {
return nil, fmt.Errorf("request events: %w", iresp.Error)
}
sc := bufio.NewScanner(resp.Body)
+ ch := make(chan Event)
go func() {
for sc.Scan() {
var ev Event
if err := json.Unmarshal(sc.Bytes(), &ev); err != nil {
- ch <- Event{Error: fmt.Errorf("decode event: %w", err)}
+ ch <- Event{Error: fmt.Errorf("decode event: %v", err)}
continue
}
ch <- ev
}
if sc.Err() != nil {
- ch <- Event{Error: fmt.Errorf("scan response: %w", err)}
+ ch <- Event{Error: fmt.Errorf("scan response: %w", sc.Err())}
}
resp.Body.Close()
close(ch)