qworker.c 41.1 KB
Newer Older
D
dapan1121 已提交
1
#include "qworker.h"
H
Haojun Liao 已提交
2
#include <common.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
int32_t qwValidateStatus(QW_FPARAMS_DEF, int8_t oriStatus, int8_t newStatus) {
D
dapan1121 已提交
15
  int32_t code = 0;
D
dapan1121 已提交
16

D
dapan1121 已提交
17 18 19 20 21 22
  if (oriStatus == newStatus) {
    QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
  }
  
  switch (oriStatus) {
    case JOB_TASK_STATUS_NULL:
D
dapan1121 已提交
23 24 25
      if (newStatus != JOB_TASK_STATUS_EXECUTING 
       && newStatus != JOB_TASK_STATUS_FAILED 
       && newStatus != JOB_TASK_STATUS_NOT_START) {
D
dapan1121 已提交
26 27 28 29 30
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_NOT_START:
D
dapan1121 已提交
31
      if (newStatus != JOB_TASK_STATUS_CANCELLED) {
D
dapan1121 已提交
32 33 34 35 36
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_EXECUTING:
D
dapan1121 已提交
37
      if (newStatus != JOB_TASK_STATUS_PARTIAL_SUCCEED 
D
dapan1121 已提交
38
       && newStatus != JOB_TASK_STATUS_SUCCEED 
D
dapan1121 已提交
39 40 41 42
       && newStatus != JOB_TASK_STATUS_FAILED 
       && newStatus != JOB_TASK_STATUS_CANCELLING 
       && newStatus != JOB_TASK_STATUS_CANCELLED 
       && newStatus != JOB_TASK_STATUS_DROPPING) {
D
dapan1121 已提交
43 44 45 46 47
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_PARTIAL_SUCCEED:
D
dapan1121 已提交
48 49 50
      if (newStatus != JOB_TASK_STATUS_EXECUTING 
       && newStatus != JOB_TASK_STATUS_SUCCEED
       && newStatus != JOB_TASK_STATUS_CANCELLED) {
D
dapan1121 已提交
51 52 53 54 55
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }
      
      break;
    case JOB_TASK_STATUS_SUCCEED:
D
dapan1121 已提交
56 57 58 59 60 61
      if (newStatus != JOB_TASK_STATUS_CANCELLED
       && newStatus != JOB_TASK_STATUS_DROPPING) {
        QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      }

      break;
D
dapan1121 已提交
62 63 64 65 66 67 68 69
    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 已提交
70 71 72 73
    case JOB_TASK_STATUS_DROPPING:
      QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
      break;
      
D
dapan1121 已提交
74
    default:
D
dapan1121 已提交
75
      QW_TASK_ELOG("invalid task status:%d", oriStatus);
D
dapan1121 已提交
76 77 78
      return TSDB_CODE_QRY_APP_ERROR;
  }

D
dapan1121 已提交
79
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
80

D
dapan1121 已提交
81
_return:
D
dapan1121 已提交
82

D
dapan1121 已提交
83 84
  QW_TASK_ELOG("invalid task status update from %d to %d", oriStatus, newStatus);
  QW_RET(code);
D
dapan1121 已提交
85 86
}

D
dapan1121 已提交
87
int32_t qwSetTaskStatus(QW_FPARAMS_DEF, SQWTaskStatus *task, int8_t status) {
D
dapan1121 已提交
88
  int32_t code = 0;
D
dapan1121 已提交
89
  int8_t origStatus = 0;
D
dapan1121 已提交
90 91 92 93 94 95 96 97

  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 已提交
98
    }
D
dapan1121 已提交
99 100 101 102
    
    QW_TASK_DLOG("task status updated from %d to %d", origStatus, status);

    break;
D
dapan1121 已提交
103
  }
D
dapan1121 已提交
104
  
D
dapan1121 已提交
105 106 107 108
  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
109
int32_t qwAddSchedulerImpl(QW_FPARAMS_DEF, int32_t rwType, SQWSchStatus **sch) {
D
dapan1121 已提交
110
  SQWSchStatus newSch = {0};
D
dapan1121 已提交
111 112
  newSch.tasksHash = taosHashInit(mgmt->cfg.maxSchTaskNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_NO_LOCK);
  if (NULL == newSch.tasksHash) {
D
dapan1121 已提交
113 114
    QW_SCH_ELOG("taosHashInit %d failed", mgmt->cfg.maxSchTaskNum);
    QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
115 116
  }

D
dapan1121 已提交
117 118 119 120 121
  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 已提交
122
      
D
dapan1121 已提交
123 124 125
      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 已提交
126
    }
D
dapan1121 已提交
127 128

    taosHashCleanup(newSch.tasksHash);
D
dapan1121 已提交
129
  }
D
dapan1121 已提交
130
  QW_UNLOCK(QW_WRITE, &mgmt->schLock);
D
dapan1121 已提交
131

D
dapan1121 已提交
132
  return TSDB_CODE_SUCCESS;  
D
dapan1121 已提交
133 134
}

D
dapan1121 已提交
135
int32_t qwAcquireSchedulerImpl(QW_FPARAMS_DEF, int32_t rwType, SQWSchStatus **sch, int32_t nOpt) {
D
dapan1121 已提交
136 137 138 139 140 141 142
  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 已提交
143
        QW_ERR_RET(qwAddSchedulerImpl(QW_FPARAMS(), rwType, sch));
D
dapan1121 已提交
144 145 146 147 148 149 150

        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 已提交
151 152
        QW_TASK_ELOG("unknown notExistOpt:%d", nOpt);
        QW_ERR_RET(TSDB_CODE_QRY_APP_ERROR);
D
dapan1121 已提交
153
      }
D
dapan1121 已提交
154
    }
D
dapan1121 已提交
155 156

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

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
162
int32_t qwAcquireAddScheduler(QW_FPARAMS_DEF, int32_t rwType, SQWSchStatus **sch) {
D
dapan1121 已提交
163
  return qwAcquireSchedulerImpl(QW_FPARAMS(), rwType, sch, QW_NOT_EXIST_ADD);
D
dapan1121 已提交
164 165
}

D
dapan1121 已提交
166
int32_t qwAcquireScheduler(QW_FPARAMS_DEF, int32_t rwType, SQWSchStatus **sch) {
D
dapan1121 已提交
167
  return qwAcquireSchedulerImpl(QW_FPARAMS(), rwType, sch, QW_NOT_EXIST_RET_ERR);
D
dapan1121 已提交
168 169
}

D
dapan1121 已提交
170
void qwReleaseScheduler(int32_t rwType, SQWorkerMgmt *mgmt) {
D
dapan1121 已提交
171 172 173
  QW_UNLOCK(rwType, &mgmt->schLock);
}

D
dapan1121 已提交
174

D
dapan1121 已提交
175
int32_t qwAcquireTaskStatus(QW_FPARAMS_DEF, int32_t rwType, SQWSchStatus *sch, SQWTaskStatus **task) {
D
dapan1121 已提交
176 177 178 179 180 181 182 183 184 185 186 187 188 189 190
  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 已提交
191
int32_t qwAddTaskStatusImpl(QW_FPARAMS_DEF, SQWSchStatus *sch, int32_t rwType, int32_t status, SQWTaskStatus **task) {
D
dapan1121 已提交
192 193
  int32_t code = 0;

D
dapan1121 已提交
194 195
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
D
dapan1121 已提交
196

D
dapan1121 已提交
197 198 199 200 201 202 203 204
  SQWTaskStatus ntask = {0};
  ntask.status = status;

  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 已提交
205 206
      if (rwType && task) {
        QW_RET(qwAcquireTaskStatus(QW_FPARAMS(), rwType, sch, task));
D
dapan1121 已提交
207
      } else {
D
dapan1121 已提交
208 209
        QW_TASK_ELOG("task status already exist, id:%s", id);
        QW_ERR_RET(TSDB_CODE_QRY_TASK_ALREADY_EXIST);
D
dapan1121 已提交
210 211
      }
    } else {
D
dapan1121 已提交
212 213
      QW_TASK_ELOG("taosHashPut to tasksHash failed, code:%x", code);
      QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
214 215 216
    }
  }
  QW_UNLOCK(QW_WRITE, &sch->tasksLock);
D
dapan1121 已提交
217

