qworker.c 46.2 KB
Newer Older
D
dapan1121 已提交
1
#include "qworker.h"
S
common  
Shengliang Guan 已提交
2
#include "tcommon.h"
3
#include "executor.h"
D
dapan1121 已提交
4
#include "planner.h"
H
Haojun Liao 已提交
5 6
#include "query.h"
#include "qworkerInt.h"
D
dapan1121 已提交
7
#include "qworkerMsg.h"
H
Haojun Liao 已提交
8
#include "tmsg.h"
9
#include "tname.h"
D
dapan1121 已提交
10
#include "dataSinkMgt.h"
D
dapan1121 已提交
11

12 13
SQWDebug gQWDebug = {0};

D
dapan1121 已提交
14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34
char *qwPhaseStr(int32_t phase) {
  switch (phase) {
    case QW_PHASE_PRE_QUERY:
      return "PRE_QUERY";
    case QW_PHASE_POST_QUERY:
      return "POST_QUERY";
    case QW_PHASE_PRE_FETCH:
      return "PRE_FETCH";
    case QW_PHASE_POST_FETCH:
      return "POST_FETCH";
    case QW_PHASE_PRE_CQUERY:
      return "PRE_CQUERY";
    case QW_PHASE_POST_CQUERY:
      return "POST_CQUERY";
    default:
      break;
  }

  return "UNKNOWN";
}

