#ifdef _KERNEL
#include <sys/cdefs.h>
__KERNEL_RCSID(0, "$NetBSD: npf_worker.c,v 1.10 2020/08/27 18:49:36 riastradh Exp $");
#include <sys/param.h>
#include <sys/types.h>
#include <sys/mutex.h>
#include <sys/kmem.h>
#include <sys/kernel.h>
#include <sys/kthread.h>
#include <sys/cprng.h>
#endif
#include "npf_impl.h"
typedef struct npf_worker {
kmutex_t lock;
kcondvar_t cv;
kcondvar_t exit_cv;
bool exit;
LIST_HEAD(, npf) instances;
unsigned worker_count;
lwp_t * worker[];
} npf_workerinfo_t;
#define NPF_GC_MINWAIT (10)
#define NPF_GC_MAXWAIT (10 * 1000)
#define WFLAG_ACTIVE 0x01
#define WFLAG_INITED 0x02
#define WFLAG_REMOVE 0x04
static void npf_worker(void *) __dead;
static npf_workerinfo_t * worker_info __read_mostly;
int
npf_worker_sysinit(unsigned nworkers)
{
const size_t len = offsetof(npf_workerinfo_t, worker[nworkers]);
npf_workerinfo_t *winfo;
KASSERT(worker_info == NULL);
if (!nworkers) {
return 0;
}
winfo = kmem_zalloc(len, KM_SLEEP);
winfo->worker_count = nworkers;
mutex_init(&winfo->lock, MUTEX_DEFAULT, IPL_SOFTNET);
cv_init(&winfo->exit_cv, "npfgcx");
cv_init(&winfo->cv, "npfgcw");
LIST_INIT(&winfo->instances);
worker_info = winfo;
for (unsigned i = 0; i < nworkers; i++) {
if (kthread_create(PRI_NONE, KTHREAD_MPSAFE | KTHREAD_MUSTJOIN,
NULL, npf_worker, winfo, &winfo->worker[i], "npfgc%u", i)) {
npf_worker_sysfini();
return ENOMEM;
}
}
return 0;
}
void
npf_worker_sysfini(void)
{
npf_workerinfo_t *winfo = worker_info;
unsigned nworkers;
if (!winfo) {
return;
}
mutex_enter(&winfo->lock);
winfo->exit = true;
cv_broadcast(&winfo->cv);
mutex_exit(&winfo->lock);
nworkers = winfo->worker_count;
for (unsigned i = 0; i < nworkers; i++) {
lwp_t *worker;
if ((worker = winfo->worker[i]) != NULL) {
kthread_join(worker);
}
}
cv_destroy(&winfo->cv);
cv_destroy(&winfo->exit_cv);
mutex_destroy(&winfo->lock);
kmem_free(winfo, offsetof(npf_workerinfo_t, worker[nworkers]));
worker_info = NULL;
}
int
npf_worker_addfunc(npf_t *npf, npf_workfunc_t work)
{
KASSERTMSG(npf->worker_flags == 0,
"the task must be added before the npf_worker_enlist() call");
for (unsigned i = 0; i < NPF_MAX_WORKS; i++) {
if (npf->worker_funcs[i] == NULL) {
npf->worker_funcs[i] = work;
return 0;
}
}
return -1;
}
void
npf_worker_signal(npf_t *npf)
{
npf_workerinfo_t *winfo = worker_info;
if ((npf->worker_flags & WFLAG_ACTIVE) == 0) {
return;
}
KASSERT(winfo != NULL);
mutex_enter(&winfo->lock);
cv_signal(&winfo->cv);
mutex_exit(&winfo->lock);
}
void
npf_worker_enlist(npf_t *npf)
{
npf_workerinfo_t *winfo = worker_info;
KASSERT(npf->worker_flags == 0);
if (!winfo) {
return;
}
mutex_enter(&winfo->lock);
LIST_INSERT_HEAD(&winfo->instances, npf, worker_entry);
npf->worker_flags |= WFLAG_ACTIVE;
mutex_exit(&winfo->lock);
}
void
npf_worker_discharge(npf_t *npf)
{
npf_workerinfo_t *winfo = worker_info;
if ((npf->worker_flags & WFLAG_ACTIVE) == 0) {
return;
}
KASSERT(winfo != NULL);
mutex_enter(&winfo->lock);
KASSERT(npf->worker_flags & WFLAG_ACTIVE);
npf->worker_flags |= WFLAG_REMOVE;
cv_broadcast(&winfo->cv);
while (npf->worker_flags & WFLAG_ACTIVE) {
cv_wait(&winfo->exit_cv, &winfo->lock);
}
mutex_exit(&winfo->lock);
KASSERT(npf->worker_flags == 0);
}
static void
remove_npf_instance(npf_workerinfo_t *winfo, npf_t *npf)
{
KASSERT(mutex_owned(&winfo->lock));
KASSERT(npf->worker_flags & WFLAG_ACTIVE);
KASSERT(npf->worker_flags & WFLAG_REMOVE);
if (npf->worker_flags & WFLAG_INITED) {
npfk_thread_unregister(npf);
}
LIST_REMOVE(npf, worker_entry);
npf->worker_flags = 0;
cv_broadcast(&winfo->exit_cv);
}
static unsigned
process_npf_instance(npf_workerinfo_t *winfo, npf_t *npf)
{
npf_workfunc_t work;
KASSERT(mutex_owned(&winfo->lock));
if (npf->worker_flags & WFLAG_REMOVE) {
remove_npf_instance(winfo, npf);
return NPF_GC_MAXWAIT;
}
if ((npf->worker_flags & WFLAG_INITED) == 0) {
npfk_thread_register(npf);
npf->worker_flags |= WFLAG_INITED;
}
for (unsigned i = 0; i < NPF_MAX_WORKS; i++) {
if ((work = npf->worker_funcs[i]) == NULL) {
break;
}
work(npf);
}
return MAX(MIN(npf->worker_wait_time, NPF_GC_MAXWAIT), NPF_GC_MINWAIT);
}
static void
npf_worker(void *arg)
{
npf_workerinfo_t *winfo = arg;
npf_t *npf;
mutex_enter(&winfo->lock);
for (;;) {
unsigned wait_time = NPF_GC_MAXWAIT;
npf = LIST_FIRST(&winfo->instances);
while (npf) {
npf_t *next = LIST_NEXT(npf, worker_entry);
unsigned i_wait_time = process_npf_instance(winfo, npf);
wait_time = MIN(wait_time, i_wait_time);
npf = next;
}
if (winfo->exit) {
break;
}
cv_timedwait(&winfo->cv, &winfo->lock, mstohz(wait_time));
}
mutex_exit(&winfo->lock);
KASSERTMSG(LIST_EMPTY(&winfo->instances),
"NPF instances must be discharged before the npfk_sysfini() call");
kthread_exit(0);
}