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

#include "trpc.h"
#include "query.h"
#include "tname.h"
#include "catalogInt.h"
#include "systable.h"
#include "tref.h"

int32_t ctgInitGetTbMetaTask(SCtgJob *pJob, int32_t taskIdx, SName *name) {
D
dapan1121 已提交
24
  SCtgTask task = {0};
D
dapan1121 已提交
25

D
dapan1121 已提交
26 27 28
  task.type = CTG_TASK_GET_TB_META;
  task.taskId = taskIdx;
  task.pJob = pJob;
D
dapan1121 已提交
29

D
dapan1121 已提交
30 31
  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgTbMetaCtx));
  if (NULL == task.taskCtx) {
D
dapan1121 已提交
32 33 34
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
35
  SCtgTbMetaCtx* ctx = task.taskCtx;
D
dapan1121 已提交
36 37
  ctx->pName = taosMemoryMalloc(sizeof(*name));
  if (NULL == ctx->pName) {
D
dapan1121 已提交
38
    taosMemoryFree(task.taskCtx);
D
dapan1121 已提交
39 40 41 42 43 44
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }
  
  memcpy(ctx->pName, name, sizeof(*name));
  ctx->flag = CTG_FLAG_UNKNOWN_STB;

D
dapan1121 已提交
45 46
  taosArrayPush(pJob->pTasks, &task);

D
dapan1121 已提交
47
  qDebug("QID:0x%" PRIx64 " the %d task type %s initialized, tbName:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), name->tname);
D
dapan1121 已提交
48 49 50 51 52

  return TSDB_CODE_SUCCESS;
}

int32_t ctgInitGetDbVgTask(SCtgJob *pJob, int32_t taskIdx, char *dbFName) {
D
dapan1121 已提交
53
  SCtgTask task = {0};
D
dapan1121 已提交
54

D
dapan1121 已提交
55 56 57
  task.type = CTG_TASK_GET_DB_VGROUP;
  task.taskId = taskIdx;
  task.pJob = pJob;
D
dapan1121 已提交
58

D
dapan1121 已提交
59 60
  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgDbVgCtx));
  if (NULL == task.taskCtx) {
D
dapan1121 已提交
61 62 63
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
64
  SCtgDbVgCtx* ctx = task.taskCtx;
D
dapan1121 已提交
65 66 67
  
  memcpy(ctx->dbFName, dbFName, sizeof(ctx->dbFName));

D
dapan1121 已提交
68 69
  taosArrayPush(pJob->pTasks, &task);

D
dapan1121 已提交
70
  qDebug("QID:0x%" PRIx64 " the %d task type %s initialized, dbFName:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), dbFName);
D
dapan1121 已提交
71 72 73 74 75

  return TSDB_CODE_SUCCESS;
}

int32_t ctgInitGetDbCfgTask(SCtgJob *pJob, int32_t taskIdx, char *dbFName) {
D
dapan1121 已提交
76
  SCtgTask task = {0};
D
dapan1121 已提交
77

D
dapan1121 已提交
78 79 80
  task.type = CTG_TASK_GET_DB_CFG;
  task.taskId = taskIdx;
  task.pJob = pJob;
D
dapan1121 已提交
81

D
dapan1121 已提交
82 83
  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgDbCfgCtx));
  if (NULL == task.taskCtx) {
D
dapan1121 已提交
84 85 86
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
87
  SCtgDbCfgCtx* ctx = task.taskCtx;
D
dapan1121 已提交
88 89 90
  
  memcpy(ctx->dbFName, dbFName, sizeof(ctx->dbFName));

D
dapan1121 已提交
91 92
  taosArrayPush(pJob->pTasks, &task);

D
dapan1121 已提交
93
  qDebug("QID:0x%" PRIx64 " the %d task type %s initialized, dbFName:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), dbFName);
D
dapan1121 已提交
94 95 96 97

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115
int32_t ctgInitGetDbInfoTask(SCtgJob *pJob, int32_t taskIdx, char *dbFName) {
  SCtgTask task = {0};

  task.type = CTG_TASK_GET_DB_INFO;
  task.taskId = taskIdx;
  task.pJob = pJob;

  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgDbInfoCtx));
  if (NULL == task.taskCtx) {
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

  SCtgDbInfoCtx* ctx = task.taskCtx;
  
  memcpy(ctx->dbFName, dbFName, sizeof(ctx->dbFName));

  taosArrayPush(pJob->pTasks, &task);

D
dapan1121 已提交
116
  qDebug("QID:0x%" PRIx64 " the %d task type %s initialized, dbFName:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), dbFName);
D
dapan1121 已提交
117 118 119 120 121

  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
122
int32_t ctgInitGetTbHashTask(SCtgJob *pJob, int32_t taskIdx, SName *name) {
D
dapan1121 已提交
123
  SCtgTask task = {0};
D
dapan1121 已提交
124

D
dapan1121 已提交
125 126 127
  task.type = CTG_TASK_GET_TB_HASH;
  task.taskId = taskIdx;
  task.pJob = pJob;
D
dapan1121 已提交
128

D
dapan1121 已提交
129 130
  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgTbHashCtx));
  if (NULL == task.taskCtx) {
D
dapan1121 已提交
131 132 133
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
134
  SCtgTbHashCtx* ctx = task.taskCtx;
D
dapan1121 已提交
135 136
  ctx->pName = taosMemoryMalloc(sizeof(*name));
  if (NULL == ctx->pName) {
D
dapan1121 已提交
137
    taosMemoryFree(task.taskCtx);
D
dapan1121 已提交
138 139 140 141 142 143
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }
  
  memcpy(ctx->pName, name, sizeof(*name));
  tNameGetFullDbName(ctx->pName, ctx->dbFName);

D
dapan1121 已提交
144 145
  taosArrayPush(pJob->pTasks, &task);

D
dapan1121 已提交
146
  qDebug("QID:0x%" PRIx64 " the %d task type %s initialized, tableName:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), name->tname);
D
dapan1121 已提交
147 148 149 150 151

  return TSDB_CODE_SUCCESS;
}

int32_t ctgInitGetQnodeTask(SCtgJob *pJob, int32_t taskIdx) {
D
dapan1121 已提交
152 153 154 155 156 157
  SCtgTask task = {0};

  task.type = CTG_TASK_GET_QNODE;
  task.taskId = taskIdx;
  task.pJob = pJob;
  task.taskCtx = NULL;
D
dapan1121 已提交
158

D
dapan1121 已提交
159
  taosArrayPush(pJob->pTasks, &task);
D
dapan1121 已提交
160

D
dapan1121 已提交
161
  qDebug("QID:%" PRIx64 " the %d task type %s initialized", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type));
D
dapan1121 已提交
162 163 164 165 166

  return TSDB_CODE_SUCCESS;
}

int32_t ctgInitGetIndexTask(SCtgJob *pJob, int32_t taskIdx, char *name) {
D
dapan1121 已提交
167
  SCtgTask task = {0};
D
dapan1121 已提交
168

D
dapan1121 已提交
169 170 171
  task.type = CTG_TASK_GET_INDEX;
  task.taskId = taskIdx;
  task.pJob = pJob;
D
dapan1121 已提交
172

D
dapan1121 已提交
173 174
  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgIndexCtx));
  if (NULL == task.taskCtx) {
D
dapan1121 已提交
175 176 177
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
178
  SCtgIndexCtx* ctx = task.taskCtx;
D
dapan1121 已提交
179 180 181
  
  strcpy(ctx->indexFName, name);

D
dapan1121 已提交
182 183
  taosArrayPush(pJob->pTasks, &task);

D
dapan1121 已提交
184
  qDebug("QID:%" PRIx64 " the %d task type %s initialized, indexFName:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), name);
D
dapan1121 已提交
185 186 187 188 189

  return TSDB_CODE_SUCCESS;
}

int32_t ctgInitGetUdfTask(SCtgJob *pJob, int32_t taskIdx, char *name) {
D
dapan1121 已提交
190
  SCtgTask task = {0};
D
dapan1121 已提交
191

D
dapan1121 已提交
192 193 194
  task.type = CTG_TASK_GET_UDF;
  task.taskId = taskIdx;
  task.pJob = pJob;
D
dapan1121 已提交
195

D
dapan1121 已提交
196 197
  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgUdfCtx));
  if (NULL == task.taskCtx) {
D
dapan1121 已提交
198 199 200
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
201
  SCtgUdfCtx* ctx = task.taskCtx;
D
dapan1121 已提交
202 203 204
  
  strcpy(ctx->udfName, name);

D
dapan1121 已提交
205 206
  taosArrayPush(pJob->pTasks, &task);

D
dapan1121 已提交
207
  qDebug("QID:%" PRIx64 " the %d task type %s initialized, udfName:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), name);
D
dapan1121 已提交
208 209 210 211 212

  return TSDB_CODE_SUCCESS;
}

int32_t ctgInitGetUserTask(SCtgJob *pJob, int32_t taskIdx, SUserAuthInfo *user) {
D
dapan1121 已提交
213
  SCtgTask task = {0};
D
dapan1121 已提交
214

D
dapan1121 已提交
215 216 217
  task.type = CTG_TASK_GET_USER;
  task.taskId = taskIdx;
  task.pJob = pJob;
D
dapan1121 已提交
218

D
dapan1121 已提交
219 220
  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgUserCtx));
  if (NULL == task.taskCtx) {
D
dapan1121 已提交
221 222 223
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
224
  SCtgUserCtx* ctx = task.taskCtx;
D
dapan1121 已提交
225 226 227
  
  memcpy(&ctx->user, user, sizeof(*user));

D
dapan1121 已提交
228 229
  taosArrayPush(pJob->pTasks, &task);

D
dapan1121 已提交
230
  qDebug("QID:%" PRIx64 " the %d task type %s initialized, user:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), user->user);
D
dapan1121 已提交
231 232 233 234

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263
int32_t ctgInitGetTbIndexTask(SCtgJob *pJob, int32_t taskIdx, SName *name) {
  SCtgTask task = {0};

  task.type = CTG_TASK_GET_TB_INDEX;
  task.taskId = taskIdx;
  task.pJob = pJob;

  task.taskCtx = taosMemoryCalloc(1, sizeof(SCtgTbIndexCtx));
  if (NULL == task.taskCtx) {
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

  SCtgTbIndexCtx* ctx = task.taskCtx;
  ctx->pName = taosMemoryMalloc(sizeof(*name));
  if (NULL == ctx->pName) {
    taosMemoryFree(task.taskCtx);
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }
  
  memcpy(ctx->pName, name, sizeof(*name));

  taosArrayPush(pJob->pTasks, &task);

  qDebug("QID:%" PRIx64 " the %d task type %s initialized, tbName:%s", pJob->queryId, taskIdx, ctgTaskTypeStr(task.type), name->tname);

  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348
int32_t ctgHandleForceUpdate(SCatalog* pCtg, SCtgJob *pJob, const SCatalogReq* pReq) {
  int32_t dbNum = pJob->dbCfgNum + pJob->dbVgNum + pJob->dbInfoNum;
  if (dbNum > 0) {
    if (dbNum > pJob->dbCfgNum && dbNum > pJob->dbVgNum && dbNum > pJob->dbInfoNum) {
      SHashObj* pDb = taosHashInit(dbNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_NO_LOCK);
      if (NULL == pDb) {
        CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
      }
      
      for (int32_t i = 0; i < pJob->dbVgNum; ++i) {
        char* dbFName = taosArrayGet(pReq->pDbVgroup, i);
        taosHashPut(pDb, dbFName, strlen(dbFName), dbFName, TSDB_DB_FNAME_LEN);
      }

      for (int32_t i = 0; i < pJob->dbCfgNum; ++i) {
        char* dbFName = taosArrayGet(pReq->pDbCfg, i);
        taosHashPut(pDb, dbFName, strlen(dbFName), dbFName, TSDB_DB_FNAME_LEN);
      }

      for (int32_t i = 0; i < pJob->dbInfoNum; ++i) {
        char* dbFName = taosArrayGet(pReq->pDbInfo, i);
        taosHashPut(pDb, dbFName, strlen(dbFName), dbFName, TSDB_DB_FNAME_LEN);
      }

      char* dbFName = taosHashIterate(pDb, NULL);
      while (dbFName) {
        ctgDropDbVgroupEnqueue(pCtg, dbFName, true);
        dbFName = taosHashIterate(pDb, dbFName);
      }

      taosHashCleanup(pDb);      
    } else {
      for (int32_t i = 0; i < pJob->dbVgNum; ++i) {
        char* dbFName = taosArrayGet(pReq->pDbVgroup, i);
        CTG_ERR_RET(ctgDropDbVgroupEnqueue(pCtg, dbFName, true));
      }
      
      for (int32_t i = 0; i < pJob->dbCfgNum; ++i) {
        char* dbFName = taosArrayGet(pReq->pDbCfg, i);
        CTG_ERR_RET(ctgDropDbVgroupEnqueue(pCtg, dbFName, true));
      }
      
      for (int32_t i = 0; i < pJob->dbInfoNum; ++i) {
        char* dbFName = taosArrayGet(pReq->pDbInfo, i);
        CTG_ERR_RET(ctgDropDbVgroupEnqueue(pCtg, dbFName, true));
      }
    }
  }

  int32_t tbNum = pJob->tbMetaNum + pJob->tbHashNum;
  if (tbNum > 0) {
    if (tbNum > pJob->tbMetaNum && tbNum > pJob->tbHashNum) {
      SHashObj* pTb = taosHashInit(tbNum, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_NO_LOCK);
      for (int32_t i = 0; i < pJob->tbMetaNum; ++i) {
        SName* name = taosArrayGet(pReq->pTableMeta, i);
        taosHashPut(pTb, name, sizeof(SName), name, sizeof(SName));
      }
      
      for (int32_t i = 0; i < pJob->tbHashNum; ++i) {
        SName* name = taosArrayGet(pReq->pTableHash, i);
        taosHashPut(pTb, name, sizeof(SName), name, sizeof(SName));
      }

      SName* name = taosHashIterate(pTb, NULL);
      while (name) {
        catalogRemoveTableMeta(pCtg, name);
        name = taosHashIterate(pTb, name);
      }

      taosHashCleanup(pTb);
    } else {
      for (int32_t i = 0; i < pJob->tbMetaNum; ++i) {
        SName* name = taosArrayGet(pReq->pTableMeta, i);
        catalogRemoveTableMeta(pCtg, name);
      }
      
      for (int32_t i = 0; i < pJob->tbHashNum; ++i) {
        SName* name = taosArrayGet(pReq->pTableHash, i);
        catalogRemoveTableMeta(pCtg, name);
      }
    }
  }

  return TSDB_CODE_SUCCESS;
}
D
dapan1121 已提交
349

350
int32_t ctgInitJob(CTG_PARAMS, SCtgJob** job, uint64_t reqId, const SCatalogReq* pReq, catalogCallback fp, void* param, int32_t* taskNum) {
D
dapan1121 已提交
351 352 353 354 355 356 357 358 359
  int32_t code = 0;
  int32_t tbMetaNum = (int32_t)taosArrayGetSize(pReq->pTableMeta);
  int32_t dbVgNum = (int32_t)taosArrayGetSize(pReq->pDbVgroup);
  int32_t tbHashNum = (int32_t)taosArrayGetSize(pReq->pTableHash);
  int32_t udfNum = (int32_t)taosArrayGetSize(pReq->pUdf);
  int32_t qnodeNum = pReq->qNodeRequired ? 1 : 0;
  int32_t dbCfgNum = (int32_t)taosArrayGetSize(pReq->pDbCfg);
  int32_t indexNum = (int32_t)taosArrayGetSize(pReq->pIndex);
  int32_t userNum = (int32_t)taosArrayGetSize(pReq->pUser);
D
dapan1121 已提交
360
  int32_t dbInfoNum = (int32_t)taosArrayGetSize(pReq->pDbInfo);
D
dapan1121 已提交
361
  int32_t tbIndexNum = (int32_t)taosArrayGetSize(pReq->pTableIndex);
D
dapan1121 已提交
362

D
dapan1121 已提交
363
  *taskNum = tbMetaNum + dbVgNum + udfNum + tbHashNum + qnodeNum + dbCfgNum + indexNum + userNum + dbInfoNum + tbIndexNum;
364
  if (*taskNum <= 0) {
365
    ctgDebug("Empty input for job, no need to retrieve meta, reqId:0x%" PRIx64, reqId);
366
    return TSDB_CODE_SUCCESS;
D
dapan1121 已提交
367
  }
368

D
dapan1121 已提交
369 370
  *job = taosMemoryCalloc(1, sizeof(SCtgJob));
  if (NULL == *job) {
371
    ctgError("failed to calloc, size:%d, reqId:0x%" PRIx64, (int32_t)sizeof(SCtgJob), reqId);
D
dapan1121 已提交
372 373 374 375 376 377 378 379 380
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

  SCtgJob *pJob = *job;
  
  pJob->queryId = reqId;
  pJob->userFp = fp;
  pJob->pCtg    = pCtg;
  pJob->pTrans    = pTrans;
D
dapan1121 已提交
381
  pJob->pMgmtEps    = *pMgmtEps;
D
dapan1121 已提交
382 383 384 385 386 387 388 389 390 391
  pJob->userParam = param;
  
  pJob->tbMetaNum = tbMetaNum;
  pJob->tbHashNum = tbHashNum;
  pJob->qnodeNum = qnodeNum;
  pJob->dbVgNum = dbVgNum;
  pJob->udfNum = udfNum;
  pJob->dbCfgNum = dbCfgNum;
  pJob->indexNum = indexNum;
  pJob->userNum = userNum;
D
dapan1121 已提交
392
  pJob->dbInfoNum = dbInfoNum;
D
dapan1121 已提交
393
  pJob->tbIndexNum = tbIndexNum;
394 395 396

  pJob->pTasks = taosArrayInit(*taskNum, sizeof(SCtgTask));

D
dapan1121 已提交
397
  if (NULL == pJob->pTasks) {
398
    ctgError("taosArrayInit %d tasks failed", *taskNum);
D
dapan1121 已提交
399 400 401
    CTG_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
  }

D
dapan1121 已提交
402 403 404 405
  if (pReq->forceUpdate) {
    CTG_ERR_JRET(ctgHandleForceUpdate(pCtg, pJob, pReq));
  }

D
dapan1121 已提交
406 407
  int32_t taskIdx = 0;
  for (int32_t i = 0; i < dbVgNum; ++i) {
408
    char* dbFName = taosArrayGet(pReq->pDbVgroup, i);
D
dapan1121 已提交
409 410 411 412
    CTG_ERR_JRET(ctgInitGetDbVgTask(pJob, taskIdx++, dbFName));
  }

  for (int32_t i = 0; i < dbCfgNum; ++i) {
413
    char* dbFName = taosArrayGet(pReq->pDbCfg, i);
D
dapan1121 已提交
414 415 416
    CTG_ERR_JRET(ctgInitGetDbCfgTask(pJob, taskIdx++, dbFName));
  }

D
dapan1121 已提交
417
  for (int32_t i = 0; i < dbInfoNum; ++i) {
418
    char* dbFName = taosArrayGet(pReq->pDbInfo, i);
D
dapan1121 已提交
419 420 421
    CTG_ERR_JRET(ctgInitGetDbInfoTask(pJob, taskIdx++, dbFName));
  }

D
dapan1121 已提交
422
  for (int32_t i = 0; i < tbMetaNum; ++i) {
423
    SName* name = taosArrayGet(pReq->pTableMeta, i);
D
dapan1121 已提交
424 425 426 427
    CTG_ERR_JRET(ctgInitGetTbMetaTask(pJob, taskIdx++, name));
  }

  for (int32_t i = 0; i < tbHashNum; ++i) {
428
    SName* name = taosArrayGet(pReq->pTableHash, i);
D
dapan1121 已提交
429 430 431
    CTG_ERR_JRET(ctgInitGetTbHashTask(pJob, taskIdx++, name));
  }

D
dapan1121 已提交
432 433 434 435 436
  for (int32_t i = 0; i < tbIndexNum; ++i) {
    SName* name = taosArrayGet(pReq->pTableIndex, i);
    CTG_ERR_JRET(ctgInitGetTbIndexTask(pJob, taskIdx++, name));
  }

D
dapan1121 已提交
437
  for (int32_t i = 0; i < indexNum; ++i) {
438
    char* indexName = taosArrayGet(pReq->pIndex, i);
D
dapan1121 已提交
439 440 441 442
    CTG_ERR_JRET(ctgInitGetIndexTask(pJob, taskIdx++, indexName));
  }

  for (int32_t i = 0; i < udfNum; ++i) {
443
    char* udfName = taosArrayGet(pReq->pUdf, i);
D
dapan1121 已提交
444 445 446 447
    CTG_ERR_JRET(ctgInitGetUdfTask(pJob, taskIdx++, udfName));
  }

  for (int32_t i = 0; i < userNum; ++i) {
448
    SUserAuthInfo* user = taosArrayGet(pReq->pUser, i);
D
dapan1121 已提交
449 450 451 452 453 454 455
    CTG_ERR_JRET(ctgInitGetUserTask(pJob, taskIdx++, user));
  }

  if (qnodeNum) {
    CTG_ERR_JRET(ctgInitGetQnodeTask(pJob, taskIdx++));
  }

H
Haojun Liao 已提交
456 457 458 459 460 461 462 463
  pJob->refId = taosAddRef(gCtgMgmt.jobPool, pJob);
  if (pJob->refId < 0) {
    ctgError("add job to ref failed, error: %s", tstrerror(terrno));
    CTG_ERR_JRET(terrno);
  }

  taosAcquireRef(gCtgMgmt.jobPool, pJob->refId);

D
dapan1121 已提交
464
  qDebug("QID:0x%" PRIx64 ", jobId: 0x%" PRIx64 " initialized, task num %d, forceUpdate %d", pJob->queryId, pJob->refId, *taskNum, pReq->forceUpdate);
H
Haojun Liao 已提交
465 466 467
  return TSDB_CODE_SUCCESS;


D
dapan1121 已提交
468 469 470 471 472 473 474 475
_return:
  taosMemoryFreeClear(*job);
  CTG_RET(code);
}

int32_t ctgDumpTbMetaRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pTableMeta) {
D
dapan1121 已提交
476
    pJob->jobRes.pTableMeta = taosArrayInit(pJob->tbMetaNum, sizeof(SMetaRes));
D
dapan1121 已提交
477 478 479 480 481
    if (NULL == pJob->jobRes.pTableMeta) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
482 483
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pTableMeta, &res);
D
dapan1121 已提交
484 485 486 487 488 489 490

  return TSDB_CODE_SUCCESS;
}

int32_t ctgDumpDbVgRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pDbVgroup) {
D
dapan1121 已提交
491
    pJob->jobRes.pDbVgroup = taosArrayInit(pJob->dbVgNum, sizeof(SMetaRes));
D
dapan1121 已提交
492 493 494 495 496
    if (NULL == pJob->jobRes.pDbVgroup) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
497 498
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pDbVgroup, &res);
D
dapan1121 已提交
499 500 501 502 503 504 505

  return TSDB_CODE_SUCCESS;
}

int32_t ctgDumpTbHashRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pTableHash) {
D
dapan1121 已提交
506
    pJob->jobRes.pTableHash = taosArrayInit(pJob->tbHashNum, sizeof(SMetaRes));
D
dapan1121 已提交
507 508 509 510 511
    if (NULL == pJob->jobRes.pTableHash) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
512 513
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pTableHash, &res);
D
dapan1121 已提交
514 515 516 517

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
518 519 520 521 522 523 524 525 526 527 528 529 530 531 532
int32_t ctgDumpTbIndexRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pTableIndex) {
    pJob->jobRes.pTableIndex = taosArrayInit(pJob->tbIndexNum, sizeof(SMetaRes));
    if (NULL == pJob->jobRes.pTableIndex) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pTableHash, &res);

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
533 534 535
int32_t ctgDumpIndexRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pIndex) {
D
dapan1121 已提交
536
    pJob->jobRes.pIndex = taosArrayInit(pJob->indexNum, sizeof(SMetaRes));
D
dapan1121 已提交
537 538 539 540 541
    if (NULL == pJob->jobRes.pIndex) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
542 543
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pIndex, &res);
D
dapan1121 已提交
544 545 546 547 548 549

  return TSDB_CODE_SUCCESS;
}

