scheduler.c 44.4 KB
Newer Older
H
refact  
Hongze Cheng 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
14 15
 */

D
dapan1121 已提交
16
#include "schedulerInt.h"
H
Hongze Cheng 已提交
17
#include "tmsg.h"
18
#include "query.h"
D
dapan 已提交
19
#include "catalog.h"
20

21
static SSchedulerMgmt schMgmt = {0};
D
dapan1121 已提交
22

D
dapan1121 已提交
23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52

int32_t schInitTask(SSchJob* pJob, SSchTask *pTask, SSubplan* pPlan, SSchLevel *pLevel) {
  pTask->plan   = pPlan;
  pTask->level  = pLevel;
  SCH_SET_TASK_STATUS(pTask, JOB_TASK_STATUS_NOT_START);
  pTask->taskId = atomic_add_fetch_64(&schMgmt.taskId, 1);
  pTask->execAddrs = taosArrayInit(SCH_MAX_CONDIDATE_EP_NUM, sizeof(SQueryNodeAddr));
  if (NULL == pTask->execAddrs) {
    SCH_TASK_ELOG("taosArrayInit %d exec addrs failed", SCH_MAX_CONDIDATE_EP_NUM);
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

  return TSDB_CODE_SUCCESS;
}

void schFreeTask(SSchTask* pTask) {
  if (pTask->candidateAddrs) {
    taosArrayDestroy(pTask->candidateAddrs);
  }

  // TODO NEED TO VERFY WITH ASYNC_SEND MEMORY FREE
  //tfree(pTask->msg); 

  if (pTask->children) {
    taosArrayDestroy(pTask->children);
  }

  if (pTask->parents) {
    taosArrayDestroy(pTask->parents);
  }
H
Haojun Liao 已提交
53 54 55 56

  if (pTask->execAddrs) {
    taosArrayDestroy(pTask->execAddrs);
  }
D
dapan1121 已提交
57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74
}


int32_t schValidateTaskReceivedMsgType(SSchJob *pJob, SSchTask *pTask, int32_t msgType) {
  int32_t lastMsgType = atomic_load_32(&pTask->lastMsgType);
  
  switch (msgType) {
    case TDMT_VND_CREATE_TABLE_RSP:
    case TDMT_VND_SUBMIT_RSP:
    case TDMT_VND_QUERY_RSP:
    case TDMT_VND_RES_READY_RSP:
    case TDMT_VND_FETCH_RSP:
    case TDMT_VND_DROP_TASK:
      if (lastMsgType != (msgType - 1)) {
        SCH_TASK_ELOG("rsp msg type mis-match, last sent msgType:%d, rspType:%d", lastMsgType, msgType);
        SCH_ERR_RET(TSDB_CODE_SCH_STATUS_ERROR);
      }
      
D
dapan1121 已提交
75
      if (SCH_GET_TASK_STATUS(pTask) != JOB_TASK_STATUS_EXECUTING && SCH_GET_TASK_STATUS(pTask) != JOB_TASK_STATUS_PARTIAL_SUCCEED) {
D
dapan1121 已提交
76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91
        SCH_TASK_ELOG("rsp msg conflicted with task status, status:%d, rspType:%d", SCH_GET_TASK_STATUS(pTask), msgType);
        SCH_ERR_RET(TSDB_CODE_SCH_STATUS_ERROR);
      }

      break;
    default:
      SCH_TASK_ELOG("unknown rsp msg, type:%d, status:%d", msgType, SCH_GET_TASK_STATUS(pTask));
      
      SCH_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
  }

  return TSDB_CODE_SUCCESS;
}


int32_t schCheckAndUpdateJobStatus(SSchJob *pJob, int8_t newStatus) {
D
dapan 已提交
92 93
  int32_t code = 0;

D
dapan1121 已提交
94
  int8_t oriStatus = 0;
D
dapan1121 已提交
95

D
dapan1121 已提交
96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143
  while (true) {
    oriStatus = SCH_GET_JOB_STATUS(pJob);

    if (oriStatus == newStatus) {
      SCH_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
    }
    
    switch (oriStatus) {
      case JOB_TASK_STATUS_NULL:
        if (newStatus != JOB_TASK_STATUS_NOT_START) {
          SCH_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
        }
        
        break;
      case JOB_TASK_STATUS_NOT_START:
        if (newStatus != JOB_TASK_STATUS_EXECUTING) {
          SCH_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
        }
        
        break;
      case JOB_TASK_STATUS_EXECUTING:
        if (newStatus != JOB_TASK_STATUS_PARTIAL_SUCCEED 
         && newStatus != JOB_TASK_STATUS_FAILED 
         && newStatus != JOB_TASK_STATUS_CANCELLING 
         && newStatus != JOB_TASK_STATUS_CANCELLED 
         && newStatus != JOB_TASK_STATUS_DROPPING) {
          SCH_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
        }
        
        break;
      case JOB_TASK_STATUS_PARTIAL_SUCCEED:
        if (newStatus != JOB_TASK_STATUS_FAILED 
         && newStatus != JOB_TASK_STATUS_SUCCEED
         && newStatus != JOB_TASK_STATUS_DROPPING) {
          SCH_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
        }
        
        break;
      case JOB_TASK_STATUS_SUCCEED:
      case JOB_TASK_STATUS_FAILED:
      case JOB_TASK_STATUS_CANCELLING:
        if (newStatus != JOB_TASK_STATUS_DROPPING) {
          SCH_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
        }
        
        break;
      case JOB_TASK_STATUS_CANCELLED:
      case JOB_TASK_STATUS_DROPPING:
D
dapan 已提交
144
        SCH_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
D
dapan1121 已提交
145 146 147 148
        break;
        
      default:
        SCH_JOB_ELOG("invalid job status:%d", oriStatus);
D
dapan 已提交
149
        SCH_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
D
dapan1121 已提交
150 151 152 153 154
    }

    if (oriStatus != atomic_val_compare_exchange_8(&pJob->status, oriStatus, newStatus)) {
      continue;
    }
D
dapan 已提交
155

D
dapan1121 已提交
156
    SCH_JOB_DLOG("job status updated from %d to %d", oriStatus, newStatus);
D
dapan1121 已提交
157

D
dapan1121 已提交
158 159
    break;
  }
D
dapan1121 已提交
160

D
dapan 已提交
161 162 163 164 165 166 167 168 169
  return TSDB_CODE_SUCCESS;

_return:

  SCH_JOB_ELOG("invalid job status update, from %d to %d", oriStatus, newStatus);
  SCH_ERR_RET(code);
}


D
dapan1121 已提交
170
int32_t schBuildTaskRalation(SSchJob *pJob, SHashObj *planToTask) {
D
dapan1121 已提交
171 172
  for (int32_t i = 0; i < pJob->levelNum; ++i) {
    SSchLevel *pLevel = taosArrayGet(pJob->levels, i);
D
dapan 已提交
173
    
D
dapan1121 已提交
174 175 176 177 178
    for (int32_t m = 0; m < pLevel->taskNum; ++m) {
      SSchTask *pTask = taosArrayGet(pLevel->subTasks, m);
      SSubplan *pPlan = pTask->plan;
      int32_t childNum = pPlan->pChildren ? (int32_t)taosArrayGetSize(pPlan->pChildren) : 0;
      int32_t parentNum = pPlan->pParents ? (int32_t)taosArrayGetSize(pPlan->pParents) : 0;
D
dapan1121 已提交
179 180

      if (childNum > 0) {
D
dapan1121 已提交
181 182 183 184 185
        if (pJob->levelIdx == pLevel->level) {
          SCH_JOB_ELOG("invalid query plan, lowest level, childNum:%d", childNum);
          SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
        }
        
D
dapan1121 已提交
186 187 188
        pTask->children = taosArrayInit(childNum, POINTER_BYTES);
        if (NULL == pTask->children) {
          SCH_TASK_ELOG("taosArrayInit %d children failed", childNum);
D
dapan1121 已提交
189 190 191 192 193
          SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
        }
      }

      for (int32_t n = 0; n < childNum; ++n) {
D
dapan1121 已提交
194
        SSubplan **child = taosArrayGet(pPlan->pChildren, n);
D
dapan 已提交
195
        SSchTask **childTask = taosHashGet(planToTask, child, POINTER_BYTES);
D
dapan 已提交
196
        if (NULL == childTask || NULL == *childTask) {
D
dapan1121 已提交
197
          SCH_TASK_ELOG("subplan children relationship error, level:%d, taskIdx:%d, childIdx:%d", i, m, n);
D
dapan1121 已提交
198 199 200
          SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
        }

D
dapan1121 已提交
201 202
        if (NULL == taosArrayPush(pTask->children, childTask)) {
          SCH_TASK_ELOG("taosArrayPush childTask failed, level:%d, taskIdx:%d, childIdx:%d", i, m, n);
D
dapan1121 已提交
203 204 205 206 207
          SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
        }
      }

      if (parentNum > 0) {
D
dapan1121 已提交
208 209 210 211 212
        if (0 == pLevel->level) {
          SCH_TASK_ELOG("invalid task info, level:0, parentNum:%d", parentNum);
          SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
        }
        
D
dapan1121 已提交
213 214 215
        pTask->parents = taosArrayInit(parentNum, POINTER_BYTES);
        if (NULL == pTask->parents) {
          SCH_TASK_ELOG("taosArrayInit %d parents failed", parentNum);
D
dapan1121 已提交
216 217
          SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
        }
D
dapan1121 已提交
218 219 220 221 222
      } else {
        if (0 != pLevel->level) {
          SCH_TASK_ELOG("invalid task info, level:%d, parentNum:%d", pLevel->level, parentNum);
          SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
        }
D
dapan1121 已提交
223 224 225
      }

      for (int32_t n = 0; n < parentNum; ++n) {
D
dapan1121 已提交
226
        SSubplan **parent = taosArrayGet(pPlan->pParents, n);
D
dapan 已提交
227
        SSchTask **parentTask = taosHashGet(planToTask, parent, POINTER_BYTES);
D
dapan 已提交
228
        if (NULL == parentTask || NULL == *parentTask) {
D
dapan1121 已提交
229
          SCH_TASK_ELOG("subplan parent relationship error, level:%d, taskIdx:%d, childIdx:%d", i, m, n);
D
dapan1121 已提交
230 231 232
          SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
        }

D
dapan1121 已提交
233 234
        if (NULL == taosArrayPush(pTask->parents, parentTask)) {
          SCH_TASK_ELOG("taosArrayPush parentTask failed, level:%d, taskIdx:%d, childIdx:%d", i, m, n);
D
dapan1121 已提交
235 236
          SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
        }
D
dapan1121 已提交
237 238 239
      }  

      SCH_TASK_DLOG("level:%d, parentNum:%d, childNum:%d", i, parentNum, childNum);
D
dapan1121 已提交
240 241 242
    }
  }

D
dapan1121 已提交
243 244 245
  SSchLevel *pLevel = taosArrayGet(pJob->levels, 0);
  if (pJob->attr.queryJob && pLevel->taskNum > 1) {
    SCH_JOB_ELOG("invalid query plan, level:0, taskNum:%d", pLevel->taskNum);
D
dapan 已提交
246 247 248
    SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
  }

D
dapan1121 已提交
249 250 251
  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
252 253 254 255 256 257 258

int32_t schRecordTaskSucceedNode(SSchTask *pTask) {
  SQueryNodeAddr *addr = taosArrayGet(pTask->candidateAddrs, atomic_load_8(&pTask->candidateIdx));

  assert(NULL != addr);

  pTask->succeedAddr = *addr;
259

D
dapan1121 已提交
260
  return TSDB_CODE_SUCCESS;
261 262
}

D
dapan1121 已提交
263 264 265 266 267 268 269 270

int32_t schRecordTaskExecNode(SSchJob *pJob, SSchTask *pTask, SQueryNodeAddr *addr) {
  if (NULL == taosArrayPush(pTask->execAddrs, addr)) {
    SCH_TASK_ELOG("taosArrayPush addr to execAddr list failed, errno:%d", errno);
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

  return TSDB_CODE_SUCCESS;
271
}
D
dapan1121 已提交
272

D
dapan1121 已提交
273

D
dapan1121 已提交
274
int32_t schValidateAndBuildJob(SQueryDag *pDag, SSchJob *pJob) {
D
dapan1121 已提交
275
  int32_t code = 0;
D
dapan1121 已提交
276
  pJob->queryId = pDag->queryId;
D
dapan1121 已提交
277
  
D
dapan1121 已提交
278 279
  if (pDag->numOfSubplans <= 0) {
    SCH_JOB_ELOG("invalid subplan num:%d", pDag->numOfSubplans);
D
dapan 已提交
280 281 282
    SCH_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
  }
  
D
dapan1121 已提交
283
  int32_t levelNum = (int32_t)taosArrayGetSize(pDag->pSubplans);
284
  if (levelNum <= 0) {
D
dapan1121 已提交
285
    SCH_JOB_ELOG("invalid level num:%d", levelNum);
D
dapan1121 已提交
286
    SCH_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
287 288
  }

D
dapan1121 已提交
289 290
  SHashObj *planToTask = taosHashInit(SCHEDULE_DEFAULT_TASK_NUMBER, taosGetDefaultHashFunction(POINTER_BYTES == sizeof(int64_t) ? TSDB_DATA_TYPE_BIGINT : TSDB_DATA_TYPE_INT), false, HASH_NO_LOCK);
  if (NULL == planToTask) {
D
dapan1121 已提交
291
    SCH_JOB_ELOG("taosHashInit %d failed", SCHEDULE_DEFAULT_TASK_NUMBER);
D
dapan1121 已提交
292 293
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }
294

D
dapan1121 已提交
295 296
  pJob->levels = taosArrayInit(levelNum, sizeof(SSchLevel));
  if (NULL == pJob->levels) {
D
dapan1121 已提交
297
    SCH_JOB_ELOG("taosArrayInit %d failed", levelNum);
D
dapan1121 已提交
298
    SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
299 300
  }

D
dapan1121 已提交
301 302
  pJob->levelNum = levelNum;
  pJob->levelIdx = levelNum - 1;
303

D
dapan1121 已提交
304
  pJob->subPlans = pDag->pSubplans;
305

D
dapan 已提交
306
  SSchLevel level = {0};
D
dapan1121 已提交
307 308
  SArray *plans = NULL;
  int32_t taskNum = 0;
D
dapan 已提交
309
  SSchLevel *pLevel = NULL;
310

D
dapan1121 已提交
311
  level.status = JOB_TASK_STATUS_NOT_START;
D
dapan1121 已提交
312

313
  for (int32_t i = 0; i < levelNum; ++i) {
D
dapan1121 已提交
314
    if (NULL == taosArrayPush(pJob->levels, &level)) {
D
dapan1121 已提交
315
      SCH_JOB_ELOG("taosArrayPush level failed, level:%d", i);
D
dapan1121 已提交
316 317 318
      SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
    }

D
dapan1121 已提交
319
    pLevel = taosArrayGet(pJob->levels, i);
D
dapan1121 已提交
320 321
  
    pLevel->level = i;
D
dapan1121 已提交
322
    
D
dapan1121 已提交
323
    plans = taosArrayGetP(pDag->pSubplans, i);
D
dapan1121 已提交
324 325
    if (NULL == plans) {
      SCH_JOB_ELOG("empty level plan, level:%d", i);
D
dapan1121 已提交
326
      SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
327 328
    }

D
dapan1121 已提交
329 330 331
    taskNum = (int32_t)taosArrayGetSize(plans);
    if (taskNum <= 0) {
      SCH_JOB_ELOG("invalid level plan number:%d, level:%d", taskNum, i);
D
dapan1121 已提交
332
      SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
D
dapan1121 已提交
333 334
    }

D
dapan1121 已提交
335
    pLevel->taskNum = taskNum;
D
dapan1121 已提交
336
    
D
dapan1121 已提交
337
    pLevel->subTasks = taosArrayInit(taskNum, sizeof(SSchTask));
D
dapan1121 已提交
338
    if (NULL == pLevel->subTasks) {
D
dapan1121 已提交
339
      SCH_JOB_ELOG("taosArrayInit %d failed", taskNum);
D
dapan1121 已提交
340
      SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
341 342
    }
    
D
dapan1121 已提交
343 344
    for (int32_t n = 0; n < taskNum; ++n) {
      SSubplan *plan = taosArrayGetP(plans, n);
D
dapan 已提交
345

D
dapan1121 已提交
346 347 348 349 350
      SCH_SET_JOB_TYPE(&pJob->attr, plan->type);

      SSchTask  task = {0};
      SSchTask *pTask = &task;
      
D
dapan1121 已提交
351
      SCH_ERR_JRET(schInitTask(pJob, &task, plan, pLevel));
D
dapan1121 已提交
352
      
D
dapan1121 已提交
353
      void *p = taosArrayPush(pLevel->subTasks, &task);
D
dapan1121 已提交
354
      if (NULL == p) {
D
dapan1121 已提交
355
        SCH_TASK_ELOG("taosArrayPush task to level failed, level:%d, taskIdx:%d", pLevel->level, n);
D
dapan1121 已提交
356 357 358 359
        SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
      }
      
      if (0 != taosHashPut(planToTask, &plan, POINTER_BYTES, &p, POINTER_BYTES)) {
D
dapan1121 已提交
360
        SCH_TASK_ELOG("taosHashPut to planToTaks failed, taskIdx:%d", n);
D
dapan1121 已提交
361 362
        SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
      }
363

D
dapan1121 已提交
364 365
      SCH_TASK_DLOG("task initialized, level:%d", pLevel->level);
    }
D
dapan1121 已提交
366

D
dapan1121 已提交
367
    SCH_JOB_DLOG("level initialized, taskNum:%d", taskNum);
D
dapan1121 已提交
368
  }
D
dapan1121 已提交
369 370

  SCH_ERR_JRET(schBuildTaskRalation(pJob, planToTask));
D
dapan1121 已提交
371 372

_return:
D
dapan1121 已提交
373 374 375 376
  if (planToTask) {
    taosHashCleanup(planToTask);
  }

D
dapan1121 已提交
377
  SCH_RET(code);
378 379
}

D
dapan1121 已提交
380 381
int32_t schSetTaskCandidateAddrs(SSchJob *pJob, SSchTask *pTask) {
  if (NULL != pTask->candidateAddrs) {
D
dapan 已提交
382 383 384
    return TSDB_CODE_SUCCESS;
  }

D
dapan1121 已提交
385 386 387 388
  pTask->candidateIdx = 0;
  pTask->candidateAddrs = taosArrayInit(SCH_MAX_CONDIDATE_EP_NUM, sizeof(SQueryNodeAddr));
  if (NULL == pTask->candidateAddrs) {
    SCH_TASK_ELOG("taosArrayInit %d condidate addrs failed", SCH_MAX_CONDIDATE_EP_NUM);
D
dapan1121 已提交
389 390 391
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
392 393 394
  if (pTask->plan->execNode.numOfEps > 0) {
    if (NULL == taosArrayPush(pTask->candidateAddrs, &pTask->plan->execNode)) {
      SCH_TASK_ELOG("taosArrayPush execNode to candidate addrs failed, errno:%d", errno);
D
dapan1121 已提交
395 396 397
      SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
    }

D
dapan1121 已提交
398
    SCH_TASK_DLOG("use execNode from plan as candidate addr, numOfEps:%d", pTask->plan->execNode.numOfEps);
D
dapan1121 已提交
399

D
dapan1121 已提交
400 401 402 403
    return TSDB_CODE_SUCCESS;
  }

  int32_t addNum = 0;
D
dapan 已提交
404
  int32_t nodeNum = 0;
405
  if (pJob->nodeList) {
D
dapan 已提交
406
    nodeNum = taosArrayGetSize(pJob->nodeList);
D
dapan1121 已提交
407
    
D
dapan 已提交
408 409 410 411 412 413 414
    for (int32_t i = 0; i < nodeNum && addNum < SCH_MAX_CONDIDATE_EP_NUM; ++i) {
      SQueryNodeAddr *naddr = taosArrayGet(pJob->nodeList, i);
      
      if (NULL == taosArrayPush(pTask->candidateAddrs, &pTask->plan->execNode)) {
        SCH_TASK_ELOG("taosArrayPush execNode to candidate addrs failed, addNum:%d, errno:%d", addNum, errno);
        SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
      }
D
dapan1121 已提交
415
    }
D
dapan1121 已提交
416 417
  }

D
dapan1121 已提交
418 419 420 421 422
  if (addNum <= 0) {
    SCH_TASK_ELOG("no available execNode as condidate addr, nodeNum:%d", nodeNum);
    return TSDB_CODE_QRY_INVALID_INPUT;
  }

D
dapan1121 已提交
423 424
/*
  for (int32_t i = 0; i < job->dataSrcEps.numOfEps && addNum < SCH_MAX_CONDIDATE_EP_NUM; ++i) {
D
dapan1121 已提交
425 426 427 428
    strncpy(epSet->fqdn[epSet->numOfEps], job->dataSrcEps.fqdn[i], sizeof(job->dataSrcEps.fqdn[i]));
    epSet->port[epSet->numOfEps] = job->dataSrcEps.port[i];
    
    ++epSet->numOfEps;
D
dapan 已提交
429
  }
D
dapan1121 已提交
430
*/
D
dapan 已提交
431 432

  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
433
}
D
dapan1121 已提交
434

D
dapan1121 已提交
435
int32_t schPushTaskToExecList(SSchJob *pJob, SSchTask *pTask) {
D
dapan1121 已提交
436 437 438
  int32_t code = taosHashPut(pJob->execTasks, &pTask->taskId, sizeof(pTask->taskId), &pTask, POINTER_BYTES);
  if (0 != code) {
    if (HASH_NODE_EXIST(code)) {
D
dapan1121 已提交
439
      SCH_TASK_ELOG("task already in execTask list, code:%x", code);
D
dapan1121 已提交
440 441 442
      SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
    }
    
D
dapan1121 已提交
443
    SCH_TASK_ELOG("taosHashPut task to execTask list failed, errno:%d", errno);
D
dapan 已提交
444 445 446
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
447
  SCH_TASK_DLOG("task added to execTask list, numOfTasks:%d", taosHashGetSize(pJob->execTasks));
448

D
dapan 已提交
449 450 451
  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
452 453 454
int32_t schMoveTaskToSuccList(SSchJob *pJob, SSchTask *pTask, bool *moved) {
  if (0 != taosHashRemove(pJob->execTasks, &pTask->taskId, sizeof(pTask->taskId))) {
    SCH_TASK_WLOG("remove task from execTask list failed, may not exist, status:%d", SCH_GET_TASK_STATUS(pTask));
D
dapan 已提交
455 456
  }

D
dapan1121 已提交
457 458 459 460 461 462 463 464 465 466
  int32_t code = taosHashPut(pJob->succTasks, &pTask->taskId, sizeof(pTask->taskId), &pTask, POINTER_BYTES);
  if (0 != code) {
    if (HASH_NODE_EXIST(code)) {
      *moved = true;
      
      SCH_TASK_ELOG("task already in succTask list, status:%d", SCH_GET_TASK_STATUS(pTask));
      SCH_ERR_RET(TSDB_CODE_SCH_STATUS_ERROR);
    }
  
    SCH_TASK_ELOG("taosHashPut task to succTask list failed, errno:%d", errno);
D
dapan1121 已提交
467 468 469 470
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

  *moved = true;
D
dapan1121 已提交
471 472

  SCH_TASK_DLOG("task moved to succTask list, numOfTasks:%d", taosHashGetSize(pJob->succTasks));
D
dapan1121 已提交
473 474 475 476
  
  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
477 478 479 480 481
int32_t schMoveTaskToFailList(SSchJob *pJob, SSchTask *pTask, bool *moved) {
  *moved = false;
  
  if (0 != taosHashRemove(pJob->execTasks, &pTask->taskId, sizeof(pTask->taskId))) {
    SCH_TASK_WLOG("remove task from execTask list failed, may not exist, status:%d", SCH_GET_TASK_STATUS(pTask));
D
dapan1121 已提交
482 483
  }

D
dapan1121 已提交
484 485 486 487 488 489 490 491 492 493
  int32_t code = taosHashPut(pJob->failTasks, &pTask->taskId, sizeof(pTask->taskId), &pTask, POINTER_BYTES);
  if (0 != code) {
    if (HASH_NODE_EXIST(code)) {
      *moved = true;
      
      SCH_TASK_WLOG("task already in failTask list, status:%d", SCH_GET_TASK_STATUS(pTask));
      SCH_ERR_RET(TSDB_CODE_SCH_STATUS_ERROR);
    }
    
    SCH_TASK_ELOG("taosHashPut task to failTask list failed, errno:%d", errno);
D
dapan 已提交
494 495 496 497
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

  *moved = true;
D
dapan1121 已提交
498 499

  SCH_TASK_DLOG("task moved to failTask list, numOfTasks:%d", taosHashGetSize(pJob->failTasks));
D
dapan 已提交
500 501 502 503
  
  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530

int32_t schMoveTaskToExecList(SSchJob *pJob, SSchTask *pTask, bool *moved) {
  if (0 != taosHashRemove(pJob->succTasks, &pTask->taskId, sizeof(pTask->taskId))) {
    SCH_TASK_WLOG("remove task from succTask list failed, may not exist, status:%d", SCH_GET_TASK_STATUS(pTask));
  }

  int32_t code = taosHashPut(pJob->execTasks, &pTask->taskId, sizeof(pTask->taskId), &pTask, POINTER_BYTES);
  if (0 != code) {
    if (HASH_NODE_EXIST(code)) {
      *moved = true;
      
      SCH_TASK_ELOG("task already in execTask list, status:%d", SCH_GET_TASK_STATUS(pTask));
      SCH_ERR_RET(TSDB_CODE_SCH_STATUS_ERROR);
    }
  
    SCH_TASK_ELOG("taosHashPut task to execTask list failed, errno:%d", errno);
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

  *moved = true;

  SCH_TASK_DLOG("task moved to execTask list, numOfTasks:%d", taosHashGetSize(pJob->execTasks));
  
  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
531
int32_t schTaskCheckAndSetRetry(SSchJob *job, SSchTask *task, int32_t errCode, bool *needRetry) {
D
dapan1121 已提交
532 533
  // TODO set retry or not based on task type/errCode/retry times/job status/available eps...
  // TODO if needRetry, set task retry info
D
dapan1121 已提交
534 535
  // TODO set condidateIdx
  // TODO record failed but tried task
D
dapan 已提交
536

D
dapan1121 已提交
537 538 539
  *needRetry = false;

  return TSDB_CODE_SUCCESS;
D
dapan 已提交
540 541
}

D
dapan 已提交
542

D
dapan1121 已提交
543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563

// Note: no more error processing, handled in function internal
int32_t schProcessOnJobFailure(SSchJob *pJob, int32_t errCode) {
  // if already FAILED, no more processing
  SCH_ERR_RET(schCheckAndUpdateJobStatus(pJob, JOB_TASK_STATUS_FAILED));
  
  if (errCode) {
    atomic_store_32(&pJob->errCode, errCode);
  }

  if (atomic_load_8(&pJob->userFetch) || ((!SCH_JOB_NEED_FETCH(&pJob->attr)) && pJob->attr.syncSchedule)) {
    tsem_post(&pJob->rspSem);
  }

  SCH_ERR_RET(atomic_load_32(&pJob->errCode));

  assert(0);
}



D
dapan1121 已提交
564
// Note: no more error processing, handled in function internal
D
dapan1121 已提交
565
int32_t schFetchFromRemote(SSchJob *pJob) {
D
dapan 已提交
566 567
  int32_t code = 0;
  
D
dapan1121 已提交
568 569
  if (atomic_val_compare_exchange_32(&pJob->remoteFetch, 0, 1) != 0) {
    SCH_JOB_ELOG("prior fetching not finished, remoteFetch:%d", atomic_load_32(&pJob->remoteFetch));
D
dapan 已提交
570 571 572
    return TSDB_CODE_SUCCESS;
  }

D
dapan1121 已提交
573 574 575 576 577 578 579
  void *res = atomic_load_ptr(&pJob->res);
  if (res) {
    atomic_val_compare_exchange_32(&pJob->remoteFetch, 1, 0);

    SCH_JOB_DLOG("res already fetched, res:%p", res);
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
580 581

  SCH_ERR_JRET(schBuildAndSendMsg(pJob, pJob->fetchTask, &pJob->resNode, TDMT_VND_FETCH));
D
dapan 已提交
582 583 584 585

  return TSDB_CODE_SUCCESS;
  
_return:
D
dapan1121 已提交
586 587

  atomic_val_compare_exchange_32(&pJob->remoteFetch, 1, 0);
D
dapan 已提交
588

D
dapan1121 已提交
589 590
  schProcessOnJobFailure(pJob, code);

D
dapan 已提交
591
  return code;
D
dapan 已提交
592 593
}

D
dapan 已提交
594

D
dapan1121 已提交
595 596 597
// Note: no more error processing, handled in function internal
int32_t schProcessOnJobPartialSuccess(SSchJob *pJob) {
  int32_t code = 0;
D
dapan 已提交
598
  
D
dapan1121 已提交
599
  SCH_ERR_RET(schCheckAndUpdateJobStatus(pJob, JOB_TASK_STATUS_PARTIAL_SUCCEED));
D
dapan 已提交
600

D
dapan1121 已提交
601
  if ((!SCH_JOB_NEED_FETCH(&pJob->attr)) && pJob->attr.syncSchedule) {
D
dapan 已提交
602
    tsem_post(&pJob->rspSem);
D
dapan 已提交
603
  }
D
dapan1121 已提交
604 605 606 607
  
  if (atomic_load_8(&pJob->userFetch)) {
    SCH_ERR_JRET(schFetchFromRemote(pJob));
  }
D
dapan 已提交
608

D
dapan 已提交
609
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
610 611 612 613 614 615

_return:

  SCH_ERR_RET(schProcessOnJobFailure(pJob, code));

  SCH_RET(code);
D
dapan 已提交
616 617
}

D
dapan1121 已提交
618

D
dapan1121 已提交
619
int32_t schProcessOnDataFetched(SSchJob *job) {
D
dapan 已提交
620 621 622
  atomic_val_compare_exchange_32(&job->remoteFetch, 1, 0);

  tsem_post(&job->rspSem);
D
dapan 已提交
623 624 625
}


D
dapan1121 已提交
626 627 628 629

// Note: no more error processing, handled in function internal
int32_t schProcessOnTaskFailure(SSchJob *pJob, SSchTask *pTask, int32_t errCode) {
  bool needRetry = false;
D
dapan 已提交
630
  bool moved = false;
D
dapan1121 已提交
631
  int32_t taskDone = 0;
D
dapan1121 已提交
632
  int32_t code = 0;
D
dapan1121 已提交
633 634

  SCH_TASK_DLOG("taskOnFailure, code:%x", errCode);
D
dapan 已提交
635
  
D
dapan1121 已提交
636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666
  SCH_ERR_JRET(schTaskCheckAndSetRetry(pJob, pTask, errCode, &needRetry));
  
  if (!needRetry) {
    SCH_TASK_ELOG("task failed and no more retry, code:%x", errCode);

    if (SCH_GET_TASK_STATUS(pTask) == JOB_TASK_STATUS_EXECUTING) {
      code = schMoveTaskToFailList(pJob, pTask, &moved);
      if (code && moved) {
        SCH_ERR_RET(errCode);
      }
    }

    SCH_SET_TASK_STATUS(pTask, JOB_TASK_STATUS_FAILED);
    
    if (SCH_TASK_NEED_WAIT_ALL(pTask)) {
      SCH_LOCK(SCH_WRITE, &pTask->level->lock);
      pTask->level->taskFailed++;
      taskDone = pTask->level->taskSucceed + pTask->level->taskFailed;
      SCH_UNLOCK(SCH_WRITE, &pTask->level->lock);

      atomic_store_32(&pJob->errCode, errCode);
      
      if (taskDone < pTask->level->taskNum) {
        SCH_TASK_DLOG("not all tasks done, done:%d, all:%d", taskDone, pTask->level->taskNum);
        SCH_ERR_RET(errCode);
      }
    }
  } else {
    // Note: no more error processing, already handled
    SCH_ERR_RET(schLaunchTask(pJob, pTask));
    
D
dapan 已提交
667 668
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
669

D
dapan1121 已提交
670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690
_return:

  SCH_ERR_RET(schProcessOnJobFailure(pJob, errCode));

  SCH_ERR_RET(errCode);
}


// Note: no more error processing, handled in function internal
int32_t schProcessOnTaskSuccess(SSchJob *pJob, SSchTask *pTask) {
  bool moved = false;
  int32_t code = 0;
  SSchTask *pErrTask = pTask;

  SCH_TASK_DLOG("taskOnSuccess, status:%d", SCH_GET_TASK_STATUS(pTask));
  
  code = schMoveTaskToSuccList(pJob, pTask, &moved);
  if (code && moved) {
    SCH_ERR_RET(code);
  }

D
dapan1121 已提交
691
  SCH_SET_TASK_STATUS(pTask, JOB_TASK_STATUS_PARTIAL_SUCCEED);
D
dapan1121 已提交
692 693

  SCH_ERR_JRET(schRecordTaskSucceedNode(pTask));
D
dapan 已提交
694
  
D
dapan1121 已提交
695
  int32_t parentNum = pTask->parents ? (int32_t)taosArrayGetSize(pTask->parents) : 0;
D
dapan 已提交
696
  if (parentNum == 0) {
D
dapan1121 已提交
697 698
    int32_t taskDone = 0;
    
D
dapan1121 已提交
699 700 701 702 703
    if (SCH_TASK_NEED_WAIT_ALL(pTask)) {
      SCH_LOCK(SCH_WRITE, &pTask->level->lock);
      pTask->level->taskSucceed++;
      taskDone = pTask->level->taskSucceed + pTask->level->taskFailed;
      SCH_UNLOCK(SCH_WRITE, &pTask->level->lock);
D
dapan1121 已提交
704
      
D
dapan1121 已提交
705 706
      if (taskDone < pTask->level->taskNum) {
        SCH_TASK_ELOG("wait all tasks, done:%d, all:%d", taskDone, pTask->level->taskNum);
D
dapan1121 已提交
707
        
D
dapan1121 已提交
708
        return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
709 710
      } else if (taskDone > pTask->level->taskNum) {
        assert(0);
D
dapan1121 已提交
711 712
      }

D
dapan1121 已提交
713
      if (pTask->level->taskFailed > 0) {
D
dapan1121 已提交
714 715 716
        SCH_RET(schProcessOnJobFailure(pJob, 0));
      } else {
        SCH_RET(schProcessOnJobPartialSuccess(pJob));
D
dapan1121 已提交
717 718
      }
    } else {
D
dapan1121 已提交
719
      pJob->resNode = pTask->succeedAddr;
D
dapan1121 已提交
720
    }
D
dapan 已提交
721

D
dapan1121 已提交
722
    pJob->fetchTask = pTask;
D
dapan1121 已提交
723 724 725 726 727

    code = schMoveTaskToExecList(pJob, pTask, &moved);
    if (code && moved) {
      SCH_ERR_RET(code);
    }
D
dapan1121 已提交
728
    
D
dapan1121 已提交
729
    SCH_ERR_RET(schProcessOnJobPartialSuccess(pJob));
D
dapan 已提交
730 731 732 733

    return TSDB_CODE_SUCCESS;
  }

D
dapan1121 已提交
734
/*
D
dapan 已提交
735
  if (SCH_IS_DATA_SRC_TASK(task) && job->dataSrcEps.numOfEps < SCH_MAX_CONDIDATE_EP_NUM) {
D
dapan 已提交
736 737 738 739 740
    strncpy(job->dataSrcEps.fqdn[job->dataSrcEps.numOfEps], task->execAddr.fqdn, sizeof(task->execAddr.fqdn));
    job->dataSrcEps.port[job->dataSrcEps.numOfEps] = task->execAddr.port;

    ++job->dataSrcEps.numOfEps;
  }
D
dapan1121 已提交
741
*/
D
dapan 已提交
742

D
dapan 已提交
743
  for (int32_t i = 0; i < parentNum; ++i) {
D
dapan1121 已提交
744
    SSchTask *par = *(SSchTask **)taosArrayGet(pTask->parents, i);
D
dapan1121 已提交
745 746
    pErrTask = par;
    
D
dapan1121 已提交
747
    atomic_add_fetch_32(&par->childReady, 1);
D
dapan 已提交
748

D
dapan1121 已提交
749
    SCH_LOCK(SCH_WRITE, &par->lock);
D
dapan1121 已提交
750
    qSetSubplanExecutionNode(par->plan, pTask->plan->id.templateId, &pTask->succeedAddr);
D
dapan1121 已提交
751
    SCH_UNLOCK(SCH_WRITE, &par->lock);
D
dapan 已提交
752 753
    
    if (SCH_TASK_READY_TO_LUNCH(par)) {
D
dapan1121 已提交
754
      SCH_ERR_RET(schLaunchTask(pJob, par));
D
dapan 已提交
755 756 757 758 759
    }
  }

  return TSDB_CODE_SUCCESS;

D
dapan1121 已提交
760
_return:
D
dapan 已提交
761

D
dapan1121 已提交
762
  SCH_ERR_RET(schProcessOnTaskFailure(pJob, pErrTask, code));
D
dapan 已提交
763

D
dapan1121 已提交
764
  SCH_ERR_RET(code);
D
dapan 已提交
765 766
}

D
dapan1121 已提交
767
int32_t schHandleResponseMsg(SSchJob *pJob, SSchTask *pTask, int32_t msgType, char *msg, int32_t msgSize, int32_t rspCode) {
D
dapan1121 已提交
768
  int32_t code = 0;
H
Haojun Liao 已提交
769

D
dapan1121 已提交
770 771
  SCH_ERR_JRET(schValidateTaskReceivedMsgType(pJob, pTask, msgType));

D
dapan1121 已提交
772
  switch (msgType) {
H
Haojun Liao 已提交
773
    case TDMT_VND_CREATE_TABLE_RSP: {
D
dapan1121 已提交
774 775
        if (rspCode != TSDB_CODE_SUCCESS) {
          SCH_ERR_RET(schProcessOnTaskFailure(pJob, pTask, rspCode));
H
Haojun Liao 已提交
776
        }
D
dapan1121 已提交
777 778
        
        SCH_ERR_RET(schProcessOnTaskSuccess(pJob, pTask));
D
dapan1121 已提交
779

D
dapan1121 已提交
780 781
        break;
      }
D
dapan1121 已提交
782
    case TDMT_VND_SUBMIT_RSP: {
D
dapan1121 已提交
783 784 785 786 787 788 789 790 791 792 793
        #if 0 //TODO OPEN THIS
        SShellSubmitRspMsg *rsp = (SShellSubmitRspMsg *)msg;

        if (rspCode != TSDB_CODE_SUCCESS || NULL == msg || rsp->code != TSDB_CODE_SUCCESS) {
          SCH_ERR_RET(schProcessOnTaskFailure(pJob, pTask, rspCode));
        }

        pJob->resNumOfRows += rsp->affectedRows;
        #else
        if (rspCode != TSDB_CODE_SUCCESS) {
          SCH_ERR_RET(schProcessOnTaskFailure(pJob, pTask, rspCode));
D
dapan1121 已提交
794
        }
D
dapan1121 已提交
795 796 797 798
        #endif

        SCH_ERR_RET(schProcessOnTaskSuccess(pJob, pTask));

D
dapan1121 已提交
799
        break;
D
dapan1121 已提交
800
      }
D
dapan1121 已提交
801
    case TDMT_VND_QUERY_RSP: {
D
dapan1121 已提交
802 803
        SQueryTableRsp *rsp = (SQueryTableRsp *)msg;
        
D
dapan1121 已提交
804
        if (rspCode != TSDB_CODE_SUCCESS || NULL == msg || rsp->code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
805
          SCH_ERR_RET(schProcessOnTaskFailure(pJob, pTask, rspCode));
D
dapan1121 已提交
806
        }
D
dapan1121 已提交
807 808 809
        
        SCH_ERR_JRET(schBuildAndSendMsg(pJob, pTask, NULL, TDMT_VND_RES_READY));
        
D
dapan1121 已提交
810 811
        break;
      }
D
dapan1121 已提交
812
    case TDMT_VND_RES_READY_RSP: {
D
dapan1121 已提交
813 814
        SResReadyRsp *rsp = (SResReadyRsp *)msg;
        
D
dapan1121 已提交
815 816
        if (rspCode != TSDB_CODE_SUCCESS || NULL == msg || rsp->code != TSDB_CODE_SUCCESS) {
          SCH_ERR_RET(schProcessOnTaskFailure(pJob, pTask, rsp->code));
D
dapan1121 已提交
817
        }
D
dapan1121 已提交
818 819 820
        
        SCH_ERR_RET(schProcessOnTaskSuccess(pJob, pTask));
        
D
dapan1121 已提交
821 822
        break;
      }
D
dapan1121 已提交
823
    case TDMT_VND_FETCH_RSP: {
D
dapan1121 已提交
824 825
        SRetrieveTableRsp *rsp = (SRetrieveTableRsp *)msg;

D
dapan1121 已提交
826 827
        if (rspCode != TSDB_CODE_SUCCESS || NULL == msg) {
          SCH_ERR_RET(schProcessOnTaskFailure(pJob, pTask, rspCode));
D
dapan1121 已提交
828
        }
D
dapan1121 已提交
829 830 831 832 833

        atomic_store_ptr(&pJob->res, rsp);
        atomic_store_32(&pJob->resNumOfRows, rsp->numOfRows);
        
        SCH_ERR_JRET(schProcessOnDataFetched(pJob));
D
dapan1121 已提交
834 835
        
        break;
D
dapan1121 已提交
836
      }
D
dapan1121 已提交
837
    case TDMT_VND_DROP_TASK: {
D
dapan1121 已提交
838 839 840
        // SHOULD NEVER REACH HERE
        assert(0);
        break;
D
dapan1121 已提交
841
      }
D
dapan1121 已提交
842
    default:
D
dapan1121 已提交
843 844 845
      SCH_TASK_ELOG("unknown rsp msg, type:%d, status:%d", msgType, SCH_GET_TASK_STATUS(pTask));
      
      SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
D
dapan1121 已提交
846 847 848 849 850
  }

  return TSDB_CODE_SUCCESS;

_return:
D
dapan1121 已提交
851 852 853 854

  SCH_ERR_RET(schProcessOnTaskFailure(pJob, pTask, code));
  
  SCH_RET(code);
D
dapan1121 已提交
855 856
}

D
dapan 已提交
857

D
dapan1121 已提交
858
int32_t schHandleCallback(void* param, const SDataBuf* pMsg, int32_t msgType, int32_t rspCode) {
D
dapan1121 已提交
859 860
  int32_t code = 0;
  SSchCallbackParam *pParam = (SSchCallbackParam *)param;
D
dapan1121 已提交
861 862
  SSchJob *pJob = NULL;
  SSchTask *pTask = NULL;
D
dapan1121 已提交
863
  
D
dapan1121 已提交
864 865
  SSchJob **job = taosHashGet(schMgmt.jobs, &pParam->queryId, sizeof(pParam->queryId));
  if (NULL == job || NULL == (*job)) {
D
dapan1121 已提交
866
    qError("QID:%"PRIx64" taosHashGet queryId not exist, may be dropped", pParam->queryId);
D
dapan1121 已提交
867 868 869
    SCH_ERR_JRET(TSDB_CODE_SCH_INTERNAL_ERROR);
  }

D
dapan1121 已提交
870 871 872 873 874
  pJob = *job;

  atomic_add_fetch_32(&pJob->ref, 1);

  int32_t s = taosHashGetSize(pJob->execTasks);
D
dapan1121 已提交
875 876 877 878
  if (s <= 0) {
    qError("QID:%"PRIx64",TID:%"PRIx64" no task in execTask list", pParam->queryId, pParam->taskId);
    SCH_ERR_JRET(TSDB_CODE_SCH_INTERNAL_ERROR);
  }
H
Haojun Liao 已提交
879

D
dapan1121 已提交
880
  SSchTask **task = taosHashGet(pJob->execTasks, &pParam->taskId, sizeof(pParam->taskId));
D
dapan1121 已提交
881
  if (NULL == task || NULL == (*task)) {
D
dapan1121 已提交
882
    qError("QID:%"PRIx64",TID:%"PRIx64" taosHashGet taskId not exist", pParam->queryId, pParam->taskId);
D
dapan1121 已提交
883 884
    SCH_ERR_JRET(TSDB_CODE_SCH_INTERNAL_ERROR);
  }
D
dapan1121 已提交
885 886 887 888

  pTask = *task;

  SCH_TASK_DLOG("rsp msg received, type:%d, code:%x", msgType, rspCode);
D
dapan1121 已提交
889
  
D
dapan1121 已提交
890
  SCH_ERR_JRET(schHandleResponseMsg(pJob, pTask, msgType, pMsg->pData, pMsg->len, rspCode));
D
dapan1121 已提交
891

H
Haojun Liao 已提交
892
_return:
D
dapan1121 已提交
893 894 895 896 897

  if (pJob) {
    atomic_sub_fetch_32(&pJob->ref, 1);
  }

D
dapan1121 已提交
898 899 900 901
  tfree(param);
  SCH_RET(code);
}

D
dapan1121 已提交
902
int32_t schHandleSubmitCallback(void* param, const SDataBuf* pMsg, int32_t code) {
D
dapan1121 已提交
903 904
  return schHandleCallback(param, pMsg, TDMT_VND_SUBMIT_RSP, code);
}
H
Haojun Liao 已提交
905

D
dapan1121 已提交
906
int32_t schHandleCreateTableCallback(void* param, const SDataBuf* pMsg, int32_t code) {
H
Haojun Liao 已提交
907 908 909
  return schHandleCallback(param, pMsg, TDMT_VND_CREATE_TABLE_RSP, code);
}

D
dapan1121 已提交
910
int32_t schHandleQueryCallback(void* param, const SDataBuf* pMsg, int32_t code) {
D
dapan1121 已提交
911 912
  return schHandleCallback(param, pMsg, TDMT_VND_QUERY_RSP, code);
}
H
Haojun Liao 已提交
913

D
dapan1121 已提交
914
int32_t schHandleFetchCallback(void* param, const SDataBuf* pMsg, int32_t code) {
D
dapan1121 已提交
915 916
  return schHandleCallback(param, pMsg, TDMT_VND_FETCH_RSP, code);
}
H
Haojun Liao 已提交
917

D
dapan1121 已提交
918
int32_t schHandleReadyCallback(void* param, const SDataBuf* pMsg, int32_t code) {
D
dapan1121 已提交
919 920
  return schHandleCallback(param, pMsg, TDMT_VND_RES_READY_RSP, code);
}
H
Haojun Liao 已提交
921

D
dapan1121 已提交
922
int32_t schHandleDropCallback(void* param, const SDataBuf* pMsg, int32_t code) {  
D
dapan1121 已提交
923
  SSchCallbackParam *pParam = (SSchCallbackParam *)param;
D
dapan1121 已提交
924
  qDebug("QID:%"PRIx64",TID:%"PRIx64" drop task rsp received, code:%x", pParam->queryId, pParam->taskId, code);
D
dapan1121 已提交
925 926
}

D
dapan1121 已提交
927
int32_t schGetCallbackFp(int32_t msgType, __async_send_cb_fn_t *fp) {
D
dapan1121 已提交
928
  switch (msgType) {
H
Haojun Liao 已提交
929 930 931
    case TDMT_VND_CREATE_TABLE:
      *fp = schHandleCreateTableCallback;
      break;
D
dapan1121 已提交
932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947
    case TDMT_VND_SUBMIT: 
      *fp = schHandleSubmitCallback;
      break;
    case TDMT_VND_QUERY: 
      *fp = schHandleQueryCallback;
      break;
    case TDMT_VND_RES_READY: 
      *fp = schHandleReadyCallback;
      break;
    case TDMT_VND_FETCH: 
      *fp = schHandleFetchCallback;
      break;
    case TDMT_VND_DROP_TASK:
      *fp = schHandleDropCallback;
      break;
    default:
D
dapan1121 已提交
948
      qError("unknown msg type for callback, msgType:%d", msgType);
D
dapan1121 已提交
949 950 951 952 953 954 955
      SCH_ERR_RET(TSDB_CODE_QRY_APP_ERROR);
  }

  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
956
int32_t schAsyncSendMsg(void *transport, SEpSet* epSet, uint64_t qId, uint64_t tId, int32_t msgType, void *msg, uint32_t msgSize) {
D
dapan1121 已提交
957 958 959
  int32_t code = 0;
  SMsgSendInfo* pMsgSendInfo = calloc(1, sizeof(SMsgSendInfo));
  if (NULL == pMsgSendInfo) {
D
dapan 已提交
960
    qError("QID:%"PRIx64 ",TID:%"PRIx64 " calloc %d failed", qId, tId, (int32_t)sizeof(SMsgSendInfo));
D
dapan1121 已提交
961 962 963 964 965
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

  SSchCallbackParam *param = calloc(1, sizeof(SSchCallbackParam));
  if (NULL == param) {
D
dapan 已提交
966
    qError("QID:%"PRIx64 ",TID:%"PRIx64 " calloc %d failed", qId, tId, (int32_t)sizeof(SSchCallbackParam));
D
dapan1121 已提交
967 968 969 970 971 972 973 974 975 976 977 978 979 980 981 982
    SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

  __async_send_cb_fn_t fp = NULL;
  SCH_ERR_JRET(schGetCallbackFp(msgType, &fp));

  param->queryId = qId;
  param->taskId = tId;

  pMsgSendInfo->param = param;
  pMsgSendInfo->msgInfo.pData = msg;
  pMsgSendInfo->msgInfo.len = msgSize;
  pMsgSendInfo->msgType = msgType;
  pMsgSendInfo->fp = fp;
  
  int64_t  transporterId = 0;
D
dapan 已提交
983
  
D
dapan1121 已提交
984 985 986 987 988
  code = asyncSendMsgToServer(transport, epSet, &transporterId, pMsgSendInfo);
  if (code) {
    qError("QID:%"PRIx64 ",TID:%"PRIx64 " asyncSendMsgToServer failed, code:%x", qId, tId, code);
    SCH_ERR_JRET(code);
  }
D
dapan1121 已提交
989 990 991 992
  
  return TSDB_CODE_SUCCESS;

_return:
D
dapan 已提交
993
  
D
dapan1121 已提交
994 995 996 997 998 999
  tfree(param);
  tfree(pMsgSendInfo);

  SCH_RET(code);
}

D
dapan1121 已提交
1000
void schConvertAddrToEpSet(SQueryNodeAddr *addr, SEpSet *epSet) {
D
dapan1121 已提交
1001 1002 1003 1004 1005 1006 1007 1008 1009
  epSet->inUse = addr->inUse;
  epSet->numOfEps = addr->numOfEps;
  
  for (int8_t i = 0; i < epSet->numOfEps; ++i) {
    strncpy(epSet->fqdn[i], addr->epAddr[i].fqdn, sizeof(addr->epAddr[i].fqdn));
    epSet->port[i] = addr->epAddr[i].port;
  }
}

D
dapan1121 已提交
1010
int32_t schBuildAndSendMsg(SSchJob *pJob, SSchTask *pTask, SQueryNodeAddr *addr, int32_t msgType) {
D
dapan1121 已提交
1011 1012 1013
  uint32_t msgSize = 0;
  void *msg = NULL;
  int32_t code = 0;
D
dapan1121 已提交
1014
  bool isCandidateAddr = false;
D
dapan 已提交
1015
  SEpSet epSet;
D
dapan1121 已提交
1016 1017 1018 1019 1020 1021 1022

  if (NULL == addr) {
    addr = taosArrayGet(pTask->candidateAddrs, atomic_load_8(&pTask->candidateIdx));
    
    isCandidateAddr = true;
  }

D
dapan 已提交
1023
  schConvertAddrToEpSet(addr, &epSet);
D
dapan1121 已提交
1024 1025
  
  switch (msgType) {
H
Haojun Liao 已提交
1026
    case TDMT_VND_CREATE_TABLE:
D
dapan1121 已提交
1027
    case TDMT_VND_SUBMIT: {
D
dapan1121 已提交
1028 1029
      msgSize = pTask->msgLen;
      msg = pTask->msg;
D
dapan1121 已提交
1030 1031 1032
      break;
    }
    case TDMT_VND_QUERY: {
D
dapan1121 已提交
1033
      msgSize = sizeof(SSubQueryMsg) + pTask->msgLen;
D
dapan1121 已提交
1034 1035
      msg = calloc(1, msgSize);
      if (NULL == msg) {
D
dapan 已提交
1036
        SCH_TASK_ELOG("calloc %d failed", msgSize);
D
dapan1121 已提交
1037 1038 1039 1040
        SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
      }

      SSubQueryMsg *pMsg = msg;
D
dapan1121 已提交
1041

D
dapan 已提交
1042 1043
      pMsg->header.vgId = htonl(addr->nodeId);
      
D
dapan1121 已提交
1044
      pMsg->sId = htobe64(schMgmt.sId);
D
dapan1121 已提交
1045 1046 1047 1048
      pMsg->queryId = htobe64(pJob->queryId);
      pMsg->taskId = htobe64(pTask->taskId);
      pMsg->contentLen = htonl(pTask->msgLen);
      memcpy(pMsg->msg, pTask->msg, pTask->msgLen);
D
dapan1121 已提交
1049 1050 1051
      break;
    }    
    case TDMT_VND_RES_READY: {
S
Shengliang Guan 已提交
1052
      msgSize = sizeof(SResReadyReq);
D
dapan1121 已提交
1053 1054
      msg = calloc(1, msgSize);
      if (NULL == msg) {
D
dapan 已提交
1055
        SCH_TASK_ELOG("calloc %d failed", msgSize);
D
dapan1121 已提交
1056 1057 1058
        SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
      }

S
Shengliang Guan 已提交
1059
      SResReadyReq *pMsg = msg;
D
dapan1121 已提交
1060
      
D
dapan 已提交
1061 1062
      pMsg->header.vgId = htonl(addr->nodeId);  
      
D
dapan1121 已提交
1063
      pMsg->sId = htobe64(schMgmt.sId);      
D
dapan1121 已提交
1064 1065
      pMsg->queryId = htobe64(pJob->queryId);
      pMsg->taskId = htobe64(pTask->taskId);      
D
dapan1121 已提交
1066 1067 1068
      break;
    }
    case TDMT_VND_FETCH: {
S
Shengliang Guan 已提交
1069
      msgSize = sizeof(SResFetchReq);
D
dapan1121 已提交
1070 1071
      msg = calloc(1, msgSize);
      if (NULL == msg) {
D
dapan 已提交
1072
        SCH_TASK_ELOG("calloc %d failed", msgSize);
D
dapan1121 已提交
1073 1074 1075
        SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
      }
    
S
Shengliang Guan 已提交
1076
      SResFetchReq *pMsg = msg;
D
dapan1121 已提交
1077
      
D
dapan 已提交
1078 1079
      pMsg->header.vgId = htonl(addr->nodeId);  
      
D
dapan1121 已提交
1080
      pMsg->sId = htobe64(schMgmt.sId);      
D
dapan1121 已提交
1081 1082
      pMsg->queryId = htobe64(pJob->queryId);
      pMsg->taskId = htobe64(pTask->taskId);      
D
dapan1121 已提交
1083 1084 1085
      break;
    }
    case TDMT_VND_DROP_TASK:{
S
Shengliang Guan 已提交
1086
      msgSize = sizeof(STaskDropReq);
D
dapan1121 已提交
1087 1088
      msg = calloc(1, msgSize);
      if (NULL == msg) {
D
dapan 已提交
1089
        SCH_TASK_ELOG("calloc %d failed", msgSize);
D
dapan1121 已提交
1090 1091 1092
        SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
      }
    
S
Shengliang Guan 已提交
1093
      STaskDropReq *pMsg = msg;
D
dapan1121 已提交
1094
      
D
dapan 已提交
1095 1096
      pMsg->header.vgId = htonl(addr->nodeId);   
      
D
dapan1121 已提交
1097
      pMsg->sId = htobe64(schMgmt.sId);      
D
dapan1121 已提交
1098 1099
      pMsg->queryId = htobe64(pJob->queryId);
      pMsg->taskId = htobe64(pTask->taskId);      
D
dapan1121 已提交
1100 1101 1102
      break;
    }
    default:
D
dapan1121 已提交
1103
      SCH_TASK_ELOG("unknown msg type to send, msgType:%d", msgType);
D
dapan1121 已提交
1104 1105 1106 1107
      SCH_ERR_RET(TSDB_CODE_SCH_INTERNAL_ERROR);
      break;
  }

D
dapan1121 已提交
1108
  atomic_store_32(&pTask->lastMsgType, msgType);
D
dapan1121 已提交
1109

D
dapan1121 已提交
1110
  SCH_ERR_JRET(schAsyncSendMsg(pJob->transport, &epSet, pJob->queryId, pTask->taskId, msgType, msg, msgSize));
D
dapan1121 已提交
1111

D
dapan1121 已提交
1112 1113 1114 1115
  if (isCandidateAddr) {
    SCH_ERR_RET(schRecordTaskExecNode(pJob, pTask, addr));
  }
  
D
dapan1121 已提交
1116 1117 1118 1119
  return TSDB_CODE_SUCCESS;

_return:

D
dapan1121 已提交
1120 1121
  atomic_store_32(&pTask->lastMsgType, -1);

D
dapan1121 已提交
1122 1123 1124
  tfree(msg);
  SCH_RET(code);
}
D
dapan1121 已提交
1125

D
dapan1121 已提交
1126 1127 1128 1129 1130 1131 1132 1133 1134
static FORCE_INLINE bool schJobNeedToStop(SSchJob *pJob, int8_t *pStatus) {
  int8_t status = SCH_GET_JOB_STATUS(pJob);
  if (pStatus) {
    *pStatus = status;
  }

  return (status == JOB_TASK_STATUS_FAILED || status == JOB_TASK_STATUS_CANCELLED 
       || status == JOB_TASK_STATUS_CANCELLING || status == JOB_TASK_STATUS_DROPPING);
}
D
dapan1121 已提交
1135

D
dapan1121 已提交
1136 1137

// Note: no more error processing, handled in function internal
D
dapan1121 已提交
1138
int32_t schLaunchTask(SSchJob *pJob, SSchTask *pTask) {
D
dapan1121 已提交
1139 1140 1141 1142 1143
  int8_t status = 0;
  int32_t code = 0;
  
  if (schJobNeedToStop(pJob, &status)) {
    SCH_TASK_ELOG("no need to launch task cause of job status, job status:%d", status);
D
dapan1121 已提交
1144 1145 1146 1147 1148
    
    code = atomic_load_32(&pJob->errCode);
    SCH_ERR_RET(code);
    
    SCH_RET(TSDB_CODE_SCH_STATUS_ERROR);
D
dapan1121 已提交
1149
  }
D
dapan1121 已提交
1150
  
D
dapan1121 已提交
1151
  SSubplan *plan = pTask->plan;
D
dapan1121 已提交
1152

D
dapan1121 已提交
1153 1154 1155 1156 1157 1158
  if (NULL == pTask->msg) {
    code = qSubPlanToString(plan, &pTask->msg, &pTask->msgLen);
    if (TSDB_CODE_SUCCESS != code || NULL == pTask->msg || pTask->msgLen <= 0) {
      SCH_TASK_ELOG("subplanToString error, code:%x, msg:%p, len:%d", code, pTask->msg, pTask->msgLen);
      SCH_ERR_JRET(code);
    }
D
dapan1121 已提交
1159
  }
D
dapan1121 已提交
1160 1161
  
  SCH_ERR_JRET(schSetTaskCandidateAddrs(pJob, pTask));
D
dapan1121 已提交
1162

H
Haojun Liao 已提交
1163
  // NOTE: race condition: the task should be put into the hash table before send msg to server
D
dapan 已提交
1164 1165
  if (SCH_GET_TASK_STATUS(pTask) != JOB_TASK_STATUS_EXECUTING) {
    SCH_ERR_JRET(schPushTaskToExecList(pJob, pTask));
D
dapan1121 已提交
1166

D
dapan 已提交
1167 1168
    SCH_SET_TASK_STATUS(pTask, JOB_TASK_STATUS_EXECUTING);
  }
D
dapan1121 已提交
1169

D
dapan1121 已提交
1170
  SCH_ERR_JRET(schBuildAndSendMsg(pJob, pTask, NULL, plan->msgType));
D
dapan1121 已提交
1171
  
D
dapan1121 已提交
1172
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
1173 1174 1175

_return:

D
dapan1121 已提交
1176
  SCH_ERR_RET(schProcessOnTaskFailure(pJob, pTask, code));
D
dapan1121 已提交
1177 1178
  
  SCH_RET(code);
D
dapan1121 已提交
1179 1180
}

D
dapan1121 已提交
1181 1182
int32_t schLaunchJob(SSchJob *pJob) {
  SSchLevel *level = taosArrayGet(pJob->levels, pJob->levelIdx);
D
dapan1121 已提交
1183 1184

  SCH_ERR_RET(schCheckAndUpdateJobStatus(pJob, JOB_TASK_STATUS_EXECUTING));
D
dapan1121 已提交
1185
  
D
dapan 已提交
1186
  for (int32_t i = 0; i < level->taskNum; ++i) {
D
dapan1121 已提交
1187 1188
    SSchTask *pTask = taosArrayGet(level->subTasks, i);
    SCH_ERR_RET(schLaunchTask(pJob, pTask));
D
dapan1121 已提交
1189
  }
D
dapan 已提交
1190
  
D
dapan1121 已提交
1191
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
1192 1193
}

D
dapan1121 已提交
1194 1195 1196 1197 1198
void schDropTaskOnExecutedNode(SSchJob *pJob, SSchTask *pTask) {
  if (NULL == pTask->execAddrs) {
    SCH_TASK_DLOG("no exec address, status:%d", SCH_GET_TASK_STATUS(pTask));
    return;
  }
H
Haojun Liao 已提交
1199

D
dapan1121 已提交
1200 1201 1202
  int32_t size = (int32_t)taosArrayGetSize(pTask->execAddrs);
  
  if (size <= 0) {
D
dapan1121 已提交
1203
    SCH_TASK_DLOG("task has no exec address, no need to drop it, status:%d", SCH_GET_TASK_STATUS(pTask));
D
dapan1121 已提交
1204 1205
    return;
  }
H
Haojun Liao 已提交
1206

D
dapan1121 已提交
1207 1208 1209
  SQueryNodeAddr *addr = NULL;
  for (int32_t i = 0; i < size; ++i) {
    addr = (SQueryNodeAddr *)taosArrayGet(pTask->execAddrs, i);
D
dapan1121 已提交
1210

D
dapan1121 已提交
1211 1212
    schBuildAndSendMsg(pJob, pTask, addr, TDMT_VND_DROP_TASK);
  }
D
dapan1121 已提交
1213 1214

  SCH_TASK_DLOG("task has %d exec address", size);
D
dapan1121 已提交
1215 1216 1217 1218
}

void schDropTaskInHashList(SSchJob *pJob, SHashObj *list) {
  void *pIter = taosHashIterate(list, NULL);
D
dapan1121 已提交
1219
  while (pIter) {
D
dapan1121 已提交
1220
    SSchTask *pTask = *(SSchTask **)pIter;
H
Haojun Liao 已提交
1221

D
dapan1121 已提交
1222 1223
    if (!SCH_TASK_NO_NEED_DROP(pTask)) {
      schDropTaskOnExecutedNode(pJob, pTask);
H
Haojun Liao 已提交
1224
    }
D
dapan1121 已提交
1225 1226 1227 1228
    
    pIter = taosHashIterate(list, pIter);
  } 
}
H
Haojun Liao 已提交
1229

D
dapan1121 已提交
1230 1231 1232 1233
void schDropJobAllTasks(SSchJob *pJob) {
  schDropTaskInHashList(pJob, pJob->execTasks);
  schDropTaskInHashList(pJob, pJob->succTasks);
  schDropTaskInHashList(pJob, pJob->failTasks);
D
dapan1121 已提交
1234
}
1235

1236
int32_t schExecJobImpl(void *transport, SArray *nodeList, SQueryDag* pDag, struct SSchJob** job, bool syncSchedule) {
D
dapan1121 已提交
1237
  if (nodeList && taosArrayGetSize(nodeList) <= 0) {
D
dapan1121 已提交
1238
    qInfo("QID:%"PRIx64" input nodeList is empty", pDag->queryId);
D
dapan1121 已提交
1239 1240
  }

D
dapan1121 已提交
1241
  int32_t code = 0;
D
dapan1121 已提交
1242 1243
  SSchJob *pJob = calloc(1, sizeof(SSchJob));
  if (NULL == pJob) {
D
dapan1121 已提交
1244
    qError("QID:%"PRIx64" calloc %d failed", pDag->queryId, (int32_t)sizeof(SSchJob));
D
dapan1121 已提交
1245
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
1246 1247
  }

D
dapan1121 已提交
1248 1249 1250
  pJob->attr.syncSchedule = syncSchedule;
  pJob->transport = transport;
  pJob->nodeList = nodeList;
D
dapan 已提交
1251

D
dapan1121 已提交
1252
  SCH_ERR_JRET(schValidateAndBuildJob(pDag, pJob));
D
dapan1121 已提交
1253

D
dapan1121 已提交
1254 1255 1256
  pJob->execTasks = taosHashInit(pDag->numOfSubplans, taosGetDefaultHashFunction(TSDB_DATA_TYPE_UBIGINT), false, HASH_ENTRY_LOCK);
  if (NULL == pJob->execTasks) {
    SCH_JOB_ELOG("taosHashInit %d execTasks failed", pDag->numOfSubplans);
D
dapan 已提交
1257 1258 1259
    SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
1260 1261 1262
  pJob->succTasks = taosHashInit(pDag->numOfSubplans, taosGetDefaultHashFunction(TSDB_DATA_TYPE_UBIGINT), false, HASH_ENTRY_LOCK);
  if (NULL == pJob->succTasks) {
    SCH_JOB_ELOG("taosHashInit %d succTasks failed", pDag->numOfSubplans);
D
dapan 已提交
1263 1264
    SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }
D
dapan 已提交
1265

D
dapan1121 已提交
1266 1267 1268
  pJob->failTasks = taosHashInit(pDag->numOfSubplans, taosGetDefaultHashFunction(TSDB_DATA_TYPE_UBIGINT), false, HASH_ENTRY_LOCK);
  if (NULL == pJob->failTasks) {
    SCH_JOB_ELOG("taosHashInit %d failTasks failed", pDag->numOfSubplans);
D
dapan1121 已提交
1269 1270 1271
    SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
1272
  tsem_init(&pJob->rspSem, 0, 0);
D
dapan1121 已提交
1273

D
dapan1121 已提交
1274
  code = taosHashPut(schMgmt.jobs, &pJob->queryId, sizeof(pJob->queryId), &pJob, POINTER_BYTES);
D
dapan1121 已提交
1275 1276
  if (0 != code) {
    if (HASH_NODE_EXIST(code)) {
D
dapan1121 已提交
1277
      SCH_JOB_ELOG("job already exist, isQueryJob:%d", pJob->attr.queryJob);
D
dapan1121 已提交
1278 1279
      SCH_ERR_JRET(TSDB_CODE_QRY_INVALID_INPUT);
    } else {
D
dapan1121 已提交
1280 1281
      SCH_JOB_ELOG("taosHashPut job failed, errno:%d", errno);
      SCH_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
1282
    }
D
dapan1121 已提交
1283 1284
  }

D
dapan1121 已提交
1285
  pJob->status = JOB_TASK_STATUS_NOT_START;
D
dapan1121 已提交
1286
  
D
dapan1121 已提交
1287
  SCH_ERR_JRET(schLaunchJob(pJob));
1288

D
dapan1121 已提交
1289
  *(SSchJob **)job = pJob;
D
dapan1121 已提交
1290
  
D
dapan 已提交
1291
  if (syncSchedule) {
D
dapan1121 已提交
1292
    SCH_JOB_DLOG("will wait for rsp now, job status:%d", SCH_GET_JOB_STATUS(pJob));
D
dapan1121 已提交
1293
    tsem_wait(&pJob->rspSem);
D
dapan1121 已提交
1294 1295
  }

D
dapan1121 已提交
1296
  SCH_JOB_DLOG("job exec done, job status:%d", SCH_GET_JOB_STATUS(pJob));
D
dapan1121 已提交
1297

D
dapan1121 已提交
1298
  return TSDB_CODE_SUCCESS;
1299

D
dapan1121 已提交
1300
_return:
D
dapan1121 已提交
1301

D
dapan1121 已提交
1302 1303 1304
  *(SSchJob **)job = NULL;
  
  scheduleFreeJob(pJob);
D
dapan1121 已提交
1305 1306
  
  SCH_RET(code);
1307
}
D
dapan1121 已提交
1308

D
dapan1121 已提交
1309 1310 1311 1312 1313 1314 1315
int32_t schCancelJob(SSchJob *pJob) {
  //TODO

  //TODO MOVE ALL TASKS FROM EXEC LIST TO FAIL LIST

}

D
dapan1121 已提交
1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336 1337 1338

int32_t schedulerInit(SSchedulerCfg *cfg) {
  if (schMgmt.jobs) {
    qError("scheduler already initialized");
    return TSDB_CODE_QRY_INVALID_INPUT;
  }

  if (cfg) {
    schMgmt.cfg = *cfg;
    
    if (schMgmt.cfg.maxJobNum == 0) {
      schMgmt.cfg.maxJobNum = SCHEDULE_DEFAULT_JOB_NUMBER;
    }
  } else {
    schMgmt.cfg.maxJobNum = SCHEDULE_DEFAULT_JOB_NUMBER;
  }

  schMgmt.jobs = taosHashInit(schMgmt.cfg.maxJobNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_UBIGINT), false, HASH_ENTRY_LOCK);
  if (NULL == schMgmt.jobs) {
    qError("init schduler jobs failed, num:%u", schMgmt.cfg.maxJobNum);
    SCH_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
1339
  if (taosGetSystemUUID((char *)&schMgmt.sId, sizeof(schMgmt.sId))) {
D
dapan1121 已提交
1340 1341 1342 1343
    qError("generate schdulerId failed, errno:%d", errno);
    SCH_ERR_RET(TSDB_CODE_QRY_SYS_ERROR);
  }

D
dapan1121 已提交
1344
  qInfo("scheduler %"PRIx64" initizlized, maxJob:%u", schMgmt.sId, schMgmt.cfg.maxJobNum);
D
dapan1121 已提交
1345 1346 1347 1348
  
  return TSDB_CODE_SUCCESS;
}

1349
int32_t scheduleExecJob(void *transport, SArray *nodeList, SQueryDag* pDag, struct SSchJob** pJob, SQueryResult *pRes) {
H
Haojun Liao 已提交
1350
  if (NULL == transport || NULL == pDag || NULL == pDag->pSubplans || NULL == pJob || NULL == pRes) {
D
dapan1121 已提交
1351 1352 1353
    SCH_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
  }

D
dapan1121 已提交
1354 1355
  SSchJob *job = NULL;

D
dapan1121 已提交
1356
  SCH_ERR_RET(schExecJobImpl(transport, nodeList, pDag, &job, true));
D
dapan1121 已提交
1357 1358

  *pJob = job;
D
dapan1121 已提交
1359

D
dapan 已提交
1360
  pRes->code = atomic_load_32(&job->errCode);
D
dapan1121 已提交
1361 1362
  pRes->numOfRows = job->resNumOfRows;
  
D
dapan1121 已提交
1363 1364 1365
  return TSDB_CODE_SUCCESS;
}

1366
int32_t scheduleAsyncExecJob(void *transport, SArray *nodeList, SQueryDag* pDag, struct SSchJob** pJob) {
1367
  if (NULL == transport || NULL == pDag || NULL == pDag->pSubplans || NULL == pJob) {
D
dapan1121 已提交
1368 1369 1370
    SCH_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
  }

H
Haojun Liao 已提交
1371
  SCH_ERR_RET(schExecJobImpl(transport, nodeList, pDag, pJob, false));
D
dapan1121 已提交
1372
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
1373 1374
}

1375 1376
int32_t scheduleFetchRows(SSchJob *pJob, void** pData) {
  if (NULL == pJob || NULL == pData) {
D
dapan1121 已提交
1377
    SCH_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
D
dapan 已提交
1378
  }
D
dapan 已提交
1379
  int32_t code = 0;
D
dapan 已提交
1380

D
dapan1121 已提交
1381 1382 1383 1384 1385 1386
  int8_t status = SCH_GET_JOB_STATUS(pJob);
  if (status == JOB_TASK_STATUS_DROPPING) {
    SCH_JOB_ELOG("job is dropping, status:%d", status);
    return TSDB_CODE_SCH_STATUS_ERROR;
  }

D
dapan1121 已提交
1387 1388 1389 1390 1391
  atomic_add_fetch_32(&pJob->ref, 1);
  
  if (!SCH_JOB_NEED_FETCH(&pJob->attr)) {
    SCH_JOB_ELOG("no need to fetch data, status:%d", SCH_GET_JOB_STATUS(pJob));
    atomic_sub_fetch_32(&pJob->ref, 1);
D
dapan1121 已提交
1392 1393 1394
    SCH_ERR_RET(TSDB_CODE_QRY_APP_ERROR);
  }

D
dapan1121 已提交
1395 1396 1397 1398
  if (atomic_val_compare_exchange_8(&pJob->userFetch, 0, 1) != 0) {
    SCH_JOB_ELOG("prior fetching not finished, userFetch:%d", atomic_load_8(&pJob->userFetch));
    atomic_sub_fetch_32(&pJob->ref, 1);
    SCH_ERR_RET(TSDB_CODE_QRY_APP_ERROR);
D
dapan 已提交
1399 1400
  }

D
dapan1121 已提交
1401
  if (status == JOB_TASK_STATUS_FAILED) {
D
dapan1121 已提交
1402
    *pData = atomic_load_ptr(&pJob->res);
D
dapan1121 已提交
1403 1404 1405
    atomic_store_ptr(&pJob->res, NULL);
    SCH_ERR_JRET(atomic_load_32(&pJob->errCode));
  } else if (status == JOB_TASK_STATUS_SUCCEED) {
D
dapan1121 已提交
1406
    *pData = atomic_load_ptr(&pJob->res);
D
dapan1121 已提交
1407 1408 1409 1410
    atomic_store_ptr(&pJob->res, NULL);
    goto _return;
  } else if (status == JOB_TASK_STATUS_PARTIAL_SUCCEED) {
    SCH_ERR_JRET(schFetchFromRemote(pJob));
D
dapan 已提交
1411 1412
  }

D
dapan1121 已提交
1413
  tsem_wait(&pJob->rspSem);
D
dapan 已提交
1414

D
dapan1121 已提交
1415
  status = SCH_GET_JOB_STATUS(pJob);
D
dapan 已提交
1416

D
dapan1121 已提交
1417 1418
  if (status == JOB_TASK_STATUS_FAILED) {
    code = atomic_load_32(&pJob->errCode);
D
dapan1121 已提交
1419
    SCH_ERR_JRET(code);
D
dapan 已提交
1420 1421
  }
  
D
dapan1121 已提交
1422 1423
  if (pJob->res && ((SRetrieveTableRsp *)pJob->res)->completed) {
    SCH_ERR_JRET(schCheckAndUpdateJobStatus(pJob, JOB_TASK_STATUS_SUCCEED));
D
dapan 已提交
1424 1425
  }

D
dapan1121 已提交
1426
  while (true) {
D
dapan1121 已提交
1427
    *pData = atomic_load_ptr(&pJob->res);
D
dapan1121 已提交
1428
    
D
dapan1121 已提交
1429
    if (*pData != atomic_val_compare_exchange_ptr(&pJob->res, *pData, NULL)) {
D
dapan1121 已提交
1430 1431 1432 1433 1434
      continue;
    }

    break;
  }
D
dapan 已提交
1435

D
dapan 已提交
1436
_return:
D
dapan1121 已提交
1437 1438 1439 1440

  atomic_val_compare_exchange_8(&pJob->userFetch, 1, 0);

  atomic_sub_fetch_32(&pJob->ref, 1);
D
dapan 已提交
1441

D
dapan1121 已提交
1442
  SCH_RET(code);
D
dapan 已提交
1443
}
D
dapan1121 已提交
1444

D
dapan1121 已提交
1445 1446
int32_t scheduleCancelJob(void *job) {
  SSchJob *pJob = (SSchJob *)job;
D
dapan1121 已提交
1447

D
dapan1121 已提交
1448
  atomic_add_fetch_32(&pJob->ref, 1);
D
dapan1121 已提交
1449

D
dapan1121 已提交
1450 1451 1452 1453 1454
  int32_t code = schCancelJob(pJob);

  atomic_sub_fetch_32(&pJob->ref, 1);

  SCH_RET(code);
D
dapan1121 已提交
1455 1456
}

D
dapan1121 已提交
1457 1458
void scheduleFreeJob(void *job) {
  if (NULL == job) {
D
dapan 已提交
1459 1460
    return;
  }
D
dapan1121 已提交
1461

D
dapan1121 已提交
1462
  SSchJob *pJob = job;
D
dapan1121 已提交
1463
  uint64_t queryId = pJob->queryId;
D
dapan1121 已提交
1464

D
dapan1121 已提交
1465 1466 1467 1468 1469
  if (SCH_GET_JOB_STATUS(pJob) > 0) {
    if (0 != taosHashRemove(schMgmt.jobs, &pJob->queryId, sizeof(pJob->queryId))) {
      SCH_JOB_ELOG("taosHashRemove job from list failed, may already freed, pJob:%p", pJob);
      return;
    }
D
dapan1121 已提交
1470

D
dapan1121 已提交
1471
    schCheckAndUpdateJobStatus(pJob, JOB_TASK_STATUS_DROPPING);
D
dapan1121 已提交
1472

D
dapan1121 已提交
1473
    SCH_JOB_DLOG("job removed from list, no further ref, ref:%d", atomic_load_32(&pJob->ref));
D
dapan1121 已提交
1474

D
dapan1121 已提交
1475 1476 1477 1478 1479 1480 1481 1482 1483
    while (true) {
      int32_t ref = atomic_load_32(&pJob->ref);
      if (0 == ref) {
        break;
      } else if (ref > 0) {
        usleep(1);
      } else {
        assert(0);
      }
D
dapan1121 已提交
1484
    }
D
dapan1121 已提交
1485

D
dapan1121 已提交
1486
    SCH_JOB_DLOG("job no ref now, status:%d", SCH_GET_JOB_STATUS(pJob));
D
dapan1121 已提交
1487

D
dapan1121 已提交
1488 1489 1490
    if (pJob->status == JOB_TASK_STATUS_EXECUTING) {
      schCancelJob(pJob);
    }
1491

D
dapan1121 已提交
1492 1493
    schDropJobAllTasks(pJob);
  }
D
dapan1121 已提交
1494 1495 1496 1497

  pJob->subPlans = NULL; // it is a reference to pDag->pSubplans
  
  int32_t numOfLevels = taosArrayGetSize(pJob->levels);
1498
  for(int32_t i = 0; i < numOfLevels; ++i) {
D
dapan1121 已提交
1499
    SSchLevel *pLevel = taosArrayGet(pJob->levels, i);
1500 1501 1502 1503

    int32_t numOfTasks = taosArrayGetSize(pLevel->subTasks);
    for(int32_t j = 0; j < numOfTasks; ++j) {
      SSchTask* pTask = taosArrayGet(pLevel->subTasks, j);
D
dapan1121 已提交
1504
      schFreeTask(pTask);
1505 1506 1507 1508
    }

    taosArrayDestroy(pLevel->subTasks);
  }
D
dapan1121 已提交
1509
  
D
dapan1121 已提交
1510 1511 1512
  taosHashCleanup(pJob->execTasks);
  taosHashCleanup(pJob->failTasks);
  taosHashCleanup(pJob->succTasks);
D
dapan1121 已提交
1513
  
D
dapan1121 已提交
1514 1515 1516
  taosArrayDestroy(pJob->levels);

  tfree(pJob->res);
D
dapan1121 已提交
1517
  
D
dapan1121 已提交
1518
  tfree(pJob);
D
dapan1121 已提交
1519 1520

  qDebug("QID:%"PRIx64" job freed", queryId);
D
dapan1121 已提交
1521 1522 1523 1524 1525 1526 1527 1528 1529
}
  
void schedulerDestroy(void) {
  if (schMgmt.jobs) {
    taosHashCleanup(schMgmt.jobs); //TODO
    schMgmt.jobs = NULL;
  }
}