-
Notifications
You must be signed in to change notification settings - Fork 376
/
func_io.go
94 lines (80 loc) · 1.98 KB
/
func_io.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
package streamutil
import (
"io"
"go.uber.org/zap"
"berty.tech/berty/v2/go/pkg/errcode"
)
func FuncReader(readFunc func() ([]byte, error), l *zap.Logger) *io.PipeReader {
in, out := io.Pipe()
go func() {
err := func() error {
for {
block, err := readFunc()
if err == io.EOF {
return nil
}
if err != nil {
return errcode.ErrStreamRead.Wrap(err)
}
if _, err := out.Write(block); err != nil {
return errcode.ErrStreamWrite.Wrap(err)
}
}
}()
ClosePipeOut(out, err, "FuncReader: close pipe out", l)
}()
return in
}
func FuncSink(buffer []byte, reader io.Reader, writeFunc func(block []byte) error) error {
for {
n, err := reader.Read(buffer)
if err == io.EOF {
return nil
}
if err != nil {
return errcode.ErrStreamRead.Wrap(err)
}
if err := writeFunc(buffer[:n]); err != nil {
return errcode.ErrStreamWrite.Wrap(err)
}
}
}
func FuncBlockTransformer(buf []byte, reader io.Reader, l *zap.Logger, transformFunc func(block []byte) ([]byte, error)) *io.PipeReader {
in, out := io.Pipe()
go func() {
err := func() error {
for {
n, readErr := io.ReadFull(reader, buf)
if readErr == io.EOF {
return nil
}
if readErr != nil && readErr != io.ErrUnexpectedEOF {
return errcode.ErrStreamRead.Wrap(readErr)
}
transformed, err := transformFunc(buf[:n])
if err != nil {
return errcode.ErrStreamTransform.Wrap(err)
}
if _, err := out.Write(transformed); err != nil {
return errcode.ErrStreamWrite.Wrap(err)
}
if readErr == io.ErrUnexpectedEOF {
return nil // last block can be smaller
}
}
}()
ClosePipeOut(out, err, "FuncBlockTransformer: close pipe out", l)
}()
return in
}
func ClosePipeOut(out *io.PipeWriter, incoming error, errPrefix string, l *zap.Logger) {
var cErr error
if incoming == nil || incoming == io.EOF {
cErr = out.Close()
} else {
cErr = out.CloseWithError(incoming)
}
if cErr != nil {
l.Error(errPrefix, zap.Error(cErr))
}
}