diff options
| -rw-r--r-- | src/ipcpd/unicast/psched.c | 56 |
1 files changed, 36 insertions, 20 deletions
diff --git a/src/ipcpd/unicast/psched.c b/src/ipcpd/unicast/psched.c index 2b535d67..01c4b066 100644 --- a/src/ipcpd/unicast/psched.c +++ b/src/ipcpd/unicast/psched.c @@ -116,6 +116,16 @@ static void dsched_untrack(struct dsched * d, d->posn[fd] = -1; } +/* Fold a deadline into the earliest pending one (0 = none yet). */ +static uint64_t dmin_fold(uint64_t dmin, + uint64_t deadline) +{ + if (dmin == 0 || deadline < dmin) + return deadline; + + return dmin; +} + static uint64_t dsched_serve(struct dsched * d, struct psched * sched, qoscube_t qc, @@ -126,31 +136,37 @@ static uint64_t dsched_serve(struct dsched * d, size_t i; int fd; time_t wait; + bool served; - for (i = 0; i < d->n; ) { - fd = d->active[i]; + /* Round-robin one packet per flow so none monopolises egress. */ + do { + served = false; - if (d->deadline[fd] > now) { - if (dmin == 0 || d->deadline[fd] < dmin) - dmin = d->deadline[fd]; - ++i; - continue; - } + for (i = 0; i < d->n; ) { + fd = d->active[i]; - if (sched->read(fd, &spb) < 0) { - dsched_untrack(d, fd); - continue; - } + if (d->deadline[fd] > now) { + dmin = dmin_fold(dmin, d->deadline[fd]); + ++i; + continue; + } - wait = sched->callback(fd, qc, spb); - if (wait == 0) - continue; + if (sched->read(fd, &spb) < 0) { + dsched_untrack(d, fd); + continue; + } - d->deadline[fd] = now + (uint64_t) wait; - if (dmin == 0 || d->deadline[fd] < dmin) - dmin = d->deadline[fd]; - ++i; - } + wait = sched->callback(fd, qc, spb); + served = true; + + if (wait > 0) { + d->deadline[fd] = now + (uint64_t) wait; + dmin = dmin_fold(dmin, d->deadline[fd]); + } + + ++i; + } + } while (served); return dmin; } |