int32_t ctgDumpQnodeRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
D
dapan1121 已提交
550 551 552 553 554 555
  if (NULL == pJob->jobRes.pQnodeList) {
    pJob->jobRes.pQnodeList = taosArrayInit(1, sizeof(SMetaRes));
    if (NULL == pJob->jobRes.pQnodeList) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }
D
dapan1121 已提交
556

D
dapan1121 已提交
557 558
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pQnodeList, &res);
D
dapan1121 已提交
559 560 561 562 563 564 565

  return TSDB_CODE_SUCCESS;
}

int32_t ctgDumpDbCfgRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pDbCfg) {
D
dapan1121 已提交
566
    pJob->jobRes.pDbCfg = taosArrayInit(pJob->dbCfgNum, sizeof(SMetaRes));
D
dapan1121 已提交
567 568 569 570 571
    if (NULL == pJob->jobRes.pDbCfg) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
572 573
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pDbCfg, &res);
D
dapan1121 已提交
574 575 576 577

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
578 579 580
int32_t ctgDumpDbInfoRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pDbInfo) {
D
dapan1121 已提交
581
    pJob->jobRes.pDbInfo = taosArrayInit(pJob->dbInfoNum, sizeof(SMetaRes));
D
dapan1121 已提交
582 583 584 585 586
    if (NULL == pJob->jobRes.pDbInfo) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
587 588
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pDbInfo, &res);
D
dapan1121 已提交
589 590 591 592

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
593 594 595
int32_t ctgDumpUdfRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pUdfList) {
D
dapan1121 已提交
596
    pJob->jobRes.pUdfList = taosArrayInit(pJob->udfNum, sizeof(SMetaRes));
D
dapan1121 已提交
597 598 599 600 601
    if (NULL == pJob->jobRes.pUdfList) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
602 603
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pUdfList, &res);
D
dapan1121 已提交
604 605 606 607 608 609 610

  return TSDB_CODE_SUCCESS;
}