D
dapan1121 已提交
35
int32_t qwValidateStatus(QW_FPARAMS_DEF, int8_t oriStatus, int8_t newStatus) {
D
dapan1121 已提交
36
  int32_t code = 0;
D
dapan1121 已提交
37

D
dapan1121 已提交
38 39 40 41 42 43
  if (oriStatus == newStatus) {
    QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
  }
  
  switch (oriStatus) {
    case JOB_TASK_STATUS_NULL:
D
dapan1121 已提交
44 45 46
      if (newStatus != JOB_TASK_STATUS_EXECUTING 
       && newStatus != JOB_TASK_STATUS_FAILED 
       && newStatus != JOB_TASK_STATUS_NOT_START) {
D
dapan1121 已提交
47 48 49 50 51
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_NOT_START:
D
dapan1121 已提交
52
      if (newStatus != JOB_TASK_STATUS_CANCELLED) {
D
dapan1121 已提交
53 54 55 56 57
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_EXECUTING:
D
dapan1121 已提交
58
      if (newStatus != JOB_TASK_STATUS_PARTIAL_SUCCEED 
D
dapan1121 已提交
59
       && newStatus != JOB_TASK_STATUS_SUCCEED 
D
dapan1121 已提交
60 61 62 63
       && newStatus != JOB_TASK_STATUS_FAILED 
       && newStatus != JOB_TASK_STATUS_CANCELLING 
       && newStatus != JOB_TASK_STATUS_CANCELLED 
       && newStatus != JOB_TASK_STATUS_DROPPING) {
D
dapan1121 已提交
64 65 66 67 68
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_PARTIAL_SUCCEED:
D
dapan1121 已提交
69 70 71
      if (newStatus != JOB_TASK_STATUS_EXECUTING 
       && newStatus != JOB_TASK_STATUS_SUCCEED
       && newStatus != JOB_TASK_STATUS_CANCELLED) {
D
dapan1121 已提交
72 73 74 75 76
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_SUCCEED:
D
dapan1121 已提交
77 78 79 80 81 82
      if (newStatus != JOB_TASK_STATUS_CANCELLED
       && newStatus != JOB_TASK_STATUS_DROPPING) {
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }

      break;
D
dapan1121 已提交
83 84 85 86 87 88 89 90
    case JOB_TASK_STATUS_FAILED:
    case JOB_TASK_STATUS_CANCELLING:
      if (newStatus != JOB_TASK_STATUS_CANCELLED) {
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_CANCELLED:
D
dapan1121 已提交
91 92 93 94
    case JOB_TASK_STATUS_DROPPING:
      QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      break;
      
D
dapan1121 已提交
95
    default:
D
dapan1121 已提交
96
      QW_TASK_ELOG("invalid task origStatus:%s", jobTaskStatusStr(oriStatus));
D
dapan1121 已提交
97 98 99
      return TSDB_CODE_QRY_APP_ERROR;
  }

D
dapan1121 已提交
100
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
101

D
dapan1121 已提交
102
_return:
D
dapan1121 已提交
103

D
dapan1121 已提交
104
  QW_TASK_ELOG("invalid task status update from %s to %s", jobTaskStatusStr(oriStatus), jobTaskStatusStr(newStatus));
D
dapan1121 已提交
105
  QW_RET(code);
D
dapan1121 已提交
106 107
}

D
dapan1121 已提交
108
int32_t qwSetTaskStatus(QW_FPARAMS_DEF, SQWTaskStatus *task, int8_t status) {
D
dapan1121 已提交
109
  int32_t code = 0;
D
dapan1121 已提交
110
  int8_t origStatus = 0;
D
dapan1121 已提交
111 112 113 114 115 116 117 118

  while (true) {
    origStatus = atomic_load_8(&task->status);
    
    QW_ERR_RET(qwValidateStatus(QW_FPARAMS(), origStatus, status));
    
    if (origStatus != atomic_val_compare_exchange_8(&task->status, origStatus, status)) {
      continue;
D
dapan1121 已提交
119
    }
D
dapan1121 已提交
120
    
D
dapan1121 已提交
121
    QW_TASK_DLOG("task status updated from %s to %s", jobTaskStatusStr(origStatus), jobTaskStatusStr(status));
D
dapan1121 已提交
122 123

    break;
D
dapan1121 已提交
124
  }
D
dapan1121 已提交
125
  
D
dapan1121 已提交
126 127 128 129
  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
130
int32_t qwAddSchedulerImpl(SQWorkerMgmt *mgmt, uint64_t sId, int32_t rwType, SQWSchStatus **sch) {
D
dapan1121 已提交
131
  SQWSchStatus newSch = {0};
D
dapan1121 已提交
132 133
  newSch.tasksHash = taosHashInit(mgmt->cfg.maxSchTaskNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_NO_LOCK);
  if (NULL == newSch.tasksHash) {
D
dapan1121 已提交
134 135
    QW_SCH_ELOG("taosHashInit %d failed", mgmt->cfg.maxSchTaskNum);
    QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
136 137
  }

D
dapan1121 已提交
138 139 140 141 142
  QW_LOCK(QW_WRITE, &mgmt->schLock);
  int32_t code = taosHashPut(mgmt->schHash, &sId, sizeof(sId), &newSch, sizeof(newSch));
  if (0 != code) {
    if (!HASH_NODE_EXIST(code)) {
      QW_UNLOCK(QW_WRITE, &mgmt->schLock);
D
dapan1121 已提交
143
      
D
dapan1121 已提交
144 145 146
      QW_SCH_ELOG("taosHashPut new sch to scheduleHash failed, errno:%d", errno);
      taosHashCleanup(newSch.tasksHash);
      QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
147
    }
D
dapan1121 已提交
148 149

    taosHashCleanup(newSch.tasksHash);
D
dapan1121 已提交
150
  }
D
dapan1121 已提交
151
  QW_UNLOCK(QW_WRITE, &mgmt->schLock);
D
dapan1121 已提交
152

D
dapan1121 已提交
153
  return TSDB_CODE_SUCCESS;  
D
dapan1121 已提交
154 155
}

D
dapan1121 已提交
156
int32_t qwAcquireSchedulerImpl(SQWorkerMgmt *mgmt, uint64_t sId, int32_t rwType, SQWSchStatus **sch, int32_t nOpt) {
D
dapan1121 已提交
157 158 159 160 161 162 163
  while (true) {
    QW_LOCK(rwType, &mgmt->schLock);
    *sch = taosHashGet(mgmt->schHash, &sId, sizeof(sId));
    if (NULL == (*sch)) {
      QW_UNLOCK(rwType, &mgmt->schLock);
      
      if (QW_NOT_EXIST_ADD == nOpt) {
D
dapan1121 已提交
164
        QW_ERR_RET(qwAddSchedulerImpl(mgmt, sId, rwType, sch));
D
dapan1121 已提交
165 166 167 168 169 170 171

        nOpt = QW_NOT_EXIST_RET_ERR;
        
        continue;
      } else if (QW_NOT_EXIST_RET_ERR == nOpt) {
        QW_RET(TSDB_CODE_QRY_SCH_NOT_EXIST);
      } else {
D
dapan1121 已提交
172
        QW_SCH_ELOG("unknown notExistOpt:%d", nOpt);
D
dapan1121 已提交
173
        QW_ERR_RET(TSDB_CODE_QRY_APP_ERROR);
D
dapan1121 已提交
174
      }
D
dapan1121 已提交
175
    }
D
dapan1121 已提交
176 177

    break;
D
dapan1121 已提交
178 179 180 181 182
  }

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
183 184
int32_t qwAcquireAddScheduler(SQWorkerMgmt *mgmt, uint64_t sId, int32_t rwType, SQWSchStatus **sch) {
  return qwAcquireSchedulerImpl(mgmt, sId, rwType, sch, QW_NOT_EXIST_ADD);
D
dapan1121 已提交
185 186
}

D
dapan1121 已提交
187 188
int32_t qwAcquireScheduler(SQWorkerMgmt *mgmt, uint64_t sId, int32_t rwType, SQWSchStatus **sch) {
  return qwAcquireSchedulerImpl(mgmt, sId, rwType, sch, QW_NOT_EXIST_RET_ERR);
D
dapan1121 已提交
189 190
}

D
dapan1121 已提交
191
void qwReleaseScheduler(int32_t rwType, SQWorkerMgmt *mgmt) {
D
dapan1121 已提交
192 193 194
  QW_UNLOCK(rwType, &mgmt->schLock);
}

D
dapan1121 已提交
195

D
dapan1121 已提交
196
int32_t qwAcquireTaskStatus(QW_FPARAMS_DEF, int32_t rwType, SQWSchStatus *sch, SQWTaskStatus **task) {
D
dapan1121 已提交
197 198 199 200 201 202 203 204 205 206 207 208 209 210 211
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);

  QW_LOCK(rwType, &sch->tasksLock);
  *task = taosHashGet(sch->tasksHash, id, sizeof(id));
  if (NULL == (*task)) {
    QW_UNLOCK(rwType, &sch->tasksLock);
    QW_ERR_RET(TSDB_CODE_QRY_TASK_NOT_EXIST);
  }

  return TSDB_CODE_SUCCESS;
}



D
dapan1121 已提交
212
int32_t qwAddTaskStatusImpl(QW_FPARAMS_DEF, SQWSchStatus *sch, int32_t rwType, int32_t status, SQWTaskStatus **task) {
D
dapan1121 已提交
213 214
  int32_t code = 0;

D
dapan1121 已提交
215 216
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
D
dapan1121 已提交
217

D
dapan1121 已提交
218 219
  SQWTaskStatus ntask = {0};
  ntask.status = status;
D
dapan1121 已提交
220
  ntask.refId = rId;
D
dapan1121 已提交
221 222 223 224 225 226

  QW_LOCK(QW_WRITE, &sch->tasksLock);
  code = taosHashPut(sch->tasksHash, id, sizeof(id), &ntask, sizeof(ntask));
  if (0 != code) {
    QW_UNLOCK(QW_WRITE, &sch->tasksLock);
    if (HASH_NODE_EXIST(code)) {
D
dapan1121 已提交
227 228
      if (rwType && task) {
        QW_RET(qwAcquireTaskStatus(QW_FPARAMS(), rwType, sch, task));
D
dapan1121 已提交
229
      } else {
D
dapan1121 已提交
230
        QW_TASK_ELOG("task status already exist, newStatus:%s", jobTaskStatusStr(status));
D
dapan1121 已提交
231
        QW_ERR_RET(TSDB_CODE_QRY_TASK_ALREADY_EXIST);
D
dapan1121 已提交
232 233
      }
    } else {
D
dapan1121 已提交
234
      QW_TASK_ELOG("taosHashPut to tasksHash failed, error:%x - %s", code, tstrerror(code));
D
dapan1121 已提交
235
      QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
236 237 238
    }
  }
  QW_UNLOCK(QW_WRITE, &sch->tasksLock);
D
dapan1121 已提交
239

D
dapan1121 已提交
240 241
  QW_TASK_DLOG("task status added, newStatus:%s", jobTaskStatusStr(status));

D
dapan1121 已提交
242 243
  if (rwType && task) {
    QW_ERR_RET(qwAcquireTaskStatus(QW_FPARAMS(), rwType, sch, task));
D
dapan1121 已提交
244 245
  }

D
dapan1121 已提交
246 247
  return TSDB_CODE_SUCCESS;
}
D
dapan1121 已提交
248

D
dapan1121 已提交
249
int32_t qwAddTaskStatus(QW_FPARAMS_DEF, int32_t status) {
D
dapan1121 已提交
250 251
  SQWSchStatus *tsch = NULL;
  int32_t code = 0;
D
dapan1121 已提交
252
  QW_ERR_RET(qwAcquireAddScheduler(mgmt, sId, QW_READ, &tsch));
D
dapan1121 已提交
253

D
dapan1121 已提交
254
  QW_ERR_JRET(qwAddTaskStatusImpl(QW_FPARAMS(), tsch, 0, status, NULL));
D
dapan1121 已提交
255 256 257 258

_return:

  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
259 260
  
  QW_RET(code);
D
dapan1121 已提交
261 262 263
}


D
dapan1121 已提交
264
int32_t qwAddAcquireTaskStatus(QW_FPARAMS_DEF, int32_t rwType, SQWSchStatus *sch, int32_t status, SQWTaskStatus **task) {
D
dapan1121 已提交
265
  return qwAddTaskStatusImpl(QW_FPARAMS(), sch, rwType, status, task);
D
dapan1121 已提交
266
}
D
dapan1121 已提交
267 268


D
dapan1121 已提交
269
void qwReleaseTaskStatus(int32_t rwType, SQWSchStatus *sch) {
D
dapan1121 已提交
270
  QW_UNLOCK(rwType, &sch->tasksLock);
D
dapan1121 已提交
271 272
}

D
dapan1121 已提交
273

D
dapan1121 已提交
274
int32_t qwAcquireTaskCtx(QW_FPARAMS_DEF, SQWTaskCtx **ctx) {
D
dapan1121 已提交
275 276 277
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
  
D
dapan1121 已提交
278
  *ctx = taosHashAcquire(mgmt->ctxHash, id, sizeof(id));
D
dapan1121 已提交
279
  if (NULL == (*ctx)) {
D
dapan1121 已提交
280
    QW_TASK_DLOG_E("task ctx not exist, may be dropped");
D
dapan1121 已提交
281
    QW_ERR_RET(TSDB_CODE_QRY_TASK_CTX_NOT_EXIST);
D
dapan1121 已提交
282 283 284 285 286
  }

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
287
int32_t qwGetTaskCtx(QW_FPARAMS_DEF, SQWTaskCtx **ctx) {
D
dapan1121 已提交
288 289 290 291 292
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
  
  *ctx = taosHashGet(mgmt->ctxHash, id, sizeof(id));
  if (NULL == (*ctx)) {
D
dapan1121 已提交
293
    QW_TASK_DLOG_E("task ctx not exist, may be dropped");
D
dapan1121 已提交
294
    QW_ERR_RET(TSDB_CODE_QRY_TASK_CTX_NOT_EXIST);
D
dapan1121 已提交
295 296 297 298 299
  }

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
300
int32_t qwAddTaskCtxImpl(QW_FPARAMS_DEF, bool acquire, SQWTaskCtx **ctx) {
D
dapan1121 已提交
301 302 303
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);

D
dapan1121 已提交
304
  SQWTaskCtx nctx = {0};
D
dapan1121 已提交
305

D
dapan1121 已提交
306
  int32_t code = taosHashPut(mgmt->ctxHash, id, sizeof(id), &nctx, sizeof(SQWTaskCtx));
D
dapan1121 已提交
307 308
  if (0 != code) {
    if (HASH_NODE_EXIST(code)) {
D
dapan1121 已提交
309
      if (acquire && ctx) {
D
dapan1121 已提交
310
        QW_RET(qwAcquireTaskCtx(QW_FPARAMS(), ctx));
D
dapan1121 已提交
311 312
      } else if (ctx) {
        QW_RET(qwGetTaskCtx(QW_FPARAMS(), ctx));
D
dapan1121 已提交
313
      } else {
D
dapan1121 已提交
314
        QW_TASK_ELOG_E("task ctx already exist");
D
dapan1121 已提交
315 316 317
        QW_ERR_RET(TSDB_CODE_QRY_TASK_ALREADY_EXIST);
      }
    } else {
D
dapan1121 已提交
318
      QW_TASK_ELOG("taosHashPut to ctxHash failed, error:%x", code);
D
dapan1121 已提交
319 320 321 322
      QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
323
  if (acquire && ctx) {
D
dapan1121 已提交
324
    QW_RET(qwAcquireTaskCtx(QW_FPARAMS(), ctx));
D
dapan1121 已提交
325 326
  } else if (ctx) {
    QW_RET(qwGetTaskCtx(QW_FPARAMS(), ctx));
D
dapan1121 已提交
327 328 329 330 331
  }

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
332
int32_t qwAddTaskCtx(QW_FPARAMS_DEF) {
D
dapan1121 已提交
333
  QW_RET(qwAddTaskCtxImpl(QW_FPARAMS(), false, NULL));
D
dapan1121 已提交
334 335
}

D
dapan1121 已提交
336
int32_t qwAddAcquireTaskCtx(QW_FPARAMS_DEF, SQWTaskCtx **ctx) {
D
dapan1121 已提交
337
  return qwAddTaskCtxImpl(QW_FPARAMS(), true, ctx);
D
dapan1121 已提交
338 339
}

D
dapan1121 已提交
340
int32_t qwAddGetTaskCtx(QW_FPARAMS_DEF, SQWTaskCtx **ctx) {
D
dapan1121 已提交
341
  return qwAddTaskCtxImpl(QW_FPARAMS(), false, ctx);
D
dapan1121 已提交
342 343 344
}


D
dapan1121 已提交
345 346 347
void qwReleaseTaskCtx(SQWorkerMgmt *mgmt, void *ctx) {
  //QW_UNLOCK(rwType, &mgmt->ctxLock);
  taosHashRelease(mgmt->ctxHash, ctx);
D
dapan1121 已提交
348 349
}

D
dapan1121 已提交
350
void qwFreeTaskHandle(QW_FPARAMS_DEF, qTaskInfo_t *taskHandle) {  
D
dapan1121 已提交
351
  // RC WARNING
D
dapan1121 已提交
352 353 354
  qTaskInfo_t otaskHandle = atomic_load_ptr(taskHandle);
  if (otaskHandle && atomic_val_compare_exchange_ptr(taskHandle, otaskHandle, NULL)) {
    qDestroyTask(otaskHandle);
D
dapan1121 已提交
355 356
  }
}
D
dapan1121 已提交
357

D
dapan1121 已提交
358 359 360 361 362
int32_t qwKillTaskHandle(QW_FPARAMS_DEF, SQWTaskCtx *ctx) {
  int32_t code = 0;
  // RC WARNING
  qTaskInfo_t taskHandle = atomic_load_ptr(&ctx->taskHandle);
  if (taskHandle && atomic_val_compare_exchange_ptr(&ctx->taskHandle, taskHandle, NULL)) {
D
dapan1121 已提交
363
    code = qAsyncKillTask(taskHandle);
D
dapan1121 已提交
364
    atomic_store_ptr(&ctx->taskHandle, taskHandle);
D
dapan1121 已提交
365
  }
D
dapan1121 已提交
366

D
dapan1121 已提交
367 368 369
  QW_RET(code);
}

D
dapan1121 已提交
370

D
dapan1121 已提交
371
void qwFreeTask(QW_FPARAMS_DEF, SQWTaskCtx *ctx) {
D
dapan1121 已提交
372
  qwFreeTaskHandle(QW_FPARAMS(), &ctx->taskHandle);
D
dapan1121 已提交
373 374 375 376
  
  if (ctx->sinkHandle) {
    dsDestroyDataSinker(ctx->sinkHandle);
    ctx->sinkHandle = NULL;
D
dapan1121 已提交
377
  }
D
dapan1121 已提交
378
}
D
dapan1121 已提交
379 380


D
dapan1121 已提交
381
// Note: NEED CTX HASH LOCKED BEFORE ENTRANCE
D
dapan1121 已提交
382
int32_t qwDropTaskCtx(QW_FPARAMS_DEF, int32_t rwType) {
D
dapan1121 已提交
383 384 385
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
  SQWTaskCtx octx;
D
dapan1121 已提交
386

D
dapan1121 已提交
387 388
  SQWTaskCtx *ctx = taosHashGet(mgmt->ctxHash, id, sizeof(id));
  if (NULL == ctx) {
D
dapan1121 已提交
389
    QW_ERR_RET(TSDB_CODE_QRY_TASK_CTX_NOT_EXIST);
D
dapan1121 已提交
390
  }
D
dapan1121 已提交
391

D
dapan1121 已提交
392
  octx = *ctx;
D
dapan1121 已提交
393

D
dapan1121 已提交
394 395 396
  atomic_store_ptr(&ctx->taskHandle, NULL);
  atomic_store_ptr(&ctx->sinkHandle, NULL);

D
dapan1121 已提交
397 398
  QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_DROP);

D
dapan1121 已提交
399 400 401 402
  if (rwType) {
    QW_UNLOCK(rwType, &ctx->lock);
  }

D
dapan1121 已提交
403
  if (taosHashRemove(mgmt->ctxHash, id, sizeof(id))) {
D
dapan1121 已提交
404
    QW_TASK_ELOG_E("taosHashRemove from ctx hash failed");    
D
dapan1121 已提交
405
    QW_ERR_RET(TSDB_CODE_QRY_TASK_CTX_NOT_EXIST);
D
dapan1121 已提交
406 407 408 409 410
  }

  if (octx.taskHandle) {
    qDestroyTask(octx.taskHandle);
  }
D
dapan1121 已提交
411

D
dapan1121 已提交
412 413 414
  if (octx.sinkHandle) {
    dsDestroyDataSinker(octx.sinkHandle);
  }
D
dapan1121 已提交
415 416

  QW_TASK_DLOG_E("task ctx dropped");
D
dapan1121 已提交
417
  
D
dapan1121 已提交
418 419 420
  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
421

D
dapan1121 已提交
422
int32_t qwDropTaskStatus(QW_FPARAMS_DEF) {
D
dapan1121 已提交
423
  SQWSchStatus *sch = NULL;
D
dapan1121 已提交
424 425 426 427 428
  SQWTaskStatus *task = NULL;
  int32_t code = 0;
  
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
D
dapan1121 已提交
429

D
dapan1121 已提交
430
  if (qwAcquireScheduler(mgmt, sId, QW_WRITE, &sch)) {
D
dapan1121 已提交
431
    QW_TASK_WLOG_E("scheduler does not exist");
D
dapan1121 已提交
432 433
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
434

D
dapan1121 已提交
435 436 437
  if (qwAcquireTaskStatus(QW_FPARAMS(), QW_WRITE, sch, &task)) {
    qwReleaseScheduler(QW_WRITE, mgmt);
    
D
dapan1121 已提交
438
    QW_TASK_WLOG_E("task does not exist");
D
dapan1121 已提交
439 440
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
441

D
dapan1121 已提交
442
  if (taosHashRemove(sch->tasksHash, id, sizeof(id))) {
D
dapan1121 已提交
443
    QW_TASK_ELOG_E("taosHashRemove task from hash failed");
D
dapan1121 已提交
444 445 446
    QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
  }

D
dapan1121 已提交
447
  QW_TASK_DLOG_E("task status dropped");
D
dapan1121 已提交
448 449 450 451 452 453 454 455 456

_return:

  qwReleaseTaskStatus(QW_WRITE, sch);
  qwReleaseScheduler(QW_WRITE, mgmt);
  
  QW_RET(code);
}

D
dapan1121 已提交
457
int32_t qwUpdateTaskStatus(QW_FPARAMS_DEF, int8_t status) {
D
dapan1121 已提交
458 459
  SQWSchStatus *sch = NULL;
  SQWTaskStatus *task = NULL;
D
dapan1121 已提交
460 461
  int32_t code = 0;

D
dapan1121 已提交
462
  QW_ERR_RET(qwAcquireScheduler(mgmt, sId, QW_READ, &sch));
D
dapan1121 已提交
463
  QW_ERR_JRET(qwAcquireTaskStatus(QW_FPARAMS(), QW_READ, sch, &task));
D
dapan1121 已提交
464

D
dapan1121 已提交
465
  QW_ERR_JRET(qwSetTaskStatus(QW_FPARAMS(), task, status));
D
dapan1121 已提交
466
  
D
dapan1121 已提交
467 468
_return:

D
dapan1121 已提交
469
  qwReleaseTaskStatus(QW_READ, sch);
D
dapan1121 已提交
470
  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
471 472 473 474

  QW_RET(code);
}

D
dapan1121 已提交
475
int32_t qwExecTask(QW_FPARAMS_DEF, SQWTaskCtx *ctx, bool *queryEnd) {
D
dapan1121 已提交
476 477 478 479
  int32_t code = 0;
  bool  qcontinue = true;
  SSDataBlock* pRes = NULL;
  uint64_t useconds = 0;
D
dapan1121 已提交
480
  int32_t i = 0;
D
dapan1121 已提交
481 482 483
  int32_t execNum = 0;
  qTaskInfo_t *taskHandle = &ctx->taskHandle; 
  DataSinkHandle sinkHandle = ctx->sinkHandle;
D
dapan1121 已提交
484 485
 
  while (true) {
H
Haojun Liao 已提交
486
    QW_TASK_DLOG("start to execTask, loopIdx:%d", i++);
D
dapan1121 已提交
487
    
D
dapan1121 已提交
488
    code = qExecTask(*taskHandle, &pRes, &useconds);
D
dapan1121 已提交
489
    if (code) {
H
Haojun Liao 已提交
490
      QW_TASK_ELOG("qExecTask failed, code:%s", tstrerror(code));
D
dapan1121 已提交
491 492 493
      QW_ERR_JRET(code);
    }

D
dapan1121 已提交
494 495
    ++execNum;

D
dapan1121 已提交
496
    if (NULL == pRes) {
D
dapan1121 已提交
497
      QW_TASK_DLOG("task query done, useconds:%"PRIu64, useconds);
D
dapan1121 已提交
498
      dsEndPut(sinkHandle, useconds);
D
dapan1121 已提交
499
      
D
dapan1121 已提交
500
      if (TASK_TYPE_TEMP == ctx->taskType) {
D
dapan1121 已提交
501 502
        qwFreeTaskHandle(QW_FPARAMS(), taskHandle);
      }
D
dapan1121 已提交
503 504 505 506

      if (queryEnd) {
        *queryEnd = true;
      }
D
dapan1121 已提交
507
      
D
dapan1121 已提交
508 509 510
      break;
    }

511 512
    ASSERT(pRes->info.rows > 0);

H
Haojun Liao 已提交
513
    SInputData inputData = {.pData = pRes};
D
dapan1121 已提交
514 515
    code = dsPutDataBlock(sinkHandle, &inputData, &qcontinue);
    if (code) {
H
Haojun Liao 已提交
516
      QW_TASK_ELOG("dsPutDataBlock failed, code:%s", tstrerror(code));
D
dapan1121 已提交
517 518
      QW_ERR_JRET(code);
    }
D
dapan1121 已提交
519 520 521

    QW_TASK_DLOG("data put into sink, rows:%d, continueExecTask:%d", pRes->info.rows, qcontinue);
    
D
dapan1121 已提交
522 523 524 525
    if (!qcontinue) {
      break;
    }

D
dapan1121 已提交
526 527 528 529 530
    if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_READY) && execNum >= QW_DEFAULT_SHORT_RUN_TIMES) {
      break;
    }

    if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_FETCH)) {
D
dapan1121 已提交
531 532
      break;
    }
D
dapan1121 已提交
533 534 535 536 537 538
  }

_return:

  QW_RET(code);
}
D
dapan1121 已提交
539

D
dapan1121 已提交
540
int32_t qwGenerateSchHbRsp(SQWorkerMgmt *mgmt, SQWSchStatus *sch, SQWHbInfo *hbInfo) {
D
dapan1121 已提交
541 542 543 544 545
  int32_t taskNum = 0;

  QW_LOCK(QW_READ, &sch->tasksLock);
  
  taskNum = taosHashGetSize(sch->tasksHash);
D
dapan1121 已提交
546 547 548

  hbInfo->rsp.taskStatus = taosArrayInit(taskNum, sizeof(STaskStatus));
  if (NULL == hbInfo->rsp.taskStatus) {
D
dapan1121 已提交
549
    QW_UNLOCK(QW_READ, &sch->tasksLock);
D
dapan1121 已提交
550
    QW_ELOG("taosArrayInit taskStatus failed, num:%d", taskNum);
D
dapan1121 已提交
551 552 553
    return TSDB_CODE_QRY_OUT_OF_MEMORY;
  }

D
dapan1121 已提交
554 555 556
  hbInfo->connection = sch->hbConnection;
  hbInfo->rsp.seqId = -1;

D
dapan1121 已提交
557 558 559
  void *key = NULL;
  size_t keyLen = 0;
  int32_t i = 0;
D
dapan1121 已提交
560
  STaskStatus status = {0};
D
dapan1121 已提交
561 562 563 564

  void *pIter = taosHashIterate(sch->tasksHash, NULL);
  while (pIter) {
    SQWTaskStatus *taskStatus = (SQWTaskStatus *)pIter;
D
dapan1121 已提交
565
    key = taosHashGetKey(pIter, &keyLen);
D
dapan1121 已提交
566 567 568

    //TODO GET EXECUTOR API TO GET MORE INFO

D
dapan1121 已提交
569 570 571 572 573
    QW_GET_QTID(key, status.queryId, status.taskId);
    status.status = taskStatus->status;
    status.refId = taskStatus->refId;
    
    taosArrayPush(hbInfo->rsp.taskStatus, &status);
D
dapan1121 已提交
574 575 576 577 578 579 580 581 582 583
    
    ++i;
    pIter = taosHashIterate(sch->tasksHash, pIter);
  }  

  QW_UNLOCK(QW_READ, &sch->tasksLock);

  return TSDB_CODE_SUCCESS;
}

D
dapan 已提交
584 585

int32_t qwGetResFromSink(QW_FPARAMS_DEF, SQWTaskCtx *ctx, int32_t *dataLen, void **rspMsg, SOutputData *pOutput) {
D
dapan1121 已提交
586 587 588
  int32_t len = 0;
  SRetrieveTableRsp *rsp = NULL;
  bool queryEnd = false;
D
dapan1121 已提交
589 590
  int32_t code = 0;

D
dapan1121 已提交
591 592 593 594 595 596 597 598 599 600 601 602 603 604
  if (ctx->emptyRes) {
    QW_TASK_DLOG("query empty result, query end, phase:%d", ctx->phase);
    
    QW_ERR_RET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_SUCCEED));
    
    QW_ERR_RET(qwMallocFetchRsp(len, &rsp));      
    
    *rspMsg = rsp;
    *dataLen = 0;
    pOutput->queryEnd = true;
    
    return TSDB_CODE_SUCCESS;
  }

D
dapan1121 已提交
605
  dsGetDataLength(ctx->sinkHandle, &len, &queryEnd);
D
dapan1121 已提交
606

D
dapan1121 已提交
607 608 609 610
  if (len < 0) {
    QW_TASK_ELOG("invalid length from dsGetDataLength, length:%d", len);
    QW_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
  }
D
dapan1121 已提交
611

D
dapan1121 已提交
612 613
  if (len == 0) {
    if (queryEnd) {
D
dapan 已提交
614
      code = dsGetDataBlock(ctx->sinkHandle, pOutput);
D
dapan1121 已提交
615 616 617 618 619
      if (code) {
        QW_TASK_ELOG("dsGetDataBlock failed, code:%x", code);
        QW_ERR_RET(code);
      }
    
D
dapan1121 已提交
620
      QW_TASK_DLOG("no data in sink and query end, phase:%d", ctx->phase);
D
dapan1121 已提交
621 622
      
      QW_ERR_RET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_SUCCEED));
D
dapan1121 已提交
623

L
Liu Jicong 已提交
624
      QW_ERR_RET(qwMallocFetchRsp(len, &rsp));
D
dapan1121 已提交
625
      *rspMsg = rsp;
D
dapan 已提交
626 627
      *dataLen = 0;
      
D
dapan1121 已提交
628 629
      return TSDB_CODE_SUCCESS;
    }
D
dapan1121 已提交
630 631

    pOutput->bufStatus = DS_BUF_EMPTY;
D
dapan1121 已提交
632
    
D
dapan1121 已提交
633 634
    QW_TASK_DLOG("no res data in sink, need response later, queryEnd:%d", queryEnd);

D
dapan1121 已提交
635
    return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
636
  }
D
dapan1121 已提交
637 638


D
dapan1121 已提交
639 640
  // Got data from sink

D
dapan 已提交
641
  *dataLen = len;
D
dapan1121 已提交
642 643

  QW_TASK_DLOG("task got data in sink, dataLength:%d", len);
D
dapan1121 已提交
644
  
D
dapan1121 已提交
645
  QW_ERR_RET(qwMallocFetchRsp(len, &rsp));
D
dapan 已提交
646
  *rspMsg = rsp;
D
dapan1121 已提交
647
  
D
dapan 已提交
648 649
  pOutput->pData = rsp->data;
  code = dsGetDataBlock(ctx->sinkHandle, pOutput);
D
dapan1121 已提交
650 651 652 653
  if (code) {
    QW_TASK_ELOG("dsGetDataBlock failed, code:%x", code);
    QW_ERR_RET(code);
  }
D
dapan1121 已提交
654

D
dapan1121 已提交
655
  if (DS_BUF_EMPTY == pOutput->bufStatus && pOutput->queryEnd) {
D
dapan1121 已提交
656
    QW_TASK_DLOG_E("task all data fetched, done");
D
dapan1121 已提交
657
    QW_ERR_RET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_SUCCEED));
D
dapan1121 已提交
658 659
  }

D
dapan1121 已提交
660
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
661 662
}

D
dapan1121 已提交
663
int32_t qwHandlePrePhaseEvents(QW_FPARAMS_DEF, int8_t phase, SQWPhaseInput *input, SQWPhaseOutput *output) {
D
dapan1121 已提交
664 665 666 667
  int32_t code = 0;
  int8_t status = 0;
  SQWTaskCtx *ctx = NULL;
  bool locked = false;
D
dapan1121 已提交
668 669
  void *dropConnection = NULL;
  void *cancelConnection = NULL;
D
dapan1121 已提交
670

D
dapan1121 已提交
671
  QW_TASK_DLOG("start to handle event at phase %s", qwPhaseStr(phase));
D
dapan1121 已提交
672 673 674 675 676 677 678 679 680 681

  if (QW_PHASE_PRE_QUERY == phase) {
    QW_ERR_JRET(qwAddAcquireTaskCtx(QW_FPARAMS(), &ctx));
  } else {
    QW_ERR_JRET(qwAcquireTaskCtx(QW_FPARAMS(), &ctx));
  }
  
  QW_LOCK(QW_WRITE, &ctx->lock);
  locked = true;

D
dapan1121 已提交
682 683
  switch (phase) {
    case QW_PHASE_PRE_QUERY: {
D
dapan1121 已提交
684 685
      atomic_store_8(&ctx->phase, phase);

D
dapan1121 已提交
686
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL) || QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
687
        QW_TASK_ELOG("task already cancelled/dropped at wrong phase %s", qwPhaseStr(phase));
D
dapan1121 已提交
688
        
D
dapan1121 已提交
689
        output->needStop = true;
D
dapan1121 已提交
690 691 692
        output->rspCode = TSDB_CODE_QRY_TASK_STATUS_ERROR;
        break;
      }
D
dapan1121 已提交
693

D
dapan1121 已提交
694
      if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
695
        QW_ERR_JRET(qwDropTaskStatus(QW_FPARAMS()));              
D
dapan1121 已提交
696
        QW_ERR_JRET(qwDropTaskCtx(QW_FPARAMS(), QW_WRITE));
D
dapan1121 已提交
697

D
dapan1121 已提交
698
        output->needStop = true;
D
dapan1121 已提交
699
        output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;
D
dapan1121 已提交
700
        QW_SET_RSP_CODE(ctx, output->rspCode);
D
dapan1121 已提交
701
        dropConnection = ctx->dropConnection;
D
dapan1121 已提交
702 703
        
        // Note: ctx freed, no need to unlock it
D
dapan 已提交
704 705 706
        locked = false;      

        break;
D
dapan1121 已提交
707 708 709 710
      } else if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_CANCEL)) {
        QW_ERR_JRET(qwAddTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_CANCELLED));
        
        QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL);
