summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorKyren223 <ulmliad223@gmail.com>2024-10-15 18:17:46 +0300
committerKyren223 <ulmliad223@gmail.com>2024-10-15 18:17:46 +0300
commitc5dba4360acff23a89f301235401ddfe19ab99a8 (patch)
treeaf968d5fe4f655930dd07bc90a77e3ae1e83068a
parent77dca0e1dfd609851a4c80ebb769b21838ec7093 (diff)
fix: channel reader blocks until EOF or an error was received
-rw-r--r--pkg/util/io.go9
-rw-r--r--pkg/util/io_test.go24
2 files changed, 30 insertions, 3 deletions
diff --git a/pkg/util/io.go b/pkg/util/io.go
index 5c56822..8c8725b 100644
--- a/pkg/util/io.go
+++ b/pkg/util/io.go
@@ -5,6 +5,8 @@ import (
"io"
)
+const bufferSize = 512
+
type ChannelReader struct {
Out <-chan []byte
Err <-chan error
@@ -17,6 +19,7 @@ func NewChannelReader(ctx context.Context, reader io.Reader) ChannelReader {
go func(in chan<- []byte, inErr chan<- error) {
defer close(in)
defer close(inErr)
+ data := make([]byte, bufferSize)
outer:
for {
select {
@@ -24,12 +27,12 @@ func NewChannelReader(ctx context.Context, reader io.Reader) ChannelReader {
inErr <- ctx.Err()
break outer
default:
- data, err := io.ReadAll(reader)
- if err != nil {
+ n, err := reader.Read(data)
+ if err != nil && err != io.EOF {
inErr <- err
break outer
}
- in <- data
+ in <- data[:n]
}
}
}(outCh, errCh)
diff --git a/pkg/util/io_test.go b/pkg/util/io_test.go
new file mode 100644
index 0000000..9d63e5d
--- /dev/null
+++ b/pkg/util/io_test.go
@@ -0,0 +1,24 @@
+package util
+
+import (
+ "bytes"
+ "context"
+ "testing"
+ "time"
+)
+
+func TestChannelReader(t *testing.T) {
+ ctx, cancel := context.WithTimeout(context.Background(), time.Second)
+ defer cancel()
+ b := []byte{1, 2, 3, 4, 5, 6, 7, 8}
+ reader := NewChannelReader(ctx, bytes.NewReader(b))
+
+ select {
+ case data := <-reader.Out:
+ if !bytes.Equal(b, data) {
+ t.Errorf("TestChannelReader() %v != %v", b, data)
+ }
+ case err := <-reader.Err:
+ t.Errorf("TestChannelReader() err = %v", err)
+ }
+}