-
Notifications
You must be signed in to change notification settings - Fork 4
Expand file tree
/
Copy pathbatch_error.go
More file actions
124 lines (107 loc) · 2.68 KB
/
Copy pathbatch_error.go
File metadata and controls
124 lines (107 loc) · 2.68 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
package treblle
import (
"sync"
"time"
)
// BatchErrorCollector handles batch collection and transmission of errors
type BatchErrorCollector struct {
mu sync.Mutex
errors []ErrorInfo
batchSize int
flushInterval time.Duration
done chan struct{}
wg sync.WaitGroup
}
// NewBatchErrorCollector creates a new BatchErrorCollector with specified batch size and flush interval
func NewBatchErrorCollector(batchSize int, flushInterval time.Duration) *BatchErrorCollector {
if batchSize <= 0 {
batchSize = 100 // default batch size
}
if flushInterval <= 0 {
flushInterval = 5 * time.Second // default flush interval
}
collector := &BatchErrorCollector{
errors: make([]ErrorInfo, 0, batchSize),
batchSize: batchSize,
flushInterval: flushInterval,
done: make(chan struct{}),
}
go collector.periodicFlush()
return collector
}
// Add adds an error to the batch
func (b *BatchErrorCollector) Add(err ErrorInfo) {
b.mu.Lock()
defer b.mu.Unlock()
b.errors = append(b.errors, err)
if len(b.errors) >= b.batchSize {
b.flush()
}
}
// flush sends the current batch of errors to Treblle
func (b *BatchErrorCollector) flush() {
if len(b.errors) == 0 {
return
}
// Create a copy of errors to send
errorsCopy := make([]ErrorInfo, len(b.errors))
copy(errorsCopy, b.errors)
// Clear the current batch
b.errors = b.errors[:0]
// Send errors asynchronously
b.wg.Add(1)
go func(errors []ErrorInfo) {
defer b.wg.Done()
// Create metadata for batch transmission
meta := MetaData{
ApiKey: Config.APIKey,
ProjectID: Config.ProjectID,
Version: Config.SDKVersion,
Sdk: Config.SDKName,
Data: DataInfo{
Server: Config.serverInfo,
Language: Config.languageInfo,
Request: RequestInfo{}, // Empty request info for batch errors
Response: ResponseInfo{}, // Empty response info for batch errors
Errors: errors,
},
}
// Send to Treblle
sendToTreblle(meta)
}(errorsCopy)
}
// periodicFlush periodically flushes the error batch based on the flush interval
func (b *BatchErrorCollector) periodicFlush() {
ticker := time.NewTicker(b.flushInterval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
b.mu.Lock()
b.flush()
b.mu.Unlock()
case <-b.done:
return
}
}
}
// Close stops the periodic flushing and flushes any remaining errors
func (b *BatchErrorCollector) Close() {
b.mu.Lock()
defer b.mu.Unlock()
select {
case <-b.done:
// Channel already closed
return
default:
close(b.done)
b.flush()
b.wg.Wait()
}
}
// Flush sends any pending errors to Treblle immediately
func (b *BatchErrorCollector) Flush() {
b.mu.Lock()
defer b.mu.Unlock()
b.flush()
}