D
dapan1121 已提交
711
        output->needStop = true;                
D
dapan1121 已提交
712
        output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;
D
dapan1121 已提交
713 714 715
        QW_SET_RSP_CODE(ctx, output->rspCode);
        
        cancelConnection = ctx->cancelConnection;
D
dapan 已提交
716 717

        break;
D
dapan1121 已提交
718
      }
D
dapan1121 已提交
719

D
dapan1121 已提交
720
      if (ctx->rspCode) {
D
dapan1121 已提交
721
        QW_TASK_ELOG("task already failed at phase %s, error:%x - %s", qwPhaseStr(phase), ctx->rspCode, tstrerror(ctx->rspCode));
D
dapan1121 已提交
722
        output->needStop = true;
D
dapan1121 已提交
723 724
        output->rspCode = ctx->rspCode;        
        QW_ERR_JRET(output->rspCode);
D
dapan1121 已提交
725
      }
D
dapan1121 已提交
726

D
dapan1121 已提交
727 728
      if (!output->needStop) {
        QW_ERR_JRET(qwAddTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_EXECUTING));
D
dapan1121 已提交
729 730 731
      }
      break;
    }
D
dapan1121 已提交
732
    case QW_PHASE_PRE_FETCH: {
D
dapan1121 已提交
733
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
734
        QW_TASK_WLOG("task already dropped, phase:%s", qwPhaseStr(phase));
D
dapan1121 已提交
735 736 737 738
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_DROPPED);
      }
