Skip to content

Commit e0f40b1

Browse files
committed
Change buffer to be managed by another structure
1 parent 38301f4 commit e0f40b1

2 files changed

Lines changed: 68 additions & 32 deletions

File tree

buffer.go

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
1+
package fluent
2+
3+
import (
4+
"sync"
5+
)
6+
7+
type buffer struct {
8+
buf []byte
9+
mu sync.Mutex
10+
Dirty chan struct{}
11+
}
12+
13+
func newBuffer() buffer {
14+
return buffer{
15+
buf: []byte{},
16+
Dirty: make(chan struct{}),
17+
}
18+
}
19+
20+
func (buffer *buffer) Add(raw []byte) {
21+
buffer.mu.Lock()
22+
defer buffer.mu.Unlock()
23+
24+
buffer.buf = append(buffer.buf, raw...)
25+
go func() {
26+
buffer.Dirty <- struct{}{}
27+
}()
28+
}
29+
30+
func (buffer *buffer) Remove() []byte {
31+
buffer.mu.Lock()
32+
defer buffer.mu.Unlock()
33+
34+
if len(buffer.buf) == 0 {
35+
return nil
36+
}
37+
38+
data := buffer.buf
39+
buffer.buf = buffer.buf[:0]
40+
return data
41+
}
42+
43+
func (buffer *buffer) Back(raw []byte) {
44+
buffer.mu.Lock()
45+
defer buffer.mu.Unlock()
46+
47+
buffer.buf = append(buffer.buf, raw...)
48+
}

logger.go

Lines changed: 20 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -34,22 +34,19 @@ func withDefaultConfig(c Config) Config {
3434
}
3535

3636
type Logger struct {
37-
conf Config
38-
conn io.WriteCloser
39-
bmu sync.Mutex
40-
cmu sync.Mutex
41-
buf []byte
42-
wg sync.WaitGroup
43-
done chan struct{}
44-
dirty chan struct{}
37+
conf Config
38+
conn io.WriteCloser
39+
mu sync.Mutex
40+
buf buffer
41+
wg sync.WaitGroup
42+
done chan struct{}
4543
}
4644

4745
func NewLogger(c Config) (*Logger, error) {
4846
logger := &Logger{
49-
conf: withDefaultConfig(c),
50-
buf: []byte{},
51-
done: make(chan struct{}),
52-
dirty: make(chan struct{}),
47+
conf: withDefaultConfig(c),
48+
buf: newBuffer(),
49+
done: make(chan struct{}),
5350
}
5451
if err := logger.connect(); err != nil {
5552
return nil, err
@@ -74,15 +71,7 @@ func (logger *Logger) PostWithTime(tag string, t time.Time, obj interface{}) err
7471
if err := enc.Encode(record); err != nil {
7572
return err
7673
}
77-
raw := buf.Bytes()
78-
79-
logger.bmu.Lock()
80-
logger.buf = append(logger.buf, raw...)
81-
logger.bmu.Unlock()
82-
83-
go func() {
84-
logger.dirty <- struct{}{}
85-
}()
74+
logger.buf.Add(buf.Bytes())
8675
return nil
8776
}
8877

@@ -92,8 +81,8 @@ func (logger *Logger) Close() error {
9281
}
9382

9483
func (logger *Logger) connect() error {
95-
logger.cmu.Lock()
96-
defer logger.cmu.Unlock()
84+
logger.mu.Lock()
85+
defer logger.mu.Unlock()
9786

9887
if logger.conn != nil {
9988
return nil
@@ -116,8 +105,8 @@ func (logger *Logger) connect() error {
116105
}
117106

118107
func (logger *Logger) disconnect() error {
119-
logger.cmu.Lock()
120-
defer logger.cmu.Unlock()
108+
logger.mu.Lock()
109+
defer logger.mu.Unlock()
121110

122111
if logger.conn == nil {
123112
return nil
@@ -130,26 +119,25 @@ func (logger *Logger) disconnect() error {
130119
const maxWriteAttempts = 3
131120

132121
func (logger *Logger) send() error {
133-
logger.bmu.Lock()
134-
defer logger.bmu.Unlock()
135-
136-
data := logger.buf
122+
data := logger.buf.Remove()
137123
if len(data) == 0 {
138124
return nil
139125
}
126+
140127
var err error
141128
for i := 0; i < maxWriteAttempts; i++ {
142129
err = logger.connect()
143130
if err == nil {
144131
_, err := logger.conn.Write(data)
145132
if err == nil {
146-
logger.buf = logger.buf[:0]
147133
break
148134
}
149135
}
150136
logger.disconnect()
151137
}
152-
138+
if err != nil {
139+
logger.buf.Back(data)
140+
}
153141
return err
154142
}
155143

@@ -163,7 +151,7 @@ func (logger *Logger) start() {
163151
case <-logger.done:
164152
logger.send()
165153
return
166-
case <-logger.dirty:
154+
case <-logger.buf.Dirty:
167155
case <-ticker.C:
168156
}
169157
logger.send()

0 commit comments

Comments
 (0)