int32_t ctgDumpUserRes(SCtgTask* pTask) {
  SCtgJob* pJob = pTask->pJob;
  if (NULL == pJob->jobRes.pUser) {
D
dapan1121 已提交
611
    pJob->jobRes.pUser = taosArrayInit(pJob->userNum, sizeof(SMetaRes));
D
dapan1121 已提交
612 613 614 615 616
    if (NULL == pJob->jobRes.pUser) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
  }

D
dapan1121 已提交
617 618
  SMetaRes res = {.code = pTask->code, .pRes = pTask->res};
  taosArrayPush(pJob->jobRes.pUser, &res);
D
dapan1121 已提交
619 620 621 622 623 624 625 626

  return TSDB_CODE_SUCCESS;
}

int32_t ctgHandleTaskEnd(SCtgTask* pTask, int32_t rspCode) {
  SCtgJob* pJob = pTask->pJob;
  int32_t code = 0;

627
  qDebug("QID:0x%" PRIx64 " task %d end with rsp %s", pJob->queryId, pTask->taskId, tstrerror(rspCode));
D
dapan1121 已提交
628

D
dapan1121 已提交
629
  pTask->code = rspCode;
D
dapan1121 已提交
630 631 632

  int32_t taskDone = atomic_add_fetch_32(&pJob->taskDone, 1);
  if (taskDone < taosArrayGetSize(pJob->pTasks)) {
633
    qDebug("QID:0x%" PRIx64 " task done: %d, total: %d", pJob->queryId, taskDone, (int32_t)taosArrayGetSize(pJob->pTasks));
D
dapan1121 已提交
634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657
    return TSDB_CODE_SUCCESS;
  }

  CTG_ERR_JRET(ctgMakeAsyncRes(pJob));

_return:

  qDebug("QID:%" PRIx64 " user callback with rsp %s", pJob->queryId, tstrerror(code));

  (*pJob->userFp)(&pJob->jobRes, pJob->userParam, code);

  taosRemoveRef(gCtgMgmt.jobPool, pJob->refId);
  
  CTG_RET(code);
}

int32_t ctgHandleGetTbMetaRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  SCtgDBCache *dbCache = NULL;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  SCtgTbMetaCtx* ctx = (SCtgTbMetaCtx*)pTask->taskCtx;
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
658
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735

  switch (reqType) {
    case TDMT_MND_USE_DB: {
      SUseDbOutput* pOut = (SUseDbOutput*)pTask->msgCtx.out;

      SVgroupInfo vgInfo = {0};
      CTG_ERR_JRET(ctgGetVgInfoFromHashValue(pCtg, pOut->dbVgroup, ctx->pName, &vgInfo));

      ctgDebug("will refresh tbmeta, not supposed to be stb, tbName:%s, flag:%d", tNameGetTableName(ctx->pName), ctx->flag);

      CTG_ERR_JRET(ctgGetTbMetaFromVnode(CTG_PARAMS_LIST(), ctx->pName, &vgInfo, NULL, pTask));

      return TSDB_CODE_SUCCESS;
    }
    case TDMT_MND_TABLE_META: {
      STableMetaOutput* pOut = (STableMetaOutput*)pTask->msgCtx.out;
      
      if (CTG_IS_META_NULL(pOut->metaType)) {
        if (CTG_FLAG_IS_STB(ctx->flag)) {
          char dbFName[TSDB_DB_FNAME_LEN] = {0};
          tNameGetFullDbName(ctx->pName, dbFName);
          
          CTG_ERR_RET(ctgAcquireVgInfoFromCache(pCtg, dbFName, &dbCache));
          if (NULL != dbCache) {
            SVgroupInfo vgInfo = {0};
            CTG_ERR_JRET(ctgGetVgInfoFromHashValue(pCtg, dbCache->vgInfo, ctx->pName, &vgInfo));
          
            ctgDebug("will refresh tbmeta, not supposed to be stb, tbName:%s, flag:%d", tNameGetTableName(ctx->pName), ctx->flag);
          
            CTG_ERR_JRET(ctgGetTbMetaFromVnode(CTG_PARAMS_LIST(), ctx->pName, &vgInfo, NULL, pTask));

            ctgReleaseVgInfo(dbCache);
            ctgReleaseDBCache(pCtg, dbCache);
          } else {
            SBuildUseDBInput input = {0};
          
            tstrncpy(input.db, dbFName, tListLen(input.db));
            input.vgVersion = CTG_DEFAULT_INVALID_VERSION;
          
            CTG_ERR_JRET(ctgGetDBVgInfoFromMnode(pCtg, pTrans, pMgmtEps, &input, NULL, pTask));
          }

          return TSDB_CODE_SUCCESS;
        }
        
        ctgError("no tbmeta got, tbNmae:%s", tNameGetTableName(ctx->pName));
        ctgRemoveTbMetaFromCache(pCtg, ctx->pName, false);
        
        CTG_ERR_JRET(CTG_ERR_CODE_TABLE_NOT_EXIST);
      }

      if (pTask->msgCtx.lastOut) {
        TSWAP(pTask->msgCtx.out, pTask->msgCtx.lastOut);
        STableMetaOutput* pLastOut = (STableMetaOutput*)pTask->msgCtx.out;
        TSWAP(pLastOut->tbMeta, pOut->tbMeta);
      }
      
      break;
    }
    case TDMT_VND_TABLE_META: {
      STableMetaOutput* pOut = (STableMetaOutput*)pTask->msgCtx.out;
      
      if (CTG_IS_META_NULL(pOut->metaType)) {
        ctgError("no tbmeta got, tbNmae:%s", tNameGetTableName(ctx->pName));
        ctgRemoveTbMetaFromCache(pCtg, ctx->pName, false);
        CTG_ERR_JRET(CTG_ERR_CODE_TABLE_NOT_EXIST);
      }

      if (CTG_FLAG_IS_STB(ctx->flag)) {
        break;
      }
      
      if (CTG_IS_META_TABLE(pOut->metaType) && TSDB_SUPER_TABLE == pOut->tbMeta->tableType) {
        ctgDebug("will continue to refresh tbmeta since got stb, tbName:%s", tNameGetTableName(ctx->pName));
      
        taosMemoryFreeClear(pOut->tbMeta);
        
D
dapan1121 已提交
736
        CTG_RET(ctgGetTbMetaFromMnode(CTG_PARAMS_LIST(), ctx->pName, NULL, pTask));
D
dapan1121 已提交
737 738 739 740 741 742 743 744
      } else if (CTG_IS_META_BOTH(pOut->metaType)) {
        int32_t exist = 0;
        if (!CTG_FLAG_IS_FORCE_UPDATE(ctx->flag)) {
          CTG_ERR_JRET(ctgTbMetaExistInCache(pCtg, pOut->dbFName, pOut->tbName, &exist));
        }
      
        if (0 == exist) {
          TSWAP(pTask->msgCtx.lastOut, pTask->msgCtx.out);
D
dapan1121 已提交
745
          CTG_RET(ctgGetTbMetaFromMnodeImpl(CTG_PARAMS_LIST(), pOut->dbFName, pOut->tbName, NULL, pTask));
D
dapan1121 已提交
746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804
        } else {
          taosMemoryFreeClear(pOut->tbMeta);
          
          SET_META_TYPE_CTABLE(pOut->metaType); 
        }
      }
      break;
    }
    default:
      ctgError("invalid reqType %d", reqType);
      CTG_ERR_JRET(TSDB_CODE_INVALID_MSG);
      break;
  }

  STableMetaOutput* pOut = (STableMetaOutput*)pTask->msgCtx.out;

  ctgUpdateTbMetaToCache(pCtg, pOut, false);

  if (CTG_IS_META_BOTH(pOut->metaType)) {
    memcpy(pOut->tbMeta, &pOut->ctbMeta, sizeof(pOut->ctbMeta));
  } else if (CTG_IS_META_CTABLE(pOut->metaType)) {
    SName stbName = *ctx->pName;
    strcpy(stbName.tname, pOut->tbName);
    SCtgTbMetaCtx stbCtx = {0};
    stbCtx.flag = ctx->flag;
    stbCtx.pName = &stbName;
    
    CTG_ERR_JRET(ctgReadTbMetaFromCache(pCtg, &stbCtx, &pOut->tbMeta));
    if (NULL == pOut->tbMeta) {
      ctgDebug("stb no longer exist, stbName:%s", stbName.tname);
      CTG_ERR_JRET(ctgRelaunchGetTbMetaTask(pTask));

      return TSDB_CODE_SUCCESS;
    }

    memcpy(pOut->tbMeta, &pOut->ctbMeta, sizeof(pOut->ctbMeta));
  }

  TSWAP(pTask->res, pOut->tbMeta);

