提交 597fe3ce 编写于 作者: M Michal Privoznik

daemon: Create priority workers pool

This patch annotates APIs with low or high priority.
In low set MUST be all APIs which might eventually access monitor
(and thus block indefinitely). Other APIs may be marked as high
priority. However, some must be (e.g. domainDestroy).

For high priority calls (HPC), there are some high priority workers
(HPW) created in the pool. HPW can execute only HPC, although normal
worker can process any call regardless priority. Therefore, only those
APIs which are guaranteed to end in reasonable small amount of time
can be marked as HPC.

The size of this HPC pool is static, because HPC are expected to end
quickly, therefore jobs assigned to this pool will be served quickly.
It can be configured in libvirtd.conf via prio_workers variable.
Default is set to 5.

To mark API with low or high priority, append priority:{low|high} to
it's comment in src/remote/remote_protocol.x. This is similar to
autogen|skipgen. If not marked, the generator assumes low as default.
上级 63379890
......@@ -57,6 +57,7 @@ module Libvirtd =
| int_entry "max_clients"
| int_entry "max_requests"
| int_entry "max_client_requests"
| int_entry "prio_workers"
let logging_entry = int_entry "log_level"
| str_entry "log_filters"
......
......@@ -134,6 +134,8 @@ struct daemonConfig {
int max_workers;
int max_clients;
int prio_workers;
int max_requests;
int max_client_requests;
......@@ -886,6 +888,8 @@ daemonConfigNew(bool privileged ATTRIBUTE_UNUSED)
data->max_workers = 20;
data->max_clients = 20;
data->prio_workers = 5;
data->max_requests = 20;
data->max_client_requests = 5;
......@@ -1042,6 +1046,8 @@ daemonConfigLoad(struct daemonConfig *data,
GET_CONF_INT (conf, filename, max_workers);
GET_CONF_INT (conf, filename, max_clients);
GET_CONF_INT (conf, filename, prio_workers);
GET_CONF_INT (conf, filename, max_requests);
GET_CONF_INT (conf, filename, max_client_requests);
......@@ -1441,6 +1447,7 @@ int main(int argc, char **argv) {
config->auth_unix_ro == REMOTE_AUTH_POLKIT;
if (!(srv = virNetServerNew(config->min_workers,
config->max_workers,
config->prio_workers,
config->max_clients,
config->mdns_adv ? config->mdns_name : NULL,
use_polkit_dbus,
......
......@@ -257,6 +257,12 @@
#min_workers = 5
#max_workers = 20
# The number of priority workers. If all workers from above
# pool will stuck, some calls marked as high priority
# (notably domainDestroy) can be executed in this pool.
#prio_workers = 5
# Total global limit on concurrent RPC calls. Should be
# at least as large as max_workers. Beyond this, RPC requests
# will be read into memory and queued. This directly impact
......
......@@ -663,7 +663,7 @@ qemudStartup(int privileged) {
virHashForEach(qemu_driver->domains.objs, qemuDomainSnapshotLoad,
qemu_driver->snapshotDir);
qemu_driver->workerPool = virThreadPoolNew(0, 1, processWatchdogEvent, qemu_driver);
qemu_driver->workerPool = virThreadPoolNew(0, 1, 0, processWatchdogEvent, qemu_driver);
if (!qemu_driver->workerPool)
goto error;
......
......@@ -571,7 +571,7 @@ qemuProcessHandleWatchdog(qemuMonitorPtr mon ATTRIBUTE_UNUSED,
* deleted before handling watchdog event is finished.
*/
virDomainObjRef(vm);
if (virThreadPoolSendJob(driver->workerPool, wdEvent) < 0) {
if (virThreadPoolSendJob(driver->workerPool, 0, wdEvent) < 0) {
if (virDomainObjUnref(vm) == 0)
vm = NULL;
VIR_FREE(wdEvent);
......
......@@ -52,9 +52,14 @@ const QEMU_PROGRAM = 0x20008087;
const QEMU_PROTOCOL_VERSION = 1;
enum qemu_procedure {
/* Each function must have a two-word comment. The first word is
/* Each function must have a three-word comment. The first word is
* whether gendispatch.pl handles daemon, the second whether
* it handles src/remote. */
QEMU_PROC_MONITOR_COMMAND = 1, /* skipgen skipgen */
QEMU_PROC_DOMAIN_ATTACH = 2 /* autogen autogen */
* it handles src/remote.
* The last argument describes priority of API. There are two accepted
* values: low, high; Each API that might eventually access hypervisor's
* monitor (and thus block) MUST fall into low priority. However, there
* are some exceptions to this rule, e.g. domainDestroy. Other APIs MAY
* be marked as high priority. If in doubt, it's safe to choose low. */
QEMU_PROC_MONITOR_COMMAND = 1, /* skipgen skipgen priority:low */
QEMU_PROC_DOMAIN_ATTACH = 2 /* autogen autogen priority:low */
};
此差异已折叠。
......@@ -146,12 +146,13 @@ while (<PROTOCOL>) {
}
if ($opt_b or $opt_k) {
if (!($flags =~ m/^\s*\/\*\s*(\S+)\s+(\S+)\s*(.*)\*\/\s*$/)) {
if (!($flags =~ m/^\s*\/\*\s*(\S+)\s+(\S+)\s*(\|.*)?\s+(priority:(\S+))?\s*\*\/\s*$/)) {
die "invalid generator flags for ${procprefix}_PROC_${name}"
}
my $genmode = $opt_b ? $1 : $2;
my $genflags = $3;
my $priority = defined $5 ? $5 : "low";
if ($genmode eq "autogen") {
push(@autogen, $ProcName);
......@@ -171,6 +172,16 @@ while (<PROTOCOL>) {
} else {
$calls{$name}->{streamflag} = "none";
}
# for now, we distinguish only two levels of prioroty:
# low (0) and high (1)
if ($priority eq "high") {
$calls{$name}->{priority} = 1;
} elsif ($priority eq "low") {
$calls{$name}->{priority} = 0;
} else {
die "invalid priority ${priority} for ${procprefix}_PROC_${name}"
}
}
$calls[$id] = $calls{$name};
......@@ -260,6 +271,7 @@ if ($opt_d) {
print "$_:\n";
print " name $calls{$_}->{name} ($calls{$_}->{ProcName})\n";
print " $calls{$_}->{args} -> $calls{$_}->{ret}\n";
print " priority -> $calls{$_}->{priority}\n";
}
}
......@@ -935,7 +947,7 @@ elsif ($opt_b) {
print "virNetServerProgramProc ${structprefix}Procs[] = {\n";
for ($id = 0 ; $id <= $#calls ; $id++) {
my ($comment, $name, $argtype, $arglen, $argfilter, $retlen, $retfilter);
my ($comment, $name, $argtype, $arglen, $argfilter, $retlen, $retfilter, $priority);
if (defined $calls[$id] && !$calls[$id]->{msg}) {
$comment = "/* Method $calls[$id]->{ProcName} => $id */";
......@@ -958,7 +970,9 @@ elsif ($opt_b) {
$retfilter = "xdr_void";
}
print "{ $comment\n ${name},\n $arglen,\n (xdrproc_t)$argfilter,\n $retlen,\n (xdrproc_t)$retfilter,\n true \n},\n";
$priority = defined $calls[$id]->{priority} ? $calls[$id]->{priority} : 0;
print "{ $comment\n ${name},\n $arglen,\n (xdrproc_t)$argfilter,\n $retlen,\n (xdrproc_t)$retfilter,\n true,\n $priority\n},\n";
}
print "};\n";
print "size_t ${structprefix}NProcs = ARRAY_CARDINALITY(${structprefix}Procs);\n";
......
......@@ -64,6 +64,7 @@ typedef virNetServerJob *virNetServerJobPtr;
struct _virNetServerJob {
virNetServerClientPtr client;
virNetMessagePtr msg;
virNetServerProgramPtr prog;
};
struct _virNetServer {
......@@ -128,61 +129,41 @@ static void virNetServerHandleJob(void *jobOpaque, void *opaque)
{
virNetServerPtr srv = opaque;
virNetServerJobPtr job = jobOpaque;
virNetServerProgramPtr prog = NULL;
size_t i;
virNetServerLock(srv);
VIR_DEBUG("server=%p client=%p message=%p",
srv, job->client, job->msg);
for (i = 0 ; i < srv->nprograms ; i++) {
if (virNetServerProgramMatches(srv->programs[i], job->msg)) {
prog = srv->programs[i];
break;
}
}
VIR_DEBUG("server=%p client=%p message=%p prog=%p",
srv, job->client, job->msg, job->prog);
if (!prog) {
if (virNetServerProgramUnknownError(job->client,
job->msg,
&job->msg->header) < 0)
goto error;
else
goto cleanup;
}
virNetServerProgramRef(prog);
virNetServerUnlock(srv);
if (virNetServerProgramDispatch(prog,
if (virNetServerProgramDispatch(job->prog,
srv,
job->client,
job->msg) < 0)
goto error;
virNetServerLock(srv);
cleanup:
virNetServerProgramFree(prog);
virNetServerProgramFree(job->prog);
virNetServerUnlock(srv);
virNetServerClientFree(job->client);
VIR_FREE(job);
return;
error:
virNetServerProgramFree(job->prog);
virNetMessageFree(job->msg);
virNetServerClientClose(job->client);
goto cleanup;
virNetServerClientFree(job->client);
VIR_FREE(job);
}
static int virNetServerDispatchNewMessage(virNetServerClientPtr client,
virNetMessagePtr msg,
void *opaque)
{
virNetServerPtr srv = opaque;
virNetServerJobPtr job;
int ret;
virNetServerProgramPtr prog = NULL;
unsigned int priority = 0;
size_t i;
int ret = -1;
VIR_DEBUG("server=%p client=%p message=%p",
srv, client, msg);
......@@ -196,8 +177,29 @@ static int virNetServerDispatchNewMessage(virNetServerClientPtr client,
job->msg = msg;
virNetServerLock(srv);
if ((ret = virThreadPoolSendJob(srv->workers, job)) < 0)
for (i = 0 ; i < srv->nprograms ; i++) {
if (virNetServerProgramMatches(srv->programs[i], job->msg)) {
prog = srv->programs[i];
break;
}
}
if (!prog) {
virNetServerProgramUnknownError(client, msg, &msg->header);
goto cleanup;
}
virNetServerProgramRef(prog);
job->prog = prog;
priority = virNetServerProgramGetPriority(prog, msg->header.proc);
ret = virThreadPoolSendJob(srv->workers, priority, job);
cleanup:
if (ret < 0) {
VIR_FREE(job);
virNetServerProgramFree(prog);
}
virNetServerUnlock(srv);
return ret;
......@@ -274,6 +276,7 @@ static void virNetServerFatalSignal(int sig, siginfo_t *siginfo ATTRIBUTE_UNUSED
virNetServerPtr virNetServerNew(size_t min_workers,
size_t max_workers,
size_t priority_workers,
size_t max_clients,
const char *mdnsGroupName,
bool connectDBus ATTRIBUTE_UNUSED,
......@@ -290,6 +293,7 @@ virNetServerPtr virNetServerNew(size_t min_workers,
srv->refs = 1;
if (!(srv->workers = virThreadPoolNew(min_workers, max_workers,
priority_workers,
virNetServerHandleJob,
srv)))
goto error;
......
......@@ -39,6 +39,7 @@ typedef int (*virNetServerClientInitHook)(virNetServerPtr srv,
virNetServerPtr virNetServerNew(size_t min_workers,
size_t max_workers,
size_t priority_workers,
size_t max_clients,
const char *mdnsGroupName,
bool connectDBus,
......
......@@ -108,6 +108,17 @@ static virNetServerProgramProcPtr virNetServerProgramGetProc(virNetServerProgram
return &prog->procs[procedure];
}
unsigned int
virNetServerProgramGetPriority(virNetServerProgramPtr prog,
int procedure)
{
virNetServerProgramProcPtr proc = virNetServerProgramGetProc(prog, procedure);
if (!proc)
return 0;
return proc->priority;
}
static int
virNetServerProgramSendError(unsigned program,
......
......@@ -53,6 +53,7 @@ struct _virNetServerProgramProc {
size_t ret_len;
xdrproc_t ret_filter;
bool needAuth;
unsigned int priority;
};
virNetServerProgramPtr virNetServerProgramNew(unsigned program,
......@@ -63,6 +64,9 @@ virNetServerProgramPtr virNetServerProgramNew(unsigned program,
int virNetServerProgramGetID(virNetServerProgramPtr prog);
int virNetServerProgramGetVersion(virNetServerProgramPtr prog);
unsigned int virNetServerProgramGetPriority(virNetServerProgramPtr prog,
int procedure);
void virNetServerProgramRef(virNetServerProgramPtr prog);
int virNetServerProgramMatches(virNetServerProgramPtr prog,
......
......@@ -37,7 +37,9 @@ typedef struct _virThreadPoolJob virThreadPoolJob;
typedef virThreadPoolJob *virThreadPoolJobPtr;
struct _virThreadPoolJob {
virThreadPoolJobPtr prev;
virThreadPoolJobPtr next;
unsigned int priority;
void *data;
};
......@@ -47,7 +49,8 @@ typedef virThreadPoolJobList *virThreadPoolJobListPtr;
struct _virThreadPoolJobList {
virThreadPoolJobPtr head;
virThreadPoolJobPtr *tail;
virThreadPoolJobPtr tail;
virThreadPoolJobPtr firstPrio;
};
......@@ -57,6 +60,7 @@ struct _virThreadPool {
virThreadPoolJobFunc jobFunc;
void *jobOpaque;
virThreadPoolJobList jobList;
size_t jobQueueDepth;
virMutex mutex;
virCond cond;
......@@ -66,33 +70,75 @@ struct _virThreadPool {
size_t freeWorkers;
size_t nWorkers;
virThreadPtr workers;
size_t nPrioWorkers;
virThreadPtr prioWorkers;
virCond prioCond;
};
struct virThreadPoolWorkerData {
virThreadPoolPtr pool;
virCondPtr cond;
bool priority;
};
static void virThreadPoolWorker(void *opaque)
{
virThreadPoolPtr pool = opaque;
struct virThreadPoolWorkerData *data = opaque;
virThreadPoolPtr pool = data->pool;
virCondPtr cond = data->cond;
bool priority = data->priority;
virThreadPoolJobPtr job = NULL;
VIR_FREE(data);
virMutexLock(&pool->mutex);
while (1) {
while (!pool->quit &&
!pool->jobList.head) {
((!priority && !pool->jobList.head) ||
(priority && !pool->jobList.firstPrio))) {
if (!priority)
pool->freeWorkers++;
if (virCondWait(&pool->cond, &pool->mutex) < 0) {
if (virCondWait(cond, &pool->mutex) < 0) {
if (!priority)
pool->freeWorkers--;
goto out;
}
if (!priority)
pool->freeWorkers--;
}
if (pool->quit)
break;
virThreadPoolJobPtr job = pool->jobList.head;
pool->jobList.head = pool->jobList.head->next;
job->next = NULL;
if (pool->jobList.tail == &job->next)
pool->jobList.tail = &pool->jobList.head;
if (priority) {
job = pool->jobList.firstPrio;
} else {
job = pool->jobList.head;
}
if (job == pool->jobList.firstPrio) {
virThreadPoolJobPtr tmp = job->next;
while (tmp) {
if (tmp->priority) {
break;
}
tmp = tmp->next;
}
pool->jobList.firstPrio = tmp;
}
if (job->prev)
job->prev->next = job->next;
else
pool->jobList.head = job->next;
if (job->next)
job->next->prev = job->prev;
else
pool->jobList.tail = job->prev;
pool->jobQueueDepth--;
virMutexUnlock(&pool->mutex);
(pool->jobFunc)(job->data, pool->jobOpaque);
......@@ -101,19 +147,24 @@ static void virThreadPoolWorker(void *opaque)
}
out:
if (priority)
pool->nPrioWorkers--;
else
pool->nWorkers--;
if (pool->nWorkers == 0)
if (pool->nWorkers == 0 && pool->nPrioWorkers==0)
virCondSignal(&pool->quit_cond);
virMutexUnlock(&pool->mutex);
}
virThreadPoolPtr virThreadPoolNew(size_t minWorkers,
size_t maxWorkers,
size_t prioWorkers,
virThreadPoolJobFunc func,
void *opaque)
{
virThreadPoolPtr pool;
size_t i;
struct virThreadPoolWorkerData *data = NULL;
if (minWorkers > maxWorkers)
minWorkers = maxWorkers;
......@@ -123,8 +174,7 @@ virThreadPoolPtr virThreadPoolNew(size_t minWorkers,
return NULL;
}
pool->jobList.head = NULL;
pool->jobList.tail = &pool->jobList.head;
pool->jobList.tail = pool->jobList.head = NULL;
pool->jobFunc = func;
pool->jobOpaque = opaque;
......@@ -141,18 +191,51 @@ virThreadPoolPtr virThreadPoolNew(size_t minWorkers,
pool->maxWorkers = maxWorkers;
for (i = 0; i < minWorkers; i++) {
if (VIR_ALLOC(data) < 0) {
virReportOOMError();
goto error;
}
data->pool = pool;
data->cond = &pool->cond;
if (virThreadCreate(&pool->workers[i],
true,
virThreadPoolWorker,
pool) < 0) {
data) < 0) {
goto error;
}
pool->nWorkers++;
}
if (prioWorkers) {
if (virCondInit(&pool->prioCond) < 0)
goto error;
if (VIR_ALLOC_N(pool->prioWorkers, prioWorkers) < 0)
goto error;
for (i = 0; i < prioWorkers; i++) {
if (VIR_ALLOC(data) < 0) {
virReportOOMError();
goto error;
}
data->pool = pool;
data->cond = &pool->prioCond;
data->priority = true;
if (virThreadCreate(&pool->prioWorkers[i],
true,
virThreadPoolWorker,
data) < 0) {
goto error;
}
pool->nPrioWorkers++;
}
}
return pool;
error:
VIR_FREE(data);
virThreadPoolFree(pool);
return NULL;
......@@ -161,17 +244,22 @@ error:
void virThreadPoolFree(virThreadPoolPtr pool)
{
virThreadPoolJobPtr job;
bool priority = false;
if (!pool)
return;
virMutexLock(&pool->mutex);
pool->quit = true;
if (pool->nWorkers > 0) {
if (pool->nWorkers > 0)
virCondBroadcast(&pool->cond);
ignore_value(virCondWait(&pool->quit_cond, &pool->mutex));
if (pool->nPrioWorkers > 0) {
priority = true;
virCondBroadcast(&pool->prioCond);
}
ignore_value(virCondWait(&pool->quit_cond, &pool->mutex));
while ((job = pool->jobList.head)) {
pool->jobList.head = pool->jobList.head->next;
VIR_FREE(job);
......@@ -182,10 +270,19 @@ void virThreadPoolFree(virThreadPoolPtr pool)
virMutexDestroy(&pool->mutex);
ignore_value(virCondDestroy(&pool->quit_cond));
ignore_value(virCondDestroy(&pool->cond));
if (priority) {
VIR_FREE(pool->prioWorkers);
ignore_value(virCondDestroy(&pool->prioCond));
}
VIR_FREE(pool);
}
/*
* @priority - job priority
* Return: 0 on success, -1 otherwise
*/
int virThreadPoolSendJob(virThreadPoolPtr pool,
unsigned int priority,
void *jobData)
{
virThreadPoolJobPtr job;
......@@ -194,7 +291,7 @@ int virThreadPoolSendJob(virThreadPoolPtr pool,
if (pool->quit)
goto error;
if (pool->freeWorkers == 0 &&
if (pool->freeWorkers - pool->jobQueueDepth <= 0 &&
pool->nWorkers < pool->maxWorkers) {
if (VIR_EXPAND_N(pool->workers, pool->nWorkers, 1) < 0) {
virReportOOMError();
......@@ -216,13 +313,26 @@ int virThreadPoolSendJob(virThreadPoolPtr pool,
}
job->data = jobData;
job->next = NULL;
*pool->jobList.tail = job;
pool->jobList.tail = &(*pool->jobList.tail)->next;
job->priority = priority;
job->prev = pool->jobList.tail;
if (pool->jobList.tail)
pool->jobList.tail->next = job;
pool->jobList.tail = job;
if (!pool->jobList.head)
pool->jobList.head = job;
if (priority && !pool->jobList.firstPrio)
pool->jobList.firstPrio = job;
pool->jobQueueDepth++;
virCondSignal(&pool->cond);
virMutexUnlock(&pool->mutex);
if (priority)
virCondSignal(&pool->prioCond);
virMutexUnlock(&pool->mutex);
return 0;
error:
......
......@@ -35,12 +35,14 @@ typedef void (*virThreadPoolJobFunc)(void *jobdata, void *opaque);
virThreadPoolPtr virThreadPoolNew(size_t minWorkers,
size_t maxWorkers,
size_t prioWorkers,
virThreadPoolJobFunc func,
void *opaque) ATTRIBUTE_NONNULL(3);
void *opaque) ATTRIBUTE_NONNULL(4);
void virThreadPoolFree(virThreadPoolPtr pool);
int virThreadPoolSendJob(virThreadPoolPtr pool,
unsigned int priority,
void *jobdata) ATTRIBUTE_NONNULL(1)
ATTRIBUTE_RETURN_CHECK;
......
Markdown is supported
0% .
You are about to add 0 people to the discussion. Proceed with caution.
先完成此消息的编辑!
想要评论请 注册