clientImpl.c 33.2 KB
Newer Older
1

2
#include "../../libs/scheduler/inc/schedulerInt.h"
3
#include "clientInt.h"
4
#include "clientLog.h"
5 6 7
#include "parser.h"
#include "planner.h"
#include "scheduler.h"
8 9
#include "tdef.h"
#include "tep.h"
10
#include "tglobal.h"
11
#include "tmsgtype.h"
12 13
#include "tnote.h"
#include "tpagedfile.h"
14
#include "tref.h"
X
Xiaoyu Wang 已提交
15

16
#define CHECK_CODE_GOTO(expr, label) \
H
Haojun Liao 已提交
17 18
  do {                               \
    int32_t code = expr;             \
X
Xiaoyu Wang 已提交
19
    if (TSDB_CODE_SUCCESS != code) { \
H
Haojun Liao 已提交
20
      terrno = code;                 \
21
      goto label;                    \
H
Haojun Liao 已提交
22
    }                                \
X
Xiaoyu Wang 已提交
23
  } while (0)
24

25
static int32_t initEpSetFromCfg(const char *firstEp, const char *secondEp, SCorEpSet *pEpSet);
26 27
static SMsgSendInfo* buildConnectMsg(SRequestObj *pRequest);
static void destroySendMsgInfo(SMsgSendInfo* pMsgBody);
28
static void setQueryResultFromRsp(SReqResultInfo* pResultInfo, const SRetrieveTableRsp* pRsp);
H
Haojun Liao 已提交
29

30
static bool stringLengthCheck(const char* str, size_t maxsize) {
31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47
  if (str == NULL) {
    return false;
  }

  size_t len = strlen(str);
  if (len <= 0 || len > maxsize) {
    return false;
  }

  return true;
}

static bool validateUserName(const char* user) {
  return stringLengthCheck(user, TSDB_USER_LEN - 1);
}

static bool validatePassword(const char* passwd) {
48
  return stringLengthCheck(passwd, TSDB_PASSWORD_LEN - 1);
49 50 51 52 53 54
}

static bool validateDbName(const char* db) {
  return stringLengthCheck(db, TSDB_DB_NAME_LEN - 1);
}

55 56 57 58 59 60
static char* getClusterKey(const char* user, const char* auth, const char* ip, int32_t port) {
  char key[512] = {0};
  snprintf(key, sizeof(key), "%s:%s:%s:%d", user, auth, ip, port);
  return strdup(key);
}

H
Haojun Liao 已提交
61
static STscObj* taosConnectImpl(const char *user, const char *auth, const char *db, __taos_async_fn_t fp, void *param, SAppInstInfo* pAppInfo);
62
static void setResSchemaInfo(SReqResultInfo* pResInfo, const SSchema* pSchema, int32_t numOfCols);
63 64

TAOS *taos_connect_internal(const char *ip, const char *user, const char *pass, const char *auth, const char *db, uint16_t port) {
65 66 67 68
  if (taos_init() != TSDB_CODE_SUCCESS) {
    return NULL;
  }

69 70 71 72 73
  if (!validateUserName(user)) {
    terrno = TSDB_CODE_TSC_INVALID_USER_LENGTH;
    return NULL;
  }

74
  char localDb[TSDB_DB_NAME_LEN] = {0};
75 76 77 78 79 80
  if (db != NULL) {
    if(!validateDbName(db)) {
      terrno = TSDB_CODE_TSC_INVALID_DB_LENGTH;
      return NULL;
    }

81 82
    tstrncpy(localDb, db, sizeof(localDb));
    strdequote(localDb);
83 84
  }

85
  char secretEncrypt[TSDB_PASSWORD_LEN + 1] = {0};
86 87 88 89 90 91
  if (auth == NULL) {
    if (!validatePassword(pass)) {
      terrno = TSDB_CODE_TSC_INVALID_PASS_LENGTH;
      return NULL;
    }

92
    taosEncryptPass_c((uint8_t *)pass, strlen(pass), secretEncrypt);
93 94 95 96
  } else {
    tstrncpy(secretEncrypt, auth, tListLen(secretEncrypt));
  }

97
  SCorEpSet epSet = {0};
98 99 100 101 102 103 104 105 106 107 108 109 110 111
  if (ip) {
    if (initEpSetFromCfg(ip, NULL, &epSet) < 0) {
      return NULL;
    }

    if (port) {
      epSet.epSet.port[0] = port;
    }
  } else {
    if (initEpSetFromCfg(tsFirst, tsSecond, &epSet) < 0) {
      return NULL;
    }
  }

112
  char* key = getClusterKey(user, secretEncrypt, ip, port);
H
Haojun Liao 已提交
113
  SAppInstInfo** pInst = NULL;
114

H
Haojun Liao 已提交
115 116 117
  pthread_mutex_lock(&appInfo.mutex);

  pInst = taosHashGet(appInfo.pInstMap, key, strlen(key));
118
  if (pInst == NULL) {
H
Haojun Liao 已提交
119 120
    SAppInstInfo* p = calloc(1, sizeof(struct SAppInstInfo));
    p->mgmtEp       = epSet;
H
Haojun Liao 已提交
121
    p->pTransporter = openTransporter(user, secretEncrypt, tsNumOfCores);
L
Liu Jicong 已提交
122
    p->pAppHbMgr = appHbMgrInit(p);
H
Haojun Liao 已提交
123
    taosHashPut(appInfo.pInstMap, key, strlen(key), &p, POINTER_BYTES);
124

H
Haojun Liao 已提交
125
    pInst = &p;
126 127
  }

H
Haojun Liao 已提交
128 129
  pthread_mutex_unlock(&appInfo.mutex);

H
Haojun Liao 已提交
130
  tfree(key);
H
Haojun Liao 已提交
131
  return taosConnectImpl(user, &secretEncrypt[0], localDb, NULL, NULL, *pInst);
132 133
}

X
Xiaoyu Wang 已提交
134 135 136
int32_t buildRequest(STscObj *pTscObj, const char *sql, int sqlLen, SRequestObj** pRequest) {
  *pRequest = createRequest(pTscObj, NULL, NULL, TSDB_SQL_SELECT);
  if (*pRequest == NULL) {
137
    tscError("failed to malloc sqlObj");
X
Xiaoyu Wang 已提交
138
    return TSDB_CODE_TSC_OUT_OF_MEMORY;
139 140
  }

X
Xiaoyu Wang 已提交
141 142 143 144 145
  (*pRequest)->sqlstr = malloc(sqlLen + 1);
  if ((*pRequest)->sqlstr == NULL) {
    tscError("0x%"PRIx64" failed to prepare sql string buffer", (*pRequest)->self);
    (*pRequest)->msgBuf = strdup("failed to prepare sql string buffer");
    return TSDB_CODE_TSC_OUT_OF_MEMORY;
146 147
  }

X
Xiaoyu Wang 已提交
148 149 150
  strntolower((*pRequest)->sqlstr, sql, (int32_t)sqlLen);
  (*pRequest)->sqlstr[sqlLen] = 0;
  (*pRequest)->sqlLen = sqlLen;
151

H
Haojun Liao 已提交
152
  tscDebugL("0x%"PRIx64" SQL: %s, reqId:0x%"PRIx64, (*pRequest)->self, (*pRequest)->sqlstr, (*pRequest)->requestId);
X
Xiaoyu Wang 已提交
153 154
  return TSDB_CODE_SUCCESS;
}
155