_return:

  if (dbCache) {
    ctgReleaseVgInfo(dbCache);
    ctgReleaseDBCache(pCtg, dbCache);
  }

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}

int32_t ctgHandleGetDbVgRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  SCtgDbVgCtx* ctx = (SCtgDbVgCtx*)pTask->taskCtx;
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
805
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
806 807 808 809 810 811 812

  switch (reqType) {
    case TDMT_MND_USE_DB: {
      SUseDbOutput* pOut = (SUseDbOutput*)pTask->msgCtx.out;

      CTG_ERR_JRET(ctgGenerateVgList(pCtg, pOut->dbVgroup->vgHash, (SArray**)&pTask->res));
      
D
dapan1121 已提交
813
      CTG_ERR_JRET(ctgUpdateVgroupEnqueue(pCtg, ctx->dbFName, pOut->dbId, pOut->dbVgroup, false));
D
dapan1121 已提交
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
      pOut->dbVgroup = NULL;

      break;
    }
    default:
      ctgError("invalid reqType %d", reqType);
      CTG_ERR_JRET(TSDB_CODE_INVALID_MSG);
      break;
  }


_return:

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}

int32_t ctgHandleGetTbHashRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  SCtgTbHashCtx* ctx = (SCtgTbHashCtx*)pTask->taskCtx;
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
839
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
840 841 842 843 844 845 846 847 848 849 850 851

  switch (reqType) {
    case TDMT_MND_USE_DB: {
      SUseDbOutput* pOut = (SUseDbOutput*)pTask->msgCtx.out;

      pTask->res = taosMemoryMalloc(sizeof(SVgroupInfo));
      if (NULL == pTask->res) {
        CTG_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
      }

      CTG_ERR_JRET(ctgGetVgInfoFromHashValue(pCtg, pOut->dbVgroup, ctx->pName, (SVgroupInfo*)pTask->res));
      
D
dapan1121 已提交
852
      CTG_ERR_JRET(ctgUpdateVgroupEnqueue(pCtg, ctx->dbFName, pOut->dbId, pOut->dbVgroup, false));
D
dapan1121 已提交
853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870
      pOut->dbVgroup = NULL;

      break;
    }
    default:
      ctgError("invalid reqType %d", reqType);
      CTG_ERR_JRET(TSDB_CODE_INVALID_MSG);
      break;
  }