D
dapan1121 已提交
218 219
  if (rwType && task) {
    QW_ERR_RET(qwAcquireTaskStatus(QW_FPARAMS(), rwType, sch, task));
D
dapan1121 已提交
220 221
  }

D
dapan1121 已提交
222 223
  return TSDB_CODE_SUCCESS;
}
D
dapan1121 已提交
224

D
dapan1121 已提交
225
int32_t qwAddTaskStatus(QW_FPARAMS_DEF, int32_t status) {
D
dapan1121 已提交
226 227
  SQWSchStatus *tsch = NULL;
  int32_t code = 0;
D
dapan1121 已提交
228
  QW_ERR_RET(qwAcquireAddScheduler(QW_FPARAMS(), QW_READ, &tsch));
D
dapan1121 已提交
229

D
dapan1121 已提交
230
  QW_ERR_JRET(qwAddTaskStatusImpl(QW_FPARAMS(), tsch, 0, status, NULL));
D
dapan1121 已提交
231 232 233 234

_return:

  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
235 236
  
  QW_RET(code);
D
dapan1121 已提交
237 238 239
}


D
dapan1121 已提交
240
int32_t qwAddAcquireTaskStatus(QW_FPARAMS_DEF, int32_t rwType, SQWSchStatus *sch, int32_t status, SQWTaskStatus **task) {
D
dapan1121 已提交
241
  return qwAddTaskStatusImpl(QW_FPARAMS(), sch, rwType, status, task);
D
dapan1121 已提交
242
}
D
dapan1121 已提交
243 244


D
dapan1121 已提交
245
void qwReleaseTaskStatus(int32_t rwType, SQWSchStatus *sch) {
D
dapan1121 已提交
246
  QW_UNLOCK(rwType, &sch->tasksLock);
D
dapan1121 已提交
247 248
}

D
dapan1121 已提交
249

D
dapan1121 已提交
250
int32_t qwAcquireTaskCtx(QW_FPARAMS_DEF, SQWTaskCtx **ctx) {
D
dapan1121 已提交
251 252 253
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
  
D
dapan1121 已提交
254 255
  //QW_LOCK(rwType, &mgmt->ctxLock);
  *ctx = taosHashAcquire(mgmt->ctxHash, id, sizeof(id));
D
dapan1121 已提交
256
  if (NULL == (*ctx)) {
D
dapan1121 已提交
257
    //QW_UNLOCK(rwType, &mgmt->ctxLock);
D
dapan1121 已提交
258
    QW_TASK_DLOG_E("task ctx not exist, may be dropped");
D
dapan1121 已提交
259
    QW_ERR_RET(TSDB_CODE_QRY_TASK_CTX_NOT_EXIST);
D
dapan1121 已提交
260 261 262 263 264
  }

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
265
int32_t qwGetTaskCtx(QW_FPARAMS_DEF, SQWTaskCtx **ctx) {
D
dapan1121 已提交
266 267 268 269 270
  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 已提交
271
    QW_TASK_DLOG_E("task ctx not exist, may be dropped");
D
dapan1121 已提交
272
    QW_ERR_RET(TSDB_CODE_QRY_TASK_CTX_NOT_EXIST);
D
dapan1121 已提交
273 274 275 276 277
  }

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
278
int32_t qwAddTaskCtxImpl(QW_FPARAMS_DEF, bool acquire, int32_t status, SQWTaskCtx **ctx) {
D
dapan1121 已提交
279 280 281
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);

D
dapan1121 已提交
282
  SQWTaskCtx nctx = {0};
D
dapan1121 已提交
283

D
dapan1121 已提交
284
  //QW_LOCK(QW_WRITE, &mgmt->ctxLock);
D
dapan1121 已提交
285
  int32_t code = taosHashPut(mgmt->ctxHash, id, sizeof(id), &nctx, sizeof(SQWTaskCtx));
D
dapan1121 已提交
286
  if (0 != code) {
D
dapan1121 已提交
287
    //QW_UNLOCK(QW_WRITE, &mgmt->ctxLock);
D
dapan1121 已提交
288 289
    
    if (HASH_NODE_EXIST(code)) {
D
dapan1121 已提交
290
      if (acquire && ctx) {
D
dapan1121 已提交
291
        QW_RET(qwAcquireTaskCtx(QW_FPARAMS(), ctx));
D
dapan1121 已提交
292 293
      } else if (ctx) {
        QW_RET(qwGetTaskCtx(QW_FPARAMS(), ctx));
D
dapan1121 已提交
294 295 296 297 298 299 300 301 302
      } else {
        QW_TASK_ELOG("task ctx already exist, id:%s", id);
        QW_ERR_RET(TSDB_CODE_QRY_TASK_ALREADY_EXIST);
      }
    } else {
      QW_TASK_ELOG("taosHashPut to ctxHash failed, code:%x", code);
      QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
    }
  }
D
dapan1121 已提交
303
  //QW_UNLOCK(QW_WRITE, &mgmt->ctxLock);
D
dapan1121 已提交
304

D
dapan1121 已提交
305
  if (acquire && ctx) {
D
dapan1121 已提交
306
    QW_RET(qwAcquireTaskCtx(QW_FPARAMS(), ctx));
D
dapan1121 已提交
307 308
  } else if (ctx) {
    QW_RET(qwGetTaskCtx(QW_FPARAMS(), ctx));
D
dapan1121 已提交
309 310 311 312 313
  }

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
314
int32_t qwAddTaskCtx(QW_FPARAMS_DEF) {
D
dapan1121 已提交
315
  QW_RET(qwAddTaskCtxImpl(QW_FPARAMS(), false, 0, NULL));
D
dapan1121 已提交
316 317 318
}


D
dapan1121 已提交
319

D
dapan1121 已提交
320
int32_t qwAddAcquireTaskCtx(QW_FPARAMS_DEF, SQWTaskCtx **ctx) {
D
dapan1121 已提交
321
  return qwAddTaskCtxImpl(QW_FPARAMS(), true, 0, ctx);
D
dapan1121 已提交
322 323
}

D
dapan1121 已提交
324
int32_t qwAddGetTaskCtx(QW_FPARAMS_DEF, SQWTaskCtx **ctx) {
D
dapan1121 已提交
325
  return qwAddTaskCtxImpl(QW_FPARAMS(), false, 0, ctx);
D
dapan1121 已提交
326 327 328
}


D
dapan1121 已提交
329 330 331
void qwReleaseTaskCtx(SQWorkerMgmt *mgmt, void *ctx) {
  //QW_UNLOCK(rwType, &mgmt->ctxLock);
  taosHashRelease(mgmt->ctxHash, ctx);
D
dapan1121 已提交
332 333
}

D
dapan1121 已提交
334
void qwFreeTaskHandle(QW_FPARAMS_DEF, qTaskInfo_t *taskHandle) {  
D
dapan1121 已提交
335
  // RC WARNING
D
dapan1121 已提交
336 337 338
  qTaskInfo_t otaskHandle = atomic_load_ptr(taskHandle);
  if (otaskHandle && atomic_val_compare_exchange_ptr(taskHandle, otaskHandle, NULL)) {
    qDestroyTask(otaskHandle);
D
dapan1121 已提交
339 340
  }
}
D
dapan1121 已提交
341

D
dapan1121 已提交
342 343 344 345 346
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 已提交
347
    code = qAsyncKillTask(taskHandle);
D
dapan1121 已提交
348
    atomic_store_ptr(&ctx->taskHandle, taskHandle);
D
dapan1121 已提交
349
  }
D
dapan1121 已提交
350

D
dapan1121 已提交
351 352 353
  QW_RET(code);
}

D
dapan1121 已提交
354

D
dapan1121 已提交
355
void qwFreeTask(QW_FPARAMS_DEF, SQWTaskCtx *ctx) {
D
dapan1121 已提交
356
  qwFreeTaskHandle(QW_FPARAMS(), &ctx->taskHandle);
D
dapan1121 已提交
357 358 359 360
  
  if (ctx->sinkHandle) {
    dsDestroyDataSinker(ctx->sinkHandle);
    ctx->sinkHandle = NULL;
D
dapan1121 已提交
361
  }
D
dapan1121 已提交
362
}
D
dapan1121 已提交
363 364