X
Xiaoyu Wang 已提交
156
int32_t parseSql(SRequestObj* pRequest, SQueryNode** pQuery) {
H
Haojun Liao 已提交
157 158
  STscObj* pTscObj = pRequest->pTscObj;

X
Xiaoyu Wang 已提交
159
  SParseContext cxt = {
H
Haojun Liao 已提交
160
    .requestId = pRequest->requestId,
H
Haojun Liao 已提交
161
    .acctId    = pTscObj->acctId,
H
Haojun Liao 已提交
162
    .db        = getDbOfConnection(pTscObj),
H
Haojun Liao 已提交
163 164 165 166
    .pSql      = pRequest->sqlstr,
    .sqlLen    = pRequest->sqlLen,
    .pMsg      = pRequest->msgBuf,
    .msgLen    = ERROR_MSG_BUF_DEFAULT_SIZE,
H
Haojun Liao 已提交
167
    .pTransporter = pTscObj->pAppInfo->pTransporter,
X
Xiaoyu Wang 已提交
168
  };
H
Haojun Liao 已提交
169

H
Haojun Liao 已提交
170 171
  cxt.mgmtEpSet = getEpSet_s(&pTscObj->pAppInfo->mgmtEp);
  int32_t code = catalogGetHandle(pTscObj->pAppInfo->clusterId, &cxt.pCatalog);
H
Haojun Liao 已提交
172
  if (code != TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
173
    tfree(cxt.db);
H
Haojun Liao 已提交
174 175 176 177
    return code;
  }

  code = qParseQuerySql(&cxt, pQuery);
178

H
Haojun Liao 已提交
179
  tfree(cxt.db);
X
Xiaoyu Wang 已提交
180 181
  return code;
}
H
Haojun Liao 已提交
182

X
Xiaoyu Wang 已提交
183 184 185
int32_t execDdlQuery(SRequestObj* pRequest, SQueryNode* pQuery) {
  SDclStmtInfo* pDcl = (SDclStmtInfo*)pQuery;
  pRequest->type = pDcl->msgType;
H
Haojun Liao 已提交
186
  pRequest->body.requestMsg = (SDataBuf){.pData = pDcl->pMsg, .len = pDcl->msgLen};
X
Xiaoyu Wang 已提交
187

H
Haojun Liao 已提交
188
  STscObj* pTscObj = pRequest->pTscObj;
189
  SMsgSendInfo* pSendMsg = buildMsgInfoImpl(pRequest);
X
Xiaoyu Wang 已提交
190

191
  int64_t transporterId = 0;
192 193 194 195
  if (pDcl->msgType == TDMT_VND_CREATE_TABLE || pDcl->msgType == TDMT_VND_SHOW_TABLES) {
    if (pDcl->msgType == TDMT_VND_SHOW_TABLES) {
      SShowReqInfo* pShowReqInfo = &pRequest->body.showInfo;
      if (pShowReqInfo->pArray == NULL) {
H
Haojun Liao 已提交
196
        pShowReqInfo->currentIndex = 0;  // set the first vnode/ then iterate the next vnode
197 198 199
        pShowReqInfo->pArray = pDcl->pExtension;
      }
    }
H
Haojun Liao 已提交
200
    asyncSendMsgToServer(pTscObj->pAppInfo->pTransporter, &pDcl->epSet, &transporterId, pSendMsg);
X
Xiaoyu Wang 已提交
201
  } else {
H
Haojun Liao 已提交
202
    asyncSendMsgToServer(pTscObj->pAppInfo->pTransporter, &pDcl->epSet, &transporterId, pSendMsg);
203 204
  }

X
Xiaoyu Wang 已提交
205 206 207
  tsem_wait(&pRequest->body.rspSem);
  return TSDB_CODE_SUCCESS;
}
X
Xiaoyu Wang 已提交
208

209
int32_t getPlan(SRequestObj* pRequest, SQueryNode* pQueryNode, SQueryDag** pDag, SArray* pNodeList) {
210
  pRequest->type = pQueryNode->type;
211

212 213 214
  SSchema* pSchema = NULL;
  int32_t  numOfCols = 0;
  int32_t code = qCreateQueryDag(pQueryNode, pDag, &pSchema, &numOfCols, pNodeList, pRequest->requestId);
215 216 217 218 219
  if (code != 0) {
    return code;
  }

  if (pQueryNode->type == TSDB_SQL_SELECT) {
220
    setResSchemaInfo(&pRequest->body.resInfo, pSchema, numOfCols);
221
    tfree(pSchema);
222
    pRequest->type = TDMT_VND_QUERY;
D
dapan1121 已提交
223 224
  } else {
    tfree(pSchema);
225 226 227
  }

  return code;
X
Xiaoyu Wang 已提交
228 229
}

230 231
void setResSchemaInfo(SReqResultInfo* pResInfo, const SSchema* pSchema, int32_t numOfCols) {
  assert(pSchema != NULL && numOfCols > 0);
232

233 234
  pResInfo->numOfCols = numOfCols;
  pResInfo->fields = calloc(numOfCols, sizeof(pSchema[0]));
235 236

  for (int32_t i = 0; i < pResInfo->numOfCols; ++i) {
237 238 239
    pResInfo->fields[i].bytes = pSchema[i].bytes;
    pResInfo->fields[i].type  = pSchema[i].type;
    tstrncpy(pResInfo->fields[i].name, pSchema[i].name, tListLen(pResInfo->fields[i].name));
240
  }
X
Xiaoyu Wang 已提交
241 242
}

243
int32_t scheduleQuery(SRequestObj* pRequest, SQueryDag* pDag, SArray* pNodeList) {
244
  if (TSDB_SQL_INSERT == pRequest->type || TSDB_SQL_CREATE_TABLE == pRequest->type) {
X
Xiaoyu Wang 已提交
245
    SQueryResult res = {.code = 0, .numOfRows = 0, .msgSize = ERROR_MSG_BUF_DEFAULT_SIZE, .msg = pRequest->msgBuf};
D
dapan1121 已提交
246
    int32_t code = schedulerExecJob(pRequest->pTscObj->pAppInfo->pTransporter, NULL, pDag, &pRequest->body.pQueryJob, &res);
247 248 249
    if (code != TSDB_CODE_SUCCESS) {
      // handle error and retry
    } else {
250
      if (pRequest->body.pQueryJob != NULL) {
D
dapan1121 已提交
251
        schedulerFreeJob(pRequest->body.pQueryJob);
252 253 254
      }
    }

H
Haojun Liao 已提交
255 256 257
    pRequest->body.resInfo.numOfRows = res.numOfRows;
    pRequest->code = res.code;
    return pRequest->code;
X
Xiaoyu Wang 已提交
258
  }
259

260
  return schedulerAsyncExecJob(pRequest->pTscObj->pAppInfo->pTransporter, pNodeList, pDag, &pRequest->body.pQueryJob);
X
Xiaoyu Wang 已提交
261
}
X
Xiaoyu Wang 已提交
262