_return:

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}

D
dapan1121 已提交
871 872 873 874 875 876 877 878 879 880 881 882 883 884
int32_t ctgHandleGetTbIndexRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  TSWAP(pTask->res, pTask->msgCtx.out);
  
_return:

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}


D
dapan1121 已提交
885 886 887 888 889 890 891 892 893 894 895 896 897
int32_t ctgHandleGetDbCfgRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  TSWAP(pTask->res, pTask->msgCtx.out);
  
_return:

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}

D
dapan1121 已提交
898 899 900 901 902
int32_t ctgHandleGetDbInfoRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  CTG_RET(TSDB_CODE_APP_ERROR);
}


D
dapan1121 已提交
903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 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 947 948 949
int32_t ctgHandleGetQnodeRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  TSWAP(pTask->res, pTask->msgCtx.out);
  
_return:

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}

int32_t ctgHandleGetIndexRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  TSWAP(pTask->res, pTask->msgCtx.out);
  
_return:

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}

int32_t ctgHandleGetUdfRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  TSWAP(pTask->res, pTask->msgCtx.out);
  
_return:

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}

int32_t ctgHandleGetUserRsp(SCtgTask* pTask, int32_t reqType, const SDataBuf *pMsg, int32_t rspCode) {
  int32_t code = 0;
  SCtgDBCache *dbCache = NULL;
  CTG_ERR_JRET(ctgProcessRspMsg(pTask->msgCtx.out, reqType, pMsg->pData, pMsg->len, rspCode, pTask->msgCtx.target));

  SCtgUserCtx* ctx = (SCtgUserCtx*)pTask->taskCtx;
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
950
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980
  bool pass = false;
  SGetUserAuthRsp* pOut = (SGetUserAuthRsp*)pTask->msgCtx.out;
  
  if (pOut->superAuth) {
    pass = true;
    goto _return;
  }

  if (pOut->createdDbs && taosHashGet(pOut->createdDbs, ctx->user.dbFName, strlen(ctx->user.dbFName))) {
    pass = true;
    goto _return;
  }

  if (ctx->user.type == AUTH_TYPE_READ && pOut->readDbs && taosHashGet(pOut->readDbs, ctx->user.dbFName, strlen(ctx->user.dbFName))) {
    pass = true;
  } else if (ctx->user.type == AUTH_TYPE_WRITE && pOut->writeDbs && taosHashGet(pOut->writeDbs, ctx->user.dbFName, strlen(ctx->user.dbFName))) {
    pass = true;
  }

_return:

  if (TSDB_CODE_SUCCESS == code) {
    pTask->res = taosMemoryCalloc(1, sizeof(bool));
    if (NULL == pTask->res) {
      code = TSDB_CODE_OUT_OF_MEMORY;
    } else {
      *(bool*)pTask->res = pass;
    }
  }

D
dapan1121 已提交
981
  ctgUpdateUserEnqueue(pCtg, pOut, false);
D
dapan1121 已提交
982
  taosMemoryFreeClear(pTask->msgCtx.out);
D
dapan1121 已提交
983 984 985 986 987 988 989 990 991

  ctgHandleTaskEnd(pTask, code);

  CTG_RET(code);
}

