/src/freeradius-server/src/lib/io/load.c
Line | Count | Source |
1 | | /* |
2 | | * This program is free software; you can redistribute it and/or modify |
3 | | * it under the terms of the GNU General Public License as published by |
4 | | * the Free Software Foundation; either version 2 of the License, or |
5 | | * (at your option) any later version. |
6 | | * |
7 | | * This program is distributed in the hope that it will be useful, |
8 | | * but WITHOUT ANY WARRANTY; without even the implied warranty of |
9 | | * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the |
10 | | * GNU General Public License for more details. |
11 | | * |
12 | | * You should have received a copy of the GNU General Public License |
13 | | * along with this program; if not, write to the Free Software |
14 | | * Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301, USA |
15 | | */ |
16 | | |
17 | | /** |
18 | | * $Id: e6ece0589d65604b550f3474a69f3bcadecbe171 $ |
19 | | * |
20 | | * @brief Load generation algorithms |
21 | | * @file io/load.c |
22 | | * |
23 | | * @copyright 2019 Network RADIUS SAS (legal@networkradius.com) |
24 | | */ |
25 | | RCSID("$Id: e6ece0589d65604b550f3474a69f3bcadecbe171 $") |
26 | | |
27 | | #include <freeradius-devel/io/load.h> |
28 | | |
29 | | /* |
30 | | * We use *inverse* numbers to avoid numerical calculation issues. |
31 | | * |
32 | | * i.e. The bad way is to take two small numbers divide them by |
33 | | * alpha / beta and then add them. That process can drop the |
34 | | * lower digits. Instead, we take two small numbers, add them, |
35 | | * and then divide the result by alpha / beta. |
36 | | */ |
37 | 0 | #define IBETA (4) |
38 | | #define IALPHA (8) |
39 | | |
40 | | #define DIFF(_rtt, _t) \ |
41 | 0 | (\ |
42 | 0 | fr_time_delta_lt(_rtt, _t) ? \ |
43 | 0 | fr_time_delta_sub(_t, _rtt) : \ |
44 | 0 | fr_time_delta_sub(_rtt, _t)\ |
45 | 0 | ) |
46 | | |
47 | | #define RTTVAR(_rtt, _rttvar, _t) \ |
48 | 0 | fr_time_delta_div(\ |
49 | 0 | fr_time_delta_add(\ |
50 | 0 | fr_time_delta_mul(_rttvar, IBETA - 1), \ |
51 | 0 | DIFF(_rtt, _t)\ |
52 | 0 | ), \ |
53 | 0 | fr_time_delta_wrap(IBETA)\ |
54 | 0 | ) |
55 | | |
56 | 0 | #define RTT(_old, _new) fr_time_delta_wrap((fr_time_delta_unwrap(_new) + (fr_time_delta_unwrap(_old) * (IALPHA - 1))) / IALPHA) |
57 | | |
58 | | typedef enum { |
59 | | FR_LOAD_STATE_INIT = 0, |
60 | | FR_LOAD_STATE_SENDING, |
61 | | FR_LOAD_STATE_GATED, |
62 | | FR_LOAD_STATE_DRAINING, |
63 | | } fr_load_state_t; |
64 | | |
65 | | struct fr_load_s { |
66 | | fr_load_state_t state; |
67 | | fr_event_list_t *el; |
68 | | fr_load_config_t const *config; |
69 | | fr_load_callback_t callback; |
70 | | fr_load_done_callback_t done; |
71 | | void *uctx; |
72 | | |
73 | | fr_load_stats_t stats; //!< sending statistics |
74 | | fr_time_t step_start; //!< when the current step started |
75 | | fr_time_t step_end; //!< when the current step will end |
76 | | int step_received; |
77 | | int sent_base; //!< stats.sent when this run started, so that |
78 | | ///< max_requests counts per-run, not cumulatively |
79 | | |
80 | | uint32_t pps; |
81 | | fr_time_delta_t delta; //!< between packets |
82 | | |
83 | | uint32_t count; |
84 | | bool header; //!< for printing statistics |
85 | | |
86 | | fr_time_t next; //!< The next time we're supposed to send a packet |
87 | | fr_timer_t *ev; |
88 | | }; |
89 | | |
90 | | fr_load_t *fr_load_generator_create(TALLOC_CTX *ctx, fr_event_list_t *el, fr_load_config_t *config, |
91 | | fr_load_callback_t callback, fr_load_done_callback_t done, |
92 | | void *uctx) |
93 | 0 | { |
94 | 0 | fr_load_t *l; |
95 | |
|
96 | 0 | l = talloc_zero(ctx, fr_load_t); |
97 | 0 | if (!l) return NULL; |
98 | | |
99 | 0 | if (!config->start_pps) config->start_pps = 1; |
100 | 0 | if (!config->milliseconds) config->milliseconds = 1000; |
101 | 0 | if (!config->parallel) config->parallel = 1; |
102 | |
|
103 | 0 | l->el = el; |
104 | 0 | l->config = config; |
105 | 0 | l->callback = callback; |
106 | 0 | l->done = done; |
107 | 0 | l->uctx = uctx; |
108 | |
|
109 | 0 | return l; |
110 | 0 | } |
111 | | |
112 | | /** Send one or more packets. |
113 | | * |
114 | | */ |
115 | | static void fr_load_generator_send(fr_load_t *l, fr_time_t now, int count) |
116 | 0 | { |
117 | 0 | int i; |
118 | | |
119 | | /* |
120 | | * Send as many packets as necessary. |
121 | | */ |
122 | 0 | l->stats.sent += count; |
123 | 0 | l->stats.last_send = now; |
124 | | |
125 | | /* |
126 | | * Run the callback AFTER we set the timer. Which makes |
127 | | * it more likely that the next timer fires on time. |
128 | | */ |
129 | 0 | for (i = 0; i < count; i++) { |
130 | 0 | l->callback(fr_time_add(now, fr_time_delta_from_nsec(i)), l->uctx); |
131 | 0 | } |
132 | 0 | } |
133 | | |
134 | | /** Stop sending new packets, and wait for outstanding replies. |
135 | | * |
136 | | * If every reply has already arrived then the test is complete right |
137 | | * now, so run the "done" callback. Otherwise completion is signalled |
138 | | * by fr_load_generator_have_reply() when the last reply comes in. |
139 | | */ |
140 | | static void load_drain(fr_load_t *l, fr_time_t now) |
141 | 0 | { |
142 | 0 | l->state = FR_LOAD_STATE_DRAINING; |
143 | |
|
144 | 0 | if (l->stats.received >= l->stats.sent) { |
145 | 0 | l->stats.end = now; |
146 | 0 | if (l->done) l->done(l->uctx); |
147 | 0 | } |
148 | 0 | } |
149 | | |
150 | | static void load_timer(fr_timer_list_t *tl, fr_time_t now, void *uctx) |
151 | 0 | { |
152 | 0 | fr_load_t *l = uctx; |
153 | 0 | fr_time_delta_t delta; |
154 | 0 | uint32_t count; |
155 | | |
156 | | /* |
157 | | * Keep track of the overall maximum backlog for the |
158 | | * duration of the entire test run. |
159 | | */ |
160 | 0 | l->stats.backlog = l->stats.sent - l->stats.received; |
161 | 0 | if (l->stats.backlog > l->stats.max_backlog) l->stats.max_backlog = l->stats.backlog; |
162 | | |
163 | | /* |
164 | | * If we're done this step, go to the next one. |
165 | | */ |
166 | 0 | if (fr_time_gteq(l->next, l->step_end)) { |
167 | 0 | l->step_start = l->next; |
168 | 0 | l->step_end = fr_time_add(l->next, l->config->duration); |
169 | 0 | l->step_received = l->stats.received; |
170 | 0 | l->pps += l->config->step; |
171 | 0 | if (l->pps > (UINT32_MAX / l->config->milliseconds)) l->pps = UINT32_MAX / l->config->milliseconds; |
172 | | |
173 | | /* |
174 | | * The ramp-up has passed max_pps. If there's no |
175 | | * reason to keep sending, stop and drain. |
176 | | * Otherwise hold the rate at max_pps, until |
177 | | * max_requests packets have been sent (checked |
178 | | * below), or forever for "unlimited". |
179 | | */ |
180 | 0 | if (l->config->max_pps && (l->pps > l->config->max_pps)) { |
181 | 0 | if (!l->config->max_requests && !l->config->unlimited) { |
182 | 0 | load_drain(l, now); |
183 | 0 | return; |
184 | 0 | } |
185 | | |
186 | 0 | l->pps = l->config->max_pps; |
187 | 0 | } |
188 | | |
189 | 0 | l->stats.pps = l->pps; |
190 | 0 | l->stats.skipped = 0; |
191 | 0 | l->delta = fr_time_delta_div(fr_time_delta_from_sec(l->config->parallel), fr_time_delta_wrap(l->pps)); |
192 | 0 | } |
193 | | |
194 | | /* |
195 | | * We don't have "pps" packets in the backlog, go send |
196 | | * some more. We scale the backlog by 1000 milliseconds |
197 | | * per second. Then, multiple the PPS by the number of |
198 | | * milliseconds of backlog we want to keep. |
199 | | * |
200 | | * If the backlog is smaller than packets/s * |
201 | | * milliseconds of backlog, then keep sending. |
202 | | * Otherwise, switch to a gated mode where we only send |
203 | | * new packets once a reply comes in. |
204 | | */ |
205 | 0 | if (((size_t) l->stats.backlog * 1000) < ((size_t) l->pps * l->config->milliseconds)) { |
206 | 0 | uint32_t capacity; |
207 | |
|
208 | 0 | l->state = FR_LOAD_STATE_SENDING; |
209 | 0 | l->stats.blocked = false; |
210 | 0 | count = l->config->parallel; |
211 | 0 | l->stats.skipped = 0; |
212 | |
|
213 | 0 | capacity = ((l->pps * l->config->milliseconds) / 1000) - l->stats.backlog; |
214 | | |
215 | | /* |
216 | | * Limit "count" so that it doesn't overflow. |
217 | | */ |
218 | 0 | if (count > capacity) count = capacity; |
219 | |
|
220 | 0 | } else { |
221 | | |
222 | | /* |
223 | | * We have too many packets in the backlog, we're |
224 | | * gated. Don't send more packets until we have |
225 | | * a reply. |
226 | | * |
227 | | * Note that we will send *these* packets. |
228 | | */ |
229 | 0 | l->state = FR_LOAD_STATE_GATED; |
230 | 0 | l->stats.blocked = true; |
231 | 0 | count = 0; |
232 | 0 | l->stats.skipped += l->count; |
233 | 0 | } |
234 | | |
235 | | /* |
236 | | * Stop after max_requests packets, if it's set. Packets |
237 | | * sent by a previous run of the generator don't count. |
238 | | */ |
239 | 0 | if (l->config->max_requests) { |
240 | 0 | uint32_t sent = (uint32_t) (l->stats.sent - l->sent_base); |
241 | |
|
242 | 0 | if (sent >= l->config->max_requests) { |
243 | 0 | load_drain(l, now); |
244 | 0 | return; |
245 | 0 | } |
246 | | |
247 | 0 | if (count > (l->config->max_requests - sent)) count = l->config->max_requests - sent; |
248 | 0 | } |
249 | | |
250 | | /* |
251 | | * Skip timers if we're too busy. |
252 | | */ |
253 | 0 | l->next = fr_time_add(l->next, l->delta); |
254 | 0 | if (fr_time_lt(l->next, now)) { |
255 | 0 | while (fr_time_lt(l->next, now)) { |
256 | | // l->stats.skipped += l->count; |
257 | 0 | l->next = fr_time_add(l->next, l->delta); |
258 | 0 | } |
259 | 0 | } |
260 | 0 | delta = fr_time_sub(l->next, now); |
261 | | |
262 | | /* |
263 | | * Set the timer for the next packet. |
264 | | */ |
265 | 0 | if (fr_timer_in(l, tl, &l->ev, delta, false, load_timer, l) < 0) { |
266 | 0 | load_drain(l, now); |
267 | 0 | return; |
268 | 0 | } |
269 | | |
270 | 0 | if (count) fr_load_generator_send(l, now, count); |
271 | 0 | } |
272 | | |
273 | | |
274 | | /** Start the load generator. |
275 | | * |
276 | | */ |
277 | | int fr_load_generator_start(fr_load_t *l) |
278 | 0 | { |
279 | 0 | uint32_t max; |
280 | |
|
281 | 0 | l->stats.start = fr_time(); |
282 | 0 | l->step_start = l->stats.start; |
283 | 0 | l->step_end = fr_time_add(l->step_start, l->config->duration); |
284 | |
|
285 | 0 | l->state = FR_LOAD_STATE_SENDING; |
286 | 0 | l->sent_base = l->stats.sent; |
287 | |
|
288 | 0 | l->pps = l->config->start_pps; |
289 | | |
290 | | /* |
291 | | * Check for numerical overflow. We later multiply pps*milliseconds, and we don't want overflow. |
292 | | */ |
293 | 0 | max = UINT32_MAX / l->config->milliseconds; |
294 | |
|
295 | 0 | if (l->pps > max) l->pps = max; |
296 | |
|
297 | 0 | l->stats.pps = l->pps; |
298 | 0 | l->count = l->config->parallel; |
299 | |
|
300 | 0 | l->delta = fr_time_delta_div(fr_time_delta_from_sec(l->config->parallel), fr_time_delta_wrap(l->pps)); |
301 | 0 | l->next = fr_time_add(l->step_start, l->delta); |
302 | |
|
303 | 0 | if (fr_timer_in(l, l->el->tl, &l->ev, fr_time_delta_wrap(0), false, load_timer, l) < 0) return -1; |
304 | 0 | return 0; |
305 | 0 | } |
306 | | |
307 | | |
308 | | /** Stop the load generation through the simple expedient of deleting |
309 | | * the timer associated with it. |
310 | | * |
311 | | */ |
312 | | int fr_load_generator_stop(fr_load_t *l) |
313 | | { |
314 | | if (!fr_timer_armed(l->ev)) return 0; |
315 | | |
316 | | FR_TIMER_DELETE_RETURN(&l->ev); |
317 | | return 0; |
318 | | } |
319 | | |
320 | | |
321 | | /** Tell the load generator that we have a reply to a packet we sent. |
322 | | * |
323 | | */ |
324 | | fr_load_reply_t fr_load_generator_have_reply(fr_load_t *l, fr_time_t request_time) |
325 | 0 | { |
326 | 0 | fr_time_t now; |
327 | 0 | fr_time_delta_t t; |
328 | | |
329 | | /* |
330 | | * Note that the replies may come out of order with |
331 | | * respect to the request. So we can't use this reply |
332 | | * for any kind of timing. |
333 | | */ |
334 | 0 | now = fr_time(); |
335 | 0 | t = fr_time_sub(now, request_time); |
336 | |
|
337 | 0 | l->stats.rttvar = RTTVAR(l->stats.rtt, l->stats.rttvar, t); |
338 | 0 | l->stats.rtt = RTT(l->stats.rtt, t); |
339 | |
|
340 | 0 | l->stats.received++; |
341 | | |
342 | | /* |
343 | | * t is in nanoseconds. |
344 | | */ |
345 | 0 | if (fr_time_delta_lt(t, fr_time_delta_wrap(1000))) { |
346 | 0 | l->stats.times[0]++; /* < microseconds */ |
347 | 0 | } else if (fr_time_delta_lt(t, fr_time_delta_wrap(10000))) { |
348 | 0 | l->stats.times[1]++; /* microseconds */ |
349 | 0 | } else if (fr_time_delta_lt(t, fr_time_delta_wrap(100000))) { |
350 | 0 | l->stats.times[2]++; /* 10s of microseconds */ |
351 | 0 | } else if (fr_time_delta_lt(t, fr_time_delta_wrap(1000000))) { |
352 | 0 | l->stats.times[3]++; /* 100s of microseconds */ |
353 | 0 | } else if (fr_time_delta_lt(t, fr_time_delta_wrap(10000000))) { |
354 | 0 | l->stats.times[4]++; /* milliseconds */ |
355 | 0 | } else if (fr_time_delta_lt(t, fr_time_delta_wrap(100000000))) { |
356 | 0 | l->stats.times[5]++; /* 10s of milliseconds */ |
357 | 0 | } else if (fr_time_delta_lt(t, fr_time_delta_wrap(NSEC))) { |
358 | 0 | l->stats.times[6]++; /* 100s of milliseconds */ |
359 | 0 | } else { |
360 | 0 | l->stats.times[7]++; /* seconds */ |
361 | 0 | } |
362 | | |
363 | | /* |
364 | | * Still sending packets. Rely on the timer to send more |
365 | | * packets. |
366 | | */ |
367 | 0 | if (l->state == FR_LOAD_STATE_SENDING) return FR_LOAD_CONTINUE; |
368 | | |
369 | | /* |
370 | | * The send code has decided that the backlog is too |
371 | | * high. New requests are blocked until replies come in. |
372 | | * Since we have a reply, send another request. Unless |
373 | | * we've already sent max_requests packets, in which case |
374 | | * the timer will notice and start draining. |
375 | | */ |
376 | 0 | if (l->state == FR_LOAD_STATE_GATED) { |
377 | 0 | if ((l->stats.skipped > 0) && |
378 | 0 | (!l->config->max_requests || |
379 | 0 | ((uint32_t) (l->stats.sent - l->sent_base) < l->config->max_requests))) { |
380 | 0 | l->stats.skipped--; |
381 | 0 | fr_load_generator_send(l, now, 1); |
382 | 0 | } |
383 | 0 | return FR_LOAD_CONTINUE; |
384 | 0 | } |
385 | | |
386 | | /* |
387 | | * We're still sending or gated, tell the caller to |
388 | | * continue. |
389 | | */ |
390 | 0 | if (l->state != FR_LOAD_STATE_DRAINING) { |
391 | 0 | return FR_LOAD_CONTINUE; |
392 | 0 | } |
393 | | /* |
394 | | * Not yet received all replies. Wait until we have all |
395 | | * replies. |
396 | | */ |
397 | 0 | if (l->stats.received < l->stats.sent) return FR_LOAD_CONTINUE; |
398 | | |
399 | 0 | l->stats.end = now; |
400 | 0 | return FR_LOAD_DONE; |
401 | 0 | } |
402 | | |
403 | | /** Print load generator statistics in CVS format. |
404 | | * |
405 | | */ |
406 | | size_t fr_load_generator_stats_sprint(fr_load_t *l, fr_time_t now, char *buffer, size_t buflen) |
407 | 0 | { |
408 | 0 | double now_f, last_send_f; |
409 | |
|
410 | 0 | if (!l->header) { |
411 | 0 | l->header = true; |
412 | 0 | return snprintf(buffer, buflen, "\"time\",\"last_packet\",\"rtt\",\"rttvar\",\"pps\",\"pps_accepted\",\"sent\",\"received\",\"backlog\",\"max_backlog\",\"<usec\",\"us\",\"10us\",\"100us\",\"ms\",\"10ms\",\"100ms\",\"s\",\"blocked\"\n"); |
413 | 0 | } |
414 | | |
415 | | |
416 | 0 | now_f = fr_time_delta_unwrap(fr_time_sub(now, l->stats.start)) / (double)NSEC; |
417 | |
|
418 | 0 | last_send_f = fr_time_delta_unwrap(fr_time_sub(l->stats.last_send, l->stats.start)) / (double)NSEC; |
419 | | |
420 | | /* |
421 | | * Track packets/s. Since times are in nanoseconds, we |
422 | | * have to scale the counters up by NSEC. And since NSEC |
423 | | * is 1B, the calculations have to be done via 64-bit |
424 | | * numbers, and then converted to a final 32-bit counter. |
425 | | */ |
426 | 0 | if (fr_time_gt(now, l->step_start)) { |
427 | 0 | l->stats.pps_accepted = |
428 | 0 | fr_time_delta_unwrap( |
429 | 0 | fr_time_delta_div(fr_time_delta_from_sec(l->stats.received - l->step_received), |
430 | 0 | fr_time_sub(now, l->step_start)) |
431 | 0 | ); |
432 | 0 | } |
433 | |
|
434 | 0 | return snprintf(buffer, buflen, |
435 | 0 | "%f,%f," |
436 | 0 | "%" PRIu64 ",%" PRIu64 "," |
437 | 0 | "%d,%d," |
438 | 0 | "%d,%d," |
439 | 0 | "%d,%d," |
440 | 0 | "%d,%d,%d,%d,%d,%d,%d,%d," |
441 | 0 | "%d\n", |
442 | 0 | now_f, last_send_f, |
443 | 0 | fr_time_delta_unwrap(l->stats.rtt), fr_time_delta_unwrap(l->stats.rttvar), |
444 | 0 | l->stats.pps, l->stats.pps_accepted, |
445 | 0 | l->stats.sent, l->stats.received, |
446 | 0 | l->stats.backlog, l->stats.max_backlog, |
447 | 0 | l->stats.times[0], l->stats.times[1], l->stats.times[2], l->stats.times[3], |
448 | 0 | l->stats.times[4], l->stats.times[5], l->stats.times[6], l->stats.times[7], |
449 | 0 | l->stats.blocked); |
450 | 0 | } |
451 | | |
452 | | fr_load_stats_t const * fr_load_generator_stats(fr_load_t const *l) |
453 | 0 | { |
454 | 0 | return &l->stats; |
455 | 0 | } |