D
dapan1121 已提交
365
// Note: NEED CTX HASH LOCKED BEFORE ENTRANCE
D
dapan1121 已提交
366
int32_t qwDropTaskCtx(QW_FPARAMS_DEF, int32_t rwType) {
D
dapan1121 已提交
367 368 369
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
  SQWTaskCtx octx;
D
dapan1121 已提交
370

D
dapan1121 已提交
371 372
  SQWTaskCtx *ctx = taosHashGet(mgmt->ctxHash, id, sizeof(id));
  if (NULL == ctx) {
D
dapan1121 已提交
373
    QW_ERR_RET(TSDB_CODE_QRY_TASK_CTX_NOT_EXIST);
D
dapan1121 已提交
374
  }
D
dapan1121 已提交
375

D
dapan1121 已提交
376
  octx = *ctx;
D
dapan1121 已提交
377

D
dapan1121 已提交
378 379 380
  atomic_store_ptr(&ctx->taskHandle, NULL);
  atomic_store_ptr(&ctx->sinkHandle, NULL);

D
dapan1121 已提交
381 382
  QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_DROP);

D
dapan1121 已提交
383 384 385 386
  if (rwType) {
    QW_UNLOCK(rwType, &ctx->lock);
  }

D
dapan1121 已提交
387
  if (taosHashRemove(mgmt->ctxHash, id, sizeof(id))) {
D
dapan1121 已提交
388
    QW_TASK_ELOG_E("taosHashRemove from ctx hash failed");    
D
dapan1121 已提交
389
    QW_ERR_RET(TSDB_CODE_QRY_TASK_CTX_NOT_EXIST);
D
dapan1121 已提交
390 391 392 393 394
  }

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

D
dapan1121 已提交
396 397 398
  if (octx.sinkHandle) {
    dsDestroyDataSinker(octx.sinkHandle);
  }
D
dapan1121 已提交
399 400

  QW_TASK_DLOG_E("task ctx dropped");
D
dapan1121 已提交
401
  
D
dapan1121 已提交
402 403 404
  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
405

D
dapan1121 已提交
406
int32_t qwDropTaskStatus(QW_FPARAMS_DEF) {
D
dapan1121 已提交
407
  SQWSchStatus *sch = NULL;
D
dapan1121 已提交
408 409 410 411 412
  SQWTaskStatus *task = NULL;
  int32_t code = 0;
  
  char id[sizeof(qId) + sizeof(tId)] = {0};
  QW_SET_QTID(id, qId, tId);
D
dapan1121 已提交
413

D
dapan1121 已提交
414
  if (qwAcquireScheduler(QW_FPARAMS(), QW_WRITE, &sch)) {
D
dapan1121 已提交
415
    QW_TASK_WLOG_E("scheduler does not exist");
D
dapan1121 已提交
416 417
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
418

D
dapan1121 已提交
419 420 421
  if (qwAcquireTaskStatus(QW_FPARAMS(), QW_WRITE, sch, &task)) {
    qwReleaseScheduler(QW_WRITE, mgmt);
    
D
dapan1121 已提交
422
    QW_TASK_WLOG_E("task does not exist");
D
dapan1121 已提交
423 424
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
425

D
dapan1121 已提交
426
  if (taosHashRemove(sch->tasksHash, id, sizeof(id))) {
D
dapan1121 已提交
427
    QW_TASK_ELOG_E("taosHashRemove task from hash failed");
D
dapan1121 已提交
428 429 430
    QW_ERR_JRET(TSDB_CODE_QRY_APP_ERROR);
  }

D
dapan1121 已提交
431
  QW_TASK_DLOG_E("task status dropped");
D
dapan1121 已提交
432 433 434 435 436 437 438 439 440

_return:

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

D
dapan1121 已提交
441
int32_t qwUpdateTaskStatus(QW_FPARAMS_DEF, int8_t status) {
D
dapan1121 已提交
442 443
  SQWSchStatus *sch = NULL;
  SQWTaskStatus *task = NULL;
D
dapan1121 已提交
444 445
  int32_t code = 0;

D
dapan1121 已提交
446 447
  QW_ERR_RET(qwAcquireScheduler(QW_FPARAMS(), QW_READ, &sch));
  QW_ERR_JRET(qwAcquireTaskStatus(QW_FPARAMS(), QW_READ, sch, &task));
D
dapan1121 已提交
448

D
dapan1121 已提交
449
  QW_ERR_JRET(qwSetTaskStatus(QW_FPARAMS(), task, status));
D
dapan1121 已提交
450
  
D
dapan1121 已提交
451 452
_return:

D
dapan1121 已提交
453
  qwReleaseTaskStatus(QW_READ, sch);
D
dapan1121 已提交
454
  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
455 456 457 458

  QW_RET(code);
}

D
dapan1121 已提交
459
int32_t qwExecTask(QW_FPARAMS_DEF, SQWTaskCtx *ctx) {
D
dapan1121 已提交
460 461 462 463
  int32_t code = 0;
  bool  qcontinue = true;
  SSDataBlock* pRes = NULL;
  uint64_t useconds = 0;
D
dapan1121 已提交
464
  int32_t i = 0;
D
dapan1121 已提交
465 466 467
  int32_t execNum = 0;
  qTaskInfo_t *taskHandle = &ctx->taskHandle; 
  DataSinkHandle sinkHandle = ctx->sinkHandle;
D
dapan1121 已提交
468 469 470 471
 
  while (true) {
    QW_TASK_DLOG("start to execTask in executor, loopIdx:%d", i++);
    
D
dapan1121 已提交
472
    code = qExecTask(*taskHandle, &pRes, &useconds);
D
dapan1121 已提交
473 474 475 476 477
    if (code) {
      QW_TASK_ELOG("qExecTask failed, code:%x", code);
      QW_ERR_JRET(code);
    }

D
dapan1121 已提交
478 479
    ++execNum;

D
dapan1121 已提交
480
    if (NULL == pRes) {
D
dapan1121 已提交
481
      QW_TASK_DLOG("task query done, useconds:%"PRIu64, useconds);
D
dapan1121 已提交
482
      
D
dapan1121 已提交
483
      dsEndPut(sinkHandle, useconds);
D
dapan1121 已提交
484
      
D
dapan1121 已提交
485
      if (TASK_TYPE_TEMP == ctx->taskType) {
D
dapan1121 已提交
486 487
        qwFreeTaskHandle(QW_FPARAMS(), taskHandle);
      }
D
dapan1121 已提交
488
      
D
dapan1121 已提交
489 490 491 492 493 494
      break;
    }

    SInputData inputData = {.pData = pRes, .pTableRetrieveTsMap = NULL};
    code = dsPutDataBlock(sinkHandle, &inputData, &qcontinue);
    if (code) {
H
Haojun Liao 已提交
495
      QW_TASK_ELOG("dsPutDataBlock failed, code:%s", tstrerror(code));
D
dapan1121 已提交
496 497
      QW_ERR_JRET(code);
    }
D
dapan1121 已提交
498 499 500

    QW_TASK_DLOG("data put into sink, rows:%d, continueExecTask:%d", pRes->info.rows, qcontinue);
    
D
dapan1121 已提交
501 502 503 504
    if (!qcontinue) {
      break;
    }

D
dapan1121 已提交
505 506 507 508 509
    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 已提交
510 511
      break;
    }
D
dapan1121 已提交
512 513 514 515 516 517
  }

_return:

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

D
dapan 已提交
519 520

int32_t qwGetResFromSink(QW_FPARAMS_DEF, SQWTaskCtx *ctx, int32_t *dataLen, void **rspMsg, SOutputData *pOutput) {
D
dapan1121 已提交
521 522 523
  int32_t len = 0;
  SRetrieveTableRsp *rsp = NULL;
  bool queryEnd = false;
D
dapan1121 已提交
524 525
  int32_t code = 0;

D
dapan1121 已提交
526 527 528 529 530 531 532 533 534 535 536 537 538 539
  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 已提交
540
  dsGetDataLength(ctx->sinkHandle, &len, &queryEnd);
D
dapan1121 已提交
541

D
dapan1121 已提交
542 543 544 545
  if (len < 0) {
    QW_TASK_ELOG("invalid length from dsGetDataLength, length:%d", len);
    QW_ERR_RET(TSDB_CODE_QRY_INVALID_INPUT);
  }
D
dapan1121 已提交
546

D
dapan1121 已提交
547 548
  if (len == 0) {
    if (queryEnd) {
D
dapan 已提交
549
      code = dsGetDataBlock(ctx->sinkHandle, pOutput);
D
dapan1121 已提交
550 551 552 553 554
      if (code) {
        QW_TASK_ELOG("dsGetDataBlock failed, code:%x", code);
        QW_ERR_RET(code);
      }
    
D
dapan1121 已提交
555
      QW_TASK_DLOG("no data in sink and query end, phase:%d", ctx->phase);
D
dapan1121 已提交
556 557
      
      QW_ERR_RET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_SUCCEED));
D
dapan1121 已提交
558

D
dapan1121 已提交
559 560
      QW_ERR_RET(qwMallocFetchRsp(len, &rsp));      
      *rspMsg = rsp;
D
dapan 已提交
561 562
      *dataLen = 0;
      
D
dapan1121 已提交
563 564
      return TSDB_CODE_SUCCESS;
    }
D
dapan1121 已提交
565 566

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

D
dapan1121 已提交
570
    return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
571
  }  
D
dapan1121 已提交
572 573


D
dapan1121 已提交
574 575
  // Got data from sink

D
dapan 已提交
576
  *dataLen = len;
D
dapan1121 已提交
577 578

  QW_TASK_DLOG("task got data in sink, dataLength:%d", len);
D
dapan1121 已提交
579
  
D
dapan1121 已提交
580
  QW_ERR_RET(qwMallocFetchRsp(len, &rsp));
D
dapan 已提交
581
  *rspMsg = rsp;
D
dapan1121 已提交
582
  
D
dapan 已提交
583 584
  pOutput->pData = rsp->data;
  code = dsGetDataBlock(ctx->sinkHandle, pOutput);
D
dapan1121 已提交
585 586 587 588
  if (code) {
    QW_TASK_ELOG("dsGetDataBlock failed, code:%x", code);
    QW_ERR_RET(code);
  }
D
dapan1121 已提交
589

D
dapan 已提交
590 591
  queryEnd = pOutput->queryEnd;
  pOutput->queryEnd = false;
D
dapan1121 已提交
592

D
dapan 已提交
593 594
  if (DS_BUF_EMPTY == pOutput->bufStatus && queryEnd) {
    pOutput->queryEnd = true;
D
dapan1121 已提交
595
    
D
dapan1121 已提交
596 597
    QW_SCH_TASK_DLOG("task all fetched, status:%d", JOB_TASK_STATUS_SUCCEED);
    QW_ERR_RET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_SUCCEED));
D
dapan1121 已提交
598 599
  }

