#include <sys/cdefs.h>
__KERNEL_RCSID(0, "$NetBSD: rf_engine.c,v 1.53 2019/10/10 03:43:59 christos Exp $");
#include <sys/errno.h>
#include "rf_threadstuff.h"
#include "rf_dag.h"
#include "rf_engine.h"
#include "rf_etimer.h"
#include "rf_general.h"
#include "rf_dagutils.h"
#include "rf_shutdown.h"
#include "rf_raid.h"
#include "rf_kintf.h"
#include "rf_paritymap.h"
static void rf_ShutdownEngine(void *);
static void DAGExecutionThread(RF_ThreadArg_t arg);
static void rf_RaidIOThread(RF_ThreadArg_t arg);
#define DO_LOCK(_r_) \
rf_lock_mutex2((_r_)->node_queue_mutex)
#define DO_UNLOCK(_r_) \
rf_unlock_mutex2((_r_)->node_queue_mutex)
#define DO_WAIT(_r_) \
rf_wait_cond2((_r_)->node_queue_cv, (_r_)->node_queue_mutex)
#define DO_SIGNAL(_r_) \
rf_broadcast_cond2((_r_)->node_queue_cv)
static void
rf_ShutdownEngine(void *arg)
{
RF_Raid_t *raidPtr;
raidPtr = (RF_Raid_t *) arg;
rf_lock_mutex2(raidPtr->iodone_lock);
raidPtr->shutdown_raidio = 1;
rf_signal_cond2(raidPtr->iodone_cv);
while (raidPtr->shutdown_raidio)
rf_wait_cond2(raidPtr->iodone_cv, raidPtr->iodone_lock);
rf_unlock_mutex2(raidPtr->iodone_lock);
DO_LOCK(raidPtr);
raidPtr->shutdown_engine = 1;
DO_SIGNAL(raidPtr);
while (raidPtr->shutdown_engine)
DO_WAIT(raidPtr);
DO_UNLOCK(raidPtr);
rf_destroy_mutex2(raidPtr->node_queue_mutex);
rf_destroy_cond2(raidPtr->node_queue_cv);
rf_destroy_mutex2(raidPtr->iodone_lock);
rf_destroy_cond2(raidPtr->iodone_cv);
}
int
rf_ConfigureEngine(RF_ShutdownList_t **listp, RF_Raid_t *raidPtr,
RF_Config_t *cfgPtr)
{
TAILQ_INIT(&(raidPtr->iodone));
rf_init_mutex2(raidPtr->iodone_lock, IPL_VM);
rf_init_cond2(raidPtr->iodone_cv, "raidiow");
rf_init_mutex2(raidPtr->node_queue_mutex, IPL_VM);
rf_init_cond2(raidPtr->node_queue_cv, "rfnodeq");
raidPtr->node_queue = NULL;
raidPtr->dags_in_flight = 0;
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: Creating engine thread\n", raidPtr->raidid);
}
#endif
if (RF_CREATE_ENGINE_THREAD(raidPtr->engine_thread,
DAGExecutionThread, raidPtr,
"raid%d", raidPtr->raidid)) {
printf("raid%d: Unable to create engine thread\n",
raidPtr->raidid);
return (ENOMEM);
}
if (RF_CREATE_ENGINE_THREAD(raidPtr->engine_helper_thread,
rf_RaidIOThread, raidPtr,
"raidio%d", raidPtr->raidid)) {
printf("raid%d: Unable to create raidio thread\n",
raidPtr->raidid);
return (ENOMEM);
}
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: Created engine thread\n", raidPtr->raidid);
}
#endif
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: Engine thread running and waiting for events\n", raidPtr->raidid);
}
#endif
rf_ShutdownCreate(listp, rf_ShutdownEngine, raidPtr);
return (0);
}
#if 0
static int
BranchDone(RF_DagNode_t *node)
{
int i;
switch (node->status) {
case rf_wait:
RF_PANIC();
break;
case rf_fired:
return (RF_FALSE);
case rf_good:
for (i = 0; i < node->numSuccedents; i++)
if (!BranchDone(node->succedents[i]))
return RF_FALSE;
return RF_TRUE;
case rf_bad:
return (RF_TRUE);
case rf_recover:
RF_PANIC();
break;
case rf_undone:
case rf_panic:
return (RF_TRUE);
default:
RF_PANIC();
break;
}
}
#endif
static int
NodeReady(RF_DagNode_t *node)
{
int ready;
ready = RF_FALSE;
switch (node->dagHdr->status) {
case rf_enable:
case rf_rollForward:
if ((node->status == rf_wait) &&
(node->numAntecedents == node->numAntDone))
ready = RF_TRUE;
break;
case rf_rollBackward:
RF_ASSERT(node->numSuccDone <= node->numSuccedents);
RF_ASSERT(node->numSuccFired <= node->numSuccedents);
RF_ASSERT(node->numSuccFired <= node->numSuccDone);
if ((node->status == rf_good) &&
(node->numSuccDone == node->numSuccedents))
ready = RF_TRUE;
break;
default:
printf("Execution engine found illegal DAG status in NodeReady\n");
RF_PANIC();
break;
}
return (ready);
}
static void
FireNode(RF_DagNode_t *node)
{
switch (node->status) {
case rf_fired:
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: Firing node 0x%lx (%s)\n",
node->dagHdr->raidPtr->raidid,
(unsigned long) node, node->name);
}
#endif
if (node->flags & RF_DAGNODE_FLAG_YIELD) {
#if defined(__NetBSD__) && defined(_KERNEL)
#else
thread_block();
#endif
}
(*(node->doFunc)) (node);
break;
case rf_recover:
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: Firing (undo) node 0x%lx (%s)\n",
node->dagHdr->raidPtr->raidid,
(unsigned long) node, node->name);
}
#endif
if (node->flags & RF_DAGNODE_FLAG_YIELD)
#if defined(__NetBSD__) && defined(_KERNEL)
#else
thread_block();
#endif
(*(node->undoFunc)) (node);
break;
default:
RF_PANIC();
break;
}
}
static void
FireNodeArray(int numNodes, RF_DagNode_t **nodeList)
{
RF_DagStatus_t dstat;
RF_DagNode_t *node;
int i, j;
for (i = 0; i < numNodes; i++) {
node = nodeList[i];
dstat = node->dagHdr->status;
RF_ASSERT((node->status == rf_wait) ||
(node->status == rf_good));
if (NodeReady(node)) {
if ((dstat == rf_enable) ||
(dstat == rf_rollForward)) {
RF_ASSERT(node->status == rf_wait);
if (node->commitNode)
node->dagHdr->numCommits++;
node->status = rf_fired;
for (j = 0; j < node->numAntecedents; j++)
node->antecedents[j]->numSuccFired++;
} else {
RF_ASSERT(dstat == rf_rollBackward);
RF_ASSERT(node->status == rf_good);
RF_ASSERT(node->commitNode == RF_FALSE);
node->status = rf_recover;
}
}
}
for (i = 0; i < numNodes; i++) {
if ((nodeList[i]->status == rf_fired) ||
(nodeList[i]->status == rf_recover))
FireNode(nodeList[i]);
}
}
static void
FireNodeList(RF_DagNode_t *nodeList)
{
RF_DagNode_t *node, *next;
RF_DagStatus_t dstat;
int j;
if (nodeList) {
for (node = nodeList; node; node = next) {
next = node->next;
dstat = node->dagHdr->status;
RF_ASSERT((node->status == rf_wait) ||
(node->status == rf_good));
if (NodeReady(node)) {
if ((dstat == rf_enable) ||
(dstat == rf_rollForward)) {
RF_ASSERT(node->status == rf_wait);
if (node->commitNode)
node->dagHdr->numCommits++;
node->status = rf_fired;
for (j = 0; j < node->numAntecedents; j++)
node->antecedents[j]->numSuccFired++;
} else {
RF_ASSERT(dstat == rf_rollBackward);
RF_ASSERT(node->status == rf_good);
RF_ASSERT(node->commitNode == RF_FALSE);
node->status = rf_recover;
}
}
}
for (node = nodeList; node; node = next) {
next = node->next;
if ((node->status == rf_fired) ||
(node->status == rf_recover))
FireNode(node);
}
}
}
static void
PropagateResults(RF_DagNode_t *node, int context)
{
RF_DagNode_t *s, *a;
RF_Raid_t *raidPtr;
int i;
RF_DagNode_t *finishlist = NULL;
RF_DagNode_t *skiplist = NULL;
RF_DagNode_t *firelist = NULL;
RF_DagNode_t *q = NULL, *qh = NULL, *next;
int j, skipNode;
raidPtr = node->dagHdr->raidPtr;
DO_LOCK(raidPtr);
for (i = 0; i < node->numAntecedents; i++) {
a = *(node->antecedents + i);
RF_ASSERT(a->numSuccFired >= a->numSuccDone);
RF_ASSERT(a->numSuccFired <= a->numSuccedents);
a->numSuccDone++;
}
switch (node->dagHdr->status) {
case rf_enable:
case rf_rollForward:
for (i = 0; i < node->numSuccedents; i++) {
s = *(node->succedents + i);
RF_ASSERT(s->status == rf_wait);
(s->numAntDone)++;
if (s->numAntDone == s->numAntecedents) {
if (s->doFunc == rf_NullNodeFunc) {
s->next = finishlist;
finishlist = s;
} else {
skipNode = RF_FALSE;
for (j = 0; j < s->numAntecedents; j++)
if ((s->antType[j] == rf_trueData) && (s->antecedents[j]->status == rf_bad))
skipNode = RF_TRUE;
if (skipNode) {
s->next = skiplist;
skiplist = s;
} else
if (context != RF_INTR_CONTEXT) {
s->next = firelist;
firelist = s;
} else {
RF_ASSERT(NodeReady(s));
if (q) {
q->next = s;
q = s;
} else {
qh = q = s;
qh->next = NULL;
}
}
}
}
}
if (q) {
q->next = raidPtr->node_queue;
raidPtr->node_queue = qh;
DO_SIGNAL(raidPtr);
}
DO_UNLOCK(raidPtr);
for (; skiplist; skiplist = next) {
next = skiplist->next;
skiplist->status = rf_skipped;
for (i = 0; i < skiplist->numAntecedents; i++) {
skiplist->antecedents[i]->numSuccFired++;
}
if (skiplist->commitNode) {
skiplist->dagHdr->numCommits++;
}
rf_FinishNode(skiplist, context);
}
for (; finishlist; finishlist = next) {
next = finishlist->next;
finishlist->status = rf_good;
for (i = 0; i < finishlist->numAntecedents; i++) {
finishlist->antecedents[i]->numSuccFired++;
}
if (finishlist->commitNode)
finishlist->dagHdr->numCommits++;
rf_FinishNode(finishlist, context);
}
FireNodeList(firelist);
break;
case rf_rollBackward:
for (i = 0; i < node->numAntecedents; i++) {
a = *(node->antecedents + i);
RF_ASSERT(a->status == rf_good);
RF_ASSERT(a->numSuccDone <= a->numSuccedents);
RF_ASSERT(a->numSuccDone <= a->numSuccFired);
if (a->numSuccDone == a->numSuccFired) {
if (a->undoFunc == rf_NullNodeFunc) {
a->next = finishlist;
finishlist = a;
} else {
if (context != RF_INTR_CONTEXT) {
a->next = firelist;
firelist = a;
} else {
RF_ASSERT(NodeReady(a));
if (q) {
q->next = a;
q = a;
} else {
qh = q = a;
qh->next = NULL;
}
}
}
}
}
if (q) {
q->next = raidPtr->node_queue;
raidPtr->node_queue = qh;
DO_SIGNAL(raidPtr);
}
DO_UNLOCK(raidPtr);
for (; finishlist; finishlist = next) {
next = finishlist->next;
finishlist->status = rf_good;
rf_FinishNode(finishlist, context);
}
FireNodeList(firelist);
break;
default:
printf("Engine found illegal DAG status in PropagateResults()\n");
RF_PANIC();
break;
}
}
static void
ProcessNode(RF_DagNode_t *node, int context)
{
#if RF_DEBUG_ENGINE
RF_Raid_t *raidPtr;
raidPtr = node->dagHdr->raidPtr;
#endif
switch (node->status) {
case rf_good:
break;
case rf_bad:
if ((node->dagHdr->numCommits > 0) ||
(node->dagHdr->numCommitNodes == 0)) {
node->dagHdr->status = rf_rollForward;
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: node (%s) returned fail, rolling forward\n", raidPtr->raidid, node->name);
}
#endif
} else {
node->dagHdr->status = rf_rollBackward;
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: node (%s) returned fail, rolling backward\n", raidPtr->raidid, node->name);
}
#endif
}
break;
case rf_undone:
break;
case rf_panic:
printf("UNDO of a node failed!!!\n");
break;
default:
printf("node finished execution with an illegal status!!!\n");
RF_PANIC();
break;
}
PropagateResults(node, context);
}
void
rf_FinishNode(RF_DagNode_t *node, int context)
{
node->dagHdr->numNodesCompleted++;
ProcessNode(node, context);
}
int
rf_DispatchDAG(RF_DagHeader_t *dag, void (*cbFunc) (void *),
void *cbArg)
{
RF_Raid_t *raidPtr;
raidPtr = dag->raidPtr;
#if RF_ACC_TRACE > 0
if (dag->tracerec) {
RF_ETIMER_START(dag->tracerec->timer);
}
#endif
#if DEBUG
#if RF_DEBUG_VALIDATE_DAG
if (rf_engineDebug || rf_validateDAGDebug) {
if (rf_ValidateDAG(dag))
RF_PANIC();
}
#endif
#endif
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: Entering DispatchDAG\n", raidPtr->raidid);
}
#endif
raidPtr->dags_in_flight++;
dag->cbFunc = cbFunc;
dag->cbArg = cbArg;
dag->numNodesCompleted = 0;
dag->status = rf_enable;
FireNodeArray(dag->numSuccedents, dag->succedents);
return (1);
}
static void
DAGExecutionThread(RF_ThreadArg_t arg)
{
RF_DagNode_t *nd, *local_nq, *term_nq, *fire_nq;
RF_Raid_t *raidPtr;
raidPtr = (RF_Raid_t *) arg;
#if RF_DEBUG_ENGINE
if (rf_engineDebug) {
printf("raid%d: Engine thread is running\n", raidPtr->raidid);
}
#endif
DO_LOCK(raidPtr);
while (!raidPtr->shutdown_engine) {
while (raidPtr->node_queue != NULL) {
local_nq = raidPtr->node_queue;
fire_nq = NULL;
term_nq = NULL;
raidPtr->node_queue = NULL;
DO_UNLOCK(raidPtr);
while (local_nq) {
nd = local_nq;
local_nq = local_nq->next;
switch (nd->dagHdr->status) {
case rf_enable:
case rf_rollForward:
if (nd->numSuccedents == 0) {
nd->next = term_nq;
term_nq = nd;
} else {
nd->next = fire_nq;
fire_nq = nd;
}
break;
case rf_rollBackward:
if (nd->numAntecedents == 0) {
nd->next = term_nq;
term_nq = nd;
} else {
nd->next = fire_nq;
fire_nq = nd;
}
break;
default:
RF_PANIC();
break;
}
}
while (term_nq) {
nd = term_nq;
term_nq = term_nq->next;
nd->next = NULL;
(nd->dagHdr->cbFunc) (nd->dagHdr->cbArg);
raidPtr->dags_in_flight--;
}
FireNodeList(fire_nq);
DO_LOCK(raidPtr);
}
while (!raidPtr->shutdown_engine &&
raidPtr->node_queue == NULL) {
DO_WAIT(raidPtr);
}
}
raidPtr->shutdown_engine = 0;
DO_SIGNAL(raidPtr);
DO_UNLOCK(raidPtr);
kthread_exit(0);
}
static void
rf_RaidIOThread(RF_ThreadArg_t arg)
{
RF_Raid_t *raidPtr;
RF_DiskQueueData_t *req;
raidPtr = (RF_Raid_t *) arg;
rf_lock_mutex2(raidPtr->iodone_lock);
while (!raidPtr->shutdown_raidio) {
if (TAILQ_EMPTY(&(raidPtr->iodone)) &&
rf_buf_queue_check(raidPtr)) {
rf_wait_cond2(raidPtr->iodone_cv, raidPtr->iodone_lock);
}
if (raidPtr->parity_map != NULL) {
rf_unlock_mutex2(raidPtr->iodone_lock);
rf_paritymap_checkwork(raidPtr->parity_map);
rf_lock_mutex2(raidPtr->iodone_lock);
}
while ((req = TAILQ_FIRST(&(raidPtr->iodone))) != NULL) {
TAILQ_REMOVE(&(raidPtr->iodone), req, iodone_entries);
rf_unlock_mutex2(raidPtr->iodone_lock);
rf_DiskIOComplete(req->queue, req, req->error);
(req->CompleteFunc) (req->argument, req->error);
rf_lock_mutex2(raidPtr->iodone_lock);
}
rf_unlock_mutex2(raidPtr->iodone_lock);
raidstart(raidPtr);
rf_lock_mutex2(raidPtr->iodone_lock);
}
raidPtr->shutdown_raidio = 0;
rf_signal_cond2(raidPtr->iodone_cv);
rf_unlock_mutex2(raidPtr->iodone_lock);
kthread_exit(0);
}