(stream io.ReadCloser, nBuffers, bufferSize uint)
| 820 | } |
| 821 | |
| 822 | func makeBufferedNetworkReader(stream io.ReadCloser, nBuffers, bufferSize uint) *bufferedNetworkReader { |
| 823 | br := bufferedNetworkReader{ |
| 824 | stream: stream, |
| 825 | emptyBuffer: make(chan *bufferedNetworkReaderBuffer, nBuffers), |
| 826 | readyBuffer: make(chan *bufferedNetworkReaderBuffer, nBuffers), |
| 827 | terminate: make(chan bool), |
| 828 | } |
| 829 | |
| 830 | go func() { |
| 831 | handleBufferedNetworkReader(&br) |
| 832 | }() |
| 833 | |
| 834 | for range nBuffers { |
| 835 | b := bufferedNetworkReaderBuffer{ |
| 836 | data: make([]byte, bufferSize), |
| 837 | } |
| 838 | br.emptyBuffer <- &b |
| 839 | } |
| 840 | |
| 841 | return &br |
| 842 | } |
| 843 | |
| 844 | type signalCloseReader struct { |
| 845 | closed chan struct{} |
no test coverage detected
searching dependent graphs…