D
dapan1121 已提交
600
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
601 602
}

D
dapan1121 已提交
603
int32_t qwHandlePrePhaseEvents(QW_FPARAMS_DEF, int8_t phase, SQWPhaseInput *input, SQWPhaseOutput *output) {
D
dapan1121 已提交
604 605 606 607
  int32_t code = 0;
  int8_t status = 0;
  SQWTaskCtx *ctx = NULL;
  bool locked = false;
D
dapan1121 已提交
608 609
  void *dropConnection = NULL;
  void *cancelConnection = NULL;
D
dapan1121 已提交
610

D
dapan1121 已提交
611
  QW_SCH_TASK_DLOG("start to handle event at phase %d", phase);
D
dapan1121 已提交
612

D
dapan1121 已提交
613 614 615 616 617 618 619 620 621 622 623
  output->needStop = false;

  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 已提交
624 625
  switch (phase) {
    case QW_PHASE_PRE_QUERY: {
D
dapan1121 已提交
626 627
      atomic_store_8(&ctx->phase, phase);

D
dapan1121 已提交
628 629
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL) || QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_DROP)) {
        QW_TASK_ELOG("task already cancelled/dropped at wrong phase, phase:%d", phase);
D
dapan1121 已提交
630
        
D
dapan1121 已提交
631
        output->needStop = true;
D
dapan1121 已提交
632 633 634
        output->rspCode = TSDB_CODE_QRY_TASK_STATUS_ERROR;
        break;
      }
D
dapan1121 已提交
635

D
dapan1121 已提交
636
      if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
637
        QW_ERR_JRET(qwDropTaskStatus(QW_FPARAMS()));              
D
dapan1121 已提交
638
        QW_ERR_JRET(qwDropTaskCtx(QW_FPARAMS(), QW_WRITE));
D
dapan1121 已提交
639

D
dapan1121 已提交
640
        output->needStop = true;
D
dapan1121 已提交
641
        output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;
D
dapan1121 已提交
642
        QW_SET_RSP_CODE(ctx, output->rspCode);
D
dapan1121 已提交
643
        dropConnection = ctx->dropConnection;
D
dapan1121 已提交
644 645
        
        // Note: ctx freed, no need to unlock it
D
dapan 已提交
646 647 648
        locked = false;      

        break;
D
dapan1121 已提交
649 650 651 652
      } 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 已提交
653
        output->needStop = true;                
D
dapan1121 已提交
654
        output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;
D
dapan1121 已提交
655 656 657
        QW_SET_RSP_CODE(ctx, output->rspCode);
        
        cancelConnection = ctx->cancelConnection;
D
dapan 已提交
658 659

        break;
D
dapan1121 已提交
660
      }
D
dapan1121 已提交
661

D
dapan1121 已提交
662 663
      if (ctx->rspCode) {
        QW_TASK_ELOG("task already failed at wrong phase, code:%x, phase:%d", ctx->rspCode, phase);
D
dapan1121 已提交
664
        output->needStop = true;
D
dapan1121 已提交
665 666
        output->rspCode = ctx->rspCode;        
        QW_ERR_JRET(output->rspCode);
D
dapan1121 已提交
667
      }
D
dapan1121 已提交
668

D
dapan1121 已提交
669 670
      if (!output->needStop) {
        QW_ERR_JRET(qwAddTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_EXECUTING));
D
dapan1121 已提交
671 672 673
      }
      break;
    }
D
dapan1121 已提交
674
    case QW_PHASE_PRE_FETCH: {
D
dapan1121 已提交
675 676 677 678 679 680
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_DROP)) {
        QW_TASK_WLOG("task already dropped, phase:%d", phase);
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_DROPPED);
      }
D
dapan1121 已提交
681 682 683 684 685 686
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL)) {
        QW_TASK_WLOG("task already cancelled, phase:%d", phase);
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
      }
D
dapan1121 已提交
687

D
dapan1121 已提交
688
      if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
689
        QW_TASK_ELOG("drop event at wrong phase, phase:%d", phase);
D
dapan1121 已提交
690
        output->needStop = true;
D
dapan1121 已提交
691 692
        output->rspCode = TSDB_CODE_QRY_TASK_STATUS_ERROR;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
D
dapan1121 已提交
693
      } else if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_CANCEL)) {
D
dapan1121 已提交
694
        QW_TASK_ELOG("cancel event at wrong phase, phase:%d", phase);
D
dapan1121 已提交
695
        output->needStop = true;
D
dapan1121 已提交
696 697 698 699
        output->rspCode = TSDB_CODE_QRY_TASK_STATUS_ERROR;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
      }

D
dapan1121 已提交
700 701 702 703 704 705 706
      if (ctx->rspCode) {
        QW_TASK_ELOG("task already failed, code:%x, phase:%d", ctx->rspCode, phase);
        output->needStop = true;
        output->rspCode = ctx->rspCode;        
        QW_ERR_JRET(output->rspCode);
      }

D
dapan1121 已提交
707 708 709 710 711 712 713 714 715 716 717 718
      if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_FETCH)) {
        QW_TASK_WLOG("last fetch not finished, phase:%d", phase);
        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)) {
        QW_TASK_ELOG("query rsp are not ready, phase:%d", phase);
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_TASK_MSG_ERROR;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_MSG_ERROR);
D
dapan1121 已提交
719 720 721
      }
      break;
    }    
D
dapan1121 已提交
722
    case QW_PHASE_PRE_CQUERY: {
D
dapan1121 已提交
723 724
      atomic_store_8(&ctx->phase, phase);

D
dapan1121 已提交
725 726 727 728 729 730 731
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_CANCEL)) {
        QW_TASK_WLOG("task already cancelled, phase:%d", phase);
        output->needStop = true;
        output->rspCode = TSDB_CODE_QRY_TASK_CANCELLED;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_CANCELLED);
      }

