diff options
Diffstat (limited to 'pkg/util')
| -rw-r--r-- | pkg/util/io.go | 57 | ||||
| -rw-r--r-- | pkg/util/io_test.go | 59 |
2 files changed, 0 insertions, 116 deletions
diff --git a/pkg/util/io.go b/pkg/util/io.go deleted file mode 100644 index 5c489ed..0000000 --- a/pkg/util/io.go +++ /dev/null @@ -1,57 +0,0 @@ -package util - -import ( - "context" - "io" -) - -const bufferSize = 512 - -type ChannelReader struct { - Out <-chan []byte - Err <-chan error -} - -func NewChannelReader(ctx context.Context, reader io.Reader) ChannelReader { - outCh := make(chan []byte) - errCh := make(chan error) - - go func(in chan<- []byte, inErr chan<- error) { - defer close(in) - defer close(inErr) - buffer := make([]byte, 2 * bufferSize) - isD1 := true - d1 := buffer[:bufferSize] - d2 := buffer[bufferSize:] - outer: - for { - select { - case <-ctx.Done(): - 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 { - if err != io.EOF { - inErr <- err - } - break outer - } - in <- data[:n] - } - } - }(outCh, errCh) - - return ChannelReader{ - Out: outCh, - Err: errCh, - } -} diff --git a/pkg/util/io_test.go b/pkg/util/io_test.go deleted file mode 100644 index f05b833..0000000 --- a/pkg/util/io_test.go +++ /dev/null @@ -1,59 +0,0 @@ -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("%v != %v", b, data) - } - case err := <-reader.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 - } - } -} |
