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
|
type FixSizeLargeBuff struct {
buf []byte
}
const Megabit = 1024 * 1024
func NewFixSizeLargeBuff() *FixSizeLargeBuff {
return &FixSizeLargeBuff{buf: make([]byte, 0, Megabit)}
}
func (f *FixSizeLargeBuff) Avail() int {
return Megabit - len(f.buf)
}
func (f *FixSizeLargeBuff) Reset() {
f.buf = f.buf[:0]
}
func (f *FixSizeLargeBuff) Append(p []byte) (int, error) {
if f.Avail() < len(p) {
return 0, fmt.Errorf("no avail free bytes")
}
f.buf = append(f.buf, p...)
return len(p), nil
}
type SimpleAsyncWriter struct {
data chan *FixSizeLargeBuff
curbuff *FixSizeLargeBuff
buffpool sync.Pool
wt io.Writer
lock sync.Mutex
wg sync.WaitGroup
ct *time.Ticker
last time.Time
active chan struct{}
}
func NewSimpleAsyncWriter(w io.Writer, limit int) *SimpleAsyncWriter {
ret := &SimpleAsyncWriter{
data: make(chan *FixSizeLargeBuff, limit),
buffpool: sync.Pool{New: func() any {
return NewFixSizeLargeBuff()
}},
wt: w,
lock: sync.Mutex{},
active: make(chan struct{}),
ct: time.NewTicker(1 * time.Second),
}
ret.addCount()
go ret.poller()
return ret
}
func (s *SimpleAsyncWriter) addCount() {
s.wg.Add(1)
}
func (s *SimpleAsyncWriter) Write(p []byte) (int, error) {
select {
case <-s.active:
return 0, ErrorWriteAsyncerIsClosed
default:
}
s.last = time.Now()
s.lock.Lock()
defer s.lock.Unlock()
select {
case <-s.active:
return 0, ErrorWriteAsyncerIsClosed
case <-s.ct.C:
if s.curbuff.Avail() > 0 && time.Since(s.last) > 5*time.Second {
s.data <- s.curbuff
s.curbuff = s.buffpool.Get().(*FixSizeLargeBuff)
}
default:
if s.curbuff == nil {
s.curbuff = s.buffpool.Get().(*FixSizeLargeBuff)
}
if len(p) > s.curbuff.Avail() {
s.data <- s.curbuff
s.curbuff = s.buffpool.Get().(*FixSizeLargeBuff)
}
}
if n, err := s.curbuff.Append(p); err != nil {
return n, err
}
return len(p), nil
}
func (s *SimpleAsyncWriter) poller() {
defer func() {
for i := len(s.data); i > 0; i-- {
d := <-s.data
s.wt.Write(d.buf)
}
if s.curbuff.Avail() > 0 {
s.wt.Write(s.curbuff.buf)
}
close(s.data)
s.data = nil
s.ct.Stop()
s.wg.Done()
}()
for {
select {
case <-s.active:
goto outer
case d := <-s.data:
s.wt.Write(d.buf)
d.Reset()
s.buffpool.Put(d)
}
}
outer:
}
func (s *SimpleAsyncWriter) Stop() {
s.active <- struct{}{}
s.wg.Wait()
}
|