D
dapan1121 已提交
739
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL)) {
D
dapan1121 已提交
740
        QW_TASK_WLOG("task already cancelled, phase:%s", qwPhaseStr(phase));
D
dapan1121 已提交
741 742 743 744
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
      }
D
dapan1121 已提交
745

D
dapan1121 已提交
746
      if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
747
        QW_TASK_ELOG("drop event at wrong phase %s", qwPhaseStr(phase));
D
dapan1121 已提交
748
        output->needStop = true;
D
dapan1121 已提交
749 750
        output->rspCode = TSDB_CODE_QRY_TASK_STATUS_ERROR;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
D
dapan1121 已提交
751
      } else if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_CANCEL)) {
D
dapan1121 已提交
752
        QW_TASK_ELOG("cancel event at wrong phase %s", qwPhaseStr(phase));
D
dapan1121 已提交
753
        output->needStop = true;
D
dapan1121 已提交
754 755 756 757
        output->rspCode = TSDB_CODE_QRY_TASK_STATUS_ERROR;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
      }

D
dapan1121 已提交
758
      if (ctx->rspCode) {
D
dapan1121 已提交
759
        QW_TASK_ELOG("task already failed at phase %s, error:%x - %s", qwPhaseStr(phase), ctx->rspCode, tstrerror(ctx->rspCode));
D
dapan1121 已提交
760 761 762 763 764
        output->needStop = true;
        output->rspCode = ctx->rspCode;        
        QW_ERR_JRET(output->rspCode);
      }

D
dapan1121 已提交
765
      if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_FETCH)) {
D
dapan1121 已提交
766
        QW_TASK_WLOG("last fetch not finished, phase:%s", qwPhaseStr(phase));
D
dapan1121 已提交
767 768 769 770 771 772
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_DUPLICATTED_OPERATION;        
        QW_ERR_JRET(TSDB_CODE_QRY_DUPLICATTED_OPERATION);
      }

      if (!QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_READY)) {
D
dapan1121 已提交
773
        QW_TASK_ELOG("query rsp are not ready, phase:%s", qwPhaseStr(phase));
D
dapan1121 已提交
774 775 776
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_TASK_MSG_ERROR;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_MSG_ERROR);
D
dapan1121 已提交
777 778 779
      }
      break;
    }    
D
dapan1121 已提交
780
    case QW_PHASE_PRE_CQUERY: {
D
dapan1121 已提交
781 782
      atomic_store_8(&ctx->phase, phase);

D
dapan1121 已提交
783
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL)) {
D
dapan1121 已提交
784
        QW_TASK_WLOG("task already cancelled, phase:%s", qwPhaseStr(phase));
D
dapan1121 已提交
785 786 787 788 789
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
      }

D
dapan1121 已提交
790
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
791
        QW_TASK_WLOG("task already dropped, phase:%s", qwPhaseStr(phase));
D
dapan1121 已提交
792
        output->needStop = true;
D
dapan1121 已提交
793 794
        output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_DROPPED);
D
dapan1121 已提交
795
      }
D
dapan1121 已提交
796

D
dapan1121 已提交
797
      if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
798 799 800 801
        QW_ERR_JRET(qwDropTaskStatus(QW_FPARAMS()));  
        QW_ERR_JRET(qwDropTaskCtx(QW_FPARAMS(), QW_WRITE));
        
        output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;
D
dapan1121 已提交
802
        output->needStop = true;
D
dapan1121 已提交
803 804 805 806 807 808 809
        QW_SET_RSP_CODE(ctx, output->rspCode);
        dropConnection = ctx->dropConnection;
        
        // Note: ctx freed, no need to unlock it
        locked = false;            
      
        QW_ERR_JRET(output->rspCode);
D
dapan1121 已提交
810
      } else if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_CANCEL)) {
D
dapan1121 已提交
811 812 813 814 815 816 817 818 819 820 821
        QW_ERR_JRET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_CANCELLED));
        qwFreeTask(QW_FPARAMS(), ctx);
        
        QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL);
      
        output->needStop = true;        
        output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;
        QW_SET_RSP_CODE(ctx, output->rspCode);
        cancelConnection = ctx->cancelConnection;
      
        QW_ERR_JRET(output->rspCode);
D
dapan1121 已提交
822
      }
D
dapan1121 已提交
823

D
dapan1121 已提交
824

D
dapan1121 已提交
825
      if (ctx->rspCode) {
D
dapan1121 已提交
826
        QW_TASK_ELOG("task already failed at phase %s, error:%x - %s", qwPhaseStr(phase), ctx->rspCode, tstrerror(ctx->rspCode));
D
dapan1121 已提交
827 828 829 830
        output->needStop = true;
        output->rspCode = ctx->rspCode;        
        QW_ERR_JRET(output->rspCode);
      }
D
dapan1121 已提交
831 832
      break;
    }    
D
dapan1121 已提交
833
  }
D
dapan1121 已提交
834

D
dapan1121 已提交
835
_return:
D
dapan1121 已提交
836

D
dapan1121 已提交
837 838 839 840 841 842 843 844
  if (ctx) {
    if (output->rspCode) {
      QW_UPDATE_RSP_CODE(ctx, output->rspCode);
    }
    
    if (locked) {
      QW_UNLOCK(QW_WRITE, &ctx->lock);
    }
D
dapan1121 已提交
845

D
dapan1121 已提交
846 847
    qwReleaseTaskCtx(mgmt, ctx);
  }
D
dapan1121 已提交
848

D
dapan1121 已提交
849 850 851 852 853
  if (code) {
    output->needStop = true;
    if (TSDB_CODE_SUCCESS == output->rspCode) {
      output->rspCode = code;
    }
D
dapan1121 已提交
854 855
  }

D
dapan1121 已提交
856 857
  if (dropConnection) {
    qwBuildAndSendDropRsp(dropConnection, output->rspCode);    
D
dapan1121 已提交
858
    QW_TASK_DLOG("drop msg rsped, code:%x - %s", output->rspCode, tstrerror(output->rspCode));
D
dapan1121 已提交
859
  }
D
dapan1121 已提交
860

D
dapan1121 已提交
861 862
  if (cancelConnection) {
    qwBuildAndSendCancelRsp(cancelConnection, output->rspCode);    
D
dapan1121 已提交
863
    QW_TASK_DLOG("cancel msg rsped, code:%x - %s", output->rspCode, tstrerror(output->rspCode));
D
dapan1121 已提交
864 865
  }

D
dapan1121 已提交
866
  QW_TASK_DLOG("end to handle event at phase %s", qwPhaseStr(phase));
D
dapan1121 已提交
867 868 869 870 871 872 873 874 875 876 877 878 879 880

  QW_RET(code);
}