L
Liu Jicong 已提交
263

L
Liu Jicong 已提交
264
typedef struct SMqClientVg {
L
Liu Jicong 已提交
265
  // statistics
L
Liu Jicong 已提交
266
  int64_t pollCnt;
L
Liu Jicong 已提交
267 268 269 270 271 272
  // offset
  int64_t committedOffset;
  int64_t currentOffset;
  //connection info
  int32_t vgId;
  SEpSet  epSet;
L
Liu Jicong 已提交
273 274 275 276 277 278 279 280 281 282
} SMqClientVg;

typedef struct SMqClientTopic {
  // subscribe info
  int32_t sqlLen;
  char*   sql;
  char*   topicName;
  int64_t topicId;
  int32_t nextVgIdx;
  SArray* vgs;    //SArray<SMqClientVg>
L
Liu Jicong 已提交
283 284 285 286 287 288
} SMqClientTopic;

typedef struct tmq_resp_err_t {
  int32_t code;
} tmq_resp_err_t;

L
Liu Jicong 已提交
289 290
typedef struct tmq_topic_vgroup_t {
  char*   topic;
L
Liu Jicong 已提交
291
  int32_t vgId;
L
Liu Jicong 已提交
292 293 294 295 296 297 298
  int64_t commitOffset;
} tmq_topic_vgroup_t;

typedef struct tmq_topic_vgroup_list_t {
  int32_t cnt;
  int32_t size;
  tmq_topic_vgroup_t* elems;
L
Liu Jicong 已提交
299 300 301 302
} tmq_topic_vgroup_list_t;

typedef void (tmq_commit_cb(tmq_t*, tmq_resp_err_t, tmq_topic_vgroup_list_t*, void* param));

L
Liu Jicong 已提交
303 304 305
struct tmq_conf_t {
  char           clientId[256];
  char           groupId[256];
L
Liu Jicong 已提交
306 307 308
  char*          ip;
  uint16_t       port;
  tmq_commit_cb* commit_cb;
L
Liu Jicong 已提交
309 310 311 312 313 314 315 316
};

tmq_conf_t* tmq_conf_new() {
  tmq_conf_t* conf = calloc(1, sizeof(tmq_conf_t));
  return conf;
}

int32_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
L
Liu Jicong 已提交
317
  if (strcmp(key, "group.id") == 0) {
L
Liu Jicong 已提交
318 319
    strcpy(conf->groupId, value);
  }
L
Liu Jicong 已提交
320
  if (strcmp(key, "client.id") == 0) {
L
Liu Jicong 已提交
321 322 323 324
    strcpy(conf->clientId, value);
  }
  return 0;
}
L
Liu Jicong 已提交
325 326 327 328

struct tmq_t {
  char           groupId[256];
  char           clientId[256];
L
Liu Jicong 已提交
329 330
  int64_t        consumerId;
  int64_t        status;
L
Liu Jicong 已提交
331 332
  STscObj*       pTscObj;
  tmq_commit_cb* commit_cb;
L
Liu Jicong 已提交
333 334 335 336 337 338 339 340 341 342 343 344 345 346 347
  int32_t        nextTopicIdx;
  SArray*        clientTopics;  //SArray<SMqClientTopic>
};

tmq_t* taos_consumer_new(void* conn, tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
  tmq_t* pTmq = calloc(sizeof(tmq_t), 1);
  if (pTmq == NULL) {
    return NULL;
  }
  pTmq->pTscObj = (STscObj*)conn;
  pTmq->status = 0;
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
  pTmq->commit_cb = conf->commit_cb;
  pTmq->consumerId = generateRequestId() & ((uint64_t)-1 >> 1);
L
Liu Jicong 已提交
348
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
L
Liu Jicong 已提交
349 350 351 352 353 354 355
  return pTmq;
}

struct tmq_list_t {
  int32_t cnt;
  int32_t tot;
  char*   elems[];
L
Liu Jicong 已提交
356
};
L
Liu Jicong 已提交
357 358 359 360 361 362 363 364 365 366 367 368
tmq_list_t* tmq_list_new() {
  tmq_list_t *ptr = malloc(sizeof(tmq_list_t) + 8 * sizeof(char*));
  if (ptr == NULL) {
    return ptr;
  }
  ptr->cnt = 0;
  ptr->tot = 8;
  return ptr;
}

int32_t tmq_list_append(tmq_list_t* ptr, char* src) {
  if (ptr->cnt >= ptr->tot-1) return -1;
L
Liu Jicong 已提交
369
  ptr->elems[ptr->cnt] = strdup(src);
L
Liu Jicong 已提交
370 371 372 373 374 375 376 377 378 379 380
  ptr->cnt++;
  return 0;
}


TAOS_RES* tmq_subscribe(tmq_t* tmq, tmq_list_t* topic_list) {
  SRequestObj *pRequest = NULL;
  tmq->status = 1;
  int32_t sz = topic_list->cnt;
  tmq->clientTopics = taosArrayInit(sz, sizeof(void*));
  for (int i = 0; i < sz; i++) {
L
Liu Jicong 已提交
381 382 383
    char* topicName = topic_list->elems[i];

    SName name = {0};
L
Liu Jicong 已提交
384
    char* dbName = getDbOfConnection(tmq->pTscObj);
L
Liu Jicong 已提交
385 386 387 388 389 390 391 392 393 394 395 396 397
    tNameSetDbName(&name, tmq->pTscObj->acctId, dbName, strlen(dbName));     
    tNameFromString(&name, topicName, T_NAME_TABLE);

    char* topicFname = calloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFname == NULL) {

    }
    tNameExtractFullName(&name, topicFname);
    tscDebug("subscribe topic: %s", topicFname);
    taosArrayPush(tmq->clientTopics, &topicFname); 
    /*SMqClientTopic topic = {*/
      /*.*/
    /*};*/
L
Liu Jicong 已提交
398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414
  }
  SCMSubscribeReq req;
  req.topicNum = taosArrayGetSize(tmq->clientTopics);
  req.consumerId = tmq->consumerId;
  req.consumerGroup = strdup(tmq->groupId);
  req.topicNames = tmq->clientTopics;

  int tlen = tSerializeSCMSubscribeReq(NULL, &req);
  void* buf = malloc(tlen);
  if(buf == NULL) {
    goto _return;
  }

  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);
  /*printf("formatted: %s\n", dagStr);*/

L
Liu Jicong 已提交
415
  pRequest = createRequest(tmq->pTscObj, NULL, NULL, TDMT_MND_SUBSCRIBE);
L
Liu Jicong 已提交
416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440
  if (pRequest == NULL) {
    tscError("failed to malloc sqlObj");
  }

  pRequest->body.requestMsg = (SDataBuf){ .pData = buf, .len = tlen };

  SMsgSendInfo* body = buildMsgInfoImpl(pRequest);
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, body);

  tsem_wait(&pRequest->body.rspSem);

_return:
  if (body != NULL) {
    destroySendMsgInfo(body);
  }

  if (pRequest != NULL && terrno != TSDB_CODE_SUCCESS) {
    pRequest->code = terrno;
  }

  return pRequest;
}
L
Liu Jicong 已提交
441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457

