Hunter0x7c7
2022-08-11 a82f9cb69f63aaeba40c024960deda7d75b9fece
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
package kcp
 
import (
    "io"
    "sync"
 
    "github.com/v2fly/v2ray-core/v5/common/buf"
    "github.com/v2fly/v2ray-core/v5/common/retry"
)
 
type SegmentWriter interface {
    Write(seg Segment) error
}
 
type SimpleSegmentWriter struct {
    sync.Mutex
    buffer *buf.Buffer
    writer io.Writer
}
 
func NewSegmentWriter(writer io.Writer) SegmentWriter {
    return &SimpleSegmentWriter{
        writer: writer,
        buffer: buf.New(),
    }
}
 
func (w *SimpleSegmentWriter) Write(seg Segment) error {
    w.Lock()
    defer w.Unlock()
 
    w.buffer.Clear()
    rawBytes := w.buffer.Extend(seg.ByteSize())
    seg.Serialize(rawBytes)
    _, err := w.writer.Write(w.buffer.Bytes())
    return err
}
 
type RetryableWriter struct {
    writer SegmentWriter
}
 
func NewRetryableWriter(writer SegmentWriter) SegmentWriter {
    return &RetryableWriter{
        writer: writer,
    }
}
 
func (w *RetryableWriter) Write(seg Segment) error {
    return retry.Timed(5, 100).On(func() error {
        return w.writer.Write(seg)
    })
}