int32_t qwHandlePostPhaseEvents(QW_FPARAMS_DEF, int8_t phase, SQWPhaseInput *input, SQWPhaseOutput *output) {
  int32_t code = 0;
  int8_t status = 0;
  SQWTaskCtx *ctx = NULL;
  bool locked = false;
  void *readyConnection = NULL;
  void *dropConnection = NULL;
  void *cancelConnection = NULL;

D
dapan1121 已提交
881
  QW_TASK_DLOG("start to handle event at phase %s", qwPhaseStr(phase));
D
dapan1121 已提交
882 883 884 885 886 887 888 889 890

  output->needStop = false;
  
  QW_ERR_JRET(qwAcquireTaskCtx(QW_FPARAMS(), &ctx));
  
  QW_LOCK(QW_WRITE, &ctx->lock);
  locked = true;     

  if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
891
    QW_TASK_WLOG("task already dropped, phase:%s", qwPhaseStr(phase));
D
dapan1121 已提交
892 893 894 895 896 897
    output->needStop = true;
    output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;        
    QW_ERR_JRET(TSDB_CODE_QRY_TASK_DROPPED);
  }
  
  if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL)) {
D
dapan1121 已提交
898
    QW_TASK_WLOG("task already cancelled, phase:%s", qwPhaseStr(phase));
D
dapan1121 已提交
899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916
    output->needStop = true;
    output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;        
    QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
  }

  if (input->code) {
    output->rspCode = input->code;
  }

  if (QW_PHASE_POST_QUERY == phase) {
    if (NULL == ctx->taskHandle && NULL == ctx->sinkHandle) {
      ctx->emptyRes = true;
    }
    
    if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_READY)) {
      readyConnection = ctx->readyConnection;
      QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_READY);
    }
D
dapan1121 已提交
917 918
  }

D
dapan1121 已提交
919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946
  if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
    QW_ERR_JRET(qwDropTaskStatus(QW_FPARAMS()));  
    QW_ERR_JRET(qwDropTaskCtx(QW_FPARAMS(), QW_WRITE));
    
    output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;
    output->needStop = true;
    QW_SET_RSP_CODE(ctx, output->rspCode);
    dropConnection = ctx->dropConnection;
    
    // Note: ctx freed, no need to unlock it
    locked = false;            

    QW_ERR_JRET(output->rspCode);
  } else if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_CANCEL)) {
    QW_ERR_JRET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_CANCELLED));
    qwFreeTask(QW_FPARAMS(), ctx);
    
    QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL);

    output->needStop = true;        
    output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;
    QW_SET_RSP_CODE(ctx, output->rspCode);
    cancelConnection = ctx->cancelConnection;

    QW_ERR_JRET(output->rspCode);
  }

  if (ctx->rspCode) {
D
dapan1121 已提交
947
    QW_TASK_ELOG("task already failed at phase %s, error:%x - %s", qwPhaseStr(phase), ctx->rspCode, tstrerror(ctx->rspCode));
D
dapan1121 已提交
948 949 950 951 952 953 954 955 956 957 958
    output->needStop = true;
    output->rspCode = ctx->rspCode;        
    QW_ERR_JRET(output->rspCode);
  }      

  if (QW_PHASE_POST_QUERY == phase && (!output->needStop)) {      
    QW_ERR_JRET(qwUpdateTaskStatus(QW_FPARAMS(), input->taskStatus));
  }

_return:

D
dapan1121 已提交
959
  if (ctx) {
D
dapan1121 已提交
960 961 962 963
    if (output->rspCode) {
      QW_UPDATE_RSP_CODE(ctx, output->rspCode);
    }

D
dapan1121 已提交
964 965 966
    if (QW_PHASE_POST_FETCH != phase) {
      atomic_store_8(&ctx->phase, phase);
    }
D
dapan1121 已提交
967 968 969 970 971
    
    if (locked) {
      QW_UNLOCK(QW_WRITE, &ctx->lock);
    }
    
D
dapan1121 已提交
972
    qwReleaseTaskCtx(mgmt, ctx);
D
dapan1121 已提交
973 974
  }

D
dapan1121 已提交
975 976 977 978 979 980 981
  if (code) {
    output->needStop = true;
    if (TSDB_CODE_SUCCESS == output->rspCode) {
      output->rspCode = code;
    }
  }

D
dapan1121 已提交
982 983
  if (readyConnection) {
    qwBuildAndSendReadyRsp(readyConnection, output->rspCode);    
D
dapan1121 已提交
984
    QW_TASK_DLOG("ready msg rsped, code:%x - %s", output->rspCode, tstrerror(output->rspCode));
D
dapan1121 已提交
985 986 987 988
  }

  if (dropConnection) {
    qwBuildAndSendDropRsp(dropConnection, output->rspCode);    
D
dapan1121 已提交
989
    QW_TASK_DLOG("drop msg rsped, code:%x - %s", output->rspCode, tstrerror(output->rspCode));
D
dapan1121 已提交
990 991 992 993
  }

  if (cancelConnection) {
    qwBuildAndSendCancelRsp(cancelConnection, output->rspCode);    
D
dapan1121 已提交
994
    QW_TASK_DLOG("cancel msg rsped, code:%x - %s", output->rspCode, tstrerror(output->rspCode));
D
dapan1121 已提交
995 996
  }

D
dapan1121 已提交
997
  QW_TASK_DLOG("end to handle event at phase %s", qwPhaseStr(phase));
D
dapan1121 已提交
998

D
dapan1121 已提交
999 1000 1001 1002
  QW_RET(code);
}


D
dapan1121 已提交
1003
int32_t qwProcessQuery(QW_FPARAMS_DEF, SQWMsg *qwMsg, int8_t taskType) {
D
dapan1121 已提交
1004 1005 1006 1007 1008 1009 1010
  int32_t code = 0;
  bool queryRsped = false;
  bool needStop = false;
  struct SSubplan *plan = NULL;
  int32_t rspCode = 0;
  SQWPhaseInput input = {0};
  SQWPhaseOutput output = {0};
D
dapan1121 已提交
1011 1012
  qTaskInfo_t pTaskInfo = NULL;
  DataSinkHandle sinkHandle = NULL;
D
dapan1121 已提交
1013
  SQWTaskCtx *ctx = NULL;
D
dapan1121 已提交
1014

D
dapan1121 已提交
1015
  QW_ERR_JRET(qwHandlePrePhaseEvents(QW_FPARAMS(), QW_PHASE_PRE_QUERY, &input, &output));
D
dapan1121 已提交
1016

D
dapan1121 已提交
1017 1018
  needStop = output.needStop;
  code = output.rspCode;
D
dapan1121 已提交
1019
  
D
dapan1121 已提交
1020
  if (needStop) {
D
dapan1121 已提交
1021
    QW_TASK_DLOG("task need stop, phase:%d", QW_PHASE_PRE_QUERY);
D
dapan1121 已提交
1022
    QW_ERR_JRET(code);
D
dapan1121 已提交
1023
  }
D
dapan1121 已提交
1024 1025 1026 1027

  QW_ERR_JRET(qwGetTaskCtx(QW_FPARAMS(), &ctx));
  
  atomic_store_8(&ctx->taskType, taskType);
D
dapan1121 已提交
1028
  
D
dapan1121 已提交
1029 1030
  code = qStringToSubplan(qwMsg->msg, &plan);
  if (TSDB_CODE_SUCCESS != code) {
H
Haojun Liao 已提交
1031
    QW_TASK_ELOG("task string to subplan failed, code:%s", tstrerror(code));
D
dapan1121 已提交
1032
    QW_ERR_JRET(code);
D
dapan1121 已提交
1033
  }
D
dapan1121 已提交
1034
  
D
dapan1121 已提交
1035
  code = qCreateExecTask(qwMsg->node, 0, tId, (struct SSubplan *)plan, &pTaskInfo, &sinkHandle);
D
dapan1121 已提交
1036
  if (code) {
H
Haojun Liao 已提交
1037
    QW_TASK_ELOG("qCreateExecTask failed, code:%s", tstrerror(code));
D
dapan1121 已提交
1038
    QW_ERR_JRET(code);
D
dapan1121 已提交
1039
  }
D
dapan1121 已提交
1040

H
Haojun Liao 已提交
1041
  if (NULL == sinkHandle || NULL == pTaskInfo) {
D
dapan1121 已提交
1042 1043 1044 1045 1046
    QW_TASK_ELOG("create task result error, taskHandle:%p, sinkHandle:%p", pTaskInfo, sinkHandle);
    QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
  }

  //TODO OPTIMIZE EMTYP RESULT QUERY RSP TO AVOID FURTHER FETCH
D
dapan1121 已提交
1047
  
D
dapan1121 已提交
1048
  QW_ERR_JRET(qwBuildAndSendQueryRsp(qwMsg->connection, code));
D
dapan1121 已提交
1049
  QW_TASK_DLOG("query msg rsped, code:%d", code);
D
dapan1121 已提交
1050

D
dapan1121 已提交
1051
  queryRsped = true;
D
dapan1121 已提交
1052

D
dapan1121 已提交
1053 1054 1055
  atomic_store_ptr(&ctx->taskHandle, pTaskInfo);
  atomic_store_ptr(&ctx->sinkHandle, sinkHandle);

D
dapan1121 已提交
1056
  if (pTaskInfo && sinkHandle) {
D
dapan1121 已提交
1057
    QW_ERR_JRET(qwExecTask(QW_FPARAMS(), ctx, NULL));
D
dapan1121 已提交
1058
  }
D
dapan1121 已提交
1059
  
D
dapan1121 已提交
1060 1061
_return:

D
dapan1121 已提交
1062
  if (code) {
D
dapan1121 已提交
1063
    rspCode = code;
D
dapan1121 已提交
1064
  }
D
dapan1121 已提交
1065 1066
  
  if (!queryRsped) {
D
dapan1121 已提交
1067
    qwBuildAndSendQueryRsp(qwMsg->connection, rspCode);
D
dapan1121 已提交
1068
    QW_TASK_DLOG("query msg rsped, code:%x", rspCode);
D
dapan1121 已提交
1069 1070
  }

D
dapan1121 已提交
1071
  input.code = rspCode;
D
dapan1121 已提交
1072
  input.taskStatus = rspCode ? JOB_TASK_STATUS_FAILED : JOB_TASK_STATUS_PARTIAL_SUCCEED;
D
dapan1121 已提交
1073
  
D
dapan1121 已提交
1074
  QW_ERR_RET(qwHandlePostPhaseEvents(QW_FPARAMS(), QW_PHASE_POST_QUERY, &input, &output));
D
dapan1121 已提交
1075 1076
  
  QW_RET(rspCode);
D
dapan1121 已提交
1077 1078
}

D
dapan1121 已提交
1079
int32_t qwProcessReady(QW_FPARAMS_DEF, SQWMsg *qwMsg) {
D
dapan1121 已提交
1080 1081
  int32_t code = 0;
  SQWTaskCtx *ctx = NULL;
D
dapan1121 已提交
1082
  int8_t phase = 0;
D
dapan1121 已提交
1083 1084
  bool needRsp = false;
  int32_t rspCode = 0;
D
dapan1121 已提交
1085 1086

  QW_ERR_JRET(qwAcquireTaskCtx(QW_FPARAMS(), &ctx));
D
dapan1121 已提交
1087

D
dapan1121 已提交
1088 1089
  QW_LOCK(QW_WRITE, &ctx->lock);

D
dapan1121 已提交
1090 1091 1092 1093 1094 1095
  if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL) || QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_DROP) ||
      QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_CANCEL) || QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
    QW_TASK_WLOG("task already cancelled/dropped, phase:%d", phase);
    QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
  }

