1 | const got = require('got')
|
2 | const url = require('url')
|
3 | const exitHook = require('async-exit-hook')
|
4 |
|
5 | const { logproto } = require('./proto')
|
6 | const protoHelpers = require('./proto/helpers')
|
7 | let snappy = false
|
8 |
|
9 |
|
10 |
|
11 |
|
12 |
|
13 |
|
14 | class Batcher {
|
15 | loadSnappy () {
|
16 | return require('snappy')
|
17 | }
|
18 | |
19 |
|
20 |
|
21 |
|
22 |
|
23 |
|
24 | constructor (options) {
|
25 |
|
26 | this.options = options
|
27 |
|
28 |
|
29 | this.url = new url.URL(this.options.host + '/api/prom/push').toString()
|
30 |
|
31 |
|
32 | this.interval = this.options.interval
|
33 | ? Number(this.options.interval) * 1000
|
34 | : 5000
|
35 | this.circuitBreakerInterval = 60000
|
36 |
|
37 |
|
38 | this.batch = {
|
39 | streams: []
|
40 | }
|
41 |
|
42 |
|
43 | if (!this.options.json) {
|
44 | try {
|
45 | snappy = this.loadSnappy()
|
46 | } catch (error) {
|
47 | this.options.json = true
|
48 | }
|
49 | if (!snappy) {
|
50 | this.options.json = true
|
51 | }
|
52 | }
|
53 |
|
54 |
|
55 | this.contentType = 'application/x-protobuf'
|
56 | if (this.options.json) {
|
57 | this.contentType = 'application/json'
|
58 | }
|
59 |
|
60 |
|
61 | this.options.batching && this.run()
|
62 |
|
63 | exitHook(callback => {
|
64 | this.sendBatchToLoki()
|
65 | .then(() => callback())
|
66 | .catch(() => callback())
|
67 | })
|
68 | }
|
69 |
|
70 | |
71 |
|
72 |
|
73 |
|
74 |
|
75 |
|
76 | wait (duration) {
|
77 | return new Promise(resolve => {
|
78 | setTimeout(resolve, duration)
|
79 | })
|
80 | }
|
81 |
|
82 | |
83 |
|
84 |
|
85 |
|
86 |
|
87 |
|
88 | async pushLogEntry (logEntry) {
|
89 | const noTimestamp =
|
90 | logEntry && logEntry.entries && logEntry.entries[0].ts === undefined
|
91 |
|
92 | if (this.options.replaceTimestamp || noTimestamp) {
|
93 | logEntry.entries[0].ts = Date.now()
|
94 | }
|
95 |
|
96 |
|
97 | if (!this.options.json) {
|
98 | logEntry = protoHelpers.createProtoTimestamps(logEntry)
|
99 | }
|
100 |
|
101 |
|
102 | if (this.options.batching !== undefined && !this.options.batching) {
|
103 | await this.sendBatchToLoki(logEntry)
|
104 | } else {
|
105 | const { streams } = this.batch
|
106 |
|
107 |
|
108 | const match = streams.findIndex(
|
109 | stream => stream.labels === logEntry.labels
|
110 | )
|
111 |
|
112 | if (match > -1) {
|
113 |
|
114 | logEntry.entries.forEach(entry => {
|
115 | streams[match].entries.push(entry)
|
116 | })
|
117 | } else {
|
118 |
|
119 | streams.push(logEntry)
|
120 | }
|
121 | }
|
122 | }
|
123 |
|
124 | |
125 |
|
126 |
|
127 | clearBatch () {
|
128 | this.batch.streams = []
|
129 | }
|
130 |
|
131 | |
132 |
|
133 |
|
134 |
|
135 |
|
136 |
|
137 |
|
138 | sendBatchToLoki (logEntry) {
|
139 | return new Promise((resolve, reject) => {
|
140 |
|
141 | if (this.batch.streams.length === 0 && !logEntry) {
|
142 | resolve()
|
143 | } else {
|
144 | let reqBody
|
145 |
|
146 |
|
147 | if (this.options.json) {
|
148 | if (logEntry !== undefined) {
|
149 |
|
150 | reqBody = JSON.stringify({ streams: [logEntry] })
|
151 | } else {
|
152 |
|
153 | reqBody = JSON.stringify(reqBody)
|
154 | }
|
155 | } else {
|
156 | try {
|
157 | let batch
|
158 | if (logEntry !== undefined) {
|
159 |
|
160 | batch = { streams: [logEntry] }
|
161 | } else {
|
162 | batch = this.batch
|
163 | }
|
164 |
|
165 |
|
166 | const err = logproto.PushRequest.verify(batch)
|
167 |
|
168 |
|
169 | if (err) reject(err)
|
170 |
|
171 |
|
172 | const message = logproto.PushRequest.create(batch)
|
173 |
|
174 |
|
175 | const buffer = logproto.PushRequest.encode(message).finish()
|
176 |
|
177 |
|
178 | reqBody = snappy.compressSync(buffer)
|
179 | } catch (err) {
|
180 | reject(err)
|
181 | }
|
182 | }
|
183 |
|
184 |
|
185 | got
|
186 | .post(this.url, {
|
187 | body: reqBody,
|
188 | headers: {
|
189 | 'content-type': this.contentType
|
190 | }
|
191 | })
|
192 | .then(res => {
|
193 |
|
194 | logEntry === undefined && this.clearBatch()
|
195 | resolve()
|
196 | })
|
197 | .catch(err => {
|
198 |
|
199 | this.options.clearOnError && this.clearBatch()
|
200 | reject(err)
|
201 | })
|
202 | }
|
203 | })
|
204 | }
|
205 |
|
206 | |
207 |
|
208 |
|
209 |
|
210 |
|
211 |
|
212 | async run () {
|
213 | while (true) {
|
214 | try {
|
215 | await this.sendBatchToLoki()
|
216 | if (this.interval === this.circuitBreakerInterval) {
|
217 | this.interval = Number(this.options.interval) * 1000
|
218 | }
|
219 | } catch (e) {
|
220 | this.interval = this.circuitBreakerInterval
|
221 | }
|
222 | await this.wait(this.interval)
|
223 | }
|
224 | }
|
225 | }
|
226 |
|
227 | module.exports = Batcher
|