void tmq_conf_set_offset_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb) {
  conf->commit_cb = cb;
}

SArray* tmqGetConnInfo(SClientHbKey connKey, void* param) {
  tmq_t* pTmq = (void*)param;
  SArray* pArray = taosArrayInit(0, sizeof(SKv));
  if (pArray == NULL) {
    return NULL;
  }
  SKv kv = {0};
  kv.key = malloc(256);
  if (kv.key == NULL) {
    taosArrayDestroy(pArray);
    return NULL;
  }
L
Liu Jicong 已提交
458 459 460 461 462
  strcpy(kv.key, "mq-tmp");
  kv.keyLen = strlen("mq-tmp") + 1;
  SMqHbMsg* pMqHb = malloc(sizeof(SMqHbMsg));
  if (pMqHb == NULL) {
    return pArray;
L
Liu Jicong 已提交
463
  }
L
Liu Jicong 已提交
464 465 466 467 468
  pMqHb->consumerId = connKey.connId;
  SArray* clientTopics = pTmq->clientTopics;
  int sz = taosArrayGetSize(clientTopics);
  for (int i = 0; i < sz; i++) {
    SMqClientTopic* pCTopic = taosArrayGet(clientTopics, i);
L
Liu Jicong 已提交
469 470 471 472
    /*if (pCTopic->vgId == -1) {*/
      /*pMqHb->status = 1;*/
      /*break;*/
    /*}*/
L
Liu Jicong 已提交
473 474 475
  }
  kv.value = pMqHb;
  kv.valueLen = sizeof(SMqHbMsg);
L
Liu Jicong 已提交
476
  taosArrayPush(pArray, &kv);
L
Liu Jicong 已提交
477 478

  return pArray;
L
Liu Jicong 已提交
479 480 481 482 483 484 485 486 487 488 489 490 491 492 493
}

tmq_t* tmqCreateConsumerImpl(TAOS* conn, tmq_conf_t* conf) {
  tmq_t* pTmq = malloc(sizeof(tmq_t));
  if (pTmq == NULL) {
    return NULL;
  }
  strcpy(pTmq->groupId, conf->groupId);
  strcpy(pTmq->clientId, conf->clientId);
  pTmq->pTscObj = (STscObj*)conn;
  pTmq->pTscObj->connType = HEARTBEAT_TYPE_MQ;

  return pTmq;
}

L
Liu Jicong 已提交
494
TAOS_RES *taos_create_topic(TAOS* taos, const char* topicName, const char* sql, int sqlLen) {
495 496 497 498
  STscObj     *pTscObj = (STscObj*)taos;
  SRequestObj *pRequest = NULL;
  SQueryNode  *pQueryNode = NULL;
  char        *pStr = NULL;
L
Liu Jicong 已提交
499 500

  terrno = TSDB_CODE_SUCCESS;
501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519
  if (taos == NULL || topicName == NULL || sql == NULL) {
    tscError("invalid parameters for creating topic, connObj:%p, topic name:%s, sql:%s", taos, topicName, sql);
    terrno = TSDB_CODE_TSC_INVALID_INPUT;
    goto _return;
  }

  if (strlen(topicName) >= TSDB_TOPIC_NAME_LEN) {
    tscError("topic name too long, max length:%d", TSDB_TOPIC_NAME_LEN - 1);
    terrno = TSDB_CODE_TSC_INVALID_INPUT;
    goto _return;
  }

  if (sqlLen > tsMaxSQLStringLen) {
    tscError("sql string exceeds max length:%d", tsMaxSQLStringLen);
    terrno = TSDB_CODE_TSC_EXCEED_SQL_LIMIT;
    goto _return;
  }

  tscDebug("start to create topic, %s", topicName);
L
Liu Jicong 已提交
520 521

  CHECK_CODE_GOTO(buildRequest(pTscObj, sql, sqlLen, &pRequest), _return);
522
  CHECK_CODE_GOTO(parseSql(pRequest, &pQueryNode), _return);
L
Liu Jicong 已提交
523

H
Haojun Liao 已提交
524 525 526
  SQueryStmtInfo* pQueryStmtInfo = (SQueryStmtInfo* ) pQueryNode;
  pQueryStmtInfo->info.continueQuery = true;

527
  // todo check for invalid sql statement and return with error code
L
Liu Jicong 已提交
528

529 530 531
  SSchema *schema = NULL;
  int32_t  numOfCols = 0;
  CHECK_CODE_GOTO(qCreateQueryDag(pQueryNode, &pRequest->body.pDag, &schema, &numOfCols, NULL, pRequest->requestId), _return);
L
Liu Jicong 已提交
532

533 534 535
  pStr = qDagToString(pRequest->body.pDag);
  if(pStr == NULL) {
    goto _return;
L
Liu Jicong 已提交
536
  }
537

H
Haojun Liao 已提交
538 539
  printf("%s\n", pStr);

540 541 542 543 544 545 546 547 548 549
  // The topic should be related to a database that the queried table is belonged to.
  SName name = {0};
  char dbName[TSDB_DB_FNAME_LEN] = {0};
  tNameGetFullDbName(&((SQueryStmtInfo*) pQueryNode)->pTableMetaInfo[0]->name, dbName);

  tNameFromString(&name, dbName, T_NAME_ACCT|T_NAME_DB);
  tNameFromString(&name, topicName, T_NAME_TABLE);

  char topicFname[TSDB_TOPIC_FNAME_LEN] = {0};
  tNameExtractFullName(&name, topicFname);
L
Liu Jicong 已提交
550 551

  SCMCreateTopicReq req = {
552
    .name         = (char*) topicFname,
L
Liu Jicong 已提交
553
    .igExists     = 1,
554 555 556
    .physicalPlan = (char*) pStr,
    .sql          = (char*) sql,
    .logicalPlan  = "no logic plan",
L
Liu Jicong 已提交
557 558
  };

L
Liu Jicong 已提交
559 560 561 562 563
  int tlen = tSerializeSCMCreateTopicReq(NULL, &req);
  void* buf = malloc(tlen);
  if(buf == NULL) {
    goto _return;
  }
564

L
Liu Jicong 已提交
565 566 567
  void* abuf = buf;
  tSerializeSCMCreateTopicReq(&abuf, &req);
  /*printf("formatted: %s\n", dagStr);*/
L
Liu Jicong 已提交
568 569

  pRequest->body.requestMsg = (SDataBuf){ .pData = buf, .len = tlen };
570
  pRequest->type = TDMT_MND_CREATE_TOPIC;
L
Liu Jicong 已提交
571

572
  SMsgSendInfo* body = buildMsgInfoImpl(pRequest);
573
  SEpSet epSet = getEpSet_s(&pTscObj->pAppInfo->mgmtEp);
L
Liu Jicong 已提交
574 575

  int64_t transporterId = 0;
576
  asyncSendMsgToServer(pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, body);
L
Liu Jicong 已提交
577 578 579 580

  tsem_wait(&pRequest->body.rspSem);

_return:
581 582 583 584 585
  qDestroyQuery(pQueryNode);
  if (body != NULL) {
    destroySendMsgInfo(body);
  }

L
Liu Jicong 已提交
586 587 588
  if (pRequest != NULL && terrno != TSDB_CODE_SUCCESS) {
    pRequest->code = terrno;
  }
589

L
Liu Jicong 已提交
590 591 592
  return pRequest;
}