int32_t ctgAsyncRefreshTbMeta(SCtgTask *pTask) {
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
992
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013
  int32_t code = 0;
  SCtgTbMetaCtx* ctx = (SCtgTbMetaCtx*)pTask->taskCtx;

  if (CTG_FLAG_IS_SYS_DB(ctx->flag)) {
    ctgDebug("will refresh sys db tbmeta, tbName:%s", tNameGetTableName(ctx->pName));

    CTG_RET(ctgGetTbMetaFromMnodeImpl(CTG_PARAMS_LIST(), (char *)ctx->pName->dbname, (char *)ctx->pName->tname, NULL, pTask));
  }

  if (CTG_FLAG_IS_STB(ctx->flag)) {
    ctgDebug("will refresh tbmeta, supposed to be stb, tbName:%s", tNameGetTableName(ctx->pName));

    // if get from mnode failed, will not try vnode
    CTG_RET(ctgGetTbMetaFromMnode(CTG_PARAMS_LIST(), ctx->pName, NULL, pTask));
  }

  SCtgDBCache *dbCache = NULL;
  char dbFName[TSDB_DB_FNAME_LEN] = {0};
  tNameGetFullDbName(ctx->pName, dbFName);
  
  CTG_ERR_RET(ctgAcquireVgInfoFromCache(pCtg, dbFName, &dbCache));
D
dapan1121 已提交
1014
  if (dbCache) {
D
dapan1121 已提交
1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042
    SVgroupInfo vgInfo = {0};
    CTG_ERR_RET(ctgGetVgInfoFromHashValue(pCtg, dbCache->vgInfo, ctx->pName, &vgInfo));

    ctgDebug("will refresh tbmeta, not supposed to be stb, tbName:%s, flag:%d", tNameGetTableName(ctx->pName), ctx->flag);

    CTG_ERR_JRET(ctgGetTbMetaFromVnode(CTG_PARAMS_LIST(), ctx->pName, &vgInfo, NULL, pTask));
  } else {
    SBuildUseDBInput input = {0};

    tstrncpy(input.db, dbFName, tListLen(input.db));
    input.vgVersion = CTG_DEFAULT_INVALID_VERSION;

    CTG_ERR_JRET(ctgGetDBVgInfoFromMnode(pCtg, pTrans, pMgmtEps, &input, NULL, pTask));
  }

_return:

  if (dbCache) {
    ctgReleaseVgInfo(dbCache);
    ctgReleaseDBCache(pCtg, dbCache);
  }

  CTG_RET(code);
}

int32_t ctgLaunchGetTbMetaTask(SCtgTask *pTask) {
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
1043
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059

  CTG_ERR_RET(ctgGetTbMetaFromCache(CTG_PARAMS_LIST(), (SCtgTbMetaCtx*)pTask->taskCtx, (STableMeta**)&pTask->res));
  if (pTask->res) {
    CTG_ERR_RET(ctgHandleTaskEnd(pTask, 0));
    return TSDB_CODE_SUCCESS;
  }

  CTG_ERR_RET(ctgAsyncRefreshTbMeta(pTask));

  return TSDB_CODE_SUCCESS;
}

int32_t ctgLaunchGetDbVgTask(SCtgTask *pTask) {
  int32_t code = 0;
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
1060
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091
  SCtgDBCache *dbCache = NULL;
  SCtgDbVgCtx* pCtx = (SCtgDbVgCtx*)pTask->taskCtx;
  
  CTG_ERR_RET(ctgAcquireVgInfoFromCache(pCtg, pCtx->dbFName, &dbCache));
  if (NULL != dbCache) {
    CTG_ERR_JRET(ctgGenerateVgList(pCtg, dbCache->vgInfo->vgHash, (SArray**)&pTask->res));
    
    CTG_ERR_JRET(ctgHandleTaskEnd(pTask, 0));
  } else {
    SBuildUseDBInput input = {0};
    
    tstrncpy(input.db, pCtx->dbFName, tListLen(input.db));
    input.vgVersion = CTG_DEFAULT_INVALID_VERSION;
    
    CTG_ERR_RET(ctgGetDBVgInfoFromMnode(pCtg, pTrans, pMgmtEps, &input, NULL, pTask));
  }

_return:

  if (dbCache) {
    ctgReleaseVgInfo(dbCache);
    ctgReleaseDBCache(pCtg, dbCache);
  }

  CTG_RET(code);
}

int32_t ctgLaunchGetTbHashTask(SCtgTask *pTask) {
  int32_t code = 0;
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
1092
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123
  SCtgDBCache *dbCache = NULL;
  SCtgTbHashCtx* pCtx = (SCtgTbHashCtx*)pTask->taskCtx;
  
  CTG_ERR_RET(ctgAcquireVgInfoFromCache(pCtg, pCtx->dbFName, &dbCache));
  if (NULL != dbCache) {
    pTask->res = taosMemoryMalloc(sizeof(SVgroupInfo));
    if (NULL == pTask->res) {
      CTG_ERR_JRET(TSDB_CODE_OUT_OF_MEMORY);
    }
    CTG_ERR_JRET(ctgGetVgInfoFromHashValue(pCtg, dbCache->vgInfo, pCtx->pName, (SVgroupInfo*)pTask->res));
    
    CTG_ERR_JRET(ctgHandleTaskEnd(pTask, 0));
  } else {
    SBuildUseDBInput input = {0};
    
    tstrncpy(input.db, pCtx->dbFName, tListLen(input.db));
    input.vgVersion = CTG_DEFAULT_INVALID_VERSION;
    
    CTG_ERR_RET(ctgGetDBVgInfoFromMnode(pCtg, pTrans, pMgmtEps, &input, NULL, pTask));
  }

_return:

  if (dbCache) {
    ctgReleaseVgInfo(dbCache);
    ctgReleaseDBCache(pCtg, dbCache);
  }

  CTG_RET(code);
}

D
dapan1121 已提交
1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135
int32_t ctgLaunchGetTbIndexTask(SCtgTask *pTask) {
  int32_t code = 0;
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
  SCtgTbIndexCtx* pCtx = (SCtgTbIndexCtx*)pTask->taskCtx;
  
  CTG_ERR_RET(ctgGetTbIndexFromMnode(CTG_PARAMS_LIST(), pCtx->pName, NULL, pTask));
  return TSDB_CODE_SUCCESS;
}


D
dapan1121 已提交
1136 1137 1138
int32_t ctgLaunchGetQnodeTask(SCtgTask *pTask) {
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
1139
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
1140 1141 1142 1143 1144 1145 1146 1147

  CTG_ERR_RET(ctgGetQnodeListFromMnode(CTG_PARAMS_LIST(), NULL, pTask));
  return TSDB_CODE_SUCCESS;
}