D
dapan1121 已提交
732 733
      if (QW_IS_EVENT_PROCESSED(ctx, QW_EVENT_DROP)) {
        QW_TASK_WLOG("task already dropped, phase:%d", phase);
D
dapan1121 已提交
734
        output->needStop = true;
D
dapan1121 已提交
735 736
        output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;        
        QW_ERR_JRET(TSDB_CODE_QRY_TASK_DROPPED);
D
dapan1121 已提交
737
      }
D
dapan1121 已提交
738

D
dapan1121 已提交
739
      if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_DROP)) {
D
dapan1121 已提交
740 741 742 743
        QW_ERR_JRET(qwDropTaskStatus(QW_FPARAMS()));  
        QW_ERR_JRET(qwDropTaskCtx(QW_FPARAMS(), QW_WRITE));
        
        output->rspCode = TSDB_CODE_QRY_TASK_DROPPED;
D
dapan1121 已提交
744
        output->needStop = true;
D
dapan1121 已提交
745 746 747 748 749 750 751
        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 已提交
752
      } else if (QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_CANCEL)) {
D
dapan1121 已提交
753 754 755 756 757 758 759 760 761 762 763
        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 已提交
764
      }
D
dapan1121 已提交
765

D
dapan1121 已提交
766

D
dapan1121 已提交
767 768 769 770 771 772
      if (ctx->rspCode) {
        QW_TASK_ELOG("task already failed, code:%x, phase:%d", ctx->rspCode, phase);
        output->needStop = true;
        output->rspCode = ctx->rspCode;        
        QW_ERR_JRET(output->rspCode);
      }
D
dapan1121 已提交
773 774
      break;
    }    
D
dapan1121 已提交
775
  }
D
dapan1121 已提交
776

D
dapan1121 已提交
777
_return:
D
dapan1121 已提交
778

D
dapan1121 已提交
779 780 781 782 783 784 785 786
  if (ctx) {
    if (output->rspCode) {
      QW_UPDATE_RSP_CODE(ctx, output->rspCode);
    }
    
    if (locked) {
      QW_UNLOCK(QW_WRITE, &ctx->lock);
    }
D
dapan1121 已提交
787

D
dapan1121 已提交
788 789
    qwReleaseTaskCtx(mgmt, ctx);
  }
D
dapan1121 已提交
790

D
dapan1121 已提交
791 792 793 794 795
  if (code) {
    output->needStop = true;
    if (TSDB_CODE_SUCCESS == output->rspCode) {
      output->rspCode = code;
    }
D
dapan1121 已提交
796 797
  }

D
dapan1121 已提交
798 799 800 801
  if (dropConnection) {
    qwBuildAndSendDropRsp(dropConnection, output->rspCode);    
    QW_TASK_DLOG("drop msg rsped, code:%x", output->rspCode);
  }
D
dapan1121 已提交
802

D
dapan1121 已提交
803 804 805
  if (cancelConnection) {
    qwBuildAndSendCancelRsp(cancelConnection, output->rspCode);    
    QW_TASK_DLOG("cancel msg rsped, code:%x", output->rspCode);
D
dapan1121 已提交
806 807
  }

D
dapan1121 已提交
808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858
  QW_SCH_TASK_DLOG("end to handle event at phase %d", phase);

  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;

  QW_SCH_TASK_DLOG("start to handle event at phase %d", phase);

  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)) {
    QW_TASK_WLOG("task already dropped, phase:%d", phase);
    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)) {
    QW_TASK_WLOG("task already cancelled, phase:%d", phase);
    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 已提交
859 860
  }

D
dapan1121 已提交
861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900
  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) {
    QW_TASK_ELOG("task failed, code:%x, phase:%d", ctx->rspCode, phase);
    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 已提交
901
  if (ctx) {
D
dapan1121 已提交
902 903 904 905
    if (output->rspCode) {
      QW_UPDATE_RSP_CODE(ctx, output->rspCode);
    }

D
dapan1121 已提交
906 907 908
    if (QW_PHASE_POST_FETCH != phase) {
      atomic_store_8(&ctx->phase, phase);
    }
D
dapan1121 已提交
909 910 911 912 913
    
    if (locked) {
      QW_UNLOCK(QW_WRITE, &ctx->lock);
    }
    
D
dapan1121 已提交
914
    qwReleaseTaskCtx(mgmt, ctx);
D
dapan1121 已提交
915 916
  }

D
dapan1121 已提交
917 918 919 920 921 922 923
  if (code) {
    output->needStop = true;
    if (TSDB_CODE_SUCCESS == output->rspCode) {
      output->rspCode = code;
    }
  }

D
dapan1121 已提交
924 925 926 927 928 929 930 931 932 933 934 935 936 937 938
  if (readyConnection) {
    qwBuildAndSendReadyRsp(readyConnection, output->rspCode);    
    QW_TASK_DLOG("ready msg rsped, code:%x", output->rspCode);
  }

  if (dropConnection) {
    qwBuildAndSendDropRsp(dropConnection, output->rspCode);    
    QW_TASK_DLOG("drop msg rsped, code:%x", output->rspCode);
  }

  if (cancelConnection) {
    qwBuildAndSendCancelRsp(cancelConnection, output->rspCode);    
    QW_TASK_DLOG("cancel msg rsped, code:%x", output->rspCode);
  }

D
dapan1121 已提交
939
  QW_SCH_TASK_DLOG("end to handle event at phase %d", phase);
D
dapan1121 已提交
940

D
dapan1121 已提交
941 942 943 944
  QW_RET(code);
}


