diff options
Diffstat (limited to 'pkg')
| -rw-r--r-- | pkg/util/io.go | 19 | ||||
| -rw-r--r-- | pkg/util/io_test.go | 39 |
2 files changed, 53 insertions, 5 deletions
diff --git a/pkg/util/io.go b/pkg/util/io.go index 8c8725b..5c489ed 100644 --- a/pkg/util/io.go +++ b/pkg/util/io.go @@ -19,7 +19,10 @@ 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) + buffer := make([]byte, 2 * bufferSize) + isD1 := true + d1 := buffer[:bufferSize] + d2 := buffer[bufferSize:] outer: for { select { @@ -27,9 +30,19 @@ func NewChannelReader(ctx context.Context, reader io.Reader) ChannelReader { inErr <- ctx.Err() break outer default: + var data []byte + if isD1 { + data = d1 + } else { + data = d2 + } + isD1 = !isD1 + n, err := reader.Read(data) - if err != nil && err != io.EOF { - inErr <- err + if err != nil { + if err != io.EOF { + inErr <- err + } break outer } in <- data[:n] diff --git a/pkg/util/io_test.go b/pkg/util/io_test.go index 9d63e5d..f05b833 100644 --- a/pkg/util/io_test.go +++ b/pkg/util/io_test.go @@ -16,9 +16,44 @@ func TestChannelReader(t *testing.T) { select { case data := <-reader.Out: if !bytes.Equal(b, data) { - t.Errorf("TestChannelReader() %v != %v", b, data) + t.Errorf("%v != %v", b, data) } case err := <-reader.Err: - t.Errorf("TestChannelReader() err = %v", err) + t.Errorf("reading err: %v", err) + } +} + +func TestChannelReaderMultiPartRead(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + + b := []byte("Testing FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting Framer") + b1 := b[:512] + b2 := b[512:] + reader := NewChannelReader(ctx, bytes.NewReader(b)) + + counter := 0 +outer: + for { + select { + case data := <-reader.Out: + if counter == 0 { + if !bytes.Equal(b1, data) { + t.Errorf("%v != %v", b, data) + } + } else { + if !bytes.Equal(b2[:len(data)], data) { + t.Errorf("%v != %v", b, data) + } + } + counter++ + if counter == 2 { + break outer + } + + case err := <-reader.Err: + t.Errorf("reading err: %v", err) + break outer + } } } |