L
Liu Jicong 已提交
593
/*typedef SMqConsumeRsp tmq_message_t;*/
L
Liu Jicong 已提交
594

L
Liu Jicong 已提交
595 596 597 598
struct tmq_message_t {
  SMqConsumeRsp rsp;
};

L
Liu Jicong 已提交
599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629
int32_t tmq_poll_cb_inner(void* param, const SDataBuf* pMsg, int32_t code) {
  return 0;
}

int32_t tmq_ask_ep_cb(void* param, const SDataBuf* pMsg, int32_t code) {
  tmq_t* tmq = (tmq_t*)param;
  SMqCMGetSubEpRsp rsp;
  tDecodeSMqCMGetSubEpRsp(pMsg->pData, &rsp);
  int32_t sz = taosArrayGetSize(rsp.topics);
  // TODO: lock
  tmq->clientTopics = taosArrayInit(sz, sizeof(SMqClientTopic));
  for (int32_t i = 0; i < sz; i++) {
    SMqClientTopic topic = {0};
    SMqSubTopicEp* pTopicEp = taosArrayGet(rsp.topics, i);
    topic.topicName = strdup(pTopicEp->topic);
    int32_t vgSz = taosArrayGetSize(pTopicEp->vgs);
    topic.vgs = taosArrayInit(vgSz, sizeof(SMqClientVg));
    for (int32_t j = 0; j < vgSz; j++) {
      SMqSubVgEp* pVgEp = taosArrayGet(pTopicEp->vgs, j);
      SMqClientVg clientVg = {
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet
      };
      taosArrayPush(topic.vgs, &clientVg);
    }
    taosArrayPush(tmq->clientTopics, &topic);
  }
  // unlock
  return 0;
}

L
Liu Jicong 已提交
630 631 632 633 634 635 636 637 638
tmq_message_t* tmq_consume_poll(tmq_t* tmq, int64_t blocking_time) {
  if (tmq->clientTopics == NULL || taosArrayGetSize(tmq->clientTopics) == 0) {
    return NULL;
  }
  SRequestObj *pRequest = NULL;
  SMqConsumeReq req = {0};
  req.reqType = 1;
  req.blockingTime = blocking_time;
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
639
  tmq_message_t* tmq_message = NULL;
L
Liu Jicong 已提交
640 641
  strcpy(req.cgroup, tmq->groupId);

L
Liu Jicong 已提交
642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670
  if (taosArrayGetSize(tmq->clientTopics) == 0) {
    int32_t tlen = sizeof(SMqCMGetSubEpReq);
    SMqCMGetSubEpReq* buf = malloc(tlen);
    if (buf == NULL) {
      tscError("failed to malloc get subscribe ep buf");
    }
    buf->consumerId = htobe64(buf->consumerId);
    
    pRequest = createRequest(tmq->pTscObj, NULL, NULL, TDMT_MND_GET_SUB_EP);
    if (pRequest == NULL) {
      tscError("failed to malloc subscribe ep request");
    }

    pRequest->body.requestMsg = (SDataBuf){ .pData = buf, .len = tlen };

    SMsgSendInfo* sendInfo = buildMsgInfoImpl(pRequest);
    sendInfo->requestObjRefId = 0;
    sendInfo->param = tmq;
    sendInfo->fp = tmq_ask_ep_cb;

    SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

    int64_t transporterId = 0;
    asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);

    tsem_wait(&pRequest->body.rspSem);
  }

  SMqClientTopic* pTopic = taosArrayGetP(tmq->clientTopics, tmq->nextTopicIdx);
L
Liu Jicong 已提交
671 672 673 674 675 676 677 678 679 680
  tmq->nextTopicIdx = (tmq->nextTopicIdx + 1) % taosArrayGetSize(tmq->clientTopics);
  strcpy(req.topic, pTopic->topicName);
  int32_t nextVgIdx = pTopic->nextVgIdx;
  pTopic->nextVgIdx = (nextVgIdx + 1) % taosArrayGetSize(pTopic->vgs);
  SMqClientVg* pVg = taosArrayGet(pTopic->vgs, nextVgIdx);
  req.offset = pVg->currentOffset;

  pRequest->body.requestMsg = (SDataBuf){ .pData = &req, .len = sizeof(SMqConsumeReq) };
  pRequest->type = TDMT_VND_CONSUME;

L
Liu Jicong 已提交
681 682 683 684
  SMsgSendInfo* sendInfo = buildMsgInfoImpl(pRequest);
  sendInfo->requestObjRefId = 0;
  sendInfo->param = &tmq_message;
  sendInfo->fp = tmq_poll_cb_inner;
L
Liu Jicong 已提交
685 686

  int64_t transporterId = 0;
L
Liu Jicong 已提交
687
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);
L
Liu Jicong 已提交
688 689 690

  tsem_wait(&pRequest->body.rspSem);

L
Liu Jicong 已提交
691
  return tmq_message;
L
Liu Jicong 已提交
692 693 694 695 696 697 698 699 700 701 702 703

  /*tsem_wait(&pRequest->body.rspSem);*/

  /*if (body != NULL) {*/
    /*destroySendMsgInfo(body);*/
  /*}*/

  /*if (pRequest != NULL && terrno != TSDB_CODE_SUCCESS) {*/
    /*pRequest->code = terrno;*/
  /*}*/

  /*return pRequest;*/
L
Liu Jicong 已提交
704 705
}

L
Liu Jicong 已提交
706 707
tmq_resp_err_t* tmq_commit(tmq_t* tmq, tmq_topic_vgroup_list_t* tmq_topic_vgroup_list, int32_t async) {
  SMqConsumeReq req = {0};
L
Liu Jicong 已提交
708 709 710
  return NULL;
}

L
Liu Jicong 已提交
711 712
void tmq_message_destroy(tmq_message_t* tmq_message) {
  if (tmq_message == NULL) return;
L
Liu Jicong 已提交
713 714 715
}