D
dapan1121 已提交
945
int32_t qwProcessQuery(QW_FPARAMS_DEF, SQWMsg *qwMsg, int8_t taskType) {
D
dapan1121 已提交
946 947 948 949 950 951 952
  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 已提交
953 954
  qTaskInfo_t pTaskInfo = NULL;
  DataSinkHandle sinkHandle = NULL;
D
dapan1121 已提交
955
  SQWTaskCtx *ctx = NULL;
D
dapan1121 已提交
956

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

D
dapan1121 已提交
959 960
  needStop = output.needStop;
  code = output.rspCode;
D
dapan1121 已提交
961
  
D
dapan1121 已提交
962
  if (needStop) {
D
dapan1121 已提交
963
    QW_TASK_DLOG("task need stop, phase:%d", QW_PHASE_PRE_QUERY);
D
dapan1121 已提交
964
    QW_ERR_JRET(code);
D
dapan1121 已提交
965
  }
D
dapan1121 已提交
966 967 968 969

  QW_ERR_JRET(qwGetTaskCtx(QW_FPARAMS(), &ctx));
  
  atomic_store_8(&ctx->taskType, taskType);
D
dapan1121 已提交
970
  
D
dapan1121 已提交
971 972
  code = qStringToSubplan(qwMsg->msg, &plan);
  if (TSDB_CODE_SUCCESS != code) {
H
Haojun Liao 已提交
973
    QW_TASK_ELOG("task string to subplan failed, code:%s", tstrerror(code));
D
dapan1121 已提交
974
    QW_ERR_JRET(code);
D
dapan1121 已提交
975
  }
D
dapan1121 已提交
976 977
  
  code = qCreateExecTask(qwMsg->node, 0, (struct SSubplan *)plan, &pTaskInfo, &sinkHandle);
D
dapan1121 已提交
978
  if (code) {
H
Haojun Liao 已提交
979
    QW_TASK_ELOG("qCreateExecTask failed, code:%s", tstrerror(code));
D
dapan1121 已提交
980
    QW_ERR_JRET(code);
D
dapan1121 已提交
981
  }
D
dapan1121 已提交
982

H
Haojun Liao 已提交
983
  if (NULL == sinkHandle || NULL == pTaskInfo) {
D
dapan1121 已提交
984 985 986 987 988
    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 已提交
989
  
D
dapan1121 已提交
990 991
  QW_ERR_JRET(qwBuildAndSendQueryRsp(qwMsg->connection, code));
  QW_TASK_DLOG("query msg rsped, code:%d", code);
D
dapan1121 已提交
992

D
dapan1121 已提交
993
  queryRsped = true;
D
dapan1121 已提交
994

D
dapan1121 已提交
995 996 997
  atomic_store_ptr(&ctx->taskHandle, pTaskInfo);
  atomic_store_ptr(&ctx->sinkHandle, sinkHandle);

D
dapan1121 已提交
998
  if (pTaskInfo && sinkHandle) {
D
dapan1121 已提交
999
    QW_ERR_JRET(qwExecTask(QW_FPARAMS(), ctx));
D
dapan1121 已提交
1000
  }
D
dapan1121 已提交
1001
  
D
dapan1121 已提交
1002 1003
_return:

D
dapan1121 已提交
1004
  if (code) {
D
dapan1121 已提交
1005
    rspCode = code;
D
dapan1121 已提交
1006
  }
D
dapan1121 已提交
1007 1008
  
  if (!queryRsped) {
D
dapan1121 已提交
1009 1010
    qwBuildAndSendQueryRsp(qwMsg->connection, rspCode);
    QW_TASK_DLOG("query msg rsped, code:%x", rspCode);
D
dapan1121 已提交
1011 1012
  }

D
dapan1121 已提交
1013
  input.code = rspCode;
D
dapan1121 已提交
1014
  input.taskStatus = rspCode ? JOB_TASK_STATUS_FAILED : JOB_TASK_STATUS_PARTIAL_SUCCEED;
D
dapan1121 已提交
1015
  
D
dapan1121 已提交
1016
  QW_ERR_RET(qwHandlePostPhaseEvents(QW_FPARAMS(), QW_PHASE_POST_QUERY, &input, &output));
D
dapan1121 已提交
1017 1018
  
  QW_RET(rspCode);
D
dapan1121 已提交
1019 1020
}

D
dapan1121 已提交
1021
int32_t qwProcessReady(QW_FPARAMS_DEF, SQWMsg *qwMsg) {
D
dapan1121 已提交
1022 1023
  int32_t code = 0;
  SQWTaskCtx *ctx = NULL;
D
dapan1121 已提交
1024
  int8_t phase = 0;
D
dapan1121 已提交
1025 1026
  bool needRsp = false;
  int32_t rspCode = 0;
D
dapan1121 已提交
1027 1028

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

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

D
dapan1121 已提交
1032 1033 1034 1035 1036 1037
  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 已提交
1038 1039 1040
  phase = QW_GET_PHASE(ctx);
  
  if (phase == QW_PHASE_PRE_QUERY) {
D
dapan1121 已提交
1041
    QW_SET_EVENT_RECEIVED(ctx, QW_EVENT_READY);
D
dapan1121 已提交
1042
    ctx->readyConnection = qwMsg->connection;
D
dapan1121 已提交
1043
    QW_TASK_DLOG("ready msg not rsped, phase:%d", phase);
D
dapan1121 已提交
1044
  } else if (phase == QW_PHASE_POST_QUERY) {
D
dapan1121 已提交
1045
    QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_READY);
D
dapan1121 已提交
1046 1047
    needRsp = true;
    rspCode = ctx->rspCode;
D
dapan1121 已提交
1048 1049
  } else {
    QW_TASK_ELOG("invalid phase when got ready msg, phase:%d", phase);
D
dapan1121 已提交
1050 1051 1052 1053
    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 已提交
1054 1055 1056 1057
  }

_return:

D
dapan1121 已提交
1058
  if (code && ctx) {
D
dapan1121 已提交
1059
    QW_UPDATE_RSP_CODE(ctx, code);
D
dapan1121 已提交
1060 1061
  }

D
dapan1121 已提交
1062 1063
  if (ctx) {
    QW_UNLOCK(QW_WRITE, &ctx->lock);
D
dapan1121 已提交
1064
    qwReleaseTaskCtx(mgmt, ctx);
D
dapan1121 已提交
1065 1066
  }

D
dapan1121 已提交
1067 1068 1069 1070 1071
  if (needRsp) {
    qwBuildAndSendReadyRsp(qwMsg->connection, rspCode);
    QW_TASK_DLOG("ready msg rsped, code:%x", rspCode);
  }

D
dapan1121 已提交
1072 1073 1074 1075
  QW_RET(code);
}


D
dapan1121 已提交
1076
int32_t qwProcessCQuery(QW_FPARAMS_DEF, SQWMsg *qwMsg) {
D
dapan1121 已提交
1077
  SQWTaskCtx *ctx = NULL;
D
dapan1121 已提交
1078
  int32_t code = 0;
1079 1080 1081 1082 1083 1084 1085
  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 已提交
1086 1087
  
  do {
D
dapan1121 已提交
1088
    QW_ERR_JRET(qwHandlePrePhaseEvents(QW_FPARAMS(), QW_PHASE_PRE_CQUERY, &input, &output));
D
dapan1121 已提交
1089

D
dapan1121 已提交
1090 1091 1092 1093 1094 1095 1096
    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 已提交
1097

D
dapan1121 已提交
1098
    QW_ERR_JRET(qwGetTaskCtx(QW_FPARAMS(), &ctx));
D
dapan1121 已提交
1099

D
dapan1121 已提交
1100
    atomic_store_8(&ctx->queryInQueue, 0);
D
dapan1121 已提交
1101
    atomic_store_8(&ctx->queryContinue, 0);
D
dapan1121 已提交
1102

D
dapan1121 已提交
1103
    DataSinkHandle  sinkHandle = ctx->sinkHandle;
D
dapan1121 已提交
1104

D
dapan1121 已提交
1105
    QW_ERR_JRET(qwExecTask(QW_FPARAMS(), ctx));
D
dapan1121 已提交
1106

D
dapan1121 已提交
1107 1108 1109 1110 1111 1112
    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 已提交
1113 1114
        
        // RC WARNING
D
dapan1121 已提交
1115
        atomic_store_8(&ctx->queryContinue, 1);
1116
      }
D
dapan1121 已提交
1117 1118 1119 1120

      if (sOutput.queryEnd) {
        needStop = true;
      }
D
dapan1121 已提交
1121 1122 1123 1124 1125
      
      if (rsp) {
        qwBuildFetchRsp(rsp, &sOutput, dataLen); 
        
        QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_FETCH);            
1126
        
D
dapan1121 已提交
1127 1128
        qwBuildAndSendFetchRsp(qwMsg->connection, rsp, dataLen, code);                
        QW_TASK_DLOG("fetch msg rsped, code:%x, dataLen:%d", code, dataLen);
D
dapan1121 已提交
1129 1130
      } else {
        atomic_store_8(&ctx->queryContinue, 1);
1131 1132 1133
      }
    }

D
dapan1121 已提交
1134
  _return:
1135

D
dapan1121 已提交
1136 1137 1138 1139
    if (NULL == ctx) {
      break;
    }

D
dapan1121 已提交
1140
    if (code && QW_IS_EVENT_RECEIVED(ctx, QW_EVENT_FETCH)) {
D
dapan1121 已提交
1141
      QW_SET_EVENT_PROCESSED(ctx, QW_EVENT_FETCH);    
1142 1143 1144
      qwFreeFetchRsp(rsp);
      rsp = NULL;
      qwBuildAndSendFetchRsp(qwMsg->connection, rsp, 0, code);
D
dapan1121 已提交
1145
      QW_TASK_DLOG("fetch msg rsped, code:%x, dataLen:%d", code, 0);      
1146
    }
D
dapan1121 已提交
1147

D
dapan1121 已提交
1148 1149 1150 1151 1152 1153 1154 1155 1156
    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 已提交
1157

D
dapan1121 已提交
1158 1159
  input.code = code;
  qwHandlePostPhaseEvents(QW_FPARAMS(), QW_PHASE_POST_CQUERY, &input, &output);    
D
dapan1121 已提交
1160 1161

  QW_RET(code);
D
dapan1121 已提交
1162
}
D
dapan1121 已提交
1163

D
dapan1121 已提交
1164