int32_t ctgLaunchGetDbCfgTask(SCtgTask *pTask) {
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
1148
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
1149 1150 1151 1152 1153 1154 1155
  SCtgDbCfgCtx* pCtx = (SCtgDbCfgCtx*)pTask->taskCtx;

  CTG_ERR_RET(ctgGetDBCfgFromMnode(CTG_PARAMS_LIST(), pCtx->dbFName, NULL, pTask));

  return TSDB_CODE_SUCCESS;
}

D
dapan1121 已提交
1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190
int32_t ctgLaunchGetDbInfoTask(SCtgTask *pTask) {
  int32_t code = 0;
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
  SCtgDBCache *dbCache = NULL;
  SCtgDbInfoCtx* pCtx = (SCtgDbInfoCtx*)pTask->taskCtx;

  pTask->res = taosMemoryCalloc(1, sizeof(SDbInfo));
  if (NULL == pTask->res) {
    CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
  }

  SDbInfo* pInfo = (SDbInfo*)pTask->res;
  CTG_ERR_RET(ctgAcquireVgInfoFromCache(pCtg, pCtx->dbFName, &dbCache));
  if (NULL != dbCache) {
    pInfo->vgVer = dbCache->vgInfo->vgVersion;
    pInfo->dbId = dbCache->dbId;
    pInfo->tbNum = dbCache->vgInfo->numOfTable;
  } else {
    pInfo->vgVer = CTG_DEFAULT_INVALID_VERSION;
  }

  CTG_ERR_JRET(ctgHandleTaskEnd(pTask, 0));

_return:

  if (dbCache) {
    ctgReleaseVgInfo(dbCache);
    ctgReleaseDBCache(pCtg, dbCache);
  }

  CTG_RET(code);
}

D
dapan1121 已提交
1191 1192 1193
int32_t ctgLaunchGetIndexTask(SCtgTask *pTask) {
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
1194
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
1195 1196 1197 1198 1199 1200 1201 1202 1203 1204
  SCtgIndexCtx* pCtx = (SCtgIndexCtx*)pTask->taskCtx;

  CTG_ERR_RET(ctgGetIndexInfoFromMnode(CTG_PARAMS_LIST(), pCtx->indexFName, NULL, pTask));

  return TSDB_CODE_SUCCESS;
}

int32_t ctgLaunchGetUdfTask(SCtgTask *pTask) {
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
1205
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
1206 1207 1208 1209 1210 1211 1212 1213 1214 1215
  SCtgUdfCtx* pCtx = (SCtgUdfCtx*)pTask->taskCtx;

  CTG_ERR_RET(ctgGetUdfInfoFromMnode(CTG_PARAMS_LIST(), pCtx->udfName, NULL, pTask));

  return TSDB_CODE_SUCCESS;
}

int32_t ctgLaunchGetUserTask(SCtgTask *pTask) {
  SCatalog* pCtg = pTask->pJob->pCtg; 
  void *pTrans = pTask->pJob->pTrans; 
D
dapan1121 已提交
1216
  const SEpSet* pMgmtEps = &pTask->pJob->pMgmtEps;
D
dapan1121 已提交
1217 1218 1219 1220 1221 1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246
  SCtgUserCtx* pCtx = (SCtgUserCtx*)pTask->taskCtx;
  bool inCache = false;
  bool pass = false;
  
  CTG_ERR_RET(ctgChkAuthFromCache(pCtg, pCtx->user.user, pCtx->user.dbFName, pCtx->user.type, &inCache, &pass));
  if (inCache) {
    pTask->res = taosMemoryCalloc(1, sizeof(bool));
    if (NULL == pTask->res) {
      CTG_ERR_RET(TSDB_CODE_OUT_OF_MEMORY);
    }
    *(bool*)pTask->res = pass;
    
    CTG_ERR_RET(ctgHandleTaskEnd(pTask, 0));
    return TSDB_CODE_SUCCESS;
  }

  CTG_ERR_RET(ctgGetUserDbAuthFromMnode(CTG_PARAMS_LIST(), pCtx->user.user, NULL, pTask));

  return TSDB_CODE_SUCCESS;
}

int32_t ctgRelaunchGetTbMetaTask(SCtgTask *pTask) {
  ctgResetTbMetaTask(pTask);

  CTG_ERR_RET(ctgLaunchGetTbMetaTask(pTask));

  return TSDB_CODE_SUCCESS;
}

SCtgAsyncFps gCtgAsyncFps[] = {
D
dapan1121 已提交
1247 1248 1249 1250 1251 1252 1253 1254 1255 1256
  {ctgLaunchGetQnodeTask,   ctgHandleGetQnodeRsp,   ctgDumpQnodeRes},
  {ctgLaunchGetDbVgTask,    ctgHandleGetDbVgRsp,    ctgDumpDbVgRes},
  {ctgLaunchGetDbCfgTask,   ctgHandleGetDbCfgRsp,   ctgDumpDbCfgRes},
  {ctgLaunchGetDbInfoTask,  ctgHandleGetDbInfoRsp,  ctgDumpDbInfoRes},
  {ctgLaunchGetTbMetaTask,  ctgHandleGetTbMetaRsp,  ctgDumpTbMetaRes},
  {ctgLaunchGetTbHashTask,  ctgHandleGetTbHashRsp,  ctgDumpTbHashRes},
  {ctgLaunchGetTbIndexTask, ctgHandleGetTbIndexRsp, ctgDumpTbIndexRes},
  {ctgLaunchGetIndexTask,   ctgHandleGetIndexRsp,   ctgDumpIndexRes},
  {ctgLaunchGetUdfTask,     ctgHandleGetUdfRsp,     ctgDumpUdfRes},
  {ctgLaunchGetUserTask,    ctgHandleGetUserRsp,    ctgDumpUserRes},
D
dapan1121 已提交
1257 1258 1259 1260 1261 1262 1263 1264 1265 1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277
};

int32_t ctgMakeAsyncRes(SCtgJob *pJob) {
  int32_t code = 0;
  int32_t taskNum = taosArrayGetSize(pJob->pTasks);
  
  for (int32_t i = 0; i < taskNum; ++i) {
    SCtgTask *pTask = taosArrayGet(pJob->pTasks, i);
    CTG_ERR_RET((*gCtgAsyncFps[pTask->type].dumpResFp)(pTask));
  }

  return TSDB_CODE_SUCCESS;
}


int32_t ctgLaunchJob(SCtgJob *pJob) {
  int32_t taskNum = taosArrayGetSize(pJob->pTasks);
  
  for (int32_t i = 0; i < taskNum; ++i) {
    SCtgTask *pTask = taosArrayGet(pJob->pTasks, i);

1278
    qDebug("QID:0x%" PRIx64 " start to launch task %d", pJob->queryId, pTask->taskId);
D
dapan1121 已提交
1279 1280 1281 1282 1283 1284 1285 1286
    CTG_ERR_RET((*gCtgAsyncFps[pTask->type].launchFp)(pTask));
  }

  return TSDB_CODE_SUCCESS;
}