X
Xiaoyu Wang 已提交
716 717 718 719 720 721
TAOS_RES *taos_query_l(TAOS *taos, const char *sql, int sqlLen) {
  STscObj *pTscObj = (STscObj *)taos;
  if (sqlLen > (size_t) tsMaxSQLStringLen) {
    tscError("sql string exceeds max length:%d", tsMaxSQLStringLen);
    terrno = TSDB_CODE_TSC_EXCEED_SQL_LIMIT;
    return NULL;
722 723
  }

X
Xiaoyu Wang 已提交
724 725
  nPrintTsc("%s", sql)

H
Haojun Liao 已提交
726
  SRequestObj *pRequest = NULL;
727
  SQueryNode  *pQueryNode = NULL;
728
  SArray      *pNodeList = taosArrayInit(4, sizeof(struct SQueryNodeAddr));
X
Xiaoyu Wang 已提交
729

X
Xiaoyu Wang 已提交
730
  terrno = TSDB_CODE_SUCCESS;
X
Xiaoyu Wang 已提交
731
  CHECK_CODE_GOTO(buildRequest(pTscObj, sql, sqlLen, &pRequest), _return);
732
  CHECK_CODE_GOTO(parseSql(pRequest, &pQueryNode), _return);
H
Haojun Liao 已提交
733

734 735
  if (qIsDdlQuery(pQueryNode)) {
    CHECK_CODE_GOTO(execDdlQuery(pRequest, pQueryNode), _return);
H
Haojun Liao 已提交
736
  } else {
D
dapan1121 已提交
737

738 739
    CHECK_CODE_GOTO(getPlan(pRequest, pQueryNode, &pRequest->body.pDag, pNodeList), _return);
    CHECK_CODE_GOTO(scheduleQuery(pRequest, pRequest->body.pDag, pNodeList), _return);
X
Xiaoyu Wang 已提交
740
    pRequest->code = terrno;
X
Xiaoyu Wang 已提交
741 742 743
  }

_return:
744
  taosArrayDestroy(pNodeList);
745
  qDestroyQuery(pQueryNode);
X
Xiaoyu Wang 已提交
746
  if (NULL != pRequest && TSDB_CODE_SUCCESS != terrno) {
X
Xiaoyu Wang 已提交
747 748
    pRequest->code = terrno;
  }
749

750 751 752
  return pRequest;
}

753
int initEpSetFromCfg(const char *firstEp, const char *secondEp, SCorEpSet *pEpSet) {
754 755
  pEpSet->version = 0;

H
Haojun Liao 已提交
756
  // init mnode ip set
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
  SEpSet *mgmtEpSet   = &(pEpSet->epSet);
  mgmtEpSet->numOfEps = 0;
  mgmtEpSet->inUse    = 0;

  if (firstEp && firstEp[0] != 0) {
    if (strlen(firstEp) >= TSDB_EP_LEN) {
      terrno = TSDB_CODE_TSC_INVALID_FQDN;
      return -1;
    }

    taosGetFqdnPortFromEp(firstEp, mgmtEpSet->fqdn[0], &(mgmtEpSet->port[0]));
    mgmtEpSet->numOfEps++;
  }

  if (secondEp && secondEp[0] != 0) {
    if (strlen(secondEp) >= TSDB_EP_LEN) {
      terrno = TSDB_CODE_TSC_INVALID_FQDN;
      return -1;
    }

    taosGetFqdnPortFromEp(secondEp, mgmtEpSet->fqdn[mgmtEpSet->numOfEps], &(mgmtEpSet->port[mgmtEpSet->numOfEps]));
    mgmtEpSet->numOfEps++;
  }

  if (mgmtEpSet->numOfEps == 0) {
    terrno = TSDB_CODE_TSC_INVALID_FQDN;
    return -1;
  }

  return 0;
}

H
Haojun Liao 已提交
789
STscObj* taosConnectImpl(const char *user, const char *auth, const char *db, __taos_async_fn_t fp, void *param, SAppInstInfo* pAppInfo) {
790
  STscObj *pTscObj = createTscObj(user, auth, db, pAppInfo);
791 792 793 794 795
  if (NULL == pTscObj) {
    terrno = TSDB_CODE_TSC_OUT_OF_MEMORY;
    return pTscObj;
  }

H
Hongze Cheng 已提交
796
  SRequestObj *pRequest = createRequest(pTscObj, fp, param, TDMT_MND_CONNECT);
797 798 799
  if (pRequest == NULL) {
    destroyTscObj(pTscObj);
    terrno = TSDB_CODE_TSC_OUT_OF_MEMORY;
800 801 802
    return NULL;
  }

803
  SMsgSendInfo* body = buildConnectMsg(pRequest);
804 805

  int64_t transporterId = 0;
H
Haojun Liao 已提交
806
  asyncSendMsgToServer(pTscObj->pAppInfo->pTransporter, &pTscObj->pAppInfo->mgmtEp.epSet, &transporterId, body);
807 808 809

  tsem_wait(&pRequest->body.rspSem);
  if (pRequest->code != TSDB_CODE_SUCCESS) {
S
Shengliang Guan 已提交
810
    const char *errorMsg = (pRequest->code == TSDB_CODE_RPC_FQDN_ERROR) ? taos_errstr(pRequest) : tstrerror(pRequest->code);
811 812 813 814 815 816
    printf("failed to connect to server, reason: %s\n\n", errorMsg);

    destroyRequest(pRequest);
    taos_close(pTscObj);
    pTscObj = NULL;
  } else {
H
Haojun Liao 已提交
817
    tscDebug("0x%"PRIx64" connection is opening, connId:%d, dnodeConn:%p, reqId:0x%"PRIx64, pTscObj->id, pTscObj->connId, pTscObj->pAppInfo->pTransporter, pRequest->requestId);
818 819 820 821 822 823
    destroyRequest(pRequest);
  }

  return pTscObj;
}

824 825 826 827 828 829 830
static SMsgSendInfo* buildConnectMsg(SRequestObj *pRequest) {
  SMsgSendInfo *pMsgSendInfo = calloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
    terrno = TSDB_CODE_TSC_OUT_OF_MEMORY;
    return NULL;
  }

H
Haojun Liao 已提交
831
  pMsgSendInfo->msgType         = TDMT_MND_CONNECT;
S
Shengliang Guan 已提交
832
  pMsgSendInfo->msgInfo.len     = sizeof(SConnectReq);
833 834
  pMsgSendInfo->requestObjRefId = pRequest->self;
  pMsgSendInfo->requestId       = pRequest->requestId;
D
catalog  
dapan1121 已提交
835
  pMsgSendInfo->fp              = handleRequestRspFp[TMSG_INDEX(pMsgSendInfo->msgType)];
836
  pMsgSendInfo->param           = pRequest;
837

S
Shengliang Guan 已提交
838
  SConnectReq *pConnect = calloc(1, sizeof(SConnectReq));
839
  if (pConnect == NULL) {
840
    tfree(pMsgSendInfo);
841
    terrno = TSDB_CODE_TSC_OUT_OF_MEMORY;
842
    return NULL;
843 844
  }

845 846
  STscObj *pObj = pRequest->pTscObj;

H
Haojun Liao 已提交
847
  char* db = getDbOfConnection(pObj);
848 849 850
  if (db != NULL) {
    tstrncpy(pConnect->db, db, sizeof(pConnect->db));
  }
851
  tfree(db);
852

853 854 855
  pConnect->pid = htonl(appInfo.pid);
  pConnect->startTime = htobe64(appInfo.startTime);
  tstrncpy(pConnect->app, appInfo.appName, tListLen(pConnect->app));
856

857 858
  pMsgSendInfo->msgInfo.pData = pConnect;
  return pMsgSendInfo;
859 860
}

861
static void destroySendMsgInfo(SMsgSendInfo* pMsgBody) {
H
Haojun Liao 已提交
862
  assert(pMsgBody != NULL);
863 864
  tfree(pMsgBody->msgInfo.pData);
  tfree(pMsgBody);
H
Haojun Liao 已提交
865 866
}