D
dapan1121 已提交
1096 1097 1098
  phase = QW_GET_PHASE(ctx);
  
  if (phase == QW_PHASE_PRE_QUERY) {
D
dapan1121 已提交
1099
    QW_SET_EVENT_RECEIVED(ctx, QW_EVENT_READY);
D
dapan1121 已提交
1100
    ctx->readyConnection = qwMsg->connection;
D
dapan1121 已提交
1101
    QW_TASK_DLOG("ready msg not rsped, phase:%d", phase);
D
dapan1121 已提交
1102
  } else if (phase == QW_PHASE_POST_QUERY) {
D
dapan1121 已提交
1103
    QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_READY);
D
dapan1121 已提交
1104 1105
    needRsp = true;
    rspCode = ctx->rspCode;
D
dapan1121 已提交
1106 1107
  } else {
    QW_TASK_ELOG("invalid phase when got ready msg, phase:%d", phase);
D
dapan1121 已提交
1108 1109 1110 1111
    QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_READY);
    needRsp = true;    
    rspCode = TSDB_CODE_QRY_TASK_STATUS_ERROR;
    QW_ERR_JRET(TSDB_CODE_QRY_TASK_STATUS_ERROR);
D
dapan1121 已提交
1112 1113 1114 1115
  }

_return:

D
dapan1121 已提交
1116
  if (code && ctx) {
D
dapan1121 已提交
1117
    QW_UPDATE_RSP_CODE(ctx, code);
D
dapan1121 已提交
1118 1119
  }

D
dapan1121 已提交
1120 1121
  if (ctx) {
    QW_UNLOCK(QW_WRITE, &ctx->lock);
D
dapan1121 已提交
1122
    qwReleaseTaskCtx(mgmt, ctx);
D
dapan1121 已提交
1123 1124
  }

D
dapan1121 已提交
1125 1126 1127 1128 1129
  if (needRsp) {
    qwBuildAndSendReadyRsp(qwMsg->connection, rspCode);
    QW_TASK_DLOG("ready msg rsped, code:%x", rspCode);
  }

D
dapan1121 已提交
1130 1131 1132 1133
  QW_RET(code);
}


D
dapan1121 已提交
1134
int32_t qwProcessCQuery(QW_FPARAMS_DEF, SQWMsg *qwMsg) {
D
dapan1121 已提交
1135
  SQWTaskCtx *ctx = NULL;
D
dapan1121 已提交
1136
  int32_t code = 0;
1137 1138 1139 1140 1141 1142 1143
  bool queryRsped = false;
  bool needStop = false;
  struct SSubplan *plan = NULL;
  SQWPhaseInput input = {0};
  SQWPhaseOutput output = {0};
  void *rsp = NULL;
  int32_t dataLen = 0;
D
dapan1121 已提交
1144
  bool queryEnd = false;
D
dapan1121 已提交
1145 1146
  
  do {
D
dapan1121 已提交
1147
    QW_ERR_JRET(qwHandlePrePhaseEvents(QW_FPARAMS(), QW_PHASE_PRE_CQUERY, &input, &output));
D
dapan1121 已提交
1148

D
dapan1121 已提交
1149 1150 1151 1152 1153 1154 1155
    needStop = output.needStop;
    code = output.rspCode;
    
    if (needStop) {
      QW_TASK_DLOG("task need stop, phase:%d", QW_PHASE_PRE_CQUERY);
      QW_ERR_JRET(code);
    }
D
dapan1121 已提交
1156

D
dapan1121 已提交
1157
    QW_ERR_JRET(qwGetTaskCtx(QW_FPARAMS(), &ctx));
D
dapan1121 已提交
1158

D
dapan1121 已提交
1159
    atomic_store_8(&ctx->queryInQueue, 0);
D
dapan1121 已提交
1160
    atomic_store_8(&ctx->queryContinue, 0);
D
dapan1121 已提交
1161

D
dapan1121 已提交
1162
    QW_ERR_JRET(qwExecTask(QW_FPARAMS(), ctx, &queryEnd));
D
dapan1121 已提交
1163

D
dapan1121 已提交
1164 1165 1166 1167 1168 1169
    if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_FETCH)) {
      SOutputData sOutput = {0};
      QW_ERR_JRET(qwGetResFromSink(QW_FPARAMS(), ctx, &dataLen, &rsp, &sOutput));
      
      if ((!sOutput.queryEnd) && (DS_BUF_LOW == sOutput.bufStatus || DS_BUF_EMPTY == sOutput.bufStatus)) {    
        QW_TASK_DLOG("task not end, need to continue, bufStatus:%d", sOutput.bufStatus);
D
dapan1121 已提交
1170 1171
        
        // RC WARNING
D
dapan1121 已提交
1172
        atomic_store_8(&ctx->queryContinue, 1);
1173
      }
D
dapan1121 已提交
1174 1175
      
      if (rsp) {
D
dapan1121 已提交
1176 1177
        bool qComplete = (DS_BUF_EMPTY == sOutput.bufStatus && sOutput.queryEnd);
        qwBuildFetchRsp(rsp, &sOutput, dataLen, qComplete);
D
dapan1121 已提交
1178 1179
        
        QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_FETCH);            
1180
        
D
dapan1121 已提交
1181 1182
        qwBuildAndSendFetchRsp(qwMsg->connection, rsp, dataLen, code);                
        QW_TASK_DLOG("fetch msg rsped, code:%x, dataLen:%d", code, dataLen);
D
dapan1121 已提交
1183 1184
      } else {
        atomic_store_8(&ctx->queryContinue, 1);
1185 1186 1187
      }
    }

D
dapan1121 已提交
1188 1189 1190 1191
    if (queryEnd) {
      needStop = true;
    }

D
dapan1121 已提交
1192
  _return:
1193

D
dapan1121 已提交
1194 1195 1196 1197
    if (NULL == ctx) {
      break;
    }

D
dapan1121 已提交
1198
    if (code && QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_FETCH)) {
D
dapan1121 已提交
1199
      QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_FETCH);    
1200 1201 1202
      qwFreeFetchRsp(rsp);
      rsp = NULL;
      qwBuildAndSendFetchRsp(qwMsg->connection, rsp, 0, code);
D
dapan1121 已提交
1203
      QW_TASK_DLOG("fetch msg rsped, code:%x, dataLen:%d", code, 0);      
1204
    }
D
dapan1121 已提交
1205

D
dapan1121 已提交
1206 1207 1208 1209 1210 1211 1212 1213 1214
    QW_LOCK(QW_WRITE, &ctx->lock);
    if (needStop || code || 0 == atomic_load_8(&ctx->queryContinue)) {
      atomic_store_8(&ctx->phase, 0);
      QW_UNLOCK(QW_WRITE,&ctx->lock);
      break;
    }
    
    QW_UNLOCK(QW_WRITE,&ctx->lock);
  } while (true);
D
dapan1121 已提交
1215

D
dapan1121 已提交
1216 1217
  input.code = code;
  qwHandlePostPhaseEvents(QW_FPARAMS(), QW_PHASE_POST_CQUERY, &input, &output);    
D
dapan1121 已提交
1218 1219

  QW_RET(code);
D
dapan1121 已提交
1220
}
D
dapan1121 已提交
1221

D
dapan1121 已提交
1222

D
dapan1121 已提交
1223
int32_t qwProcessFetch(QW_FPARAMS_DEF, SQWMsg *qwMsg) {
D
dapan1121 已提交
1224
  int32_t code = 0;
D
dapan1121 已提交
1225 1226
  int32_t needRsp = true;
  void *data = NULL;
1227
  int32_t sinkStatus = 0;
D
dapan1121 已提交
1228
  int32_t dataLen = 0;
1229
  bool queryEnd = false;
D
dapan1121 已提交
1230 1231 1232
  bool needStop = false;
  bool locked = false;
  SQWTaskCtx *ctx = NULL;
D
dapan1121 已提交
1233
  int8_t status = 0;
D
dapan1121 已提交
1234
  void *rsp = NULL;
D
dapan1121 已提交
1235

D
dapan1121 已提交
1236 1237
  SQWPhaseInput input = {0};
  SQWPhaseOutput output = {0};
D
dapan1121 已提交
1238

D
dapan1121 已提交
1239
  QW_ERR_JRET(qwHandlePrePhaseEvents(QW_FPARAMS(), QW_PHASE_PRE_FETCH, &input, &output));
D
dapan1121 已提交
1240 1241 1242 1243 1244
  
  needStop = output.needStop;
  code = output.rspCode;
  
  if (needStop) {
D
dapan1121 已提交
1245
    QW_TASK_DLOG("task need stop, phase:%d", QW_PHASE_PRE_FETCH);
D
dapan1121 已提交
1246
    QW_ERR_JRET(code);
D
dapan1121 已提交
1247
  }
1248

1249 1250
  QW_ERR_JRET(qwGetTaskCtx(QW_FPARAMS(), &ctx));
 
D
dapan 已提交
1251 1252
  SOutputData sOutput = {0};
  QW_ERR_JRET(qwGetResFromSink(QW_FPARAMS(), ctx, &dataLen, &rsp, &sOutput));
D
dapan1121 已提交
1253

1254 1255
  if (NULL == rsp) {
    QW_SET_EVENT_RECEIVED(ctx, QW_EVENT_FETCH);
D
dapan1121 已提交
1256
  } else {
D
dapan1121 已提交
1257 1258
    bool qComplete = (DS_BUF_EMPTY == sOutput.bufStatus && sOutput.queryEnd);
    qwBuildFetchRsp(rsp, &sOutput, dataLen, qComplete);
D
dapan1121 已提交
1259 1260
  }

D
dapan1121 已提交
1261
  if ((!sOutput.queryEnd) && (DS_BUF_LOW == sOutput.bufStatus || DS_BUF_EMPTY == sOutput.bufStatus)) {    
1262
    QW_TASK_DLOG("task not end, need to continue, bufStatus:%d", sOutput.bufStatus);
D
dapan1121 已提交
1263

D
dapan1121 已提交
1264 1265
    QW_LOCK(QW_WRITE, &ctx->lock);
    locked = true;
1266

D
dapan1121 已提交
1267
    // RC WARNING
D
dapan1121 已提交
1268
    if (QW_IS_QUERY_RUNNING(ctx)) {
D
dapan1121 已提交
1269
      atomic_store_8(&ctx->queryContinue, 1);
D
dapan1121 已提交
1270
    } else if (0 == atomic_load_8(&ctx->queryInQueue)) {
D
dapan1121 已提交
1271 1272 1273 1274
      if (!ctx->multiExec) {
        QW_ERR_JRET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_EXECUTING));      
        ctx->multiExec = true;
      }
D
dapan1121 已提交
1275

D
dapan1121 已提交
1276
      atomic_store_8(&ctx->queryInQueue, 1);
1277
      
D
dapan1121 已提交
1278
      QW_ERR_JRET(qwBuildAndSendCQueryMsg(QW_FPARAMS(), qwMsg->connection));
D
dapan1121 已提交
1279 1280

      QW_TASK_DLOG("schedule query in queue, phase:%d", ctx->phase);
1281
    }
D
dapan 已提交
1282 1283
  }
  
D
dapan1121 已提交
1284
_return:
D
dapan1121 已提交
1285

D
dapan1121 已提交
1286 1287 1288 1289 1290 1291
  if (locked) {
    QW_UNLOCK(QW_WRITE, &ctx->lock);
  }

  input.code = code;