D
dapan1121 已提交
1165
int32_t qwProcessFetch(QW_FPARAMS_DEF, SQWMsg *qwMsg) {
D
dapan1121 已提交
1166
  int32_t code = 0;
D
dapan1121 已提交
1167 1168
  int32_t needRsp = true;
  void *data = NULL;
1169
  int32_t sinkStatus = 0;
D
dapan1121 已提交
1170
  int32_t dataLen = 0;
1171
  bool queryEnd = false;
D
dapan1121 已提交
1172 1173 1174
  bool needStop = false;
  bool locked = false;
  SQWTaskCtx *ctx = NULL;
D
dapan1121 已提交
1175
  int8_t status = 0;
D
dapan1121 已提交
1176
  void *rsp = NULL;
D
dapan1121 已提交
1177

D
dapan1121 已提交
1178 1179
  SQWPhaseInput input = {0};
  SQWPhaseOutput output = {0};
D
dapan1121 已提交
1180

D
dapan1121 已提交
1181
  QW_ERR_JRET(qwHandlePrePhaseEvents(QW_FPARAMS(), QW_PHASE_PRE_FETCH, &input, &output));
D
dapan1121 已提交
1182 1183 1184 1185 1186
  
  needStop = output.needStop;
  code = output.rspCode;
  
  if (needStop) {
D
dapan1121 已提交
1187
    QW_TASK_DLOG("task need stop, phase:%d", QW_PHASE_PRE_FETCH);
D
dapan1121 已提交
1188
    QW_ERR_JRET(code);
D
dapan1121 已提交
1189
  }
1190

1191 1192
  QW_ERR_JRET(qwGetTaskCtx(QW_FPARAMS(), &ctx));
 
D
dapan 已提交
1193 1194
  SOutputData sOutput = {0};
  QW_ERR_JRET(qwGetResFromSink(QW_FPARAMS(), ctx, &dataLen, &rsp, &sOutput));
D
dapan1121 已提交
1195

1196 1197
  if (NULL == rsp) {
    QW_SET_EVENT_RECEIVED(ctx, QW_EVENT_FETCH);
D
dapan1121 已提交
1198 1199
  } else {
    qwBuildFetchRsp(rsp, &sOutput, dataLen);
D
dapan1121 已提交
1200 1201
  }

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

D
dapan1121 已提交
1205 1206
    QW_LOCK(QW_WRITE, &ctx->lock);
    locked = true;
1207

D
dapan1121 已提交
1208
    // RC WARNING
D
dapan1121 已提交
1209
    if (QW_IS_QUERY_RUNNING(ctx)) {
D
dapan1121 已提交
1210
      atomic_store_8(&ctx->queryContinue, 1);
D
dapan1121 已提交
1211
    } else if (0 == atomic_load_8(&ctx->queryInQueue)) {
D
dapan1121 已提交
1212 1213 1214 1215
      if (!ctx->multiExec) {
        QW_ERR_JRET(qwUpdateTaskStatus(QW_FPARAMS(), JOB_TASK_STATUS_EXECUTING));      
        ctx->multiExec = true;
      }
D
dapan1121 已提交
1216

D
dapan1121 已提交
1217
      atomic_store_8(&ctx->queryInQueue, 1);
1218
      
D
dapan1121 已提交
1219
      QW_ERR_JRET(qwBuildAndSendCQueryMsg(QW_FPARAMS(), qwMsg->connection));
D
dapan1121 已提交
1220 1221

      QW_TASK_DLOG("schedule query in queue, phase:%d", ctx->phase);
1222
    }
D
dapan 已提交
1223 1224
  }
  
D
dapan1121 已提交
1225
_return:
D
dapan1121 已提交
1226

D
dapan1121 已提交
1227 1228 1229 1230 1231 1232
  if (locked) {
    QW_UNLOCK(QW_WRITE, &ctx->lock);
  }

  input.code = code;

D
dapan1121 已提交
1233 1234 1235 1236 1237
  qwHandlePostPhaseEvents(QW_FPARAMS(), QW_PHASE_POST_FETCH, &input, &output);

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

D
dapan 已提交
1239 1240 1241
  if (code) {
    qwFreeFetchRsp(rsp);
    rsp = NULL;
D
dapan1121 已提交
1242 1243 1244
    dataLen = 0;
    qwBuildAndSendFetchRsp(qwMsg->connection, rsp, dataLen, code);
    QW_TASK_DLOG("fetch msg rsped, code:%x, dataLen:%d", code, dataLen);
D
dapan 已提交
1245 1246
  } else if (rsp) {
    qwBuildAndSendFetchRsp(qwMsg->connection, rsp, dataLen, code);
D
dapan1121 已提交
1247
    QW_TASK_DLOG("fetch msg rsped, code:%x, dataLen:%d", code, dataLen);
D
dapan1121 已提交
1248 1249
  }

D
dapan1121 已提交
1250 1251
  QW_RET(code);
}
D
dapan1121 已提交
1252

1253

D
dapan1121 已提交
1254
int32_t qwProcessDrop(QW_FPARAMS_DEF, SQWMsg *qwMsg) {
D
dapan1121 已提交
1255 1256
  int32_t code = 0;
  bool needRsp = false;
D
dapan1121 已提交
1257 1258 1259
  SQWTaskCtx *ctx = NULL;
  bool locked = false;

D
dapan1121 已提交
1260
  QW_ERR_JRET(qwAddAcquireTaskCtx(QW_FPARAMS(), &ctx));
D
dapan1121 已提交
1261
  
D
dapan1121 已提交
1262 1263 1264 1265 1266 1267 1268 1269 1270
  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 已提交
1271
  if (QW_IS_QUERY_RUNNING(ctx)) {
D
dapan1121 已提交
1272 1273 1274 1275 1276
    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 已提交
1277
    QW_ERR_JRET(qwDropTaskCtx(QW_FPARAMS(), QW_WRITE));
D
dapan1121 已提交
1278

D
dapan1121 已提交
1279 1280
    QW_SET_RSP_CODE(ctx, TSDB_CODE_QRY_TASK_DROPPED);

D
dapan1121 已提交
1281 1282 1283
    locked = false;
    needRsp = true;
  }
D
dapan1121 已提交
1284

D
dapan 已提交
1285 1286 1287
  if (!needRsp) {    
    ctx->dropConnection = qwMsg->connection;
    
D
dapan1121 已提交
1288 1289 1290
    QW_SET_EVENT_RECEIVED(ctx, QW_EVENT_DROP);
  }
  
D
dapan1121 已提交
1291
_return:
D
dapan1121 已提交
1292

D
dapan1121 已提交
1293
  if (code) {
D
dapan1121 已提交
1294
    QW_UPDATE_RSP_CODE(ctx, code);
D
dapan1121 已提交
1295 1296
  }

D
dapan 已提交
1297 1298 1299 1300
  if (locked) {
    QW_UNLOCK(QW_WRITE, &ctx->lock);
  }

D
dapan1121 已提交
1301
  if (ctx) {
D
dapan1121 已提交
1302
    qwReleaseTaskCtx(mgmt, ctx);
D
dapan1121 已提交
1303 1304
  }

D
dapan1121 已提交
1305 1306
  if (TSDB_CODE_SUCCESS != code || needRsp) {
    QW_ERR_RET(qwBuildAndSendDropRsp(qwMsg->connection, code));
D
dapan1121 已提交
1307 1308

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

D
dapan1121 已提交
1311
  QW_RET(code);
D
dapan1121 已提交
1312
}
D
dapan1121 已提交
1313

D
dapan1121 已提交
1314 1315 1316 1317 1318 1319
int32_t qWorkerInit(int8_t nodeType, int32_t nodeId, SQWorkerCfg *cfg, void **qWorkerMgmt, void *nodeObj, putReqToQueryQFp fp) {
  if (NULL == qWorkerMgmt || NULL == nodeObj || NULL == fp) {
    qError("invalid param to init qworker");
    QW_RET(TSDB_CODE_QRY_INVALID_INPUT);
  }
  
D
dapan1121 已提交
1320 1321 1322
  SQWorkerMgmt *mgmt = calloc(1, sizeof(SQWorkerMgmt));
  if (NULL == mgmt) {
    qError("calloc %d failed", (int32_t)sizeof(SQWorkerMgmt));
D
dapan1121 已提交
1323
    QW_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
1324 1325 1326 1327
  }

  if (cfg) {
    mgmt->cfg = *cfg;
D
dapan1121 已提交
1328
    if (0 == mgmt->cfg.maxSchedulerNum) {
D
dapan1121 已提交
1329
      mgmt->cfg.maxSchedulerNum = QW_DEFAULT_SCHEDULER_NUMBER;
D
dapan1121 已提交
1330 1331
    }
    if (0 == mgmt->cfg.maxTaskNum) {
D
dapan1121 已提交
1332
      mgmt->cfg.maxTaskNum = QW_DEFAULT_TASK_NUMBER;
D
dapan1121 已提交
1333 1334
    }
    if (0 == mgmt->cfg.maxSchTaskNum) {
D
dapan1121 已提交
1335
      mgmt->cfg.maxSchTaskNum = QW_DEFAULT_SCH_TASK_NUMBER;
D
dapan1121 已提交
1336
    }
D
dapan1121 已提交
1337
  } else {
D
dapan1121 已提交
1338 1339 1340
    mgmt->cfg.maxSchedulerNum = QW_DEFAULT_SCHEDULER_NUMBER;
    mgmt->cfg.maxTaskNum = QW_DEFAULT_TASK_NUMBER;
    mgmt->cfg.maxSchTaskNum = QW_DEFAULT_SCH_TASK_NUMBER;
D
dapan1121 已提交
1341 1342
  }

1343
  mgmt->schHash = taosHashInit(mgmt->cfg.maxSchedulerNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_UBIGINT), false, HASH_ENTRY_LOCK);
D
dapan1121 已提交
1344
  if (NULL == mgmt->schHash) {
D
dapan1121 已提交
1345
    tfree(mgmt);
D
dapan1121 已提交
1346 1347
    qError("init %d scheduler hash failed", mgmt->cfg.maxSchedulerNum);
    QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
1348 1349
  }

1350
  mgmt->ctxHash = taosHashInit(mgmt->cfg.maxTaskNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_ENTRY_LOCK);
D
dapan1121 已提交
1351
  if (NULL == mgmt->ctxHash) {
D
dapan1121 已提交
1352
    qError("init %d task ctx hash failed", mgmt->cfg.maxTaskNum);
D
dapan1121 已提交
1353 1354
    taosHashCleanup(mgmt->schHash);
    mgmt->schHash = NULL;
D
dapan1121 已提交
1355
    tfree(mgmt);
D
dapan1121 已提交
1356
    QW_ERR_RET(TSDB_CODE_QRY_OUT_OF_MEMORY);
D
dapan1121 已提交
1357 1358
  }

D
dapan1121 已提交
1359 1360 1361 1362 1363
  mgmt->nodeType = nodeType;
  mgmt->nodeId = nodeId;
  mgmt->nodeObj = nodeObj;
  mgmt->putToQueueFp = fp;

D
dapan1121 已提交
1364 1365
  *qWorkerMgmt = mgmt;

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

D
dapan1121 已提交
1368 1369
  return TSDB_CODE_SUCCESS;
}
D
dapan1121 已提交
1370 1371 1372 1373

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