867
void processMsgFromServer(void* parent, SRpcMsg* pMsg, SEpSet* pEpSet) {
868 869
  SMsgSendInfo *pSendInfo = (SMsgSendInfo *) pMsg->ahandle;
  assert(pMsg->ahandle != NULL);
870

871
  if (pSendInfo->requestObjRefId != 0) {
H
Haojun Liao 已提交
872
    SRequestObj *pRequest = (SRequestObj *)taosAcquireRef(clientReqRefPool, pSendInfo->requestObjRefId);
873
    assert(pRequest->self == pSendInfo->requestObjRefId);
874

875 876
    pRequest->metric.rsp = taosGetTimestampMs();
    pRequest->code = pMsg->code;
877

878 879 880 881 882
    STscObj *pTscObj = pRequest->pTscObj;
    if (pEpSet) {
      if (!isEpsetEqual(&pTscObj->pAppInfo->mgmtEp.epSet, pEpSet)) {
        updateEpSet_s(&pTscObj->pAppInfo->mgmtEp, pEpSet);
      }
883 884
    }

885
    /*
886 887
   * There is not response callback function for submit response.
   * The actual inserted number of points is the first number.
888
     */
889
    int32_t elapsed = pRequest->metric.rsp - pRequest->metric.start;
890
    if (pMsg->code == TSDB_CODE_SUCCESS) {
H
Haojun Liao 已提交
891
      tscDebug("0x%" PRIx64 " message:%s, code:%s rspLen:%d, elapsed:%d ms, reqId:0x%"PRIx64, pRequest->self,
892
          TMSG_INFO(pMsg->msgType), tstrerror(pMsg->code), pMsg->contLen, elapsed, pRequest->requestId);
893
    } else {
H
Haojun Liao 已提交
894
      tscError("0x%" PRIx64 " SQL cmd:%s, code:%s rspLen:%d, elapsed time:%d ms, reqId:0x%"PRIx64, pRequest->self,
895
          TMSG_INFO(pMsg->msgType), tstrerror(pMsg->code), pMsg->contLen, elapsed, pRequest->requestId);
896
    }
897

H
Haojun Liao 已提交
898
    taosReleaseRef(clientReqRefPool, pSendInfo->requestObjRefId);
899 900
  }

901 902 903 904 905 906 907 908 909 910
  SDataBuf buf = {.len = pMsg->contLen, .pData = NULL};

  if (pMsg->contLen > 0) {
    buf.pData = calloc(1, pMsg->contLen);
    if (buf.pData == NULL) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      pMsg->code = TSDB_CODE_OUT_OF_MEMORY;
    } else {
      memcpy(buf.pData, pMsg->pCont, pMsg->contLen);
    }
911 912
  }

913
  pSendInfo->fp(pSendInfo->param, &buf, pMsg->code);
914
  rpcFreeCont(pMsg->pCont);
915
  destroySendMsgInfo(pSendInfo);
916
}
917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937

TAOS *taos_connect_auth(const char *ip, const char *user, const char *auth, const char *db, uint16_t port) {
  tscDebug("try to connect to %s:%u by auth, user:%s db:%s", ip, port, user, db);
  if (user == NULL) {
    user = TSDB_DEFAULT_USER;
  }

  if (auth == NULL) {
    tscError("No auth info is given, failed to connect to server");
    return NULL;
  }

  return taos_connect_internal(ip, user, NULL, auth, db, port);
}

TAOS *taos_connect_l(const char *ip, int ipLen, const char *user, int userLen, const char *pass, int passLen, const char *db, int dbLen, uint16_t port) {
  char ipStr[TSDB_EP_LEN]      = {0};
  char dbStr[TSDB_DB_NAME_LEN] = {0};
  char userStr[TSDB_USER_LEN]  = {0};
  char passStr[TSDB_PASSWORD_LEN]   = {0};

dengyihao's avatar
dengyihao 已提交
938 939 940 941
  strncpy(ipStr,   ip,   TMIN(TSDB_EP_LEN - 1, ipLen));
  strncpy(userStr, user, TMIN(TSDB_USER_LEN - 1, userLen));
  strncpy(passStr, pass, TMIN(TSDB_PASSWORD_LEN - 1, passLen));
  strncpy(dbStr,   db,   TMIN(TSDB_DB_NAME_LEN - 1, dbLen));
942
  return taos_connect(ipStr, userStr, passStr, dbStr, port);
943 944 945 946
}