D
dapan1121 已提交
1292 1293 1294 1295 1296
  qwHandlePostPhaseEvents(QW_FPARAMS(), QW_PHASE_POST_FETCH, &input, &output);

  if (output.rspCode) {
    code = output.rspCode;
  }
D
dapan1121 已提交
1297

D
dapan 已提交
1298 1299 1300
  if (code) {
    qwFreeFetchRsp(rsp);
    rsp = NULL;
D
dapan1121 已提交
1301 1302 1303
    dataLen = 0;
    qwBuildAndSendFetchRsp(qwMsg->connection, rsp, dataLen, code);
    QW_TASK_DLOG("fetch msg rsped, code:%x, dataLen:%d", code, dataLen);
D
dapan 已提交
1304 1305
  } else if (rsp) {
    qwBuildAndSendFetchRsp(qwMsg->connection, rsp, dataLen, code);
D
dapan1121 已提交
1306
    QW_TASK_DLOG("fetch msg rsped, code:%x, dataLen:%d", code, dataLen);
D
dapan1121 已提交
1307 1308
  }

D
dapan1121 已提交
1309 1310
  QW_RET(code);
}
D
dapan1121 已提交
1311

1312

D
dapan1121 已提交
1313
int32_t qwProcessDrop(QW_FPARAMS_DEF, SQWMsg *qwMsg) {
D
dapan1121 已提交
1314 1315
  int32_t code = 0;
  bool needRsp = false;
D
dapan1121 已提交
1316 1317 1318
  SQWTaskCtx *ctx = NULL;
  bool locked = false;

D
dapan1121 已提交
1319
  QW_ERR_JRET(qwAddAcquireTaskCtx(QW_FPARAMS(), &ctx));
D
dapan1121 已提交
1320
  
D
dapan1121 已提交
1321 1322 1323 1324 1325 1326 1327 1328 1329
  QW_LOCK(QW_WRITE, &ctx->lock);

  locked = true;

  if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
    QW_TASK_WLOG("task already dropping, phase:%d", ctx->phase);
    QW_ERR_JRET(TSDB_CODE_QRY_DUPLICATTED_OPERATION);
  }

D
dapan1121 已提交
1330
  if (QW_IS_QUERY_RUNNING(ctx)) {
D
dapan1121 已提交
1331 1332 1333 1334 1335
    QW_ERR_JRET(qwKillTaskHandle(QW_FPARAMS(), ctx));
    
    QW_ERR_JRET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_DROPPING));
  } else if (ctx->phase > 0) {
    QW_ERR_JRET(qwDropTaskStatus(QW_FPARAMS()));
D
dapan1121 已提交
1336
    QW_ERR_JRET(qwDropTaskCtx(QW_FPARAMS(), QW_WRITE));
D
dapan1121 已提交
1337

D
dapan1121 已提交
1338 1339
    QW_SET_RSP_CODE(ctx, TSDB_CODE_QRY_TASK_DROPPED);

D
dapan1121 已提交
1340 1341 1342
    locked = false;
    needRsp = true;
  }
D
dapan1121 已提交
1343

D
dapan 已提交
1344 1345 1346
  if (!needRsp) {    
    ctx->dropConnection = qwMsg->connection;
    
D
dapan1121 已提交
1347 1348 1349
    QW_SET_EVENT_RECEIVED(ctx, QW_EVENT_DROP);
  }
  
D
dapan1121 已提交
1350
_return:
D
dapan1121 已提交
1351

D
dapan1121 已提交
1352
  if (code) {
D
dapan1121 已提交
1353
    QW_UPDATE_RSP_CODE(ctx, code);
D
dapan1121 已提交
1354 1355
  }

D
dapan 已提交
1356 1357 1358 1359
  if (locked) {
    QW_UNLOCK(QW_WRITE, &ctx->lock);
  }

D
dapan1121 已提交
1360
  if (ctx) {
D
dapan1121 已提交
1361
    qwReleaseTaskCtx(mgmt, ctx);
D
dapan1121 已提交
1362 1363
  }

D
dapan1121 已提交
1364 1365
  if (TSDB_CODE_SUCCESS != code || needRsp) {
    QW_ERR_RET(qwBuildAndSendDropRsp(qwMsg->connection, code));
D
dapan1121 已提交
1366 1367

    QW_TASK_DLOG("drop msg rsped, code:%x", code);
D
dapan1121 已提交
1368
  }
D
dapan1121 已提交
1369

D
dapan1121 已提交
1370
  QW_RET(code);
D
dapan1121 已提交
1371
}
D
dapan1121 已提交
1372

D
dapan1121 已提交
1373 1374 1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387
int32_t qwProcessHb(SQWorkerMgmt *mgmt, SQWMsg *qwMsg, SSchedulerHbReq *req) {
  int32_t code = 0;
  SSchedulerHbRsp rsp = {0};
  SQWSchStatus *sch = NULL;
  uint64_t seqId = 0;

  memcpy(&rsp.epId, &req->epId, sizeof(req->epId));
  
  QW_ERR_JRET(qwAcquireAddScheduler(mgmt, req->sId, QW_READ, &sch));

  atomic_store_ptr(&sch->hbConnection, qwMsg->connection);
  ++sch->hbSeqId;

  rsp.seqId = sch->hbSeqId;

D
dapan1121 已提交
1388 1389 1390
  QW_DLOG("hb connection updated, seqId:%" PRIx64 ", sId:%" PRIx64 ", nodeId:%d, fqdn:%s, port:%d, connection:%p",
    sch->hbSeqId, req->sId, req->epId.nodeId, req->epId.ep.fqdn, req->epId.ep.port, qwMsg->connection);

D
dapan1121 已提交
1391 1392 1393 1394 1395 1396 1397 1398 1399 1400 1401
  qwReleaseScheduler(QW_READ, mgmt);

_return:

  qwBuildAndSendHbRsp(qwMsg->connection, &rsp, code);
  
  QW_RET(code);
}


void qwProcessHbTimerEvent(void *param, void *tmrId) {
D
dapan1121 已提交
1402 1403
  return;
  
D
dapan1121 已提交
1404 1405 1406
  SQWorkerMgmt *mgmt = (SQWorkerMgmt *)param;
  SQWSchStatus *sch = NULL;
  int32_t taskNum = 0;
D
dapan1121 已提交
1407
  SQWHbInfo *rspList = NULL;
D
dapan1121 已提交
1408 1409 1410 1411 1412 1413 1414
  int32_t code = 0;

  QW_LOCK(QW_READ, &mgmt->schLock);

  int32_t schNum = taosHashGetSize(mgmt->schHash);
  if (schNum <= 0) {
    QW_UNLOCK(QW_READ, &mgmt->schLock);
D
dapan1121 已提交
1415 1416
    taosTmrReset(qwProcessHbTimerEvent, QW_DEFAULT_HEARTBEAT_MSEC, param, mgmt->timer, &mgmt->hbTimer);
    return;
D
dapan1121 已提交
1417 1418
  }

D
dapan1121 已提交
1419
  rspList = calloc(schNum, sizeof(SQWHbInfo));
D
dapan1121 已提交
1420 1421
  if (NULL == rspList) {
    QW_UNLOCK(QW_READ, &mgmt->schLock);
D
dapan1121 已提交
1422 1423 1424
    QW_ELOG("calloc %d SQWHbInfo failed", schNum);
    taosTmrReset(qwProcessHbTimerEvent, QW_DEFAULT_HEARTBEAT_MSEC, param, mgmt->timer, &mgmt->hbTimer);
    return;
D
dapan1121 已提交
1425 1426 1427 1428 1429 1430 1431 1432 1433 1434 1435 1436 1437 1438 1439 1440 1441 1442 1443 1444 1445 1446 1447
  }

  void *key = NULL;
  size_t keyLen = 0;
  int32_t i = 0;

  void *pIter = taosHashIterate(mgmt->schHash, NULL);
  while (pIter) {
    code = qwGenerateSchHbRsp(mgmt, (SQWSchStatus *)pIter, &rspList[i]);
    if (code) {
      taosHashCancelIterate(mgmt->schHash, pIter);
      QW_ERR_JRET(code);
    }

    ++i;
    pIter = taosHashIterate(mgmt->schHash, pIter);
  }

_return:

  QW_UNLOCK(QW_READ, &mgmt->schLock);

  for (int32_t j = 0; j < i; ++j) {
D
dapan1121 已提交
1448
    QW_DLOG("hb on connection %p, taskNum:%d", rspList[j].connection, (rspList[j].rsp.taskStatus ? (int32_t)taosArrayGetSize(rspList[j].rsp.taskStatus) : 0));
D
dapan1121 已提交
1449 1450
    qwBuildAndSendHbRsp(rspList[j].connection, &rspList[j].rsp, code);
    tFreeSSchedulerHbRsp(&rspList[j].rsp);
D
dapan1121 已提交
1451 1452 1453 1454
  }

  tfree(rspList);

D
dapan1121 已提交
1455
  taosTmrReset(qwProcessHbTimerEvent, QW_DEFAULT_HEARTBEAT_MSEC, param, mgmt->timer, &mgmt->hbTimer);  
D
dapan1121 已提交
1456 1457
}

S
Shengliang 已提交
1458 1459 1460
int32_t qWorkerInit(int8_t nodeType, int32_t nodeId, SQWorkerCfg *cfg, void **qWorkerMgmt, void *nodeObj,
                    putReqToQueryQFp fp1, sendReqToDnodeFp fp2) {
  if (NULL == qWorkerMgmt || NULL == nodeObj || NULL == fp1 || NULL == fp2) {
D
dapan1121 已提交
1461 1462 1463
    qError("invalid param to init qworker");
    QW_RET(TSDB_CODE_QRY_INVALID_INPUT);
  }
S
Shengliang 已提交
1464

D
dapan1121 已提交
1465
  int32_t code = 0;
D
dapan1121 已提交
1466 1467 1468
  SQWorkerMgmt *mgmt = calloc(1, sizeof(SQWorkerMgmt));
  if (NULL == mgmt) {
    qError("calloc %d failed", (int32_t)sizeof(SQWorkerMgmt));
D
dapan1121 已提交
1469
    QW_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
1470 1471 1472 1473
  }

  if (cfg) {
    mgmt->cfg = *cfg;
D
dapan1121 已提交
1474
    if (0 == mgmt->cfg.maxSchedulerNum) {
D
dapan1121 已提交
1475
      mgmt->cfg.maxSchedulerNum = QW_DEFAULT_SCHEDULER_NUMBER;
D
dapan1121 已提交
1476 1477
    }
    if (0 == mgmt->cfg.maxTaskNum) {
D
dapan1121 已提交
1478
      mgmt->cfg.maxTaskNum = QW_DEFAULT_TASK_NUMBER;
D
dapan1121 已提交
1479 1480
    }
    if (0 == mgmt->cfg.maxSchTaskNum) {
D
dapan1121 已提交
1481
      mgmt->cfg.maxSchTaskNum = QW_DEFAULT_SCH_TASK_NUMBER;
D
dapan1121 已提交
1482
    }
D
dapan1121 已提交
1483
  } else {
D
dapan1121 已提交
1484 1485 1486
    mgmt->cfg.maxSchedulerNum = QW_DEFAULT_SCHEDULER_NUMBER;
    mgmt->cfg.maxTaskNum = QW_DEFAULT_TASK_NUMBER;
    mgmt->cfg.maxSchTaskNum = QW_DEFAULT_SCH_TASK_NUMBER;
D
dapan1121 已提交
1487 1488
  }

1489
  mgmt->schHash = taosHashInit(mgmt->cfg.maxSchedulerNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_UBIGINT), false, HASH_ENTRY_LOCK);