D
dapan1121 已提交
1376
  SQWorkerMgmt *mgmt = *qWorkerMgmt;
D
dapan1121 已提交
1377
  
D
dapan1121 已提交
1378
  //TODO STOP ALL QUERY
D
dapan1121 已提交
1379

D
dapan1121 已提交
1380
  //TODO FREE ALL
D
dapan1121 已提交
1381

D
dapan1121 已提交
1382 1383
  tfree(*qWorkerMgmt);
}
D
dapan1121 已提交
1384

D
dapan1121 已提交
1385 1386 1387
int32_t qwGetSchTasksStatus(SQWorkerMgmt *mgmt, uint64_t sId, uint64_t qId, uint64_t tId, SSchedulerStatusRsp **rsp) {
  SQWSchStatus *sch = NULL;
  int32_t taskNum = 0;
1388

D
dapan1121 已提交
1389
/*
D
dapan1121 已提交
1390 1391 1392
  QW_ERR_RET(qwAcquireScheduler(QW_READ, mgmt, sId, &sch));
  
  sch->lastAccessTs = taosGetTimestampSec();
1393

D
dapan1121 已提交
1394 1395 1396 1397 1398 1399 1400 1401 1402 1403 1404 1405 1406
  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) {
    qError("calloc %d failed", size);
    QW_UNLOCK(QW_READ, &sch->tasksLock);
    qwReleaseScheduler(QW_READ, mgmt);
    
    return TSDB_CODE_QRY_OUT_OF_MEMORY;
  }
1407

D
dapan1121 已提交
1408 1409 1410
  void *key = NULL;
  size_t keyLen = 0;
  int32_t i = 0;
1411

D
dapan1121 已提交
1412 1413 1414 1415
  void *pIter = taosHashIterate(sch->tasksHash, NULL);
  while (pIter) {
    SQWTaskStatus *taskStatus = (SQWTaskStatus *)pIter;
    taosHashGetKey(pIter, &key, &keyLen);
D
dapan1121 已提交
1416

D
dapan1121 已提交
1417 1418 1419 1420
    QW_GET_QTID(key, (*rsp)->status[i].queryId, (*rsp)->status[i].taskId);
    (*rsp)->status[i].status = taskStatus->status;
    
    pIter = taosHashIterate(sch->tasksHash, pIter);
D
dapan1121 已提交
1421
  }  
D
dapan1121 已提交
1422

D
dapan1121 已提交
1423 1424 1425 1426
  QW_UNLOCK(QW_READ, &sch->tasksLock);
  qwReleaseScheduler(QW_READ, mgmt);

  (*rsp)->num = taskNum;
D
dapan1121 已提交
1427
*/
D
dapan1121 已提交
1428

D
dapan1121 已提交
1429 1430 1431 1432
  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
1433

D
dapan1121 已提交
1434 1435
int32_t qwUpdateSchLastAccess(SQWorkerMgmt *mgmt, uint64_t sId, uint64_t qId, uint64_t tId) {
  SQWSchStatus *sch = NULL;
D
dapan1121 已提交
1436

D
dapan1121 已提交
1437
/*
D
dapan1121 已提交
1438
  QW_ERR_RET(qwAcquireScheduler(QW_READ, mgmt, sId, &sch));
D
dapan1121 已提交
1439

D
dapan1121 已提交
1440
  sch->lastAccessTs = taosGetTimestampSec();
D
dapan1121 已提交
1441

D
dapan1121 已提交
1442
  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
1443
*/
D
dapan1121 已提交
1444 1445 1446 1447
  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
1448 1449 1450 1451
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 已提交
1452 1453

/*  
D
dapan1121 已提交
1454 1455 1456 1457
  if (qwAcquireScheduler(QW_READ, mgmt, sId, &sch)) {
    *taskStatus = JOB_TASK_STATUS_NULL;
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
1458

D
dapan1121 已提交
1459 1460 1461 1462 1463 1464
  if (qwAcquireTask(mgmt, QW_READ, sch, queryId, taskId, &task)) {
    qwReleaseScheduler(QW_READ, mgmt);
    
    *taskStatus = JOB_TASK_STATUS_NULL;
    return TSDB_CODE_SUCCESS;
  }
D
dapan1121 已提交
1465

D
dapan1121 已提交
1466
  *taskStatus = task->status;
D
dapan1121 已提交
1467

D
dapan1121 已提交
1468 1469
  qwReleaseTask(QW_READ, sch);
  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
1470
*/
D
dapan1121 已提交
1471 1472 1473 1474 1475

  QW_RET(code);
}


D
dapan1121 已提交
1476 1477 1478
int32_t qwCancelTask(SQWorkerMgmt *mgmt, uint64_t sId, uint64_t qId, uint64_t tId) {
  SQWSchStatus *sch = NULL;
  SQWTaskStatus *task = NULL;
D
dapan1121 已提交
1479
  int32_t code = 0;
D
dapan1121 已提交
1480

D
dapan1121 已提交
1481
/*
D
dapan1121 已提交
1482
  QW_ERR_RET(qwAcquireAddScheduler(QW_READ, mgmt, sId, &sch));
D
dapan1121 已提交
1483

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

D
dapan1121 已提交
1486

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

D
dapan1121 已提交
1489 1490 1491 1492 1493 1494 1495
  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 已提交
1496

D
dapan1121 已提交
1497 1498 1499 1500 1501 1502 1503 1504
    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 已提交
1505 1506
  }

D
dapan1121 已提交
1507 1508 1509 1510
  QW_UNLOCK(QW_WRITE, &task->lock);
  
  qwReleaseTask(QW_READ, sch);
  qwReleaseScheduler(QW_READ, mgmt);
D
dapan1121 已提交
1511

D
dapan1121 已提交
1512 1513 1514 1515 1516
  if (oriStatus == JOB_TASK_STATUS_EXECUTING) {
    //TODO call executer to cancel subquery async
  }
  
  return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
1517 1518

_return:
D
dapan1121 已提交
1519

D
dapan1121 已提交
1520 1521 1522 1523 1524 1525 1526 1527 1528
  if (task) {
    QW_UNLOCK(QW_WRITE, &task->lock);
    
    qwReleaseTask(QW_READ, sch);
  }

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

D
dapan1121 已提交
1531
  QW_RET(code);
D
dapan1121 已提交
1532 1533
}

D
dapan1121 已提交
1534