-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathoperator.go
More file actions
140 lines (107 loc) · 2.56 KB
/
Copy pathoperator.go
File metadata and controls
140 lines (107 loc) · 2.56 KB
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
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
package fastrpc
import (
"bufio"
"encoding/binary"
"errors"
"io"
"sync"
)
const (
DEFAULT_CHUNK_SIZE = 65536
)
type IOOperator struct {
written bool
leftLength uint64
mutex *sync.Mutex
reader *bufio.Reader
writer *bufio.Writer
}
func NewIOOperator(reader *bufio.Reader, writer *bufio.Writer, readLength uint64) *IOOperator {
return &IOOperator{
written: false,
leftLength: readLength,
mutex: new(sync.Mutex),
reader: reader,
writer: writer,
}
}
func (i *IOOperator) ReadDataLeft() uint64 {
i.mutex.Lock()
defer i.mutex.Unlock()
return i.leftLength
}
func (i *IOOperator) ReadIOStream(count int) ([]byte, error) {
i.mutex.Lock()
defer i.mutex.Unlock()
if uint64(count) > i.leftLength {
return nil, errors.New("invalid count")
}
buf, err := readSpecifiedBytes(i.reader, count)
if err != nil {
return nil, err
}
i.leftLength -= uint64(count)
return buf, nil
}
func (i *IOOperator) WriteIOFromBuffer(buf []byte) error {
i.mutex.Lock()
defer i.mutex.Unlock()
if i.written {
return nil
} else {
i.written = true
var metaDataBuffer []byte = make([]byte, 9)
binary.BigEndian.PutUint64(metaDataBuffer[1:], uint64(len(buf)))
err := writeSpecifiedBytes(i.writer, metaDataBuffer, 9)
if err != nil {
return err
}
return writeSpecifiedBytes(i.writer, buf, len(buf))
}
}
func (i *IOOperator) WriteIOFromReader(reader io.Reader, count int, chunkSize int) error {
if chunkSize > count {
return errors.New("chunkSize is not in bounds with count")
}
i.mutex.Lock()
defer i.mutex.Unlock()
if i.written {
return nil
} else {
i.written = true
var metaDataBuffer []byte = make([]byte, 9)
binary.BigEndian.PutUint64(metaDataBuffer[1:], uint64(count))
err := writeSpecifiedBytes(i.writer, metaDataBuffer, 9)
if err != nil {
return err
}
return readWriteSpecifiedBytes(reader, i.writer, count, chunkSize)
}
}
func (i *IOOperator) WriteNothing() error {
i.mutex.Lock()
defer i.mutex.Unlock()
if i.written {
return nil
} else {
i.written = true
return writeSpecifiedBytes(i.writer, make([]byte, 9), 9)
}
}
func (i *IOOperator) WriteError(message string) error {
i.mutex.Lock()
defer i.mutex.Unlock()
if i.written {
return nil
} else {
i.written = true
var metaDataBuffer []byte = make([]byte, 9)
metaDataBuffer[0] = 0b00000001
binary.BigEndian.PutUint64(metaDataBuffer[1:], uint64(len(message)))
err := writeSpecifiedBytes(i.writer, metaDataBuffer, 9)
if err != nil {
return err
}
return writeSpecifiedBytes(i.writer, []byte(message), len(message))
}
}