Coverage Report

Created: 2026-09-28 06:27

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}