File
Blob: spec/tun/pipe.go
| 1 | package tun |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "io" |
| 7 | "net" |
| 8 | "sync" |
| 9 | |
| 10 | pool "github.com/libp2p/go-buffer-pool" |
| 11 | ) |
| 12 | |
| 13 | const ( |
| 14 | BufferSize = 1024 * 16 |
| 15 | ) |
| 16 | |
| 17 | func IsTimeout(err error) bool { |
| 18 | t := errors.Is(err, context.DeadlineExceeded) |
| 19 | if t { |
| 20 | return t |
| 21 | } |
| 22 | if e, ok := err.(net.Error); ok { |
| 23 | return e.Timeout() |
| 24 | } |
| 25 | return false |
| 26 | } |
| 27 | |
| 28 | func Pipe(src, dst io.ReadWriteCloser) <-chan error { |
| 29 | err := make(chan error, 2) |
| 30 | wg := &sync.WaitGroup{} |
| 31 | wg.Add(2) |
| 32 | |
| 33 | go pipe(wg, err, src, dst) |
| 34 | go pipe(wg, err, dst, src) |
| 35 | go func() { |
| 36 | wg.Wait() |
| 37 | close(err) |
| 38 | }() |
| 39 | |
| 40 | return err |
| 41 | } |
| 42 | |
| 43 | func pipe(wg *sync.WaitGroup, errChan chan<- error, dst, src io.ReadWriteCloser) { |
| 44 | defer wg.Done() |
| 45 | |
| 46 | buf := pool.Get(BufferSize) |
| 47 | defer pool.Put(buf) |
| 48 | |
| 49 | _, err := io.CopyBuffer(src, dst, buf) |
| 50 | |
| 51 | src.Close() |
| 52 | dst.Close() |
| 53 | |
| 54 | if err != nil { |
| 55 | errChan <- err |
| 56 | } |
| 57 | } |