D
dapan1121 已提交
1490
  if (NULL == mgmt->schHash) {
D
dapan1121 已提交
1491
    tfree(mgmt);
D
dapan1121 已提交
1492
    qError("init %d scheduler hash failed", mgmt->cfg.maxSchedulerNum);
D
dapan1121 已提交
1493
    QW_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
1494 1495
  }

1496
  mgmt->ctxHash = taosHashInit(mgmt->cfg.maxTaskNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_ENTRY_LOCK);
D
dapan1121 已提交
1497
  if (NULL == mgmt->ctxHash) {
D
dapan1121 已提交
1498
    qError("init %d task ctx hash failed", mgmt->cfg.maxTaskNum);
D
dapan1121 已提交
1499 1500 1501 1502 1503 1504 1505 1506 1507
    QW_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

  mgmt->timer = taosTmrInit(0, 0, 0, "qworker");
  if (NULL == mgmt->timer) {
    qError("init timer failed, error:%s", tstrerror(terrno));
    QW_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
1508
  mgmt->hbTimer = taosTmrStart(qwProcessHbTimerEvent, QW_DEFAULT_HEARTBEAT_MSEC, mgmt, mgmt->timer);
D
dapan1121 已提交
1509 1510 1511
  if (NULL == mgmt->hbTimer) {
    qError("start hb timer failed");
    QW_ERR_JRET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
1512 1513
  }

D
dapan1121 已提交
1514 1515 1516
  mgmt->nodeType = nodeType;
  mgmt->nodeId = nodeId;
  mgmt->nodeObj = nodeObj;
S
Shengliang 已提交
1517 1518
  mgmt->putToQueueFp = fp1;
  mgmt->sendReqFp = fp2;
D
dapan1121 已提交
1519

D
dapan1121 已提交
1520 1521
  *qWorkerMgmt = mgmt;

D
dapan1121 已提交
1522 1523
  qDebug("qworker initialized for node, type:%d, id:%d, handle:%p", mgmt->nodeType, mgmt->nodeId, mgmt);

D
dapan1121 已提交
1524
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
1525 1526 1527 1528 1529 1530 1531 1532 1533 1534 1535

_return:

  taosHashCleanup(mgmt->schHash);
  taosHashCleanup(mgmt->ctxHash);

  taosTmrCleanUp(mgmt->timer);
  
  tfree(mgmt);

  QW_RET(code);
D
dapan1121 已提交
1536
}
D
dapan1121 已提交
1537 1538 1539 1540

void qWorkerDestroy(void **qWorkerMgmt) {
  if (NULL == qWorkerMgmt || NULL == *qWorkerMgmt) {
    return;
D
dapan1121 已提交
1541
  }
D
dapan 已提交
1542

D
dapan1121 已提交
1543
  SQWorkerMgmt *mgmt = *qWorkerMgmt;
D
dapan1121 已提交
1544 1545 1546

  taosTmrStopA(&mgmt->hbTimer);
  taosTmrCleanUp(mgmt->timer);
D
dapan1121 已提交
1547
  
D
dapan1121 已提交
1548
  //TODO STOP ALL QUERY
D
dapan1121 已提交
1549

D
dapan1121 已提交
1550
  //TODO FREE ALL
D
dapan1121 已提交
1551

D
dapan1121 已提交
1552 1553
  tfree(*qWorkerMgmt);
}
D
dapan1121 已提交
1554

D
dapan1121 已提交
1555
int32_t qwGetSchTasksStatus(SQWorkerMgmt *mgmt, uint64_t sId, SSchedulerStatusRsp **rsp) {
D
dapan1121 已提交
1556
/*
D
dapan1121 已提交
1557 1558
  SQWSchStatus *sch = NULL;
  int32_t taskNum = 0;
1559

D
dapan1121 已提交
1560
  QW_ERR_RET(qwAcquireScheduler(mgmt, sId, QW_READ, &sch));
D
dapan1121 已提交
1561 1562
  
  sch->lastAccessTs = taosGetTimestampSec();
1563

D
dapan1121 已提交
1564 1565 1566 1567 1568 1569 1570
  QW_LOCK(QW_READ, &sch->tasksLock);
  
  taskNum = taosHashGetSize(sch->tasksHash);
  
  int32_t size = sizeof(SSchedulerStatusRsp) + sizeof((*rsp)->status[0]) * taskNum;
  *rsp = calloc(1, size);
  if (NULL == *rsp) {
D
dapan1121 已提交
1571
    QW_SCH_ELOG("calloc %d failed", size);
D
dapan1121 已提交
1572 1573 1574 1575 1576
    QW_UNLOCK(QW_READ, &sch->tasksLock);
    qwReleaseScheduler(QW_READ, mgmt);
    
    return TSDB_CODE_QRY_OUT_OF_MEMORY;
  }
1577

D
dapan1121 已提交
1578 1579 1580
  void *key = NULL;
  size_t keyLen = 0;
  int32_t i = 0;
1581

D
dapan1121 已提交
1582 1583 1584 1585
  void *pIter = taosHashIterate(sch->tasksHash, NULL);
  while (pIter) {
    SQWTaskStatus *taskStatus = (SQWTaskStatus *)pIter;
    taosHashGetKey(pIter, &key, &keyLen);
D
dapan1121 已提交
1586

D
dapan1121 已提交
1587 1588 1589
    QW_GET_QTID(key, (*rsp)->status[i].queryId, (*rsp)->status[i].taskId);
    (*rsp)->status[i].status = taskStatus->status;
    
D
dapan1121 已提交
1590
    ++i;
D
dapan1121 已提交
1591
    pIter = taosHashIterate(sch->tasksHash, pIter);
D
dapan1121 已提交
1592
  }  
D
dapan1121 已提交
1593

D
dapan1121 已提交
1594 1595 1596 1597
  QW_UNLOCK(QW_READ, &sch->tasksLock);
  qwReleaseScheduler(QW_READ, mgmt);

  (*rsp)->num = taskNum;
D
dapan1121 已提交
1598
*/
D
dapan1121 已提交
1599 1600 1601 1602
  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
1603

D
dapan1121 已提交
1604 1605
int32_t qwUpdateSchLastAccess(SQWorkerMgmt *mgmt, uint64_t sId, uint64_t qId, uint64_t tId) {
  SQWSchStatus *sch = NULL;
D
dapan1121 已提交
1606

D
dapan1121 已提交
1607
/*
D
dapan1121 已提交
1608
  QW_ERR_RET(qwAcquireScheduler(QW_READ, mgmt, sId, &sch));
D
dapan1121 已提交
1609

D
dapan1121 已提交
1610
  sch->lastAccessTs = taosGetTimestampSec();
D
dapan1121 已提交
1611

D
dapan1121 已提交
1612
  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
1613
*/
D
dapan1121 已提交
1614 1615 1616 1617
  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
1618 1619 1620 1621
int32_t qwGetTaskStatus(SQWorkerMgmt *mgmt, uint64_t sId, uint64_t qId, uint64_t tId, int8_t *taskStatus) {
  SQWSchStatus *sch = NULL;
  SQWTaskStatus *task = NULL;
  int32_t code = 0;
D
dapan1121 已提交
1622 1623

/*  
D
dapan1121 已提交
1624 1625 1626 1627
  if (qwAcquireScheduler(QW_READ, mgmt, sId, &sch)) {
    *taskStatus = JOB_TASK_STATUS_NULL;
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
1628

D
dapan1121 已提交
1629 1630 1631 1632 1633 1634
  if (qwAcquireTask(mgmt, QW_READ, sch, queryId, taskId, &task)) {
    qwReleaseScheduler(QW_READ, mgmt);
    
    *taskStatus = JOB_TASK_STATUS_NULL;
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
1635

D
dapan1121 已提交
1636
  *taskStatus = task->status;
D
dapan1121 已提交
1637

D
dapan1121 已提交
1638 1639
  qwReleaseTask(QW_READ, sch);
  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
1640
*/
D
dapan1121 已提交
1641 1642 1643 1644 1645

  QW_RET(code);
}


D
dapan1121 已提交
1646 1647 1648
int32_t qwCancelTask(SQWorkerMgmt *mgmt, uint64_t sId, uint64_t qId, uint64_t tId) {
  SQWSchStatus *sch = NULL;
  SQWTaskStatus *task = NULL;
D
dapan1121 已提交
1649
  int32_t code = 0;
D
dapan1121 已提交
1650

D
dapan1121 已提交
1651
/*
D
dapan1121 已提交
1652
  QW_ERR_RET(qwAcquireAddScheduler(QW_READ, mgmt, sId, &sch));
D
dapan1121 已提交
1653

D
dapan1121 已提交
1654
  QW_ERR_JRET(qwAcquireAddTask(mgmt, QW_READ, sch, qId, tId, JOB_TASK_STATUS_NOT_START, &task));
D
dapan1121 已提交
1655

D
dapan1121 已提交
1656

D
dapan1121 已提交
1657
  QW_LOCK(QW_WRITE, &task->lock);
D
dapan1121 已提交
1658

D
dapan1121 已提交
1659 1660 1661 1662 1663 1664 1665
  task->cancel = true;
  
  int8_t oriStatus = task->status;
  int8_t newStatus = 0;
  
  if (task->status == JOB_TASK_STATUS_CANCELLED || task->status == JOB_TASK_STATUS_NOT_START || task->status == JOB_TASK_STATUS_CANCELLING || task->status == JOB_TASK_STATUS_DROPPING) {
    QW_UNLOCK(QW_WRITE, &task->lock);
D
dapan1121 已提交
1666

D
dapan1121 已提交
1667 1668 1669 1670 1671 1672 1673 1674
    qwReleaseTask(QW_READ, sch);
    qwReleaseScheduler(QW_READ, mgmt);
    
    return TSDB_CODE_SUCCESS;
  } else if (task->status == JOB_TASK_STATUS_FAILED || task->status == JOB_TASK_STATUS_SUCCEED || task->status == JOB_TASK_STATUS_PARTIAL_SUCCEED) {
    QW_ERR_JRET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_CANCELLED));
  } else {
    QW_ERR_JRET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_CANCELLING));
D
dapan1121 已提交
1675 1676
  }

D
dapan1121 已提交
1677 1678 1679 1680
  QW_UNLOCK(QW_WRITE, &task->lock);
  
  qwReleaseTask(QW_READ, sch);
  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
1681

D
dapan1121 已提交
1682 1683 1684 1685 1686
  if (oriStatus == JOB_TASK_STATUS_EXECUTING) {
    //TODO call executer to cancel subquery async
  }
  
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
1687 1688

_return:
D
dapan1121 已提交
1689

D
dapan1121 已提交
1690 1691 1692 1693 1694 1695 1696 1697 1698
  if (task) {
    QW_UNLOCK(QW_WRITE, &task->lock);
    
    qwReleaseTask(QW_READ, sch);
  }

  if (sch) {
    qwReleaseScheduler(QW_READ, mgmt);
  }
D
dapan1121 已提交
1699
*/
D
dapan1121 已提交
1700

D
dapan1121 已提交
1701
  QW_RET(code);
D
dapan1121 已提交
1702 1703
}

D
dapan1121 已提交
1704