void* doFetchRow(SRequestObj* pRequest) {
  assert(pRequest != NULL);
947
  SReqResultInfo* pResultInfo = &pRequest->body.resInfo;
948

H
Haojun Liao 已提交
949 950
  SEpSet epSet = {0};

H
Haojun Liao 已提交
951
  if (pResultInfo->pData == NULL || pResultInfo->current >= pResultInfo->numOfRows) {
952
    if (pRequest->type == TDMT_VND_QUERY) {
953 954 955 956 957
      // All data has returned to App already, no need to try again
      if (pResultInfo->completed) {
        return NULL;
      }

958
      SReqResultInfo* pResInfo = &pRequest->body.resInfo;
959
      int32_t code = schedulerFetchRows(pRequest->body.pQueryJob, (void **)&pResInfo->pData);
960 961 962 963
      if (code != TSDB_CODE_SUCCESS) {
        pRequest->code = code;
        return NULL;
      }
H
Haojun Liao 已提交
964

965 966 967 968
      setQueryResultFromRsp(&pRequest->body.resInfo, (SRetrieveTableRsp*)pResInfo->pData);
      tscDebug("0x%"PRIx64 " fetch results, numOfRows:%d total Rows:%"PRId64", complete:%d, reqId:0x%"PRIx64, pRequest->self, pResInfo->numOfRows,
               pResInfo->totalRows, pResInfo->completed, pRequest->requestId);

969
      if (pResultInfo->numOfRows == 0) {
D
dapan1121 已提交
970 971 972 973
        return NULL;
      }

      goto _return;
974
    } else if (pRequest->type == TDMT_MND_SHOW) {
H
Haojun Liao 已提交
975
      pRequest->type = TDMT_MND_SHOW_RETRIEVE;
H
Haojun Liao 已提交
976
      epSet = getEpSet_s(&pRequest->pTscObj->pAppInfo->mgmtEp);
H
Haojun Liao 已提交
977 978
    } else if (pRequest->type == TDMT_VND_SHOW_TABLES) {
      pRequest->type = TDMT_VND_SHOW_TABLES_FETCH;
H
Haojun Liao 已提交
979 980 981 982 983 984 985 986 987 988 989
      SShowReqInfo* pShowReqInfo = &pRequest->body.showInfo;
      SVgroupInfo* pVgroupInfo = taosArrayGet(pShowReqInfo->pArray, pShowReqInfo->currentIndex);

      epSet.numOfEps = pVgroupInfo->numOfEps;
      epSet.inUse = pVgroupInfo->inUse;

      for (int32_t i = 0; i < epSet.numOfEps; ++i) {
        strncpy(epSet.fqdn[i], pVgroupInfo->epAddr[i].fqdn, tListLen(epSet.fqdn[i]));
        epSet.port[i] = pVgroupInfo->epAddr[i].port;
      }

H
Haojun Liao 已提交
990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006
    } else if (pRequest->type == TDMT_VND_SHOW_TABLES_FETCH) {
      pRequest->type = TDMT_VND_SHOW_TABLES;
      SShowReqInfo* pShowReqInfo = &pRequest->body.showInfo;
      pShowReqInfo->currentIndex += 1;
      if (pShowReqInfo->currentIndex >= taosArrayGetSize(pShowReqInfo->pArray)) {
        return NULL;
      }

      SVgroupInfo* pVgroupInfo = taosArrayGet(pShowReqInfo->pArray, pShowReqInfo->currentIndex);
      SVShowTablesReq* pShowReq = calloc(1, sizeof(SVShowTablesReq));
      pShowReq->head.vgId = htonl(pVgroupInfo->vgId);

      pRequest->body.requestMsg.len = sizeof(SVShowTablesReq);
      pRequest->body.requestMsg.pData = pShowReq;

      SMsgSendInfo* body = buildMsgInfoImpl(pRequest);

H
Haojun Liao 已提交
1007 1008 1009 1010 1011 1012 1013 1014
      epSet.numOfEps = pVgroupInfo->numOfEps;
      epSet.inUse = pVgroupInfo->inUse;

      for (int32_t i = 0; i < epSet.numOfEps; ++i) {
        strncpy(epSet.fqdn[i], pVgroupInfo->epAddr[i].fqdn, tListLen(epSet.fqdn[i]));
        epSet.port[i] = pVgroupInfo->epAddr[i].port;
      }

H
Haojun Liao 已提交
1015 1016
      int64_t  transporterId = 0;
      STscObj *pTscObj = pRequest->pTscObj;
H
Haojun Liao 已提交
1017
      asyncSendMsgToServer(pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, body);
H
Haojun Liao 已提交
1018 1019 1020
      tsem_wait(&pRequest->body.rspSem);

      pRequest->type = TDMT_VND_SHOW_TABLES_FETCH;
H
Haojun Liao 已提交
1021 1022
    } else if (pRequest->type == TDMT_MND_SHOW_RETRIEVE && pResultInfo->pData != NULL) {
      return NULL;
H
Haojun Liao 已提交
1023
    }
H
Haojun Liao 已提交
1024

1025
    SMsgSendInfo* body = buildMsgInfoImpl(pRequest);
H
Haojun Liao 已提交
1026

1027 1028
    int64_t  transporterId = 0;
    STscObj *pTscObj = pRequest->pTscObj;
H
Haojun Liao 已提交
1029
    asyncSendMsgToServer(pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, body);
H
Haojun Liao 已提交
1030 1031 1032 1033 1034 1035 1036 1037 1038

    tsem_wait(&pRequest->body.rspSem);

    pResultInfo->current = 0;
    if (pResultInfo->numOfRows <= pResultInfo->current) {
      return NULL;
    }
  }

D
dapan1121 已提交
1039 1040
_return:

H
Haojun Liao 已提交
1041 1042 1043 1044 1045
  for(int32_t i = 0; i < pResultInfo->numOfCols; ++i) {
    pResultInfo->row[i] = pResultInfo->pCol[i] + pResultInfo->fields[i].bytes * pResultInfo->current;
    if (IS_VAR_DATA_TYPE(pResultInfo->fields[i].type)) {
      pResultInfo->length[i] = varDataLen(pResultInfo->row[i]);
      pResultInfo->row[i] = varDataVal(pResultInfo->row[i]);
1046
    }
H
Haojun Liao 已提交
1047 1048 1049 1050 1051 1052
  }

  pResultInfo->current += 1;
  return pResultInfo->row;
}

H
Haojun Liao 已提交
1053 1054 1055 1056 1057 1058 1059 1060
static void doPrepareResPtr(SReqResultInfo* pResInfo) {
  if (pResInfo->row == NULL) {
    pResInfo->row    = calloc(pResInfo->numOfCols, POINTER_BYTES);
    pResInfo->pCol   = calloc(pResInfo->numOfCols, POINTER_BYTES);
    pResInfo->length = calloc(pResInfo->numOfCols, sizeof(int32_t));
  }
}

1061
void setResultDataPtr(SReqResultInfo* pResultInfo, TAOS_FIELD* pFields, int32_t numOfCols, int32_t numOfRows) {
H
Haojun Liao 已提交
1062 1063 1064 1065 1066
  assert(numOfCols > 0 && pFields != NULL && pResultInfo != NULL);
  if (numOfRows == 0) {
    return;
  }

H
Haojun Liao 已提交
1067 1068 1069
  // todo check for the failure of malloc
  doPrepareResPtr(pResultInfo);

H
Haojun Liao 已提交
1070 1071 1072
  int32_t offset = 0;
  for (int32_t i = 0; i < numOfCols; ++i) {
    pResultInfo->length[i] = pResultInfo->fields[i].bytes;
H
Haojun Liao 已提交
1073
    pResultInfo->row[i]    = (char*) (pResultInfo->pData + offset * pResultInfo->numOfRows);
H
Haojun Liao 已提交
1074 1075
    pResultInfo->pCol[i]   = pResultInfo->row[i];
    offset += pResultInfo->fields[i].bytes;
1076
  }
S
Shengliang Guan 已提交
1077 1078
}

H
Haojun Liao 已提交
1079
char* getDbOfConnection(STscObj* pObj) {
1080 1081
  char *p = NULL;
  pthread_mutex_lock(&pObj->mutex);
1082 1083 1084 1085
  size_t len = strlen(pObj->db);
  if (len > 0) {
    p = strndup(pObj->db, tListLen(pObj->db));
  }
S
Shengliang Guan 已提交
1086

1087
  pthread_mutex_unlock(&pObj->mutex);
1088 1089 1090 1091 1092 1093 1094 1095 1096
  return p;
}

void setConnectionDB(STscObj* pTscObj, const char* db) {
  assert(db != NULL && pTscObj != NULL);
  pthread_mutex_lock(&pTscObj->mutex);
  tstrncpy(pTscObj->db, db, tListLen(pTscObj->db));
  pthread_mutex_unlock(&pTscObj->mutex);
}
S
Shengliang Guan 已提交
1097

1098
void setQueryResultFromRsp(SReqResultInfo* pResultInfo, const SRetrieveTableRsp* pRsp) {
1099 1100 1101
  assert(pResultInfo != NULL && pRsp != NULL);

  pResultInfo->pRspMsg   = (const char*) pRsp;
H
Haojun Liao 已提交
1102 1103
  pResultInfo->pData     = (void*) pRsp->data;
  pResultInfo->numOfRows = htonl(pRsp->numOfRows);
1104 1105
  pResultInfo->current   = 0;
  pResultInfo->completed = (pRsp->completed == 1);
H
Haojun Liao 已提交
1106

1107
  pResultInfo->totalRows += pResultInfo->numOfRows;
H
Haojun Liao 已提交
1108
  setResultDataPtr(pResultInfo, pResultInfo->fields, pResultInfo->numOfCols, pResultInfo->numOfRows);
L
fix  
Liu Jicong 已提交
1109
}