1 | 'use strict';
|
2 |
|
3 |
|
4 |
|
5 |
|
6 |
|
7 | const _ = require('underscore');
|
8 | const events = require('abacus-events');
|
9 |
|
10 | const keys = _.keys;
|
11 | const filter = _.filter;
|
12 | const last = _.last;
|
13 | const extend = _.extend;
|
14 | const reduce = _.reduce;
|
15 | const map = _.map;
|
16 | const find = _.find;
|
17 |
|
18 |
|
19 | const debug = require('abacus-debug')('abacus-perf');
|
20 |
|
21 |
|
22 |
|
23 | const emitter = events.emitter('call-metrics/emitter');
|
24 | const on = (e, l) => {
|
25 | emitter.on(e, l);
|
26 | };
|
27 |
|
28 |
|
29 | const calls = events.emitter('abacus-perf/calls');
|
30 |
|
31 |
|
32 | const report = (name, time, latency, err, timeout, reject, circuit) => {
|
33 | const t = time || Date.now();
|
34 | const call = () => ({
|
35 | name: name,
|
36 | time: t,
|
37 | latency: latency === undefined ? Date.now() - t : latency,
|
38 | error: err || undefined,
|
39 | timeout: timeout || 0,
|
40 | reject: reject || false,
|
41 | circuit: circuit || 'closed'
|
42 | });
|
43 | calls.emit('message', {
|
44 | metrics: {
|
45 | call: call()
|
46 | }
|
47 | });
|
48 | };
|
49 |
|
50 |
|
51 |
|
52 |
|
53 | const roll = (buckets, time, win, max, proto) => {
|
54 | const i = Math.ceil(time / (win / max));
|
55 | const current = filter(buckets, (b) => b.i > i - max);
|
56 | const all = current.length !== 0 && last(current).i === i ? current :
|
57 | current.concat([extend({}, proto(), {
|
58 | i: i
|
59 | })]);
|
60 | return last(all, max);
|
61 | };
|
62 |
|
63 |
|
64 | const rollCounts = (buckets, time) => {
|
65 | return roll(buckets, time, 10000, 10, () => ({
|
66 | i: 0,
|
67 | ok: 0,
|
68 | errors: 0,
|
69 | timeouts: 0,
|
70 | rejects: 0
|
71 | }));
|
72 | };
|
73 |
|
74 |
|
75 | const rollLatencies = (buckets, time) => {
|
76 | return roll(buckets, time, 60000, 6, () => ({
|
77 | i: 0,
|
78 | latencies: []
|
79 | }));
|
80 | };
|
81 |
|
82 |
|
83 | const rollHealth = (buckets, time, counts) => {
|
84 |
|
85 | const b = reduce(counts, (a, c) => ({
|
86 | ok: a.ok + c.ok,
|
87 | errors: a.errors + c.errors + c.timeouts + c.rejects
|
88 | }), {
|
89 | ok: 0,
|
90 | errors: 0
|
91 | });
|
92 |
|
93 |
|
94 | const health = roll(buckets, time, 2000, 4, () => ({
|
95 | i: 0,
|
96 | ok: 0,
|
97 | errors: 0
|
98 | }));
|
99 | extend(last(health), b);
|
100 | return health;
|
101 | };
|
102 |
|
103 |
|
104 |
|
105 | let accumulatedStats = {};
|
106 |
|
107 |
|
108 | const stats = (name, time, roll) => {
|
109 | const t = time || Date.now();
|
110 | const s = accumulatedStats[name] || {
|
111 | name: name,
|
112 | time: t,
|
113 | counts: [],
|
114 | latencies: [],
|
115 | health: [],
|
116 | circuit: 'closed'
|
117 | };
|
118 |
|
119 | return roll !== false ? (() => {
|
120 | const counts = rollCounts(s.counts, t);
|
121 | return {
|
122 | name: name,
|
123 | time: t,
|
124 | counts: counts,
|
125 | latencies: rollLatencies(s.latencies, t),
|
126 | health: rollHealth(s.health, t, counts),
|
127 | circuit: s.circuit
|
128 | };
|
129 | })() : {
|
130 | name: name,
|
131 | time: t,
|
132 | counts: s.counts,
|
133 | latencies: s.latencies,
|
134 | health: s.health,
|
135 | circuit: s.circuit
|
136 | };
|
137 | };
|
138 |
|
139 |
|
140 |
|
141 | const reset = (name, time) => {
|
142 | const t = time || Date.now();
|
143 |
|
144 |
|
145 |
|
146 | const astats = ((s) => {
|
147 | return s ? {
|
148 | name: name,
|
149 | time: t,
|
150 | counts: [],
|
151 | latencies: s.latencies,
|
152 | health: [],
|
153 | circuit: 'closed'
|
154 | } :
|
155 | {
|
156 | name: name,
|
157 | time: t,
|
158 | counts: [],
|
159 | latencies: [],
|
160 | health: [],
|
161 | circuit: 'closed'
|
162 | };
|
163 | })(accumulatedStats[name]);
|
164 | accumulatedStats[name] = astats;
|
165 |
|
166 |
|
167 | debug('Emitting stats for function %s', name);
|
168 | emitter.emit('message', {
|
169 | metrics: {
|
170 | stats: astats
|
171 | }
|
172 | });
|
173 | return astats;
|
174 | };
|
175 |
|
176 |
|
177 | const all = (time, roll) => {
|
178 | const t = time || Date.now();
|
179 | return map(keys(accumulatedStats), (k) => stats(k, t, roll));
|
180 | };
|
181 |
|
182 |
|
183 | const accumulateStats = (name, time, latency, err, timeout,
|
184 | reject, circuit) => {
|
185 | debug('Accumulating stats for function %s', name);
|
186 | debug('latency %d, err %s, timeout %d, reject %s, circuit %s',
|
187 | latency, err, timeout, reject, circuit);
|
188 |
|
189 |
|
190 | const s = stats(name, time, false);
|
191 |
|
192 |
|
193 | const counts = rollCounts(s ? s.counts : [], time);
|
194 | const updateCount = (c) => {
|
195 | c.ok = c.ok + (!err && !timeout && !reject ? 1 : 0);
|
196 | c.errors = c.errors + (err ? 1 : 0);
|
197 | c.timeouts = c.timeouts + (timeout ? 1 : 0);
|
198 | c.rejects = c.rejects + (reject ? 1 : 0);
|
199 | debug('%d ok, %d errors, %d timeouts, %d rejects, %d count buckets',
|
200 | c.ok, c.errors, c.timeouts, c.rejects, counts.length);
|
201 | };
|
202 | updateCount(last(counts));
|
203 |
|
204 |
|
205 |
|
206 | const latencies = rollLatencies(s ? s.latencies : [], time);
|
207 | if(!err && !timeout && !reject) {
|
208 | const updateLatency = (l) => {
|
209 | l.latencies = l.latencies.length < 100 ?
|
210 | l.latencies.concat([latency]) : l.latencies;
|
211 | debug('%d latencies, %d latencies buckets',
|
212 | l.latencies.length, latencies .length);
|
213 | };
|
214 | updateLatency(last(latencies));
|
215 | }
|
216 |
|
217 |
|
218 | const health = rollHealth(s ? s.health : [], time, counts);
|
219 | const h = last(health);
|
220 | debug('%d ok, %d errors, %d health buckets', h.ok, h.errors, health.length);
|
221 |
|
222 |
|
223 |
|
224 | const astats = {
|
225 | name: name,
|
226 | time: time,
|
227 | counts: counts,
|
228 | latencies: latencies,
|
229 | health: health,
|
230 | circuit: circuit
|
231 | };
|
232 | accumulatedStats[name] = astats;
|
233 |
|
234 |
|
235 | debug('Emitting stats for function %s', name);
|
236 | emitter.emit('message', {
|
237 | metrics: {
|
238 | stats: astats
|
239 | }
|
240 | });
|
241 | return astats;
|
242 | };
|
243 |
|
244 |
|
245 | const onMessage = (msg) => {
|
246 | if(msg.metrics) {
|
247 | debug('Received message %o', keys(msg).concat(keys(msg.metrics)));
|
248 | if(msg.metrics.call) {
|
249 |
|
250 | const c = msg.metrics.call;
|
251 | accumulateStats(
|
252 | c.name, c.time, c.latency, c.error, c.timeout, c.reject, c.circuit);
|
253 | }
|
254 | if(msg.metrics.stats) {
|
255 |
|
256 | debug('Storing stats for function %s', msg.metrics.stats.name);
|
257 | accumulatedStats[msg.metrics.stats.name] = msg.metrics.stats;
|
258 | }
|
259 | }
|
260 | };
|
261 |
|
262 |
|
263 | const healthy = (threshold) => {
|
264 |
|
265 |
|
266 | return find(all(Date.now(), false), (stat) => {
|
267 |
|
268 |
|
269 | const total = reduce(stat.health, (a, c) => ({
|
270 | requests: a.requests + a.errors + c.ok + c.errors,
|
271 | errors: a.errors + c.errors
|
272 | }), {
|
273 | requests: 0,
|
274 | errors: 0
|
275 | });
|
276 |
|
277 | const percent = 100 * (total.errors / total.requests || 1);
|
278 | debug('%s has %d% failure rate', stat.name, percent);
|
279 |
|
280 |
|
281 |
|
282 | return percent > (threshold || 5);
|
283 |
|
284 | }) ? false : true;
|
285 | }
|
286 |
|
287 | calls.on('message', onMessage);
|
288 |
|
289 |
|
290 | module.exports.report = report;
|
291 | module.exports.stats = stats;
|
292 | module.exports.reset = reset;
|
293 | module.exports.healthy = healthy;
|
294 | module.exports.all = all;
|
295 | module.exports.onMessage = onMessage;
|
296 | module.exports.on = on;
|