-
Notifications
You must be signed in to change notification settings - Fork 0
/
writer.go
124 lines (94 loc) · 2.26 KB
/
writer.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
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
package lhw
import (
"errors"
"fmt"
"strings"
"sync"
"github.com/loghole/lhw/internal"
"github.com/loghole/lhw/transport"
)
var ErrWriteFailed = errors.New("[loghole-writer] write data to queue failed")
// The url can contain secret token e.g. https://secret_token@localhost:50000
// Comma separated arrays are also supported, e.g. urlA, urlB.
// Options start with the defaults but can be overridden.
func NewWriter(url string, options ...Option) (writer *Writer, err error) {
opts := GetDefaultOptions()
opts.Servers = processURLString(url)
for _, option := range options {
if option == nil {
continue
}
if err := option(opts); err != nil {
return nil, err
}
}
writer = &Writer{
logger: opts.Logger,
queue: internal.NewQueue(opts.QueueCap),
}
writer.transport, err = transport.New(opts.transportConfig())
if err != nil {
return nil, err
}
writer.wg.Add(1)
go writer.worker()
return writer, nil
}
type Writer struct {
transport transport.Transport
queue *internal.Queue
logger Logger
wg sync.WaitGroup
}
// Write writes the data to the queue if it is not full.
func (w *Writer) Write(p []byte) (n int, err error) {
return w.write(append([]byte{}, p...))
}
// write writes the data to the queue if it is not full.
func (w *Writer) write(p []byte) (n int, err error) {
if err := w.queue.Push(p); err != nil {
return 0, fmt.Errorf("%w: %v", ErrWriteFailed, err)
}
return len(p), nil
}
// Close flushes any buffered log entries.
func (w *Writer) Close() error {
w.queue.Close()
w.wg.Wait()
return nil
}
func (w *Writer) worker() {
defer w.wg.Done()
for data := range w.queue.Read() {
if !w.transport.IsConnected() {
<-w.transport.IsReconnected()
}
w.wg.Add(1)
go w.send(data)
}
}
func (w *Writer) send(data []byte) {
defer w.wg.Done()
err := w.transport.Send(data)
if err == nil {
return
}
if w.logger != nil {
w.logger.Printf("[error] send data failed: %v", err)
}
// if sending failed, return data to queue if it is not full.
_, err = w.write(data)
if err == nil {
return
}
if w.logger != nil {
w.logger.Printf("[error] %v", err)
}
}
func processURLString(url string) []string {
urls := strings.Split(url, ",")
for idx, val := range urls {
urls[idx] = strings.TrimSpace(val)
}
return urls
}