tmq.c 104.0 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
/*
 * 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/>.
 */

L
Liu Jicong 已提交
16
#include "cJSON.h"
17 18 19
#include "clientInt.h"
#include "clientLog.h"
#include "parser.h"
H
Haojun Liao 已提交
20
#include "tdatablock.h"
21 22 23
#include "tdef.h"
#include "tglobal.h"
#include "tmsgtype.h"
X
Xiaoyu Wang 已提交
24
#include "tqueue.h"
25
#include "tref.h"
L
Liu Jicong 已提交
26 27
#include "ttimer.h"

L
Liu Jicong 已提交
28 29 30 31 32 33 34 35
int32_t tmqAskEp(tmq_t* tmq, bool async);

typedef struct {
  int8_t inited;
  tmr_h  timer;
} SMqMgmt;

static SMqMgmt tmqMgmt = {0};
36

L
Liu Jicong 已提交
37 38 39 40 41 42
typedef struct {
  int8_t  tmqRspType;
  int32_t epoch;
} SMqRspWrapper;

typedef struct {
L
Liu Jicong 已提交
43 44 45
  int8_t      tmqRspType;
  int32_t     epoch;
  SMqAskEpRsp msg;
L
Liu Jicong 已提交
46 47
} SMqAskEpRspWrapper;

L
Liu Jicong 已提交
48
struct tmq_list_t {
L
Liu Jicong 已提交
49
  SArray container;
L
Liu Jicong 已提交
50
};
L
Liu Jicong 已提交
51

L
Liu Jicong 已提交
52
struct tmq_conf_t {
53 54 55 56 57 58 59 60 61 62
  char    clientId[256];
  char    groupId[TSDB_CGROUP_LEN];
  int8_t  autoCommit;
  int8_t  resetOffset;
  int8_t  withTbName;
  int8_t  ssEnable;
  int32_t ssBatchSize;

  bool hbBgEnable;

63 64 65 66 67
  uint16_t       port;
  int32_t        autoCommitInterval;
  char*          ip;
  char*          user;
  char*          pass;
68
  tmq_commit_cb* commitCb;
L
Liu Jicong 已提交
69
  void*          commitCbUserParam;
L
Liu Jicong 已提交
70 71 72
};

struct tmq_t {
L
Liu Jicong 已提交
73
  // conf
74 75 76 77 78 79 80 81 82 83 84
  char    groupId[TSDB_CGROUP_LEN];
  char    clientId[256];
  int8_t  withTbName;
  int8_t  useSnapshot;
  int8_t  autoCommit;
  int32_t autoCommitInterval;
  int32_t resetOffsetCfg;
  int64_t consumerId;

  bool hbBgEnable;

L
Liu Jicong 已提交
85 86
  tmq_commit_cb* commitCb;
  void*          commitCbUserParam;
L
Liu Jicong 已提交
87 88 89 90

  // status
  int8_t  status;
  int32_t epoch;
L
Liu Jicong 已提交
91 92
#if 0
  int8_t  epStatus;
L
Liu Jicong 已提交
93
  int32_t epSkipCnt;
L
Liu Jicong 已提交
94
#endif
L
Liu Jicong 已提交
95 96
  int64_t pollCnt;

L
Liu Jicong 已提交
97
  // timer
98 99
  tmr_h hbLiveTimer;
  tmr_h epTimer;
L
Liu Jicong 已提交
100 101 102
  tmr_h reportTimer;
  tmr_h commitTimer;

L
Liu Jicong 已提交
103 104 105 106
  // connection
  STscObj* pTscObj;

  // container
L
Liu Jicong 已提交
107
  SArray*     clientTopics;  // SArray<SMqClientTopic>
L
Liu Jicong 已提交
108
  STaosQueue* mqueue;        // queue of rsp
L
Liu Jicong 已提交
109
  STaosQall*  qall;
L
Liu Jicong 已提交
110 111 112 113
  STaosQueue* delayedTask;  // delayed task queue for heartbeat and auto commit

  // ctl
  tsem_t rspSem;
L
Liu Jicong 已提交
114 115
};

X
Xiaoyu Wang 已提交
116 117 118 119 120 121 122 123
enum {
  TMQ_VG_STATUS__IDLE = 0,
  TMQ_VG_STATUS__WAIT,
};

enum {
  TMQ_CONSUMER_STATUS__INIT = 0,
  TMQ_CONSUMER_STATUS__READY,
124
  TMQ_CONSUMER_STATUS__NO_TOPIC,
L
Liu Jicong 已提交
125 126
};

L
Liu Jicong 已提交
127
enum {
128
  TMQ_DELAYED_TASK__ASK_EP = 1,
L
Liu Jicong 已提交
129 130 131 132
  TMQ_DELAYED_TASK__REPORT,
  TMQ_DELAYED_TASK__COMMIT,
};

L
Liu Jicong 已提交
133
typedef struct {
134 135 136
  // statistics
  int64_t pollCnt;
  // offset
L
Liu Jicong 已提交
137 138 139 140
  /*int64_t      committedOffset;*/
  /*int64_t      currentOffset;*/
  STqOffsetVal committedOffsetNew;
  STqOffsetVal currentOffsetNew;
L
Liu Jicong 已提交
141
  // connection info
142
  int32_t vgId;
X
Xiaoyu Wang 已提交
143
  int32_t vgStatus;
L
Liu Jicong 已提交
144
  int32_t vgSkipCnt;
145 146 147
  SEpSet  epSet;
} SMqClientVg;

L
Liu Jicong 已提交
148
typedef struct {
149
  // subscribe info
L
Liu Jicong 已提交
150
  char* topicName;
L
Liu Jicong 已提交
151
  char  db[TSDB_DB_FNAME_LEN];
L
Liu Jicong 已提交
152 153 154

  SArray* vgs;  // SArray<SMqClientVg>

L
Liu Jicong 已提交
155 156
  int8_t         isSchemaAdaptive;
  SSchemaWrapper schema;
157 158
} SMqClientTopic;

L
Liu Jicong 已提交
159 160 161 162 163
typedef struct {
  int8_t          tmqRspType;
  int32_t         epoch;
  SMqClientVg*    vgHandle;
  SMqClientTopic* topicHandle;
L
Liu Jicong 已提交
164
  union {
L
Liu Jicong 已提交
165 166
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
167
  };
L
Liu Jicong 已提交
168 169
} SMqPollRspWrapper;

L
Liu Jicong 已提交
170
typedef struct {
L
Liu Jicong 已提交
171 172 173
  tmq_t*  tmq;
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
174
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
175

L
Liu Jicong 已提交
176
typedef struct {
177
  tmq_t*  tmq;
L
Liu Jicong 已提交
178
  int32_t code;
L
Liu Jicong 已提交
179
  int32_t async;
X
Xiaoyu Wang 已提交
180
  tsem_t  rspSem;
181 182
} SMqAskEpCbParam;

L
Liu Jicong 已提交
183
typedef struct {
L
Liu Jicong 已提交
184 185
  tmq_t*          tmq;
  SMqClientVg*    pVg;
L
Liu Jicong 已提交
186
  SMqClientTopic* pTopic;
L
Liu Jicong 已提交
187
  int32_t         epoch;
L
Liu Jicong 已提交
188
  int32_t         vgId;
L
Liu Jicong 已提交
189
  tsem_t          rspSem;
X
Xiaoyu Wang 已提交
190
} SMqPollCbParam;
191

192
typedef struct {
L
Liu Jicong 已提交
193 194 195 196
  tmq_t* tmq;
  int8_t automatic;
  int8_t async;
  /*int8_t         freeOffsets;*/
L
Liu Jicong 已提交
197 198
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
L
Liu Jicong 已提交
199
  int32_t        rspErr;
200
  tmq_commit_cb* userCb;
L
Liu Jicong 已提交
201 202 203 204
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
  void*  userParam;
  tsem_t rspSem;
205 206 207 208 209 210 211
} SMqCommitCbParamSet;

typedef struct {
  SMqCommitCbParamSet* params;
  STqOffset*           pOffset;
} SMqCommitCbParam2;

212
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
213
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
214
  conf->withTbName = false;
L
Liu Jicong 已提交
215
  conf->autoCommit = true;
L
Liu Jicong 已提交
216
  conf->autoCommitInterval = 5000;
X
Xiaoyu Wang 已提交
217
  conf->resetOffset = TMQ_CONF__RESET_OFFSET__EARLIEAST;
218 219 220
  return conf;
}

L
Liu Jicong 已提交
221
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
222 223 224 225 226 227
  if (conf) {
    if (conf->ip) taosMemoryFree(conf->ip);
    if (conf->user) taosMemoryFree(conf->user);
    if (conf->pass) taosMemoryFree(conf->pass);
    taosMemoryFree(conf);
  }
L
Liu Jicong 已提交
228 229 230
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
231 232
  if (strcmp(key, "group.id") == 0) {
    strcpy(conf->groupId, value);
L
Liu Jicong 已提交
233
    return TMQ_CONF_OK;
234
  }
L
Liu Jicong 已提交
235

236 237
  if (strcmp(key, "client.id") == 0) {
    strcpy(conf->clientId, value);
L
Liu Jicong 已提交
238 239
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
240

L
Liu Jicong 已提交
241 242
  if (strcmp(key, "enable.auto.commit") == 0) {
    if (strcmp(value, "true") == 0) {
L
Liu Jicong 已提交
243
      conf->autoCommit = true;
L
Liu Jicong 已提交
244 245
      return TMQ_CONF_OK;
    } else if (strcmp(value, "false") == 0) {
L
Liu Jicong 已提交
246
      conf->autoCommit = false;
L
Liu Jicong 已提交
247 248 249 250
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
251
  }
L
Liu Jicong 已提交
252

L
Liu Jicong 已提交
253 254 255 256 257
  if (strcmp(key, "auto.commit.interval.ms") == 0) {
    conf->autoCommitInterval = atoi(value);
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
258 259 260 261 262 263 264 265 266 267 268 269 270 271
  if (strcmp(key, "auto.offset.reset") == 0) {
    if (strcmp(value, "none") == 0) {
      conf->resetOffset = TMQ_CONF__RESET_OFFSET__NONE;
      return TMQ_CONF_OK;
    } else if (strcmp(value, "earliest") == 0) {
      conf->resetOffset = TMQ_CONF__RESET_OFFSET__EARLIEAST;
      return TMQ_CONF_OK;
    } else if (strcmp(value, "latest") == 0) {
      conf->resetOffset = TMQ_CONF__RESET_OFFSET__LATEST;
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }
L
Liu Jicong 已提交
272

273 274
  if (strcmp(key, "msg.with.table.name") == 0) {
    if (strcmp(value, "true") == 0) {
275
      conf->withTbName = true;
L
Liu Jicong 已提交
276
      return TMQ_CONF_OK;
277
    } else if (strcmp(value, "false") == 0) {
278
      conf->withTbName = false;
L
Liu Jicong 已提交
279
      return TMQ_CONF_OK;
280 281 282 283 284
    } else {
      return TMQ_CONF_INVALID;
    }
  }

L
Liu Jicong 已提交
285
  if (strcmp(key, "experimental.snapshot.enable") == 0) {
L
Liu Jicong 已提交
286
    if (strcmp(value, "true") == 0) {
287 288 289 290 291 292 293 294 295 296 297 298 299
      conf->ssEnable = true;
      return TMQ_CONF_OK;
    } else if (strcmp(value, "false") == 0) {
      conf->ssEnable = false;
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }

  if (strcmp(key, "enable.heartbeat.background") == 0) {
    if (strcmp(value, "true") == 0) {
      conf->hbBgEnable = true;
L
Liu Jicong 已提交
300 301
      return TMQ_CONF_OK;
    } else if (strcmp(value, "false") == 0) {
302
      conf->hbBgEnable = false;
L
Liu Jicong 已提交
303 304 305 306
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
307
    return TMQ_CONF_OK;
L
Liu Jicong 已提交
308 309
  }

L
Liu Jicong 已提交
310
  if (strcmp(key, "experimental.snapshot.batch.size") == 0) {
311
    conf->ssBatchSize = atoi(value);
L
Liu Jicong 已提交
312 313 314
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
315
  if (strcmp(key, "td.connect.ip") == 0) {
L
Liu Jicong 已提交
316 317 318
    conf->ip = strdup(value);
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
319
  if (strcmp(key, "td.connect.user") == 0) {
L
Liu Jicong 已提交
320 321 322
    conf->user = strdup(value);
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
323
  if (strcmp(key, "td.connect.pass") == 0) {
L
Liu Jicong 已提交
324 325 326
    conf->pass = strdup(value);
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
327
  if (strcmp(key, "td.connect.port") == 0) {
L
Liu Jicong 已提交
328 329 330
    conf->port = atoi(value);
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
331
  if (strcmp(key, "td.connect.db") == 0) {
332
    /*conf->db = strdup(value);*/
L
Liu Jicong 已提交
333 334 335
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
336
  return TMQ_CONF_UNKNOWN;
337 338 339
}

tmq_list_t* tmq_list_new() {
L
Liu Jicong 已提交
340 341
  //
  return (tmq_list_t*)taosArrayInit(0, sizeof(void*));
342 343
}

L
Liu Jicong 已提交
344 345
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
346 347 348 349 350
  if (src == NULL || src[0] == 0) return -1;
  char* topic = strdup(src);
  if (topic[0] != '`') {
    strtolower(topic, src);
  }
L
fix  
Liu Jicong 已提交
351
  if (taosArrayPush(container, &topic) == NULL) return -1;
352 353 354
  return 0;
}

L
Liu Jicong 已提交
355
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
356
  SArray* container = &list->container;
L
Liu Jicong 已提交
357
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
358 359
}

L
Liu Jicong 已提交
360 361 362 363 364 365 366 367 368 369
int32_t tmq_list_get_size(const tmq_list_t* list) {
  const SArray* container = &list->container;
  return taosArrayGetSize(container);
}

char** tmq_list_to_c_array(const tmq_list_t* list) {
  const SArray* container = &list->container;
  return container->pData;
}

L
Liu Jicong 已提交
370 371 372 373
static int32_t tmqMakeTopicVgKey(char* dst, const char* topicName, int32_t vg) {
  return sprintf(dst, "%s:%d", topicName, vg);
}

D
dapan1121 已提交
374
int32_t tmqCommitCb2(void* param, SDataBuf* pBuf, int32_t code) {
375 376 377
  SMqCommitCbParam2*   pParam = (SMqCommitCbParam2*)param;
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
  // push into array
L
Liu Jicong 已提交
378
#if 0
379 380 381 382 383
  if (code == 0) {
    taosArrayPush(pParamSet->failedOffsets, &pParam->pOffset);
  } else {
    taosArrayPush(pParamSet->successfulOffsets, &pParam->pOffset);
  }
L
Liu Jicong 已提交
384
#endif
L
Liu Jicong 已提交
385

S
Shengliang Guan 已提交
386
  /*tscDebug("receive offset commit cb of %s on vgId:%d, offset is %" PRId64, pParam->pOffset->subKey, pParam->->vgId,
L
Liu Jicong 已提交
387 388
   * pOffset->version);*/

389
  // count down waiting rsp
L
Liu Jicong 已提交
390
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
391 392 393 394 395 396 397
  ASSERT(waitingRspNum >= 0);

  if (waitingRspNum == 0) {
    // if no more waiting rsp
    if (pParamSet->async) {
      // call async cb func
      if (pParamSet->automatic && pParamSet->tmq->commitCb) {
L
Liu Jicong 已提交
398
        pParamSet->tmq->commitCb(pParamSet->tmq, pParamSet->rspErr, pParamSet->tmq->commitCbUserParam);
399 400
      } else if (!pParamSet->automatic && pParamSet->userCb) {
        // sem post
L
Liu Jicong 已提交
401
        pParamSet->userCb(pParamSet->tmq, pParamSet->rspErr, pParamSet->userParam);
402
      }
L
Liu Jicong 已提交
403 404
    } else {
      tsem_post(&pParamSet->rspSem);
405 406
    }

L
Liu Jicong 已提交
407
#if 0
408 409
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
L
Liu Jicong 已提交
410
#endif
411 412 413 414
  }
  return 0;
}

L
Liu Jicong 已提交
415 416 417 418 419 420 421
static int32_t tmqSendCommitReq(tmq_t* tmq, SMqClientVg* pVg, SMqClientTopic* pTopic, SMqCommitCbParamSet* pParamSet) {
  STqOffset* pOffset = taosMemoryCalloc(1, sizeof(STqOffset));
  if (pOffset == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  pOffset->val = pVg->currentOffsetNew;
422

L
Liu Jicong 已提交
423 424 425 426
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pOffset->subKey, tmq->groupId, groupLen);
  pOffset->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pOffset->subKey + groupLen + 1, pTopic->topicName);
L
Liu Jicong 已提交
427

L
Liu Jicong 已提交
428 429 430 431 432 433 434 435 436 437
  int32_t len;
  int32_t code;
  tEncodeSize(tEncodeSTqOffset, pOffset, len, code);
  if (code < 0) {
    ASSERT(0);
    return -1;
  }
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
  if (buf == NULL) return -1;
  ((SMsgHead*)buf)->vgId = htonl(pVg->vgId);
L
Liu Jicong 已提交
438

L
Liu Jicong 已提交
439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTqOffset(&encoder, pOffset);

  // build param
  SMqCommitCbParam2* pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam2));
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;

  // build send info
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
    return -1;
  }
  pMsgSendInfo->msgInfo = (SDataBuf){
      .pData = buf,
      .len = sizeof(SMsgHead) + len,
      .handle = NULL,
  };

S
Shengliang Guan 已提交
461 462
  tscDebug("consumer:%" PRId64 ", commit offset of %s on vgId:%d, offset is %" PRId64, tmq->consumerId, pOffset->subKey,
           pVg->vgId, pOffset->val.version);
L
Liu Jicong 已提交
463 464 465 466 467 468 469

  // TODO: put into cb
  pVg->committedOffsetNew = pVg->currentOffsetNew;

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
470
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
L
Liu Jicong 已提交
471 472 473 474
  pMsgSendInfo->fp = tmqCommitCb2;
  pMsgSendInfo->msgType = TDMT_VND_MQ_COMMIT_OFFSET;
  // send msg

L
Liu Jicong 已提交
475 476 477
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

L
Liu Jicong 已提交
478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
  return 0;
}

int32_t tmqCommitMsgImpl(tmq_t* tmq, const TAOS_RES* msg, int8_t async, tmq_commit_cb* userCb, void* userParam) {
  char*   topic;
  int32_t vgId;
  ASSERT(msg != NULL);
  if (TD_RES_TMQ(msg)) {
    SMqRspObj* pRspObj = (SMqRspObj*)msg;
    topic = pRspObj->topic;
    vgId = pRspObj->vgId;
  } else if (TD_RES_TMQ_META(msg)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)msg;
    topic = pMetaRspObj->topic;
    vgId = pMetaRspObj->vgId;
  } else {
    return TSDB_CODE_TMQ_INVALID_MSG;
  }

  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  pParamSet->tmq = tmq;
  pParamSet->automatic = 0;
  pParamSet->async = async;
  /*pParamSet->freeOffsets = 1;*/
  pParamSet->userCb = userCb;
  pParamSet->userParam = userParam;
  tsem_init(&pParamSet->rspSem, 0, 0);

L
Liu Jicong 已提交
512 513
  int32_t code = -1;

L
Liu Jicong 已提交
514 515 516 517 518 519 520 521 522 523
  for (int32_t i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
    if (strcmp(pTopic->topicName, topic) != 0) continue;
    for (int32_t j = 0; j < taosArrayGetSize(pTopic->vgs); j++) {
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
      if (pVg->vgId != vgId) continue;

      if (pVg->currentOffsetNew.type > 0 && !tOffsetEqual(&pVg->currentOffsetNew, &pVg->committedOffsetNew)) {
        if (tmqSendCommitReq(tmq, pVg, pTopic, pParamSet) < 0) {
          goto FAIL;
L
Liu Jicong 已提交
524
        }
L
Liu Jicong 已提交
525
        goto HANDLE_RSP;
L
Liu Jicong 已提交
526 527
      }
    }
L
Liu Jicong 已提交
528
  }
L
Liu Jicong 已提交
529

L
Liu Jicong 已提交
530 531 532 533
HANDLE_RSP:
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
534 535 536
    return 0;
  }

L
Liu Jicong 已提交
537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560
  if (!async) {
    tsem_wait(&pParamSet->rspSem);
    code = pParamSet->rspErr;
    tsem_destroy(&pParamSet->rspSem);
    return code;
  } else {
    code = 0;
  }

FAIL:
  if (code != 0 && async) {
    userCb(tmq, code, userParam);
  }
  return 0;
}

int32_t tmqCommitInner2(tmq_t* tmq, const TAOS_RES* msg, int8_t automatic, int8_t async, tmq_commit_cb* userCb,
                        void* userParam) {
  int32_t code = -1;

  if (msg != NULL) {
    return tmqCommitMsgImpl(tmq, msg, async, userCb, userParam);
  }

561 562 563 564 565 566 567 568
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
  pParamSet->tmq = tmq;
  pParamSet->automatic = automatic;
  pParamSet->async = async;
L
Liu Jicong 已提交
569
  /*pParamSet->freeOffsets = 1;*/
570 571 572 573 574 575
  pParamSet->userCb = userCb;
  pParamSet->userParam = userParam;
  tsem_init(&pParamSet->rspSem, 0, 0);

  for (int32_t i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
576

S
Shengliang Guan 已提交
577
    tscDebug("consumer:%" PRId64 ", begin commit for topic %s, vgNum %d", tmq->consumerId, pTopic->topicName,
L
Liu Jicong 已提交
578 579
             (int32_t)taosArrayGetSize(pTopic->vgs));

580
    for (int32_t j = 0; j < taosArrayGetSize(pTopic->vgs); j++) {
L
Liu Jicong 已提交
581 582
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);

S
Shengliang Guan 已提交
583 584
      tscDebug("consumer:%" PRId64 ", begin commit for topic %s, vgId:%d", tmq->consumerId, pTopic->topicName,
               pVg->vgId);
L
Liu Jicong 已提交
585

L
Liu Jicong 已提交
586 587 588 589
      if (pVg->currentOffsetNew.type > 0 && !tOffsetEqual(&pVg->currentOffsetNew, &pVg->committedOffsetNew)) {
        if (tmqSendCommitReq(tmq, pVg, pTopic, pParamSet) < 0) {
          continue;
        }
590 591 592 593
      }
    }
  }

L
Liu Jicong 已提交
594 595 596 597 598 599
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
    return 0;
  }

600 601 602 603 604 605 606 607 608 609
  if (!async) {
    tsem_wait(&pParamSet->rspSem);
    code = pParamSet->rspErr;
    tsem_destroy(&pParamSet->rspSem);
  } else {
    code = 0;
  }

  if (code != 0 && async) {
    if (automatic) {
L
Liu Jicong 已提交
610
      tmq->commitCb(tmq, code, tmq->commitCbUserParam);
611
    } else {
L
Liu Jicong 已提交
612
      userCb(tmq, code, userParam);
613 614 615
    }
  }

L
Liu Jicong 已提交
616
#if 0
617 618 619 620
  if (!async) {
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
  }
L
Liu Jicong 已提交
621
#endif
622 623 624 625

  return 0;
}

626
void tmqAssignAskEpTask(void* param, void* tmrId) {
L
Liu Jicong 已提交
627
  tmq_t*  tmq = (tmq_t*)param;
628
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
629
  *pTaskType = TMQ_DELAYED_TASK__ASK_EP;
L
Liu Jicong 已提交
630
  taosWriteQitem(tmq->delayedTask, pTaskType);
631
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
632 633 634 635
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
  tmq_t*  tmq = (tmq_t*)param;
636
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
L
Liu Jicong 已提交
637 638
  *pTaskType = TMQ_DELAYED_TASK__COMMIT;
  taosWriteQitem(tmq->delayedTask, pTaskType);
639
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
640 641 642 643
}

void tmqAssignDelayedReportTask(void* param, void* tmrId) {
  tmq_t*  tmq = (tmq_t*)param;
644
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
L
Liu Jicong 已提交
645 646
  *pTaskType = TMQ_DELAYED_TASK__REPORT;
  taosWriteQitem(tmq->delayedTask, pTaskType);
647
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
648 649
}

650 651 652 653 654 655 656 657 658 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
int32_t tmqHbCb(void* param, SDataBuf* pMsg, int32_t code) {
  if (pMsg && pMsg->pData) taosMemoryFree(pMsg->pData);
  return 0;
}

void tmqSendHbReq(void* param, void* tmrId) {
  // TODO replace with ref
  tmq_t*    tmq = (tmq_t*)param;
  int64_t   consumerId = tmq->consumerId;
  int32_t   epoch = tmq->epoch;
  SMqHbReq* pReq = taosMemoryMalloc(sizeof(SMqHbReq));
  if (pReq == NULL) goto OVER;
  pReq->consumerId = consumerId;
  pReq->epoch = epoch;

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pReq);
  }
  sendInfo->msgInfo = (SDataBuf){
      .pData = pReq,
      .len = sizeof(SMqHbReq),
      .handle = NULL,
  };

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
  sendInfo->msgType = TDMT_MND_MQ_HB;

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

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

OVER:
  taosTmrReset(tmqSendHbReq, 1000, tmq, tmqMgmt.timer, &tmq->hbLiveTimer);
}

L
Liu Jicong 已提交
690 691 692 693 694 695 696 697
int32_t tmqHandleAllDelayedTask(tmq_t* tmq) {
  STaosQall* qall = taosAllocateQall();
  taosReadAllQitems(tmq->delayedTask, qall);
  while (1) {
    int8_t* pTaskType = NULL;
    taosGetQitem(qall, (void**)&pTaskType);
    if (pTaskType == NULL) break;

698
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
L
Liu Jicong 已提交
699
      tmqAskEp(tmq, true);
700
      taosTmrReset(tmqAssignAskEpTask, 1000, tmq, tmqMgmt.timer, &tmq->epTimer);
L
Liu Jicong 已提交
701
    } else if (*pTaskType == TMQ_DELAYED_TASK__COMMIT) {
L
Liu Jicong 已提交
702
      tmqCommitInner2(tmq, NULL, 1, 1, tmq->commitCb, tmq->commitCbUserParam);
L
Liu Jicong 已提交
703 704 705 706 707
      taosTmrReset(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, tmq, tmqMgmt.timer, &tmq->commitTimer);
    } else if (*pTaskType == TMQ_DELAYED_TASK__REPORT) {
    } else {
      ASSERT(0);
    }
L
Liu Jicong 已提交
708
    taosFreeQitem(pTaskType);
L
Liu Jicong 已提交
709 710 711 712 713
  }
  taosFreeQall(qall);
  return 0;
}

L
Liu Jicong 已提交
714
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
715
  SMqRspWrapper* msg = NULL;
L
Liu Jicong 已提交
716 717 718 719 720 721 722 723
  while (1) {
    taosGetQitem(tmq->qall, (void**)&msg);
    if (msg)
      taosFreeQitem(msg);
    else
      break;
  }

L
fix  
Liu Jicong 已提交
724
  msg = NULL;
L
Liu Jicong 已提交
725 726 727 728 729 730 731 732 733 734
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
    taosGetQitem(tmq->qall, (void**)&msg);
    if (msg)
      taosFreeQitem(msg);
    else
      break;
  }
}

D
dapan1121 已提交
735
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
736 737
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
L
Liu Jicong 已提交
738
  /*tmq_t* tmq = pParam->tmq;*/
L
Liu Jicong 已提交
739 740 741
  tsem_post(&pParam->rspSem);
  return 0;
}
742

L
Liu Jicong 已提交
743
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
744 745 746 747
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
748
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
749
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
750
  }
L
Liu Jicong 已提交
751
  return 0;
X
Xiaoyu Wang 已提交
752 753
}

L
Liu Jicong 已提交
754 755 756
int32_t tmq_unsubscribe(tmq_t* tmq) {
  tmq_list_t* lst = tmq_list_new();
  int32_t     rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
757 758
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
759 760
}

L
Liu Jicong 已提交
761
#if 0
762
tmq_t* tmq_consumer_new(void* conn, tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
wafwerar's avatar
wafwerar 已提交
763
  tmq_t* pTmq = taosMemoryCalloc(sizeof(tmq_t), 1);
764 765 766 767 768 769
  if (pTmq == NULL) {
    return NULL;
  }
  pTmq->pTscObj = (STscObj*)conn;
  pTmq->status = 0;
  pTmq->pollCnt = 0;
L
Liu Jicong 已提交
770
  pTmq->epoch = 0;
L
fix txn  
Liu Jicong 已提交
771
  pTmq->epStatus = 0;
L
temp  
Liu Jicong 已提交
772
  pTmq->epSkipCnt = 0;
L
Liu Jicong 已提交
773
  // set conf
774 775
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
L
Liu Jicong 已提交
776
  pTmq->autoCommit = conf->autoCommit;
777
  pTmq->commit_cb = conf->commit_cb;
L
Liu Jicong 已提交
778
  pTmq->resetOffsetCfg = conf->resetOffset;
L
Liu Jicong 已提交
779

L
Liu Jicong 已提交
780 781 782 783 784 785 786 787 788 789
  pTmq->consumerId = generateRequestId() & (((uint64_t)-1) >> 1);
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  if (pTmq->clientTopics == NULL) {
    taosMemoryFree(pTmq);
    return NULL;
  }

  pTmq->mqueue = taosOpenQueue();
  pTmq->qall = taosAllocateQall();

L
Liu Jicong 已提交
790 791
  tsem_init(&pTmq->rspSem, 0, 0);

L
Liu Jicong 已提交
792 793
  return pTmq;
}
L
Liu Jicong 已提交
794
#endif
L
Liu Jicong 已提交
795

L
Liu Jicong 已提交
796
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
L
Liu Jicong 已提交
797 798 799 800 801 802 803 804 805 806 807
  // init timer
  int8_t inited = atomic_val_compare_exchange_8(&tmqMgmt.inited, 0, 1);
  if (inited == 0) {
    tmqMgmt.timer = taosTmrInit(1000, 100, 360000, "TMQ");
    if (tmqMgmt.timer == NULL) {
      atomic_store_8(&tmqMgmt.inited, 0);
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      return NULL;
    }
  }

L
Liu Jicong 已提交
808 809
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
810
    terrno = TSDB_CODE_OUT_OF_MEMORY;
S
Shengliang Guan 已提交
811 812
    tscError("consumer %" PRId64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
             pTmq->groupId);
L
Liu Jicong 已提交
813 814
    return NULL;
  }
L
Liu Jicong 已提交
815

L
Liu Jicong 已提交
816 817 818 819 820
  const char* user = conf->user == NULL ? TSDB_DEFAULT_USER : conf->user;
  const char* pass = conf->pass == NULL ? TSDB_DEFAULT_PASS : conf->pass;

  ASSERT(user);
  ASSERT(pass);
L
Liu Jicong 已提交
821
  ASSERT(conf->groupId[0]);
L
Liu Jicong 已提交
822

L
Liu Jicong 已提交
823 824 825 826 827 828
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->qall = taosAllocateQall();
  pTmq->delayedTask = taosOpenQueue();

  if (pTmq->clientTopics == NULL || pTmq->mqueue == NULL || pTmq->qall == NULL || pTmq->delayedTask == NULL) {
L
Liu Jicong 已提交
829
    terrno = TSDB_CODE_OUT_OF_MEMORY;
S
Shengliang Guan 已提交
830 831
    tscError("consumer %" PRId64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
             pTmq->groupId);
L
Liu Jicong 已提交
832 833
    goto FAIL;
  }
L
Liu Jicong 已提交
834

L
Liu Jicong 已提交
835 836
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
837 838
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
839 840
  /*pTmq->epStatus = 0;*/
  /*pTmq->epSkipCnt = 0;*/
L
Liu Jicong 已提交
841

L
Liu Jicong 已提交
842 843 844
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
845
  pTmq->withTbName = conf->withTbName;
846
  pTmq->useSnapshot = conf->ssEnable;
L
Liu Jicong 已提交
847
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
848
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
849 850
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
851 852
  pTmq->resetOffsetCfg = conf->resetOffset;

853 854
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
855
  // assign consumerId
L
Liu Jicong 已提交
856
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
857

L
Liu Jicong 已提交
858 859
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
S
Shengliang Guan 已提交
860 861
    tscError("consumer %" PRId64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
             pTmq->groupId);
L
Liu Jicong 已提交
862 863
    goto FAIL;
  }
L
Liu Jicong 已提交
864

L
Liu Jicong 已提交
865 866 867
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
S
Shengliang Guan 已提交
868 869
    tscError("consumer %" PRId64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
             pTmq->groupId);
L
Liu Jicong 已提交
870 871 872
    tsem_destroy(&pTmq->rspSem);
    goto FAIL;
  }
L
Liu Jicong 已提交
873

874 875 876 877
  if (pTmq->hbBgEnable) {
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pTmq, tmqMgmt.timer);
  }

S
Shengliang Guan 已提交
878
  tscInfo("consumer %" PRId64 " is setup, consumer group %s", pTmq->consumerId, pTmq->groupId);
L
Liu Jicong 已提交
879

880
  return pTmq;
L
Liu Jicong 已提交
881 882 883 884 885 886 887 888

FAIL:
  if (pTmq->clientTopics) taosArrayDestroy(pTmq->clientTopics);
  if (pTmq->mqueue) taosCloseQueue(pTmq->mqueue);
  if (pTmq->delayedTask) taosCloseQueue(pTmq->delayedTask);
  if (pTmq->qall) taosFreeQall(pTmq->qall);
  taosMemoryFree(pTmq);
  return NULL;
889 890
}

L
Liu Jicong 已提交
891 892
#if 0
int32_t tmq_commit(tmq_t* tmq, const tmq_topic_vgroup_list_t* offsets, int32_t async) {
L
Liu Jicong 已提交
893
  return tmqCommitInner2(tmq, offsets, 0, async, tmq->commitCb, tmq->commitCbUserParam);
L
Liu Jicong 已提交
894
}
L
Liu Jicong 已提交
895
#endif
L
Liu Jicong 已提交
896

L
Liu Jicong 已提交
897
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
L
Liu Jicong 已提交
898 899 900 901 902
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
  SCMSubscribeReq req = {0};
  int32_t         code = -1;
903 904

  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
905
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
906
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
907
  req.topicNames = taosArrayInit(sz, sizeof(void*));
L
Liu Jicong 已提交
908
  if (req.topicNames == NULL) goto FAIL;
909

L
Liu Jicong 已提交
910 911
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
912 913

    SName name = {0};
L
Liu Jicong 已提交
914
    tNameSetDbName(&name, tmq->pTscObj->acctId, topic, strlen(topic));
915

L
Liu Jicong 已提交
916 917 918
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
919
    }
L
Liu Jicong 已提交
920
    tNameExtractFullName(&name, topicFName);
921

L
Liu Jicong 已提交
922 923 924
    tscDebug("subscribe topic: %s", topicFName);

    taosArrayPush(req.topicNames, &topicFName);
925 926
  }

L
Liu Jicong 已提交
927 928 929 930
  int32_t tlen = tSerializeSCMSubscribeReq(NULL, &req);
  buf = taosMemoryMalloc(tlen);
  if (buf == NULL) goto FAIL;

931 932 933
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

934
  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
935
  if (sendInfo == NULL) goto FAIL;
936

X
Xiaoyu Wang 已提交
937
  SMqSubscribeCbParam param = {
L
Liu Jicong 已提交
938
      .rspErr = 0,
X
Xiaoyu Wang 已提交
939 940
      .tmq = tmq,
  };
L
Liu Jicong 已提交
941

L
Liu Jicong 已提交
942 943 944
  if (tsem_init(&param.rspSem, 0, 0) != 0) goto FAIL;

  sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
945 946 947 948
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };
949

L
Liu Jicong 已提交
950 951
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
952 953
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
954 955
  sendInfo->msgType = TDMT_MND_SUBSCRIBE;

956 957 958 959 960
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
961 962 963
  // avoid double free if msg is sent
  buf = NULL;

L
Liu Jicong 已提交
964 965
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
966

L
Liu Jicong 已提交
967 968 969
  code = param.rspErr;
  if (code != 0) goto FAIL;

L
Liu Jicong 已提交
970
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == tmqAskEp(tmq, false)) {
L
fix  
Liu Jicong 已提交
971
    tscDebug("consumer not ready, retry");
L
Liu Jicong 已提交
972 973
    taosMsleep(500);
  }
974

975 976 977
  // init ep timer
  if (tmq->epTimer == NULL) {
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, tmq, tmqMgmt.timer);
978
  }
L
Liu Jicong 已提交
979 980

  // init auto commit timer
981
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
L
Liu Jicong 已提交
982 983 984
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, tmq, tmqMgmt.timer);
  }

L
Liu Jicong 已提交
985 986 987
  code = 0;
FAIL:
  if (req.topicNames != NULL) taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
988
  if (code != 0 && buf) {
L
Liu Jicong 已提交
989 990 991
    taosMemoryFree(buf);
  }
  return code;
992 993
}

L
Liu Jicong 已提交
994
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
L
Liu Jicong 已提交
995
  //
996
  conf->commitCb = cb;
L
Liu Jicong 已提交
997
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
998
}
999

D
dapan1121 已提交
1000
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1001 1002
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1003
  SMqClientTopic* pTopic = pParam->pTopic;
X
Xiaoyu Wang 已提交
1004
  tmq_t*          tmq = pParam->tmq;
L
Liu Jicong 已提交
1005 1006 1007
  int32_t         vgId = pParam->vgId;
  int32_t         epoch = pParam->epoch;
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1008
  if (code != 0) {
S
Shengliang Guan 已提交
1009
    tscWarn("msg discard from vgId:%d, epoch %d, code:%x", vgId, epoch, code);
D
dapan1121 已提交
1010
    if (pMsg->pData) taosMemoryFreeClear(pMsg->pData);
L
Liu Jicong 已提交
1011 1012 1013 1014
    if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM);
      if (pRspWrapper == NULL) {
        taosMemoryFree(pMsg->pData);
S
Shengliang Guan 已提交
1015
        tscWarn("msg discard from vgId:%d, epoch %d since out of memory", vgId, epoch);
L
Liu Jicong 已提交
1016 1017 1018 1019 1020 1021 1022 1023
        goto CREATE_MSG_FAIL;
      }
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      /*pRspWrapper->vgHandle = pVg;*/
      /*pRspWrapper->topicHandle = pTopic;*/
      taosWriteQitem(tmq->mqueue, pRspWrapper);
      tsem_post(&tmq->rspSem);
    }
L
fix txn  
Liu Jicong 已提交
1024
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1025 1026
  }

X
Xiaoyu Wang 已提交
1027 1028 1029
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
  int32_t tmqEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < tmqEpoch) {
L
Liu Jicong 已提交
1030
    // do not write into queue since updating epoch reset
S
Shengliang Guan 已提交
1031
    tscWarn("msg discard from vgId:%d since from earlier epoch, rsp epoch %d, current epoch %d", vgId, msgEpoch,
L
Liu Jicong 已提交
1032
            tmqEpoch);
1033
    tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1034
    taosMemoryFree(pMsg->pData);
X
Xiaoyu Wang 已提交
1035 1036 1037 1038
    return 0;
  }

  if (msgEpoch != tmqEpoch) {
S
Shengliang Guan 已提交
1039
    tscWarn("mismatch rsp from vgId:%d, epoch %d, current epoch %d", vgId, msgEpoch, tmqEpoch);
X
Xiaoyu Wang 已提交
1040 1041
  }

L
Liu Jicong 已提交
1042 1043 1044
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

1045
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM);
L
Liu Jicong 已提交
1046
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1047
    taosMemoryFree(pMsg->pData);
S
Shengliang Guan 已提交
1048
    tscWarn("msg discard from vgId:%d, epoch %d since out of memory", vgId, epoch);
L
fix txn  
Liu Jicong 已提交
1049
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1050
  }
L
Liu Jicong 已提交
1051

L
Liu Jicong 已提交
1052
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1053 1054
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
L
Liu Jicong 已提交
1055

L
Liu Jicong 已提交
1056
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1057 1058 1059
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqDataRsp(&decoder, &pRspWrapper->dataRsp);
wmmhello's avatar
wmmhello 已提交
1060
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1061
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1062 1063 1064
  } else {
    ASSERT(rspType == TMQ_MSG_TYPE__POLL_META_RSP);
    tDecodeSMqMetaRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pRspWrapper->metaRsp);
L
Liu Jicong 已提交
1065
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1066
  }
L
Liu Jicong 已提交
1067

L
Liu Jicong 已提交
1068
  taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1069

S
Shengliang Guan 已提交
1070 1071 1072
  tscDebug("consumer:%" PRId64 ", recv poll: vgId:%d, req offset %" PRId64 ", rsp offset %" PRId64 " type %d",
           tmq->consumerId, pVg->vgId, pRspWrapper->dataRsp.reqOffset.version, pRspWrapper->dataRsp.rspOffset.version,
           rspType);
L
fix  
Liu Jicong 已提交
1073

L
Liu Jicong 已提交
1074
  taosWriteQitem(tmq->mqueue, pRspWrapper);
1075
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1076

L
Liu Jicong 已提交
1077
  return 0;
L
fix txn  
Liu Jicong 已提交
1078
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1079
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1080 1081
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
1082
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1083
  return -1;
1084 1085
}

1086 1087 1088 1089 1090
bool tmqUpdateEp2(tmq_t* tmq, int32_t epoch, SMqAskEpRsp* pRsp) {
  bool set = false;

  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
S
Shengliang Guan 已提交
1091
  tscDebug("consumer:%" PRId64 ", update ep epoch %d to epoch %d, topic num:%d", tmq->consumerId, tmq->epoch, epoch,
1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109
           topicNumGet);

  SArray* newTopics = taosArrayInit(topicNumGet, sizeof(SMqClientTopic));
  if (newTopics == NULL) {
    return false;
  }

  SHashObj* pHash = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pHash == NULL) {
    taosArrayDestroy(newTopics);
    return false;
  }
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
  for (int32_t i = 0; i < topicNumCur; i++) {
    // find old topic
    SMqClientTopic* pTopicCur = taosArrayGet(tmq->clientTopics, i);
    if (pTopicCur->vgs) {
      int32_t vgNumCur = taosArrayGetSize(pTopicCur->vgs);
S
Shengliang Guan 已提交
1110
      tscDebug("consumer:%" PRId64 ", new vg num: %d", tmq->consumerId, vgNumCur);
1111 1112 1113
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
        sprintf(vgKey, "%s:%d", pTopicCur->topicName, pVgCur->vgId);
L
Liu Jicong 已提交
1114 1115 1116 1117
        char buf[80];
        tFormatOffset(buf, 80, &pVgCur->currentOffsetNew);
        tscDebug("consumer:%" PRId64 ", epoch %d vgId:%d vgKey is %s, offset is %s", tmq->consumerId, epoch,
                 pVgCur->vgId, vgKey, buf);
L
Liu Jicong 已提交
1118
        taosHashPut(pHash, vgKey, strlen(vgKey), &pVgCur->currentOffsetNew, sizeof(STqOffsetVal));
1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129
      }
    }
  }

  for (int32_t i = 0; i < topicNumGet; i++) {
    SMqClientTopic topic = {0};
    SMqSubTopicEp* pTopicEp = taosArrayGet(pRsp->topics, i);
    topic.schema = pTopicEp->schema;
    topic.topicName = strdup(pTopicEp->topic);
    tstrncpy(topic.db, pTopicEp->db, TSDB_DB_FNAME_LEN);

S
Shengliang Guan 已提交
1130
    tscDebug("consumer:%" PRId64 ", update topic: %s", tmq->consumerId, topic.topicName);
1131 1132 1133 1134 1135 1136

    int32_t vgNumGet = taosArrayGetSize(pTopicEp->vgs);
    topic.vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));
    for (int32_t j = 0; j < vgNumGet; j++) {
      SMqSubVgEp* pVgEp = taosArrayGet(pTopicEp->vgs, j);
      sprintf(vgKey, "%s:%d", topic.topicName, pVgEp->vgId);
L
Liu Jicong 已提交
1137 1138
      STqOffsetVal* pOffset = taosHashGet(pHash, vgKey, strlen(vgKey));
      STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
1139
      if (pOffset != NULL) {
L
Liu Jicong 已提交
1140
        offsetNew = *pOffset;
1141 1142 1143 1144
      }

      SMqClientVg clientVg = {
          .pollCnt = 0,
L
Liu Jicong 已提交
1145
          .currentOffsetNew = offsetNew,
1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168
          .vgId = pVgEp->vgId,
          .epSet = pVgEp->epSet,
          .vgStatus = TMQ_VG_STATUS__IDLE,
          .vgSkipCnt = 0,
      };
      taosArrayPush(topic.vgs, &clientVg);
      set = true;
    }
    taosArrayPush(newTopics, &topic);
  }
  if (tmq->clientTopics) taosArrayDestroy(tmq->clientTopics);
  taosHashCleanup(pHash);
  tmq->clientTopics = newTopics;

  if (taosArrayGetSize(tmq->clientTopics) == 0)
    atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__NO_TOPIC);
  else
    atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__READY);

  atomic_store_32(&tmq->epoch, epoch);
  return set;
}

L
Liu Jicong 已提交
1169
#if 0
L
Liu Jicong 已提交
1170
bool tmqUpdateEp(tmq_t* tmq, int32_t epoch, SMqAskEpRsp* pRsp) {
L
Liu Jicong 已提交
1171
  /*printf("call update ep %d\n", epoch);*/
X
Xiaoyu Wang 已提交
1172
  bool    set = false;
L
Liu Jicong 已提交
1173 1174
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
S
Shengliang Guan 已提交
1175
  tscDebug("consumer:%" PRId64 ", update ep epoch %d to epoch %d, topic num: %d", tmq->consumerId, tmq->epoch, epoch,
L
Liu Jicong 已提交
1176
           topicNumGet);
L
Liu Jicong 已提交
1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188
  SArray* newTopics = taosArrayInit(topicNumGet, sizeof(SMqClientTopic));
  if (newTopics == NULL) {
    return false;
  }
  SHashObj* pHash = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pHash == NULL) {
    taosArrayDestroy(newTopics);
    return false;
  }

  // find topic, build hash
  for (int32_t i = 0; i < topicNumGet; i++) {
X
Xiaoyu Wang 已提交
1189 1190
    SMqClientTopic topic = {0};
    SMqSubTopicEp* pTopicEp = taosArrayGet(pRsp->topics, i);
L
Liu Jicong 已提交
1191
    topic.schema = pTopicEp->schema;
L
Liu Jicong 已提交
1192
    taosHashClear(pHash);
X
Xiaoyu Wang 已提交
1193
    topic.topicName = strdup(pTopicEp->topic);
L
Liu Jicong 已提交
1194
    tstrncpy(topic.db, pTopicEp->db, TSDB_DB_FNAME_LEN);
L
Liu Jicong 已提交
1195

S
Shengliang Guan 已提交
1196
    tscDebug("consumer:%" PRId64 ", update topic: %s", tmq->consumerId, topic.topicName);
L
Liu Jicong 已提交
1197 1198 1199 1200 1201 1202
    int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
    for (int32_t j = 0; j < topicNumCur; j++) {
      // find old topic
      SMqClientTopic* pTopicCur = taosArrayGet(tmq->clientTopics, j);
      if (pTopicCur->vgs && strcmp(pTopicCur->topicName, pTopicEp->topic) == 0) {
        int32_t vgNumCur = taosArrayGetSize(pTopicCur->vgs);
S
Shengliang Guan 已提交
1203
        tscDebug("consumer:%" PRId64 ", new vg num: %d", tmq->consumerId, vgNumCur);
L
Liu Jicong 已提交
1204 1205 1206 1207
        if (vgNumCur == 0) break;
        for (int32_t k = 0; k < vgNumCur; k++) {
          SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, k);
          sprintf(vgKey, "%s:%d", topic.topicName, pVgCur->vgId);
S
Shengliang Guan 已提交
1208
          tscDebug("consumer:%" PRId64 ", epoch %d vgId:%d build %s", tmq->consumerId, epoch, pVgCur->vgId, vgKey);
L
Liu Jicong 已提交
1209 1210 1211 1212 1213 1214 1215 1216 1217
          taosHashPut(pHash, vgKey, strlen(vgKey), &pVgCur->currentOffset, sizeof(int64_t));
        }
        break;
      }
    }

    int32_t vgNumGet = taosArrayGetSize(pTopicEp->vgs);
    topic.vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));
    for (int32_t j = 0; j < vgNumGet; j++) {
X
Xiaoyu Wang 已提交
1218
      SMqSubVgEp* pVgEp = taosArrayGet(pTopicEp->vgs, j);
L
Liu Jicong 已提交
1219 1220 1221
      sprintf(vgKey, "%s:%d", topic.topicName, pVgEp->vgId);
      int64_t* pOffset = taosHashGet(pHash, vgKey, strlen(vgKey));
      int64_t  offset = pVgEp->offset;
S
Shengliang Guan 已提交
1222
      tscDebug("consumer:%" PRId64 ", (epoch %d) original offset of vgId:%d is %" PRId64, tmq->consumerId, epoch, pVgEp->vgId, offset);
L
Liu Jicong 已提交
1223 1224
      if (pOffset != NULL) {
        offset = *pOffset;
S
Shengliang Guan 已提交
1225
        tscDebug("consumer:%" PRId64 ", (epoch %d) receive offset of vgId:%d, full key is %s", tmq->consumerId, epoch, pVgEp->vgId,
1226
                 vgKey);
L
Liu Jicong 已提交
1227
      }
S
Shengliang Guan 已提交
1228
      tscDebug("consumer:%" PRId64 ", (epoch %d) offset of vgId:%d updated to %" PRId64, tmq->consumerId, epoch, pVgEp->vgId, offset);
X
Xiaoyu Wang 已提交
1229 1230
      SMqClientVg clientVg = {
          .pollCnt = 0,
L
Liu Jicong 已提交
1231
          .currentOffset = offset,
X
Xiaoyu Wang 已提交
1232 1233 1234
          .vgId = pVgEp->vgId,
          .epSet = pVgEp->epSet,
          .vgStatus = TMQ_VG_STATUS__IDLE,
L
Liu Jicong 已提交
1235
          .vgSkipCnt = 0,
X
Xiaoyu Wang 已提交
1236 1237 1238 1239
      };
      taosArrayPush(topic.vgs, &clientVg);
      set = true;
    }
L
Liu Jicong 已提交
1240
    taosArrayPush(newTopics, &topic);
X
Xiaoyu Wang 已提交
1241
  }
L
Liu Jicong 已提交
1242
  if (tmq->clientTopics) taosArrayDestroy(tmq->clientTopics);
L
Liu Jicong 已提交
1243
  taosHashCleanup(pHash);
L
Liu Jicong 已提交
1244
  tmq->clientTopics = newTopics;
1245

1246 1247 1248 1249
  if (taosArrayGetSize(tmq->clientTopics) == 0)
    atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__NO_TOPIC);
  else
    atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__READY);
1250

X
Xiaoyu Wang 已提交
1251 1252 1253
  atomic_store_32(&tmq->epoch, epoch);
  return set;
}
L
Liu Jicong 已提交
1254
#endif
X
Xiaoyu Wang 已提交
1255

D
dapan1121 已提交
1256
int32_t tmqAskEpCb(void* param, SDataBuf* pMsg, int32_t code) {
1257
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
L
Liu Jicong 已提交
1258
  tmq_t*           tmq = pParam->tmq;
L
Liu Jicong 已提交
1259
  int8_t           async = pParam->async;
L
Liu Jicong 已提交
1260
  pParam->code = code;
1261
  if (code != 0) {
S
Shengliang Guan 已提交
1262
    tscError("consumer:%" PRId64 ", get topic endpoint error, not ready, wait:%d", tmq->consumerId, pParam->async);
L
Liu Jicong 已提交
1263
    goto END;
1264
  }
L
Liu Jicong 已提交
1265

L
Liu Jicong 已提交
1266
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1267
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1268
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1269 1270
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
S
Shengliang Guan 已提交
1271
  tscDebug("consumer:%" PRId64 ", recv ep, msg epoch %d, current epoch %d", tmq->consumerId, head->epoch, epoch);
L
Liu Jicong 已提交
1272 1273
  if (head->epoch <= epoch) {
    goto END;
1274
  }
L
Liu Jicong 已提交
1275

L
Liu Jicong 已提交
1276
  if (!async) {
L
Liu Jicong 已提交
1277 1278
    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
S
Shengliang Guan 已提交
1279 1280
    /*printf("rsp epoch %" PRId64 " sz %" PRId64 "\n", rsp.epoch, rsp.topics->size);*/
    /*printf("tmq epoch %" PRId64 " sz %" PRId64 "\n", tmq->epoch, tmq->clientTopics->size);*/
L
Liu Jicong 已提交
1281
    tmqUpdateEp2(tmq, head->epoch, &rsp);
L
Liu Jicong 已提交
1282
    tDeleteSMqAskEpRsp(&rsp);
X
Xiaoyu Wang 已提交
1283
  } else {
1284
    SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM);
L
Liu Jicong 已提交
1285
    if (pWrapper == NULL) {
X
Xiaoyu Wang 已提交
1286
      terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1287 1288
      code = -1;
      goto END;
X
Xiaoyu Wang 已提交
1289
    }
L
Liu Jicong 已提交
1290 1291 1292
    pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
    pWrapper->epoch = head->epoch;
    memcpy(&pWrapper->msg, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1293
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pWrapper->msg);
L
Liu Jicong 已提交
1294

L
Liu Jicong 已提交
1295
    taosWriteQitem(tmq->mqueue, pWrapper);
1296
    tsem_post(&tmq->rspSem);
1297
  }
L
Liu Jicong 已提交
1298 1299

END:
L
Liu Jicong 已提交
1300
  /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1301
  if (!async) {
L
Liu Jicong 已提交
1302
    tsem_post(&pParam->rspSem);
L
Liu Jicong 已提交
1303 1304
  } else {
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
1305 1306
  }
  return code;
1307 1308
}

L
Liu Jicong 已提交
1309
int32_t tmqAskEp(tmq_t* tmq, bool async) {
L
Liu Jicong 已提交
1310
  int32_t code = 0;
L
Liu Jicong 已提交
1311
#if 0
L
Liu Jicong 已提交
1312
  int8_t  epStatus = atomic_val_compare_exchange_8(&tmq->epStatus, 0, 1);
L
fix txn  
Liu Jicong 已提交
1313
  if (epStatus == 1) {
L
temp  
Liu Jicong 已提交
1314
    int32_t epSkipCnt = atomic_add_fetch_32(&tmq->epSkipCnt, 1);
S
Shengliang Guan 已提交
1315
    tscTrace("consumer:%" PRId64 ", skip ask ep cnt %d", tmq->consumerId, epSkipCnt);
L
Liu Jicong 已提交
1316
    if (epSkipCnt < 5000) return 0;
L
fix txn  
Liu Jicong 已提交
1317
  }
L
temp  
Liu Jicong 已提交
1318
  atomic_store_32(&tmq->epSkipCnt, 0);
L
Liu Jicong 已提交
1319
#endif
L
Liu Jicong 已提交
1320
  int32_t      tlen = sizeof(SMqAskEpReq);
L
Liu Jicong 已提交
1321
  SMqAskEpReq* req = taosMemoryCalloc(1, tlen);
L
Liu Jicong 已提交
1322
  if (req == NULL) {
L
Liu Jicong 已提交
1323
    tscError("failed to malloc get subscribe ep buf");
L
Liu Jicong 已提交
1324
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1325
    return -1;
L
Liu Jicong 已提交
1326
  }
L
Liu Jicong 已提交
1327 1328 1329
  req->consumerId = htobe64(tmq->consumerId);
  req->epoch = htonl(tmq->epoch);
  strcpy(req->cgroup, tmq->groupId);
1330

L
Liu Jicong 已提交
1331
  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
L
Liu Jicong 已提交
1332 1333
  if (pParam == NULL) {
    tscError("failed to malloc subscribe param");
wafwerar's avatar
wafwerar 已提交
1334
    taosMemoryFree(req);
L
Liu Jicong 已提交
1335
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1336
    return -1;
L
Liu Jicong 已提交
1337 1338
  }
  pParam->tmq = tmq;
L
Liu Jicong 已提交
1339
  pParam->async = async;
X
Xiaoyu Wang 已提交
1340
  tsem_init(&pParam->rspSem, 0, 0);
1341

1342
  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
1343 1344
  if (sendInfo == NULL) {
    tsem_destroy(&pParam->rspSem);
wafwerar's avatar
wafwerar 已提交
1345 1346
    taosMemoryFree(pParam);
    taosMemoryFree(req);
L
Liu Jicong 已提交
1347
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1348 1349 1350 1351 1352 1353 1354 1355 1356 1357
    return -1;
  }

  sendInfo->msgInfo = (SDataBuf){
      .pData = req,
      .len = tlen,
      .handle = NULL,
  };

  sendInfo->requestId = generateRequestId();
L
Liu Jicong 已提交
1358 1359 1360
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqAskEpCb;
L
Liu Jicong 已提交
1361
  sendInfo->msgType = TDMT_MND_MQ_ASK_EP;
1362

L
Liu Jicong 已提交
1363 1364
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

S
Shengliang Guan 已提交
1365
  tscDebug("consumer:%" PRId64 ", ask ep", tmq->consumerId);
L
add log  
Liu Jicong 已提交
1366

L
Liu Jicong 已提交
1367 1368
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
1369

L
Liu Jicong 已提交
1370
  if (!async) {
L
Liu Jicong 已提交
1371 1372 1373 1374 1375
    tsem_wait(&pParam->rspSem);
    code = pParam->code;
    taosMemoryFree(pParam);
  }
  return code;
1376 1377
}

L
Liu Jicong 已提交
1378 1379
#if 0
int32_t tmq_seek(tmq_t* tmq, const tmq_topic_vgroup_t* offset) {
L
Liu Jicong 已提交
1380 1381 1382 1383 1384 1385 1386 1387 1388 1389 1390 1391 1392
  const SMqOffset* pOffset = &offset->offset;
  if (strcmp(pOffset->cgroup, tmq->groupId) != 0) {
    return TMQ_RESP_ERR__FAIL;
  }
  int32_t sz = taosArrayGetSize(tmq->clientTopics);
  for (int32_t i = 0; i < sz; i++) {
    SMqClientTopic* clientTopic = taosArrayGet(tmq->clientTopics, i);
    if (strcmp(clientTopic->topicName, pOffset->topicName) == 0) {
      int32_t vgSz = taosArrayGetSize(clientTopic->vgs);
      for (int32_t j = 0; j < vgSz; j++) {
        SMqClientVg* pVg = taosArrayGet(clientTopic->vgs, j);
        if (pVg->vgId == pOffset->vgId) {
          pVg->currentOffset = pOffset->offset;
L
Liu Jicong 已提交
1393
          tmqClearUnhandleMsg(tmq);
L
Liu Jicong 已提交
1394 1395 1396 1397 1398 1399 1400
          return TMQ_RESP_ERR__SUCCESS;
        }
      }
    }
  }
  return TMQ_RESP_ERR__FAIL;
}
L
Liu Jicong 已提交
1401
#endif
L
Liu Jicong 已提交
1402

1403
SMqPollReq* tmqBuildConsumeReqImpl(tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1404 1405 1406 1407 1408 1409 1410 1411 1412 1413
  /*int64_t reqOffset;*/
  /*if (pVg->currentOffset >= 0) {*/
  /*reqOffset = pVg->currentOffset;*/
  /*} else {*/
  /*if (tmq->resetOffsetCfg == TMQ_CONF__RESET_OFFSET__NONE) {*/
  /*tscError("unable to poll since no committed offset but reset offset is set to none");*/
  /*return NULL;*/
  /*}*/
  /*reqOffset = tmq->resetOffsetCfg;*/
  /*}*/
L
Liu Jicong 已提交
1414

L
Liu Jicong 已提交
1415
  SMqPollReq* pReq = taosMemoryCalloc(1, sizeof(SMqPollReq));
1416 1417 1418
  if (pReq == NULL) {
    return NULL;
  }
L
Liu Jicong 已提交
1419

L
Liu Jicong 已提交
1420 1421 1422
  /*strcpy(pReq->topic, pTopic->topicName);*/
  /*strcpy(pReq->cgroup, tmq->groupId);*/

L
Liu Jicong 已提交
1423 1424 1425 1426
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1427

1428
  pReq->withTbName = tmq->withTbName;
1429
  pReq->timeout = timeout;
L
Liu Jicong 已提交
1430
  pReq->consumerId = tmq->consumerId;
X
Xiaoyu Wang 已提交
1431
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1432 1433
  /*pReq->currentOffset = reqOffset;*/
  pReq->reqOffset = pVg->currentOffsetNew;
L
Liu Jicong 已提交
1434
  pReq->reqId = generateRequestId();
1435

L
Liu Jicong 已提交
1436 1437
  pReq->useSnapshot = tmq->useSnapshot;

1438
  pReq->head.vgId = htonl(pVg->vgId);
L
Liu Jicong 已提交
1439
  pReq->head.contLen = htonl(sizeof(SMqPollReq));
1440 1441 1442
  return pReq;
}

L
Liu Jicong 已提交
1443 1444
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1445
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1446 1447 1448 1449 1450 1451 1452 1453
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
  pRspObj->vgId = pWrapper->vgHandle->vgId;

  memcpy(&pRspObj->metaRsp, &pWrapper->metaRsp, sizeof(SMqMetaRsp));
  return pRspObj;
}

L
Liu Jicong 已提交
1454 1455 1456
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
L
Liu Jicong 已提交
1457 1458
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
L
Liu Jicong 已提交
1459
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1460
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1461
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1462

L
Liu Jicong 已提交
1463 1464
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
L
Liu Jicong 已提交
1465
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1466 1467
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1468

L
Liu Jicong 已提交
1469
  return pRspObj;
X
Xiaoyu Wang 已提交
1470 1471
}

1472
int32_t tmqPollImpl(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1473
  /*tscDebug("call poll");*/
X
Xiaoyu Wang 已提交
1474 1475 1476 1477 1478 1479
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
    for (int j = 0; j < taosArrayGetSize(pTopic->vgs); j++) {
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
      int32_t      vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
      if (vgStatus != TMQ_VG_STATUS__IDLE) {
L
Liu Jicong 已提交
1480
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
L
Liu Jicong 已提交
1481 1482
        tscTrace("consumer:%" PRId64 ", epoch %d skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch, pVg->vgId,
                 vgSkipCnt);
X
Xiaoyu Wang 已提交
1483
        continue;
L
Liu Jicong 已提交
1484
        /*if (vgSkipCnt < 10000) continue;*/
L
temp  
Liu Jicong 已提交
1485 1486 1487 1488
#if 0
        if (skipCnt < 30000) {
          continue;
        } else {
S
Shengliang Guan 已提交
1489
        tscDebug("consumer:%" PRId64 ",skip vgId:%d skip too much reset", tmq->consumerId, pVg->vgId);
L
temp  
Liu Jicong 已提交
1490 1491
        }
#endif
X
Xiaoyu Wang 已提交
1492
      }
L
Liu Jicong 已提交
1493
      atomic_store_32(&pVg->vgSkipCnt, 0);
1494
      SMqPollReq* pReq = tmqBuildConsumeReqImpl(tmq, timeout, pTopic, pVg);
X
Xiaoyu Wang 已提交
1495 1496
      if (pReq == NULL) {
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1497
        tsem_post(&tmq->rspSem);
X
Xiaoyu Wang 已提交
1498 1499
        return -1;
      }
wafwerar's avatar
wafwerar 已提交
1500
      SMqPollCbParam* pParam = taosMemoryMalloc(sizeof(SMqPollCbParam));
L
Liu Jicong 已提交
1501
      if (pParam == NULL) {
wafwerar's avatar
wafwerar 已提交
1502
        taosMemoryFree(pReq);
X
Xiaoyu Wang 已提交
1503
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1504
        tsem_post(&tmq->rspSem);
X
Xiaoyu Wang 已提交
1505 1506
        return -1;
      }
L
Liu Jicong 已提交
1507 1508
      pParam->tmq = tmq;
      pParam->pVg = pVg;
L
Liu Jicong 已提交
1509
      pParam->pTopic = pTopic;
L
Liu Jicong 已提交
1510
      pParam->vgId = pVg->vgId;
L
Liu Jicong 已提交
1511 1512
      pParam->epoch = tmq->epoch;

1513
      SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
1514
      if (sendInfo == NULL) {
wafwerar's avatar
wafwerar 已提交
1515 1516
        taosMemoryFree(pReq);
        taosMemoryFree(pParam);
L
Liu Jicong 已提交
1517
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1518
        tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1519 1520 1521 1522
        return -1;
      }

      sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1523
          .pData = pReq,
L
Liu Jicong 已提交
1524
          .len = sizeof(SMqPollReq),
X
Xiaoyu Wang 已提交
1525 1526
          .handle = NULL,
      };
L
Liu Jicong 已提交
1527
      sendInfo->requestId = pReq->reqId;
X
Xiaoyu Wang 已提交
1528
      sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1529
      sendInfo->param = pParam;
X
Xiaoyu Wang 已提交
1530
      sendInfo->fp = tmqPollCb;
L
Liu Jicong 已提交
1531
      sendInfo->msgType = TDMT_VND_CONSUME;
X
Xiaoyu Wang 已提交
1532 1533

      int64_t transporterId = 0;
L
fix  
Liu Jicong 已提交
1534
      /*printf("send poll\n");*/
L
Liu Jicong 已提交
1535

1536 1537
      char offsetFormatBuf[80];
      tFormatOffset(offsetFormatBuf, 80, &pVg->currentOffsetNew);
L
Liu Jicong 已提交
1538 1539
      tscDebug("consumer:%" PRId64 ", send poll to %s vgId:%d, epoch %d, req offset:%s, reqId:%" PRIu64,
               tmq->consumerId, pTopic->topicName, pVg->vgId, tmq->epoch, offsetFormatBuf, pReq->reqId);
S
Shengliang Guan 已提交
1540
      /*printf("send vgId:%d %" PRId64 "\n", pVg->vgId, pVg->currentOffset);*/
X
Xiaoyu Wang 已提交
1541 1542 1543 1544 1545 1546 1547 1548
      asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);
      pVg->pollCnt++;
      tmq->pollCnt++;
    }
  }
  return 0;
}

L
Liu Jicong 已提交
1549 1550
int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1551
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1552 1553
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1554
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
L
Liu Jicong 已提交
1555
      tmqUpdateEp2(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1556
      /*tmqClearUnhandleMsg(tmq);*/
X
Xiaoyu Wang 已提交
1557 1558 1559 1560 1561 1562 1563 1564 1565 1566
      *pReset = true;
    } else {
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

L
Liu Jicong 已提交
1567
void* tmqHandleAllRsp(tmq_t* tmq, int64_t timeout, bool pollIfReset) {
X
Xiaoyu Wang 已提交
1568
  while (1) {
L
Liu Jicong 已提交
1569 1570 1571
    SMqRspWrapper* rspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper == NULL) {
L
Liu Jicong 已提交
1572
      taosReadAllQitems(tmq->mqueue, tmq->qall);
L
Liu Jicong 已提交
1573 1574
      taosGetQitem(tmq->qall, (void**)&rspWrapper);
      if (rspWrapper == NULL) return NULL;
X
Xiaoyu Wang 已提交
1575 1576
    }

L
Liu Jicong 已提交
1577 1578 1579 1580 1581
    if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(rspWrapper);
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
      return NULL;
    } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1582
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1583
      /*atomic_sub_fetch_32(&tmq->readyRequest, 1);*/
1584
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
L
Liu Jicong 已提交
1585
      if (pollRspWrapper->dataRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1586
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
S
Shengliang Guan 已提交
1587
        /*printf("vgId:%d, offset %" PRId64 " up to %" PRId64 "\n", pVg->vgId, pVg->currentOffset,
L
Liu Jicong 已提交
1588
         * rspMsg->msg.rspOffset);*/
L
Liu Jicong 已提交
1589
        pVg->currentOffsetNew = pollRspWrapper->dataRsp.rspOffset;
X
Xiaoyu Wang 已提交
1590
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
L
Liu Jicong 已提交
1591
        if (pollRspWrapper->dataRsp.blockNum == 0) {
L
Liu Jicong 已提交
1592 1593
          taosFreeQitem(pollRspWrapper);
          rspWrapper = NULL;
L
temp  
Liu Jicong 已提交
1594 1595
          continue;
        }
L
Liu Jicong 已提交
1596
        // build rsp
L
Liu Jicong 已提交
1597
        SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1598
        taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1599
        return pRsp;
X
Xiaoyu Wang 已提交
1600
      } else {
L
Liu Jicong 已提交
1601 1602 1603 1604 1605 1606 1607
        tscDebug("msg discard since epoch mismatch: msg epoch %d, consumer epoch %d\n",
                 pollRspWrapper->dataRsp.head.epoch, consumerEpoch);
        taosFreeQitem(pollRspWrapper);
      }
    } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
L
Liu Jicong 已提交
1608
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1609
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
S
Shengliang Guan 已提交
1610
        /*printf("vgId:%d, offset %" PRId64 " up to %" PRId64 "\n", pVg->vgId, pVg->currentOffset,
L
Liu Jicong 已提交
1611
         * rspMsg->msg.rspOffset);*/
wmmhello's avatar
wmmhello 已提交
1612 1613
        pVg->currentOffsetNew.version = pollRspWrapper->metaRsp.rspOffset;
        pVg->currentOffsetNew.type = TMQ_OFFSET__LOG;
L
Liu Jicong 已提交
1614 1615
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1616
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1617 1618 1619 1620
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
        tscDebug("msg discard since epoch mismatch: msg epoch %d, consumer epoch %d\n",
L
Liu Jicong 已提交
1621
                 pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
L
Liu Jicong 已提交
1622
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1623 1624
      }
    } else {
L
fix  
Liu Jicong 已提交
1625
      /*printf("handle ep rsp %d\n", rspMsg->head.mqMsgType);*/
X
Xiaoyu Wang 已提交
1626
      bool reset = false;
L
Liu Jicong 已提交
1627 1628
      tmqHandleNoPollRsp(tmq, rspWrapper, &reset);
      taosFreeQitem(rspWrapper);
X
Xiaoyu Wang 已提交
1629
      if (pollIfReset && reset) {
S
Shengliang Guan 已提交
1630
        tscDebug("consumer:%" PRId64 ", reset and repoll", tmq->consumerId);
1631
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1632 1633 1634 1635 1636
      }
    }
  }
}

1637
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1638
  /*tscDebug("call poll1");*/
L
Liu Jicong 已提交
1639 1640
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1641

1642 1643 1644
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1645
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1646 1647
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1648
  }
1649
#endif
X
Xiaoyu Wang 已提交
1650

L
Liu Jicong 已提交
1651
  // in no topic status also need process delayed task
L
Liu Jicong 已提交
1652
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1653 1654 1655
    return NULL;
  }

X
Xiaoyu Wang 已提交
1656
  while (1) {
L
Liu Jicong 已提交
1657
    tmqHandleAllDelayedTask(tmq);
1658
    if (tmqPollImpl(tmq, timeout) < 0) return NULL;
L
Liu Jicong 已提交
1659

1660
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1661 1662
    if (rspObj) {
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1663 1664
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
      return NULL;
X
Xiaoyu Wang 已提交
1665
    }
1666
    if (timeout != -1) {
X
Xiaoyu Wang 已提交
1667
      int64_t endTime = taosGetTimestampMs();
1668
      int64_t leftTime = endTime - startTime;
1669
      if (leftTime > timeout) {
S
Shengliang Guan 已提交
1670
        tscDebug("consumer:%" PRId64 ", (epoch %d) timeout, no rsp", tmq->consumerId, tmq->epoch);
X
Xiaoyu Wang 已提交
1671 1672
        return NULL;
      }
1673
      tsem_timewait(&tmq->rspSem, leftTime * 1000);
L
Liu Jicong 已提交
1674 1675 1676
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
      tsem_timewait(&tmq->rspSem, 500 * 1000);
X
Xiaoyu Wang 已提交
1677 1678 1679 1680
    }
  }
}

L
Liu Jicong 已提交
1681
int32_t tmq_consumer_close(tmq_t* tmq) {
1682
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
L
Liu Jicong 已提交
1683 1684
    int32_t rsp = tmq_commit_sync(tmq, NULL);
    if (rsp != 0) {
L
Liu Jicong 已提交
1685
      return rsp;
1686 1687
    }

L
Liu Jicong 已提交
1688
    int32_t     retryCnt = 0;
1689
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
1690 1691 1692 1693 1694 1695 1696 1697 1698 1699
    while (1) {
      rsp = tmq_subscribe(tmq, lst);
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

1700
    tmq_list_destroy(lst);
1701

L
Liu Jicong 已提交
1702 1703
    /*return rsp;*/
    return 0;
L
Liu Jicong 已提交
1704
  }
1705
  // TODO: free resources
L
Liu Jicong 已提交
1706
  return 0;
1707
}
L
Liu Jicong 已提交
1708

L
Liu Jicong 已提交
1709 1710
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
1711
    return "success";
L
Liu Jicong 已提交
1712
  } else if (err == -1) {
L
Liu Jicong 已提交
1713 1714 1715
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
1716 1717
  }
}
L
Liu Jicong 已提交
1718

L
Liu Jicong 已提交
1719 1720 1721 1722
tmq_res_t tmq_get_res_type(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    return TMQ_RES_DATA;
  } else if (TD_RES_TMQ_META(res)) {
wmmhello's avatar
wmmhello 已提交
1723 1724 1725 1726
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_DELETE) {
      return TMQ_RES_DATA;
    }
L
Liu Jicong 已提交
1727 1728 1729 1730 1731 1732
    return TMQ_RES_TABLE_META;
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
1733
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
1734 1735
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
1736
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
1737 1738 1739
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
1740 1741 1742 1743 1744
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
1745 1746 1747 1748
const char* tmq_get_db_name(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
1749 1750 1751
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
L
Liu Jicong 已提交
1752 1753 1754 1755 1756
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
1757 1758 1759 1760
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
1761 1762 1763
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
L
Liu Jicong 已提交
1764 1765 1766 1767
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
1768 1769 1770 1771 1772 1773 1774 1775

const char* tmq_get_table_name(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
    }
1776
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
L
Liu Jicong 已提交
1777 1778 1779
  }
  return NULL;
}
1780

1781 1782
static char* buildCreateTableJson(SSchemaWrapper* schemaRow, SSchemaWrapper* schemaTag, char* name, int64_t id,
                                  int8_t t) {
wmmhello's avatar
wmmhello 已提交
1783 1784 1785 1786 1787 1788 1789
  char*  string = NULL;
  cJSON* json = cJSON_CreateObject();
  if (json == NULL) {
    return string;
  }
  cJSON* type = cJSON_CreateString("create");
  cJSON_AddItemToObject(json, "type", type);
wmmhello's avatar
wmmhello 已提交
1790

1791 1792 1793 1794
  //  char uid[32] = {0};
  //  sprintf(uid, "%"PRIi64, id);
  //  cJSON* id_ = cJSON_CreateString(uid);
  //  cJSON_AddItemToObject(json, "id", id_);
wmmhello's avatar
wmmhello 已提交
1795 1796 1797 1798
  cJSON* tableName = cJSON_CreateString(name);
  cJSON_AddItemToObject(json, "tableName", tableName);
  cJSON* tableType = cJSON_CreateString(t == TSDB_NORMAL_TABLE ? "normal" : "super");
  cJSON_AddItemToObject(json, "tableType", tableType);
1799 1800
  //  cJSON* version = cJSON_CreateNumber(1);
  //  cJSON_AddItemToObject(json, "version", version);
wmmhello's avatar
wmmhello 已提交
1801 1802

  cJSON* columns = cJSON_CreateArray();
1803 1804 1805 1806
  for (int i = 0; i < schemaRow->nCols; i++) {
    cJSON*   column = cJSON_CreateObject();
    SSchema* s = schemaRow->pSchema + i;
    cJSON*   cname = cJSON_CreateString(s->name);
wmmhello's avatar
wmmhello 已提交
1807 1808 1809
    cJSON_AddItemToObject(column, "name", cname);
    cJSON* ctype = cJSON_CreateNumber(s->type);
    cJSON_AddItemToObject(column, "type", ctype);
1810
    if (s->type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
1811
      int32_t length = s->bytes - VARSTR_HEADER_SIZE;
1812
      cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1813
      cJSON_AddItemToObject(column, "length", cbytes);
1814 1815 1816
    } else if (s->type == TSDB_DATA_TYPE_NCHAR) {
      int32_t length = (s->bytes - VARSTR_HEADER_SIZE) / TSDB_NCHAR_SIZE;
      cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1817 1818
      cJSON_AddItemToObject(column, "length", cbytes);
    }
wmmhello's avatar
wmmhello 已提交
1819 1820 1821 1822 1823
    cJSON_AddItemToArray(columns, column);
  }
  cJSON_AddItemToObject(json, "columns", columns);

  cJSON* tags = cJSON_CreateArray();
1824 1825 1826 1827
  for (int i = 0; schemaTag && i < schemaTag->nCols; i++) {
    cJSON*   tag = cJSON_CreateObject();
    SSchema* s = schemaTag->pSchema + i;
    cJSON*   tname = cJSON_CreateString(s->name);
wmmhello's avatar
wmmhello 已提交
1828 1829 1830
    cJSON_AddItemToObject(tag, "name", tname);
    cJSON* ttype = cJSON_CreateNumber(s->type);
    cJSON_AddItemToObject(tag, "type", ttype);
1831
    if (s->type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
1832
      int32_t length = s->bytes - VARSTR_HEADER_SIZE;
1833
      cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1834
      cJSON_AddItemToObject(tag, "length", cbytes);
1835 1836 1837
    } else if (s->type == TSDB_DATA_TYPE_NCHAR) {
      int32_t length = (s->bytes - VARSTR_HEADER_SIZE) / TSDB_NCHAR_SIZE;
      cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1838 1839
      cJSON_AddItemToObject(tag, "length", cbytes);
    }
wmmhello's avatar
wmmhello 已提交
1840 1841 1842 1843 1844 1845 1846 1847 1848
    cJSON_AddItemToArray(tags, tag);
  }
  cJSON_AddItemToObject(json, "tags", tags);

  string = cJSON_PrintUnformatted(json);
  cJSON_Delete(json);
  return string;
}

1849 1850 1851 1852
static char* buildAlterSTableJson(void* alterData, int32_t alterDataLen) {
  SMAlterStbReq req = {0};
  cJSON*        json = NULL;
  char*         string = NULL;
wmmhello's avatar
wmmhello 已提交
1853 1854 1855 1856 1857 1858 1859 1860 1861 1862 1863

  if (tDeserializeSMAlterStbReq(alterData, alterDataLen, &req) != 0) {
    goto end;
  }

  json = cJSON_CreateObject();
  if (json == NULL) {
    goto end;
  }
  cJSON* type = cJSON_CreateString("alter");
  cJSON_AddItemToObject(json, "type", type);
1864 1865
  //  cJSON* uid = cJSON_CreateNumber(id);
  //  cJSON_AddItemToObject(json, "uid", uid);
wmmhello's avatar
wmmhello 已提交
1866 1867 1868 1869 1870 1871 1872 1873 1874 1875 1876 1877
  SName name = {0};
  tNameFromString(&name, req.name, T_NAME_ACCT | T_NAME_DB | T_NAME_TABLE);
  cJSON* tableName = cJSON_CreateString(name.tname);
  cJSON_AddItemToObject(json, "tableName", tableName);
  cJSON* tableType = cJSON_CreateString("super");
  cJSON_AddItemToObject(json, "tableType", tableType);

  cJSON* alterType = cJSON_CreateNumber(req.alterType);
  cJSON_AddItemToObject(json, "alterType", alterType);
  switch (req.alterType) {
    case TSDB_ALTER_TABLE_ADD_TAG:
    case TSDB_ALTER_TABLE_ADD_COLUMN: {
1878 1879
      TAOS_FIELD* field = taosArrayGet(req.pFields, 0);
      cJSON*      colName = cJSON_CreateString(field->name);
wmmhello's avatar
wmmhello 已提交
1880 1881 1882 1883
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colType = cJSON_CreateNumber(field->type);
      cJSON_AddItemToObject(json, "colType", colType);

1884
      if (field->type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
1885
        int32_t length = field->bytes - VARSTR_HEADER_SIZE;
1886
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1887
        cJSON_AddItemToObject(json, "colLength", cbytes);
1888 1889 1890
      } else if (field->type == TSDB_DATA_TYPE_NCHAR) {
        int32_t length = (field->bytes - VARSTR_HEADER_SIZE) / TSDB_NCHAR_SIZE;
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1891 1892 1893 1894 1895
        cJSON_AddItemToObject(json, "colLength", cbytes);
      }
      break;
    }
    case TSDB_ALTER_TABLE_DROP_TAG:
1896 1897 1898
    case TSDB_ALTER_TABLE_DROP_COLUMN: {
      TAOS_FIELD* field = taosArrayGet(req.pFields, 0);
      cJSON*      colName = cJSON_CreateString(field->name);
wmmhello's avatar
wmmhello 已提交
1899 1900 1901 1902
      cJSON_AddItemToObject(json, "colName", colName);
      break;
    }
    case TSDB_ALTER_TABLE_UPDATE_TAG_BYTES:
1903 1904 1905
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_BYTES: {
      TAOS_FIELD* field = taosArrayGet(req.pFields, 0);
      cJSON*      colName = cJSON_CreateString(field->name);
wmmhello's avatar
wmmhello 已提交
1906 1907 1908
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colType = cJSON_CreateNumber(field->type);
      cJSON_AddItemToObject(json, "colType", colType);
1909
      if (field->type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
1910
        int32_t length = field->bytes - VARSTR_HEADER_SIZE;
1911
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1912
        cJSON_AddItemToObject(json, "colLength", cbytes);
1913 1914 1915
      } else if (field->type == TSDB_DATA_TYPE_NCHAR) {
        int32_t length = (field->bytes - VARSTR_HEADER_SIZE) / TSDB_NCHAR_SIZE;
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1916 1917 1918 1919 1920
        cJSON_AddItemToObject(json, "colLength", cbytes);
      }
      break;
    }
    case TSDB_ALTER_TABLE_UPDATE_TAG_NAME:
1921 1922 1923 1924
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_NAME: {
      TAOS_FIELD* oldField = taosArrayGet(req.pFields, 0);
      TAOS_FIELD* newField = taosArrayGet(req.pFields, 1);
      cJSON*      colName = cJSON_CreateString(oldField->name);
wmmhello's avatar
wmmhello 已提交
1925 1926 1927 1928 1929 1930 1931 1932 1933 1934
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colNewName = cJSON_CreateString(newField->name);
      cJSON_AddItemToObject(json, "colNewName", colNewName);
      break;
    }
    default:
      break;
  }
  string = cJSON_PrintUnformatted(json);

1935
end:
wmmhello's avatar
wmmhello 已提交
1936 1937 1938 1939 1940
  cJSON_Delete(json);
  tFreeSMAltertbReq(&req);
  return string;
}

1941
static char* processCreateStb(SMqMetaRsp* metaRsp) {
wmmhello's avatar
wmmhello 已提交
1942 1943
  SVCreateStbReq req = {0};
  SDecoder       coder;
1944
  char*          string = NULL;
wmmhello's avatar
wmmhello 已提交
1945 1946

  // decode and process req
1947
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
1948 1949 1950 1951 1952 1953 1954 1955 1956 1957
  int32_t len = metaRsp->metaRspLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);

  if (tDecodeSVCreateStbReq(&coder, &req) < 0) {
    goto _err;
  }
  string = buildCreateTableJson(&req.schemaRow, &req.schemaTag, req.name, req.suid, TSDB_SUPER_TABLE);
  tDecoderClear(&coder);
  return string;

1958
_err:
wmmhello's avatar
wmmhello 已提交
1959 1960 1961 1962
  tDecoderClear(&coder);
  return string;
}

1963
static char* processAlterStb(SMqMetaRsp* metaRsp) {
wmmhello's avatar
wmmhello 已提交
1964 1965
  SVCreateStbReq req = {0};
  SDecoder       coder;
1966
  char*          string = NULL;
wmmhello's avatar
wmmhello 已提交
1967 1968

  // decode and process req
1969
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
1970 1971 1972 1973 1974 1975 1976 1977 1978 1979
  int32_t len = metaRsp->metaRspLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);

  if (tDecodeSVCreateStbReq(&coder, &req) < 0) {
    goto _err;
  }
  string = buildAlterSTableJson(req.alterOriData, req.alterOriDataLen);
  tDecoderClear(&coder);
  return string;

1980
_err:
wmmhello's avatar
wmmhello 已提交
1981 1982 1983 1984
  tDecoderClear(&coder);
  return string;
}

1985
static char* buildCreateCTableJson(STag* pTag, char* sname, char* name, SArray* tagName, int64_t id, uint8_t tagNum) {
1986
  char*   string = NULL;
wmmhello's avatar
wmmhello 已提交
1987
  SArray* pTagVals = NULL;
1988
  cJSON*  json = cJSON_CreateObject();
wmmhello's avatar
wmmhello 已提交
1989 1990 1991 1992 1993
  if (json == NULL) {
    return string;
  }
  cJSON* type = cJSON_CreateString("create");
  cJSON_AddItemToObject(json, "type", type);
1994 1995 1996 1997
  //  char cid[32] = {0};
  //  sprintf(cid, "%"PRIi64, id);
  //  cJSON* cid_ = cJSON_CreateString(cid);
  //  cJSON_AddItemToObject(json, "id", cid_);
wmmhello's avatar
wmmhello 已提交
1998

wmmhello's avatar
wmmhello 已提交
1999 2000 2001 2002
  cJSON* tableName = cJSON_CreateString(name);
  cJSON_AddItemToObject(json, "tableName", tableName);
  cJSON* tableType = cJSON_CreateString("child");
  cJSON_AddItemToObject(json, "tableType", tableType);
wmmhello's avatar
wmmhello 已提交
2003
  cJSON* using = cJSON_CreateString(sname);
wmmhello's avatar
wmmhello 已提交
2004
  cJSON_AddItemToObject(json, "using", using);
2005 2006
  cJSON* tagNumJson = cJSON_CreateNumber(tagNum);
  cJSON_AddItemToObject(json, "tagNum", tagNumJson);
2007 2008
  //  cJSON* version = cJSON_CreateNumber(1);
  //  cJSON_AddItemToObject(json, "version", version);
wmmhello's avatar
wmmhello 已提交
2009

2010
  cJSON*  tags = cJSON_CreateArray();
wmmhello's avatar
wmmhello 已提交
2011 2012 2013 2014
  int32_t code = tTagToValArray(pTag, &pTagVals);
  if (code) {
    goto end;
  }
wmmhello's avatar
wmmhello 已提交
2015

wmmhello's avatar
wmmhello 已提交
2016 2017
  if (tTagIsJson(pTag)) {
    STag* p = (STag*)pTag;
2018
    if (p->nTag == 0) {
wmmhello's avatar
wmmhello 已提交
2019 2020
      goto end;
    }
2021 2022
    char*    pJson = parseTagDatatoJson(pTag);
    cJSON*   tag = cJSON_CreateObject();
wmmhello's avatar
wmmhello 已提交
2023 2024
    STagVal* pTagVal = taosArrayGet(pTagVals, 0);

wmmhello's avatar
wmmhello 已提交
2025 2026
    char*  ptname = taosArrayGet(tagName, 0);
    cJSON* tname = cJSON_CreateString(ptname);
wmmhello's avatar
wmmhello 已提交
2027
    cJSON_AddItemToObject(tag, "name", tname);
2028 2029
    //    cJSON* cid_ = cJSON_CreateString("");
    //    cJSON_AddItemToObject(tag, "cid", cid_);
wmmhello's avatar
wmmhello 已提交
2030 2031
    cJSON* ttype = cJSON_CreateNumber(TSDB_DATA_TYPE_JSON);
    cJSON_AddItemToObject(tag, "type", ttype);
wmmhello's avatar
wmmhello 已提交
2032
    cJSON* tvalue = cJSON_CreateString(pJson);
wmmhello's avatar
wmmhello 已提交
2033 2034
    cJSON_AddItemToObject(tag, "value", tvalue);
    cJSON_AddItemToArray(tags, tag);
wmmhello's avatar
wmmhello 已提交
2035
    taosMemoryFree(pJson);
wmmhello's avatar
wmmhello 已提交
2036 2037 2038
    goto end;
  }

2039
  for (int i = 0; i < taosArrayGetSize(pTagVals); i++) {
wmmhello's avatar
wmmhello 已提交
2040 2041 2042
    STagVal* pTagVal = (STagVal*)taosArrayGet(pTagVals, i);

    cJSON* tag = cJSON_CreateObject();
wmmhello's avatar
wmmhello 已提交
2043

wmmhello's avatar
wmmhello 已提交
2044 2045
    char*  ptname = taosArrayGet(tagName, i);
    cJSON* tname = cJSON_CreateString(ptname);
wmmhello's avatar
wmmhello 已提交
2046
    cJSON_AddItemToObject(tag, "name", tname);
2047 2048
    //    cJSON* cid = cJSON_CreateNumber(pTagVal->cid);
    //    cJSON_AddItemToObject(tag, "cid", cid);
wmmhello's avatar
wmmhello 已提交
2049 2050
    cJSON* ttype = cJSON_CreateNumber(pTagVal->type);
    cJSON_AddItemToObject(tag, "type", ttype);
wmmhello's avatar
wmmhello 已提交
2051

wmmhello's avatar
wmmhello 已提交
2052
    cJSON* tvalue = NULL;
wmmhello's avatar
wmmhello 已提交
2053
    if (IS_VAR_DATA_TYPE(pTagVal->type)) {
wmmhello's avatar
wmmhello 已提交
2054
      char* buf = taosMemoryCalloc(pTagVal->nData + 3, 1);
L
Liu Jicong 已提交
2055
      if (!buf) goto end;
wmmhello's avatar
wmmhello 已提交
2056
      dataConverToStr(buf, pTagVal->type, pTagVal->pData, pTagVal->nData, NULL);
wmmhello's avatar
wmmhello 已提交
2057 2058
      tvalue = cJSON_CreateString(buf);
      taosMemoryFree(buf);
wmmhello's avatar
wmmhello 已提交
2059
    } else {
wmmhello's avatar
wmmhello 已提交
2060 2061 2062
      double val = 0;
      GET_TYPED_DATA(val, double, pTagVal->type, &pTagVal->i64);
      tvalue = cJSON_CreateNumber(val);
wmmhello's avatar
wmmhello 已提交
2063 2064
    }

wmmhello's avatar
wmmhello 已提交
2065 2066 2067
    cJSON_AddItemToObject(tag, "value", tvalue);
    cJSON_AddItemToArray(tags, tag);
  }
wmmhello's avatar
wmmhello 已提交
2068

2069
end:
wmmhello's avatar
wmmhello 已提交
2070 2071 2072
  cJSON_AddItemToObject(json, "tags", tags);
  string = cJSON_PrintUnformatted(json);
  cJSON_Delete(json);
wmmhello's avatar
wmmhello 已提交
2073
  taosArrayDestroy(pTagVals);
wmmhello's avatar
wmmhello 已提交
2074 2075 2076
  return string;
}

2077
static char* processCreateTable(SMqMetaRsp* metaRsp) {
wmmhello's avatar
wmmhello 已提交
2078 2079
  SDecoder           decoder = {0};
  SVCreateTbBatchReq req = {0};
2080 2081
  SVCreateTbReq*     pCreateReq;
  char*              string = NULL;
wmmhello's avatar
wmmhello 已提交
2082
  // decode
2083
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2084 2085 2086 2087 2088 2089 2090 2091 2092
  int32_t len = metaRsp->metaRspLen - sizeof(SMsgHead);
  tDecoderInit(&decoder, data, len);
  if (tDecodeSVCreateTbBatchReq(&decoder, &req) < 0) {
    goto _exit;
  }

  // loop to create table
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
    pCreateReq = req.pReqs + iReq;
2093 2094
    if (pCreateReq->type == TSDB_CHILD_TABLE) {
      string = buildCreateCTableJson((STag*)pCreateReq->ctb.pTag, pCreateReq->ctb.name, pCreateReq->name,
2095
                                     pCreateReq->ctb.tagName, pCreateReq->uid, pCreateReq->ctb.tagNum);
2096 2097 2098
    } else if (pCreateReq->type == TSDB_NORMAL_TABLE) {
      string =
          buildCreateTableJson(&pCreateReq->ntb.schemaRow, NULL, pCreateReq->name, pCreateReq->uid, TSDB_NORMAL_TABLE);
wmmhello's avatar
wmmhello 已提交
2099 2100 2101 2102 2103
    }
  }

  tDecoderClear(&decoder);

2104
_exit:
wmmhello's avatar
wmmhello 已提交
2105 2106 2107 2108
  tDecoderClear(&decoder);
  return string;
}

2109 2110 2111 2112
static char* processAlterTable(SMqMetaRsp* metaRsp) {
  SDecoder     decoder = {0};
  SVAlterTbReq vAlterTbReq = {0};
  char*        string = NULL;
wmmhello's avatar
wmmhello 已提交
2113 2114

  // decode
2115
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2116 2117 2118 2119 2120 2121 2122 2123 2124 2125 2126 2127
  int32_t len = metaRsp->metaRspLen - sizeof(SMsgHead);
  tDecoderInit(&decoder, data, len);
  if (tDecodeSVAlterTbReq(&decoder, &vAlterTbReq) < 0) {
    goto _exit;
  }

  cJSON* json = cJSON_CreateObject();
  if (json == NULL) {
    goto _exit;
  }
  cJSON* type = cJSON_CreateString("alter");
  cJSON_AddItemToObject(json, "type", type);
2128 2129
  //  cJSON* uid = cJSON_CreateNumber(id);
  //  cJSON_AddItemToObject(json, "uid", uid);
wmmhello's avatar
wmmhello 已提交
2130 2131
  cJSON* tableName = cJSON_CreateString(vAlterTbReq.tbName);
  cJSON_AddItemToObject(json, "tableName", tableName);
wmmhello's avatar
wmmhello 已提交
2132
  cJSON* tableType = cJSON_CreateString(vAlterTbReq.action == TSDB_ALTER_TABLE_UPDATE_TAG_VAL ? "child" : "normal");
wmmhello's avatar
wmmhello 已提交
2133
  cJSON_AddItemToObject(json, "tableType", tableType);
wmmhello's avatar
wmmhello 已提交
2134 2135
  cJSON* alterType = cJSON_CreateNumber(vAlterTbReq.action);
  cJSON_AddItemToObject(json, "alterType", alterType);
wmmhello's avatar
wmmhello 已提交
2136 2137 2138 2139 2140 2141 2142

  switch (vAlterTbReq.action) {
    case TSDB_ALTER_TABLE_ADD_COLUMN: {
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colType = cJSON_CreateNumber(vAlterTbReq.type);
      cJSON_AddItemToObject(json, "colType", colType);
wmmhello's avatar
wmmhello 已提交
2143

2144
      if (vAlterTbReq.type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
2145
        int32_t length = vAlterTbReq.bytes - VARSTR_HEADER_SIZE;
2146
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
2147
        cJSON_AddItemToObject(json, "colLength", cbytes);
2148 2149 2150
      } else if (vAlterTbReq.type == TSDB_DATA_TYPE_NCHAR) {
        int32_t length = (vAlterTbReq.bytes - VARSTR_HEADER_SIZE) / TSDB_NCHAR_SIZE;
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
2151 2152
        cJSON_AddItemToObject(json, "colLength", cbytes);
      }
wmmhello's avatar
wmmhello 已提交
2153 2154
      break;
    }
2155
    case TSDB_ALTER_TABLE_DROP_COLUMN: {
wmmhello's avatar
wmmhello 已提交
2156 2157 2158 2159
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      break;
    }
2160
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_BYTES: {
wmmhello's avatar
wmmhello 已提交
2161 2162
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
wmmhello's avatar
wmmhello 已提交
2163
      cJSON* colType = cJSON_CreateNumber(vAlterTbReq.colModType);
wmmhello's avatar
wmmhello 已提交
2164
      cJSON_AddItemToObject(json, "colType", colType);
2165
      if (vAlterTbReq.colModType == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
2166
        int32_t length = vAlterTbReq.colModBytes - VARSTR_HEADER_SIZE;
2167
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
2168
        cJSON_AddItemToObject(json, "colLength", cbytes);
2169 2170 2171
      } else if (vAlterTbReq.colModType == TSDB_DATA_TYPE_NCHAR) {
        int32_t length = (vAlterTbReq.colModBytes - VARSTR_HEADER_SIZE) / TSDB_NCHAR_SIZE;
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
2172 2173
        cJSON_AddItemToObject(json, "colLength", cbytes);
      }
wmmhello's avatar
wmmhello 已提交
2174 2175
      break;
    }
2176
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_NAME: {
wmmhello's avatar
wmmhello 已提交
2177 2178 2179 2180 2181 2182
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colNewName = cJSON_CreateString(vAlterTbReq.colNewName);
      cJSON_AddItemToObject(json, "colNewName", colNewName);
      break;
    }
2183
    case TSDB_ALTER_TABLE_UPDATE_TAG_VAL: {
wmmhello's avatar
wmmhello 已提交
2184 2185
      cJSON* tagName = cJSON_CreateString(vAlterTbReq.tagName);
      cJSON_AddItemToObject(json, "colName", tagName);
wmmhello's avatar
wmmhello 已提交
2186

wmmhello's avatar
wmmhello 已提交
2187
      bool isNull = vAlterTbReq.isNull;
2188 2189 2190
      if (vAlterTbReq.tagType == TSDB_DATA_TYPE_JSON) {
        STag* jsonTag = (STag*)vAlterTbReq.pTagVal;
        if (jsonTag->nTag == 0) isNull = true;
wmmhello's avatar
wmmhello 已提交
2191
      }
2192
      if (!isNull) {
wmmhello's avatar
wmmhello 已提交
2193 2194 2195 2196 2197 2198 2199 2200 2201 2202 2203 2204 2205 2206 2207
        char* buf = NULL;

        if (vAlterTbReq.tagType == TSDB_DATA_TYPE_JSON) {
          ASSERT(tTagIsJson(vAlterTbReq.pTagVal) == true);
          buf = parseTagDatatoJson(vAlterTbReq.pTagVal);
        } else {
          buf = taosMemoryCalloc(vAlterTbReq.nTagVal + 1, 1);
          dataConverToStr(buf, vAlterTbReq.tagType, vAlterTbReq.pTagVal, vAlterTbReq.nTagVal, NULL);
        }

        cJSON* colValue = cJSON_CreateString(buf);
        cJSON_AddItemToObject(json, "colValue", colValue);
        taosMemoryFree(buf);
      }

wmmhello's avatar
wmmhello 已提交
2208 2209
      cJSON* isNullCJson = cJSON_CreateBool(isNull);
      cJSON_AddItemToObject(json, "colValueNull", isNullCJson);
wmmhello's avatar
wmmhello 已提交
2210 2211
      break;
    }
wmmhello's avatar
wmmhello 已提交
2212 2213 2214 2215 2216
    default:
      break;
  }
  string = cJSON_PrintUnformatted(json);

2217
_exit:
wmmhello's avatar
wmmhello 已提交
2218 2219 2220 2221
  tDecoderClear(&decoder);
  return string;
}

2222 2223 2224 2225
static char* processDropSTable(SMqMetaRsp* metaRsp) {
  SDecoder     decoder = {0};
  SVDropStbReq req = {0};
  char*        string = NULL;
wmmhello's avatar
wmmhello 已提交
2226 2227

  // decode
2228
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2229 2230 2231 2232 2233 2234 2235 2236 2237 2238 2239 2240 2241 2242 2243 2244 2245 2246 2247
  int32_t len = metaRsp->metaRspLen - sizeof(SMsgHead);
  tDecoderInit(&decoder, data, len);
  if (tDecodeSVDropStbReq(&decoder, &req) < 0) {
    goto _exit;
  }

  cJSON* json = cJSON_CreateObject();
  if (json == NULL) {
    goto _exit;
  }
  cJSON* type = cJSON_CreateString("drop");
  cJSON_AddItemToObject(json, "type", type);
  cJSON* tableName = cJSON_CreateString(req.name);
  cJSON_AddItemToObject(json, "tableName", tableName);
  cJSON* tableType = cJSON_CreateString("super");
  cJSON_AddItemToObject(json, "tableType", tableType);

  string = cJSON_PrintUnformatted(json);

2248
_exit:
wmmhello's avatar
wmmhello 已提交
2249 2250 2251 2252
  tDecoderClear(&decoder);
  return string;
}

2253 2254 2255 2256
static char* processDropTable(SMqMetaRsp* metaRsp) {
  SDecoder         decoder = {0};
  SVDropTbBatchReq req = {0};
  char*            string = NULL;
wmmhello's avatar
wmmhello 已提交
2257 2258

  // decode
2259
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2260 2261 2262 2263 2264 2265 2266 2267 2268 2269 2270 2271
  int32_t len = metaRsp->metaRspLen - sizeof(SMsgHead);
  tDecoderInit(&decoder, data, len);
  if (tDecodeSVDropTbBatchReq(&decoder, &req) < 0) {
    goto _exit;
  }

  cJSON* json = cJSON_CreateObject();
  if (json == NULL) {
    goto _exit;
  }
  cJSON* type = cJSON_CreateString("drop");
  cJSON_AddItemToObject(json, "type", type);
2272 2273 2274 2275
  //  cJSON* uid = cJSON_CreateNumber(id);
  //  cJSON_AddItemToObject(json, "uid", uid);
  //  cJSON* tableType = cJSON_CreateString("normal");
  //  cJSON_AddItemToObject(json, "tableType", tableType);
wmmhello's avatar
wmmhello 已提交
2276 2277 2278 2279 2280

  cJSON* tableNameList = cJSON_CreateArray();
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
    SVDropTbReq* pDropTbReq = req.pReqs + iReq;

wmmhello's avatar
wmmhello 已提交
2281
    cJSON* tableName = cJSON_CreateString(pDropTbReq->name);
wmmhello's avatar
wmmhello 已提交
2282 2283 2284 2285 2286 2287
    cJSON_AddItemToArray(tableNameList, tableName);
  }
  cJSON_AddItemToObject(json, "tableNameList", tableNameList);

  string = cJSON_PrintUnformatted(json);

2288
_exit:
wmmhello's avatar
wmmhello 已提交
2289 2290 2291 2292
  tDecoderClear(&decoder);
  return string;
}

2293
static int32_t taosCreateStb(TAOS* taos, void* meta, int32_t metaLen) {
wmmhello's avatar
wmmhello 已提交
2294 2295 2296
  SVCreateStbReq req = {0};
  SDecoder       coder;
  SMCreateStbReq pReq = {0};
2297 2298
  int32_t        code = TSDB_CODE_SUCCESS;
  SRequestObj*   pRequest = NULL;
wmmhello's avatar
wmmhello 已提交
2299

2300
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2301 2302 2303 2304
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2305
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2306 2307 2308 2309
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2310
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2311 2312 2313 2314 2315 2316 2317 2318
  int32_t len = metaLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);
  if (tDecodeSVCreateStbReq(&coder, &req) < 0) {
    code = TSDB_CODE_INVALID_PARA;
    goto end;
  }
  // build create stable
  pReq.pColumns = taosArrayInit(req.schemaRow.nCols, sizeof(SField));
2319
  for (int32_t i = 0; i < req.schemaRow.nCols; i++) {
wmmhello's avatar
wmmhello 已提交
2320
    SSchema* pSchema = req.schemaRow.pSchema + i;
2321
    SField   field = {.type = pSchema->type, .bytes = pSchema->bytes};
wmmhello's avatar
wmmhello 已提交
2322 2323 2324 2325
    strcpy(field.name, pSchema->name);
    taosArrayPush(pReq.pColumns, &field);
  }
  pReq.pTags = taosArrayInit(req.schemaTag.nCols, sizeof(SField));
2326
  for (int32_t i = 0; i < req.schemaTag.nCols; i++) {
wmmhello's avatar
wmmhello 已提交
2327
    SSchema* pSchema = req.schemaTag.pSchema + i;
2328
    SField   field = {.type = pSchema->type, .bytes = pSchema->bytes};
wmmhello's avatar
wmmhello 已提交
2329 2330 2331
    strcpy(field.name, pSchema->name);
    taosArrayPush(pReq.pTags, &field);
  }
2332

2333 2334
  pReq.colVer = req.schemaRow.version;
  pReq.tagVer = req.schemaTag.version;
wmmhello's avatar
wmmhello 已提交
2335 2336 2337 2338
  pReq.numOfColumns = req.schemaRow.nCols;
  pReq.numOfTags = req.schemaTag.nCols;
  pReq.commentLen = -1;
  pReq.suid = req.suid;
wmmhello's avatar
wmmhello 已提交
2339
  pReq.source = TD_REQ_FROM_TAOX;
wmmhello's avatar
wmmhello 已提交
2340
  pReq.igExists = true;
wmmhello's avatar
wmmhello 已提交
2341

2342
  STscObj* pTscObj = pRequest->pTscObj;
2343
  SName    tableName;
wmmhello's avatar
wmmhello 已提交
2344
  tNameExtractFullName(toName(pTscObj->acctId, pRequest->pDb, req.name, &tableName), pReq.name);
wmmhello's avatar
wmmhello 已提交
2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356

  SCmdMsgInfo pCmdMsg = {0};
  pCmdMsg.epSet = getEpSet_s(&pTscObj->pAppInfo->mgmtEp);
  pCmdMsg.msgType = TDMT_MND_CREATE_STB;
  pCmdMsg.msgLen = tSerializeSMCreateStbReq(NULL, 0, &pReq);
  pCmdMsg.pMsg = taosMemoryMalloc(pCmdMsg.msgLen);
  if (NULL == pCmdMsg.pMsg) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
  tSerializeSMCreateStbReq(pCmdMsg.pMsg, pCmdMsg.msgLen, &pReq);

2357
  SQuery pQuery = {0};
wmmhello's avatar
wmmhello 已提交
2358 2359 2360 2361 2362 2363
  pQuery.execMode = QUERY_EXEC_MODE_RPC;
  pQuery.pCmdMsg = &pCmdMsg;
  pQuery.msgType = pQuery.pCmdMsg->msgType;
  pQuery.stableQuery = true;

  launchQueryImpl(pRequest, &pQuery, true, NULL);
wmmhello's avatar
wmmhello 已提交
2364

L
Liu Jicong 已提交
2365 2366
  if (pRequest->code == TSDB_CODE_SUCCESS) {
    SCatalog* pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2367 2368 2369 2370
    catalogGetHandle(pTscObj->pAppInfo->clusterId, &pCatalog);
    catalogRemoveTableMeta(pCatalog, &tableName);
  }

wmmhello's avatar
wmmhello 已提交
2371 2372 2373
  code = pRequest->code;
  taosMemoryFree(pCmdMsg.pMsg);

2374
end:
wmmhello's avatar
wmmhello 已提交
2375 2376 2377 2378 2379 2380
  destroyRequest(pRequest);
  tFreeSMCreateStbReq(&pReq);
  tDecoderClear(&coder);
  return code;
}

2381
static int32_t taosDropStb(TAOS* taos, void* meta, int32_t metaLen) {
wmmhello's avatar
wmmhello 已提交
2382 2383 2384
  SVDropStbReq req = {0};
  SDecoder     coder;
  SMDropStbReq pReq = {0};
2385
  int32_t      code = TSDB_CODE_SUCCESS;
wmmhello's avatar
wmmhello 已提交
2386 2387
  SRequestObj* pRequest = NULL;

2388
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2389 2390 2391 2392
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2393
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2394 2395 2396 2397
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2398
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2399 2400 2401 2402 2403 2404 2405 2406 2407
  int32_t len = metaLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);
  if (tDecodeSVDropStbReq(&coder, &req) < 0) {
    code = TSDB_CODE_INVALID_PARA;
    goto end;
  }

  // build drop stable
  pReq.igNotExists = true;
wmmhello's avatar
wmmhello 已提交
2408
  pReq.source = TD_REQ_FROM_TAOX;
wmmhello's avatar
wmmhello 已提交
2409
  pReq.suid = req.suid;
2410 2411

  STscObj* pTscObj = pRequest->pTscObj;
wmmhello's avatar
wmmhello 已提交
2412
  SName    tableName = {0};
wmmhello's avatar
wmmhello 已提交
2413
  tNameExtractFullName(toName(pTscObj->acctId, pRequest->pDb, req.name, &tableName), pReq.name);
wmmhello's avatar
wmmhello 已提交
2414 2415 2416 2417 2418 2419 2420 2421 2422 2423 2424 2425

  SCmdMsgInfo pCmdMsg = {0};
  pCmdMsg.epSet = getEpSet_s(&pTscObj->pAppInfo->mgmtEp);
  pCmdMsg.msgType = TDMT_MND_DROP_STB;
  pCmdMsg.msgLen = tSerializeSMDropStbReq(NULL, 0, &pReq);
  pCmdMsg.pMsg = taosMemoryMalloc(pCmdMsg.msgLen);
  if (NULL == pCmdMsg.pMsg) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
  tSerializeSMDropStbReq(pCmdMsg.pMsg, pCmdMsg.msgLen, &pReq);

2426
  SQuery pQuery = {0};
wmmhello's avatar
wmmhello 已提交
2427 2428 2429 2430 2431 2432
  pQuery.execMode = QUERY_EXEC_MODE_RPC;
  pQuery.pCmdMsg = &pCmdMsg;
  pQuery.msgType = pQuery.pCmdMsg->msgType;
  pQuery.stableQuery = true;

  launchQueryImpl(pRequest, &pQuery, true, NULL);
wmmhello's avatar
wmmhello 已提交
2433

L
Liu Jicong 已提交
2434 2435
  if (pRequest->code == TSDB_CODE_SUCCESS) {
    SCatalog* pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2436 2437 2438 2439
    catalogGetHandle(pTscObj->pAppInfo->clusterId, &pCatalog);
    catalogRemoveTableMeta(pCatalog, &tableName);
  }

wmmhello's avatar
wmmhello 已提交
2440 2441 2442
  code = pRequest->code;
  taosMemoryFree(pCmdMsg.pMsg);

2443
end:
wmmhello's avatar
wmmhello 已提交
2444 2445 2446 2447 2448 2449 2450 2451 2452 2453 2454 2455
  destroyRequest(pRequest);
  tDecoderClear(&coder);
  return code;
}

typedef struct SVgroupCreateTableBatch {
  SVCreateTbBatchReq req;
  SVgroupInfo        info;
  char               dbName[TSDB_DB_NAME_LEN];
} SVgroupCreateTableBatch;

static void destroyCreateTbReqBatch(void* data) {
2456
  SVgroupCreateTableBatch* pTbBatch = (SVgroupCreateTableBatch*)data;
wmmhello's avatar
wmmhello 已提交
2457 2458 2459
  taosArrayDestroy(pTbBatch->req.pArray);
}

L
Liu Jicong 已提交
2460 2461 2462 2463 2464 2465 2466
static int32_t taosCreateTable(TAOS* taos, void* meta, int32_t metaLen) {
  SVCreateTbBatchReq req = {0};
  SDecoder           coder = {0};
  int32_t            code = TSDB_CODE_SUCCESS;
  SRequestObj*       pRequest = NULL;
  SQuery*            pQuery = NULL;
  SHashObj*          pVgroupHashmap = NULL;
wmmhello's avatar
wmmhello 已提交
2467

L
Liu Jicong 已提交
2468
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2469 2470 2471 2472
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2473
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2474 2475 2476 2477
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2478
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2479 2480 2481 2482 2483 2484 2485
  int32_t len = metaLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);
  if (tDecodeSVCreateTbBatchReq(&coder, &req) < 0) {
    code = TSDB_CODE_INVALID_PARA;
    goto end;
  }

2486 2487
  STscObj* pTscObj = pRequest->pTscObj;

L
Liu Jicong 已提交
2488 2489
  SVCreateTbReq* pCreateReq = NULL;
  SCatalog*      pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2490 2491 2492 2493 2494 2495 2496 2497 2498 2499 2500 2501 2502
  code = catalogGetHandle(pTscObj->pAppInfo->clusterId, &pCatalog);
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

  pVgroupHashmap = taosHashInit(4, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), false, HASH_NO_LOCK);
  if (NULL == pVgroupHashmap) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
  taosHashSetFreeFp(pVgroupHashmap, destroyCreateTbReqBatch);

  SRequestConnInfo conn = {.pTrans = pTscObj->pAppInfo->pTransporter,
2503 2504 2505
                           .requestId = pRequest->requestId,
                           .requestObjRefId = pRequest->self,
                           .mgmtEps = getEpSet_s(&pTscObj->pAppInfo->mgmtEp)};
wmmhello's avatar
wmmhello 已提交
2506 2507

  pRequest->tableList = taosArrayInit(req.nReqs, sizeof(SName));
wmmhello's avatar
wmmhello 已提交
2508 2509 2510 2511 2512
  // loop to create table
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
    pCreateReq = req.pReqs + iReq;

    SVgroupInfo pInfo = {0};
wmmhello's avatar
wmmhello 已提交
2513
    SName       pName = {0};
wmmhello's avatar
wmmhello 已提交
2514 2515 2516 2517 2518
    toName(pTscObj->acctId, pRequest->pDb, pCreateReq->name, &pName);
    code = catalogGetTableHashVgroup(pCatalog, &conn, &pName, &pInfo);
    if (code != TSDB_CODE_SUCCESS) {
      goto end;
    }
wmmhello's avatar
wmmhello 已提交
2519
    taosArrayPush(pRequest->tableList, &pName);
wmmhello's avatar
wmmhello 已提交
2520 2521 2522 2523 2524 2525 2526 2527 2528 2529 2530 2531 2532 2533 2534 2535 2536 2537 2538 2539 2540 2541 2542 2543 2544 2545

    SVgroupCreateTableBatch* pTableBatch = taosHashGet(pVgroupHashmap, &pInfo.vgId, sizeof(pInfo.vgId));
    if (pTableBatch == NULL) {
      SVgroupCreateTableBatch tBatch = {0};
      tBatch.info = pInfo;
      strcpy(tBatch.dbName, pRequest->pDb);

      tBatch.req.pArray = taosArrayInit(4, sizeof(struct SVCreateTbReq));
      taosArrayPush(tBatch.req.pArray, pCreateReq);

      taosHashPut(pVgroupHashmap, &pInfo.vgId, sizeof(pInfo.vgId), &tBatch, sizeof(tBatch));
    } else {  // add to the correct vgroup
      taosArrayPush(pTableBatch->req.pArray, pCreateReq);
    }
  }

  SArray* pBufArray = serializeVgroupsCreateTableBatch(pVgroupHashmap);
  if (NULL == pBufArray) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }

  pQuery = (SQuery*)nodesMakeNode(QUERY_NODE_QUERY);
  pQuery->execMode = QUERY_EXEC_MODE_SCHEDULE;
  pQuery->msgType = TDMT_VND_CREATE_TABLE;
  pQuery->stableQuery = false;
2546
  pQuery->pRoot = nodesMakeNode(QUERY_NODE_CREATE_TABLE_STMT);
wmmhello's avatar
wmmhello 已提交
2547 2548 2549 2550 2551 2552

  code = rewriteToVnodeModifyOpStmt(pQuery, pBufArray);
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

wmmhello's avatar
wmmhello 已提交
2553
  launchQueryImpl(pRequest, pQuery, true, NULL);
L
Liu Jicong 已提交
2554
  if (pRequest->code == TSDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
2555 2556 2557
    removeMeta(pTscObj, pRequest->tableList);
  }

2558
  code = pRequest->code;
wmmhello's avatar
wmmhello 已提交
2559

2560
end:
wmmhello's avatar
wmmhello 已提交
2561 2562 2563 2564 2565 2566 2567 2568 2569 2570 2571 2572 2573 2574 2575 2576 2577 2578
  taosHashCleanup(pVgroupHashmap);
  destroyRequest(pRequest);
  tDecoderClear(&coder);
  qDestroyQuery(pQuery);
  return code;
}

typedef struct SVgroupDropTableBatch {
  SVDropTbBatchReq req;
  SVgroupInfo      info;
  char             dbName[TSDB_DB_NAME_LEN];
} SVgroupDropTableBatch;

static void destroyDropTbReqBatch(void* data) {
  SVgroupDropTableBatch* pTbBatch = (SVgroupDropTableBatch*)data;
  taosArrayDestroy(pTbBatch->req.pArray);
}

L
Liu Jicong 已提交
2579 2580 2581 2582 2583 2584 2585
static int32_t taosDropTable(TAOS* taos, void* meta, int32_t metaLen) {
  SVDropTbBatchReq req = {0};
  SDecoder         coder = {0};
  int32_t          code = TSDB_CODE_SUCCESS;
  SRequestObj*     pRequest = NULL;
  SQuery*          pQuery = NULL;
  SHashObj*        pVgroupHashmap = NULL;
wmmhello's avatar
wmmhello 已提交
2586

2587
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2588 2589 2590 2591
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2592
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2593 2594 2595 2596
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2597
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2598 2599 2600 2601 2602 2603 2604
  int32_t len = metaLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);
  if (tDecodeSVDropTbBatchReq(&coder, &req) < 0) {
    code = TSDB_CODE_INVALID_PARA;
    goto end;
  }

2605 2606
  STscObj* pTscObj = pRequest->pTscObj;

L
Liu Jicong 已提交
2607 2608
  SVDropTbReq* pDropReq = NULL;
  SCatalog*    pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2609 2610 2611 2612 2613 2614 2615 2616 2617 2618 2619 2620 2621
  code = catalogGetHandle(pTscObj->pAppInfo->clusterId, &pCatalog);
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

  pVgroupHashmap = taosHashInit(4, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), false, HASH_NO_LOCK);
  if (NULL == pVgroupHashmap) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
  taosHashSetFreeFp(pVgroupHashmap, destroyDropTbReqBatch);

  SRequestConnInfo conn = {.pTrans = pTscObj->pAppInfo->pTransporter,
2622 2623 2624
                           .requestId = pRequest->requestId,
                           .requestObjRefId = pRequest->self,
                           .mgmtEps = getEpSet_s(&pTscObj->pAppInfo->mgmtEp)};
wmmhello's avatar
wmmhello 已提交
2625
  pRequest->tableList = taosArrayInit(req.nReqs, sizeof(SName));
wmmhello's avatar
wmmhello 已提交
2626 2627 2628
  // loop to create table
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
    pDropReq = req.pReqs + iReq;
wmmhello's avatar
wmmhello 已提交
2629
    pDropReq->igNotExists = true;
wmmhello's avatar
wmmhello 已提交
2630 2631

    SVgroupInfo pInfo = {0};
wmmhello's avatar
wmmhello 已提交
2632
    SName       pName = {0};
wmmhello's avatar
wmmhello 已提交
2633 2634 2635 2636 2637 2638
    toName(pTscObj->acctId, pRequest->pDb, pDropReq->name, &pName);
    code = catalogGetTableHashVgroup(pCatalog, &conn, &pName, &pInfo);
    if (code != TSDB_CODE_SUCCESS) {
      goto end;
    }

wmmhello's avatar
wmmhello 已提交
2639
    taosArrayPush(pRequest->tableList, &pName);
wmmhello's avatar
wmmhello 已提交
2640 2641 2642 2643 2644 2645 2646 2647 2648 2649 2650 2651 2652 2653 2654 2655 2656 2657 2658 2659 2660 2661 2662
    SVgroupDropTableBatch* pTableBatch = taosHashGet(pVgroupHashmap, &pInfo.vgId, sizeof(pInfo.vgId));
    if (pTableBatch == NULL) {
      SVgroupDropTableBatch tBatch = {0};
      tBatch.info = pInfo;
      tBatch.req.pArray = taosArrayInit(TARRAY_MIN_SIZE, sizeof(SVDropTbReq));
      taosArrayPush(tBatch.req.pArray, pDropReq);

      taosHashPut(pVgroupHashmap, &pInfo.vgId, sizeof(pInfo.vgId), &tBatch, sizeof(tBatch));
    } else {  // add to the correct vgroup
      taosArrayPush(pTableBatch->req.pArray, pDropReq);
    }
  }

  SArray* pBufArray = serializeVgroupsDropTableBatch(pVgroupHashmap);
  if (NULL == pBufArray) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }

  pQuery = (SQuery*)nodesMakeNode(QUERY_NODE_QUERY);
  pQuery->execMode = QUERY_EXEC_MODE_SCHEDULE;
  pQuery->msgType = TDMT_VND_DROP_TABLE;
  pQuery->stableQuery = false;
2663
  pQuery->pRoot = nodesMakeNode(QUERY_NODE_DROP_TABLE_STMT);
wmmhello's avatar
wmmhello 已提交
2664 2665 2666 2667 2668 2669

  code = rewriteToVnodeModifyOpStmt(pQuery, pBufArray);
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

wmmhello's avatar
wmmhello 已提交
2670
  launchQueryImpl(pRequest, pQuery, true, NULL);
L
Liu Jicong 已提交
2671
  if (pRequest->code == TSDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
2672 2673
    removeMeta(pTscObj, pRequest->tableList);
  }
2674
  code = pRequest->code;
wmmhello's avatar
wmmhello 已提交
2675

2676
end:
wmmhello's avatar
wmmhello 已提交
2677 2678 2679 2680 2681 2682 2683
  taosHashCleanup(pVgroupHashmap);
  destroyRequest(pRequest);
  tDecoderClear(&coder);
  qDestroyQuery(pQuery);
  return code;
}

wmmhello's avatar
wmmhello 已提交
2684 2685
// delete from db.tabl where ..       -> delete from tabl where ..
// delete from db    .tabl where ..   -> delete from tabl where ..
L
Liu Jicong 已提交
2686
// static void getTbName(char *sql){
2687 2688 2689 2690 2691 2692 2693 2694 2695 2696 2697 2698 2699 2700 2701 2702 2703 2704 2705 2706 2707 2708 2709 2710 2711 2712 2713 2714
//  char *ch = sql;
//
//  bool inBackQuote = false;
//  int8_t dotIndex = 0;
//  while(*ch != '\0'){
//    if(!inBackQuote && *ch == '`'){
//      inBackQuote = true;
//      ch++;
//      continue;
//    }
//
//    if(inBackQuote && *ch == '`'){
//      inBackQuote = false;
//      ch++;
//
//      continue;
//    }
//
//    if(!inBackQuote && *ch == '.'){
//      dotIndex ++;
//      if(dotIndex == 2){
//        memmove(sql, ch + 1, strlen(ch + 1) + 1);
//        break;
//      }
//    }
//    ch++;
//  }
//}
wmmhello's avatar
wmmhello 已提交
2715 2716

static int32_t taosDeleteData(TAOS* taos, void* meta, int32_t metaLen) {
L
Liu Jicong 已提交
2717 2718 2719
  SDeleteRes req = {0};
  SDecoder   coder = {0};
  int32_t    code = TSDB_CODE_SUCCESS;
wmmhello's avatar
wmmhello 已提交
2720 2721 2722 2723 2724 2725 2726 2727 2728 2729

  // decode and process req
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
  int32_t len = metaLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);
  if (tDecodeDeleteRes(&coder, &req) < 0) {
    code = TSDB_CODE_INVALID_PARA;
    goto end;
  }

L
Liu Jicong 已提交
2730
  //  getTbName(req.tableFName);
wmmhello's avatar
wmmhello 已提交
2731
  char sql[256] = {0};
L
Liu Jicong 已提交
2732 2733
  sprintf(sql, "delete from `%s` where `%s` >= %" PRId64 " and `%s` <= %" PRId64, req.tableFName, req.tsColName,
          req.skey, req.tsColName, req.ekey);
wmmhello's avatar
wmmhello 已提交
2734 2735
  printf("delete sql:%s\n", sql);

L
Liu Jicong 已提交
2736 2737
  TAOS_RES*    res = taos_query(taos, sql);
  SRequestObj* pRequest = (SRequestObj*)res;
wmmhello's avatar
wmmhello 已提交
2738 2739 2740 2741 2742 2743 2744 2745 2746 2747 2748
  code = pRequest->code;
  if (code == TSDB_CODE_PAR_TABLE_NOT_EXIST) {
    code = TSDB_CODE_SUCCESS;
  }
  taos_free_result(res);

end:
  tDecoderClear(&coder);
  return code;
}

2749 2750 2751 2752 2753 2754 2755 2756
static int32_t taosAlterTable(TAOS* taos, void* meta, int32_t metaLen) {
  SVAlterTbReq   req = {0};
  SDecoder       coder = {0};
  int32_t        code = TSDB_CODE_SUCCESS;
  SRequestObj*   pRequest = NULL;
  SQuery*        pQuery = NULL;
  SArray*        pArray = NULL;
  SVgDataBlocks* pVgData = NULL;
wmmhello's avatar
wmmhello 已提交
2757

2758
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2759 2760 2761 2762 2763

  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2764
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2765 2766 2767 2768
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2769
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2770 2771 2772 2773 2774 2775 2776 2777
  int32_t len = metaLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);
  if (tDecodeSVAlterTbReq(&coder, &req) < 0) {
    code = TSDB_CODE_INVALID_PARA;
    goto end;
  }

  // do not deal TSDB_ALTER_TABLE_UPDATE_OPTIONS
2778
  if (req.action == TSDB_ALTER_TABLE_UPDATE_OPTIONS) {
wmmhello's avatar
wmmhello 已提交
2779 2780 2781
    goto end;
  }

2782
  STscObj*  pTscObj = pRequest->pTscObj;
2783
  SCatalog* pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2784 2785 2786 2787 2788 2789
  code = catalogGetHandle(pTscObj->pAppInfo->clusterId, &pCatalog);
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

  SRequestConnInfo conn = {.pTrans = pTscObj->pAppInfo->pTransporter,
2790 2791 2792
                           .requestId = pRequest->requestId,
                           .requestObjRefId = pRequest->self,
                           .mgmtEps = getEpSet_s(&pTscObj->pAppInfo->mgmtEp)};
wmmhello's avatar
wmmhello 已提交
2793 2794

  SVgroupInfo pInfo = {0};
2795
  SName       pName = {0};
wmmhello's avatar
wmmhello 已提交
2796 2797 2798 2799 2800 2801 2802 2803 2804 2805 2806 2807 2808 2809 2810 2811 2812 2813 2814 2815 2816 2817 2818 2819 2820 2821 2822 2823 2824 2825 2826 2827 2828
  toName(pTscObj->acctId, pRequest->pDb, req.tbName, &pName);
  code = catalogGetTableHashVgroup(pCatalog, &conn, &pName, &pInfo);
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

  pArray = taosArrayInit(1, sizeof(void*));
  if (NULL == pArray) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }

  pVgData = taosMemoryCalloc(1, sizeof(SVgDataBlocks));
  if (NULL == pVgData) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
  pVgData->vg = pInfo;
  pVgData->pData = taosMemoryMalloc(metaLen);
  if (NULL == pVgData->pData) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
  memcpy(pVgData->pData, meta, metaLen);
  ((SMsgHead*)pVgData->pData)->vgId = htonl(pInfo.vgId);
  pVgData->size = metaLen;
  pVgData->numOfTables = 1;
  taosArrayPush(pArray, &pVgData);

  pQuery = (SQuery*)nodesMakeNode(QUERY_NODE_QUERY);
  pQuery->execMode = QUERY_EXEC_MODE_SCHEDULE;
  pQuery->msgType = TDMT_VND_ALTER_TABLE;
  pQuery->stableQuery = false;
2829
  pQuery->pRoot = nodesMakeNode(QUERY_NODE_ALTER_TABLE_STMT);
wmmhello's avatar
wmmhello 已提交
2830 2831 2832 2833 2834 2835

  code = rewriteToVnodeModifyOpStmt(pQuery, pArray);
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

wmmhello's avatar
wmmhello 已提交
2836 2837
  launchQueryImpl(pRequest, pQuery, true, NULL);

wmmhello's avatar
wmmhello 已提交
2838
  pVgData = NULL;
2839 2840 2841
  pArray = NULL;
  code = pRequest->code;
  if (code == TSDB_CODE_VND_TABLE_NOT_EXIST) {
wmmhello's avatar
wmmhello 已提交
2842
    code = TSDB_CODE_SUCCESS;
wmmhello's avatar
wmmhello 已提交
2843
  }
wmmhello's avatar
wmmhello 已提交
2844

L
Liu Jicong 已提交
2845
  if (pRequest->code == TSDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
2846
    SExecResult* pRes = &pRequest->body.resInfo.execRes;
L
Liu Jicong 已提交
2847
    if (pRes->res != NULL) {
wmmhello's avatar
wmmhello 已提交
2848 2849 2850
      code = handleAlterTbExecRes(pRes->res, pCatalog);
    }
  }
2851
end:
wmmhello's avatar
wmmhello 已提交
2852
  taosArrayDestroy(pArray);
2853
  if (pVgData) taosMemoryFreeClear(pVgData->pData);
wmmhello's avatar
wmmhello 已提交
2854 2855 2856 2857 2858 2859 2860
  taosMemoryFreeClear(pVgData);
  destroyRequest(pRequest);
  tDecoderClear(&coder);
  qDestroyQuery(pQuery);
  return code;
}

L
Liu Jicong 已提交
2861
typedef struct {
2862
  SVgroupInfo vg;
L
Liu Jicong 已提交
2863 2864
  void*       data;
} VgData;
2865 2866 2867 2868 2869 2870

static void destroyVgHash(void* data) {
  VgData* vgData = (VgData*)data;
  taosMemoryFreeClear(vgData->data);
}

L
Liu Jicong 已提交
2871 2872
int taos_write_raw_block(TAOS* taos, int rows, char* pData, const char* tbname) {
  int32_t     code = TSDB_CODE_SUCCESS;
wmmhello's avatar
wmmhello 已提交
2873
  STableMeta* pTableMeta = NULL;
L
Liu Jicong 已提交
2874
  SQuery*     pQuery = NULL;
wmmhello's avatar
wmmhello 已提交
2875 2876

  SRequestObj* pRequest = (SRequestObj*)createRequest(*(int64_t*)taos, TSDB_SQL_INSERT);
L
Liu Jicong 已提交
2877
  if (!pRequest) {
wmmhello's avatar
wmmhello 已提交
2878 2879 2880
    uError("WriteRaw:createRequest error request is null");
    code = terrno;
    goto end;
2881 2882
  }

wmmhello's avatar
wmmhello 已提交
2883 2884 2885 2886 2887 2888 2889 2890 2891 2892
  if (!pRequest->pDb) {
    uError("WriteRaw:not use db");
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }

  SName pName = {TSDB_TABLE_NAME_T, pRequest->pTscObj->acctId, {0}, {0}};
  strcpy(pName.dbname, pRequest->pDb);
  strcpy(pName.tname, tbname);

L
Liu Jicong 已提交
2893
  struct SCatalog* pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2894
  code = catalogGetHandle(pRequest->pTscObj->pAppInfo->clusterId, &pCatalog);
L
Liu Jicong 已提交
2895
  if (code != TSDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
2896 2897 2898 2899 2900 2901 2902 2903 2904 2905 2906 2907 2908 2909 2910 2911 2912 2913 2914 2915 2916 2917 2918 2919
    uError("WriteRaw: get gatlog error");
    goto end;
  }

  SRequestConnInfo conn = {0};
  conn.pTrans = pRequest->pTscObj->pAppInfo->pTransporter;
  conn.requestId = pRequest->requestId;
  conn.requestObjRefId = pRequest->self;
  conn.mgmtEps = getEpSet_s(&pRequest->pTscObj->pAppInfo->mgmtEp);

  SVgroupInfo vgData = {0};
  code = catalogGetTableHashVgroup(pCatalog, &conn, &pName, &vgData);
  if (code != TSDB_CODE_SUCCESS) {
    uError("WriteRaw:catalogGetTableHashVgroup failed. table name: %s", tbname);
    goto end;
  }

  code = catalogGetTableMeta(pCatalog, &conn, &pName, &pTableMeta);
  if (code != TSDB_CODE_SUCCESS) {
    uError("WriteRaw:catalogGetTableMeta failed. table name: %s", tbname);
    goto end;
  }
  uint64_t suid = (TSDB_NORMAL_TABLE == pTableMeta->tableType ? 0 : pTableMeta->suid);
  uint64_t uid = pTableMeta->uid;
L
Liu Jicong 已提交
2920
  int32_t  numOfCols = pTableMeta->tableInfo.numOfColumns;
wmmhello's avatar
wmmhello 已提交
2921 2922

  uint16_t fLen = 0;
L
Liu Jicong 已提交
2923 2924
  int32_t  rowSize = 0;
  int16_t  nVar = 0;
wmmhello's avatar
wmmhello 已提交
2925
  for (int i = 0; i < numOfCols; i++) {
L
Liu Jicong 已提交
2926
    SSchema* schema = pTableMeta->schema + i;
wmmhello's avatar
wmmhello 已提交
2927 2928
    fLen += TYPE_BYTES[schema->type];
    rowSize += schema->bytes;
L
Liu Jicong 已提交
2929 2930
    if (IS_VAR_DATA_TYPE(schema->type)) {
      nVar++;
wmmhello's avatar
wmmhello 已提交
2931 2932 2933 2934 2935 2936 2937 2938
    }
  }

  int32_t extendedRowSize = rowSize + TD_ROW_HEAD_LEN - sizeof(TSKEY) + nVar * sizeof(VarDataOffsetT) +
                            (int32_t)TD_BITMAP_BYTES(numOfCols - 1);
  int32_t schemaLen = 0;
  int32_t submitLen = sizeof(SSubmitBlk) + schemaLen + rows * extendedRowSize;

L
Liu Jicong 已提交
2939
  int32_t     totalLen = sizeof(SSubmitReq) + submitLen;
wmmhello's avatar
wmmhello 已提交
2940 2941
  SSubmitReq* subReq = taosMemoryCalloc(1, totalLen);
  SSubmitBlk* blk = POINTER_SHIFT(subReq, sizeof(SSubmitReq));
L
Liu Jicong 已提交
2942 2943
  void*       blkSchema = POINTER_SHIFT(blk, sizeof(SSubmitBlk));
  STSRow*     rowData = POINTER_SHIFT(blkSchema, schemaLen);
wmmhello's avatar
wmmhello 已提交
2944 2945 2946 2947 2948 2949

  SRowBuilder rb = {0};
  tdSRowInit(&rb, pTableMeta->sversion);
  tdSRowSetTpInfo(&rb, numOfCols, fLen);
  int32_t dataLen = 0;

H
Haojun Liao 已提交
2950
  char*    pStart = pData + getVersion1BlockMetaSize(pData, numOfCols);
wmmhello's avatar
wmmhello 已提交
2951 2952 2953
  int32_t* colLength = (int32_t*)pStart;
  pStart += sizeof(int32_t) * numOfCols;

L
Liu Jicong 已提交
2954
  SResultColumn* pCol = taosMemoryCalloc(numOfCols, sizeof(SResultColumn));
wmmhello's avatar
wmmhello 已提交
2955 2956 2957 2958 2959 2960 2961 2962 2963 2964 2965 2966 2967 2968 2969 2970 2971 2972

  for (int32_t i = 0; i < numOfCols; ++i) {
    if (IS_VAR_DATA_TYPE(pTableMeta->schema[i].type)) {
      pCol[i].offset = (int32_t*)pStart;
      pStart += rows * sizeof(int32_t);
    } else {
      pCol[i].nullbitmap = pStart;
      pStart += BitmapLen(rows);
    }

    pCol[i].pData = pStart;
    pStart += colLength[i];
  }

  for (int32_t j = 0; j < rows; j++) {
    tdSRowResetBuf(&rb, rowData);
    int32_t offset = 0;
    for (int32_t k = 0; k < numOfCols; k++) {
L
Liu Jicong 已提交
2973
      const SSchema* pColumn = &pTableMeta->schema[k];
wmmhello's avatar
wmmhello 已提交
2974 2975 2976 2977 2978 2979

      if (IS_VAR_DATA_TYPE(pColumn->type)) {
        if (pCol[k].offset[j] != -1) {
          char* data = pCol[k].pData + pCol[k].offset[j];
          tdAppendColValToRow(&rb, pColumn->colId, pColumn->type, TD_VTYPE_NORM, data, true, offset, k);
        } else {
H
Haojun Liao 已提交
2980

wmmhello's avatar
wmmhello 已提交
2981 2982 2983 2984 2985 2986 2987 2988 2989 2990 2991 2992 2993
          tdAppendColValToRow(&rb, pColumn->colId, pColumn->type, TD_VTYPE_NULL, NULL, false, offset, k);
        }
      } else {
        if (!colDataIsNull_f(pCol[k].nullbitmap, j)) {
          char* data = pCol[k].pData + pColumn->bytes * j;
          tdAppendColValToRow(&rb, pColumn->colId, pColumn->type, TD_VTYPE_NORM, data, true, offset, k);
        } else {
          tdAppendColValToRow(&rb, pColumn->colId, pColumn->type, TD_VTYPE_NULL, NULL, false, offset, k);
        }
      }

      offset += TYPE_BYTES[pColumn->type];
    }
2994
    tdSRowEnd(&rb);
wmmhello's avatar
wmmhello 已提交
2995 2996 2997 2998 2999 3000 3001 3002 3003 3004 3005
    int32_t rowLen = TD_ROW_LEN(rowData);
    rowData = POINTER_SHIFT(rowData, rowLen);
    dataLen += rowLen;
  }

  taosMemoryFree(pCol);

  blk->uid = htobe64(uid);
  blk->suid = htobe64(suid);
  blk->sversion = htonl(pTableMeta->sversion);
  blk->schemaLen = htonl(schemaLen);
3006
  blk->numOfRows = htonl(rows);
wmmhello's avatar
wmmhello 已提交
3007 3008 3009 3010 3011 3012 3013 3014 3015 3016 3017 3018 3019
  blk->dataLen = htonl(dataLen);
  subReq->length = sizeof(SSubmitReq) + sizeof(SSubmitBlk) + schemaLen + dataLen;
  subReq->numOfBlocks = 1;

  pQuery = (SQuery*)nodesMakeNode(QUERY_NODE_QUERY);
  if (NULL == pQuery) {
    uError("create SQuery error");
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
  pQuery->execMode = QUERY_EXEC_MODE_SCHEDULE;
  pQuery->haveResultSet = false;
  pQuery->msgType = TDMT_VND_SUBMIT;
L
Liu Jicong 已提交
3020
  pQuery->pRoot = (SNode*)nodesMakeNode(QUERY_NODE_VNODE_MODIF_STMT);
wmmhello's avatar
wmmhello 已提交
3021 3022 3023 3024 3025
  if (NULL == pQuery->pRoot) {
    uError("create pQuery->pRoot error");
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
L
Liu Jicong 已提交
3026
  SVnodeModifOpStmt* nodeStmt = (SVnodeModifOpStmt*)(pQuery->pRoot);
wmmhello's avatar
wmmhello 已提交
3027 3028 3029
  nodeStmt->payloadType = PAYLOAD_TYPE_KV;
  nodeStmt->pDataBlocks = taosArrayInit(1, POINTER_BYTES);

L
Liu Jicong 已提交
3030
  SVgDataBlocks* dst = taosMemoryCalloc(1, sizeof(SVgDataBlocks));
wmmhello's avatar
wmmhello 已提交
3031 3032 3033 3034 3035 3036 3037 3038 3039 3040 3041 3042 3043
  if (NULL == dst) {
    code = TSDB_CODE_TSC_OUT_OF_MEMORY;
    goto end;
  }
  dst->vg = vgData;
  dst->numOfTables = subReq->numOfBlocks;
  dst->size = subReq->length;
  dst->pData = (char*)subReq;
  subReq->header.vgId = htonl(dst->vg.vgId);
  subReq->version = htonl(1);
  subReq->header.contLen = htonl(subReq->length);
  subReq->length = htonl(subReq->length);
  subReq->numOfBlocks = htonl(subReq->numOfBlocks);
L
Liu Jicong 已提交
3044
  subReq = NULL;  // no need free
wmmhello's avatar
wmmhello 已提交
3045 3046 3047 3048 3049 3050 3051 3052 3053 3054 3055
  taosArrayPush(nodeStmt->pDataBlocks, &dst);

  launchQueryImpl(pRequest, pQuery, true, NULL);
  code = pRequest->code;

end:
  taosMemoryFreeClear(pTableMeta);
  qDestroyQuery(pQuery);
  return code;
}

L
Liu Jicong 已提交
3056 3057 3058 3059
static int32_t tmqWriteRaw(TAOS* taos, void* data, int32_t dataLen) {
  int32_t   code = TSDB_CODE_SUCCESS;
  SHashObj* pVgHash = NULL;
  SQuery*   pQuery = NULL;
wmmhello's avatar
wmmhello 已提交
3060
  SMqRspObj rspObj = {0};
L
Liu Jicong 已提交
3061
  SDecoder  decoder = {0};
3062 3063 3064

  terrno = TSDB_CODE_SUCCESS;
  SRequestObj* pRequest = (SRequestObj*)createRequest(*(int64_t*)taos, TSDB_SQL_INSERT);
L
Liu Jicong 已提交
3065
  if (!pRequest) {
3066 3067 3068 3069
    uError("WriteRaw:createRequest error request is null");
    return terrno;
  }

wmmhello's avatar
wmmhello 已提交
3070 3071 3072 3073 3074
  rspObj.resIter = -1;
  rspObj.resType = RES_TYPE__TMQ;

  tDecoderInit(&decoder, data, dataLen);
  code = tDecodeSMqDataRsp(&decoder, &rspObj.rsp);
L
Liu Jicong 已提交
3075
  if (code != 0) {
wmmhello's avatar
wmmhello 已提交
3076 3077 3078 3079 3080
    uError("WriteRaw:decode smqDataRsp error");
    code = TSDB_CODE_INVALID_MSG;
    goto end;
  }

3081 3082 3083 3084 3085 3086 3087 3088
  if (!pRequest->pDb) {
    uError("WriteRaw:not use db");
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }

  pVgHash = taosHashInit(16, taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT), true, HASH_NO_LOCK);
  taosHashSetFreeFp(pVgHash, destroyVgHash);
L
Liu Jicong 已提交
3089
  struct SCatalog* pCatalog = NULL;
3090
  code = catalogGetHandle(pRequest->pTscObj->pAppInfo->clusterId, &pCatalog);
L
Liu Jicong 已提交
3091
  if (code != TSDB_CODE_SUCCESS) {
3092 3093 3094 3095 3096 3097 3098 3099 3100
    uError("WriteRaw: get gatlog error");
    goto end;
  }

  SRequestConnInfo conn = {0};
  conn.pTrans = pRequest->pTscObj->pAppInfo->pTransporter;
  conn.requestId = pRequest->requestId;
  conn.requestObjRefId = pRequest->self;
  conn.mgmtEps = getEpSet_s(&pRequest->pTscObj->pAppInfo->mgmtEp);
wmmhello's avatar
wmmhello 已提交
3101 3102 3103 3104 3105 3106

  printf("raw data block num:%d\n", rspObj.rsp.blockNum);
  while (++rspObj.resIter < rspObj.rsp.blockNum) {
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(rspObj.rsp.blockData, rspObj.resIter);
    if (!rspObj.rsp.withSchema) {
      uError("WriteRaw:no schema, iter:%d", rspObj.resIter);
3107 3108
      goto end;
    }
wmmhello's avatar
wmmhello 已提交
3109 3110
    SSchemaWrapper* pSW = (SSchemaWrapper*)taosArrayGetP(rspObj.rsp.blockSchema, rspObj.resIter);
    setResSchemaInfo(&rspObj.resInfo, pSW->pSchema, pSW->nCols);
3111

wmmhello's avatar
wmmhello 已提交
3112
    code = setQueryResultFromRsp(&rspObj.resInfo, pRetrieve, false, false);
L
Liu Jicong 已提交
3113
    if (code != TSDB_CODE_SUCCESS) {
3114 3115 3116 3117 3118
      uError("WriteRaw: setQueryResultFromRsp error");
      goto end;
    }

    uint16_t fLen = 0;
L
Liu Jicong 已提交
3119 3120
    int32_t  rowSize = 0;
    int16_t  nVar = 0;
3121
    for (int i = 0; i < pSW->nCols; i++) {
L
Liu Jicong 已提交
3122
      SSchema* schema = pSW->pSchema + i;
3123 3124
      fLen += TYPE_BYTES[schema->type];
      rowSize += schema->bytes;
L
Liu Jicong 已提交
3125 3126
      if (IS_VAR_DATA_TYPE(schema->type)) {
        nVar++;
3127 3128 3129
      }
    }

wmmhello's avatar
wmmhello 已提交
3130
    int32_t rows = rspObj.resInfo.numOfRows;
3131 3132 3133 3134 3135
    int32_t extendedRowSize = rowSize + TD_ROW_HEAD_LEN - sizeof(TSKEY) + nVar * sizeof(VarDataOffsetT) +
                              (int32_t)TD_BITMAP_BYTES(pSW->nCols - 1);
    int32_t schemaLen = 0;
    int32_t submitLen = sizeof(SSubmitBlk) + schemaLen + rows * extendedRowSize;

wmmhello's avatar
wmmhello 已提交
3136
    const char* tbName = (const char*)taosArrayGetP(rspObj.rsp.blockTbName, rspObj.resIter);
L
Liu Jicong 已提交
3137
    if (!tbName) {
3138 3139 3140 3141 3142
      uError("WriteRaw: tbname is null");
      code = TSDB_CODE_TMQ_INVALID_MSG;
      goto end;
    }

3143
    printf("raw data tbname:%s\n", tbName);
3144 3145 3146 3147 3148 3149 3150 3151 3152 3153 3154 3155 3156
    SName pName = {TSDB_TABLE_NAME_T, pRequest->pTscObj->acctId, {0}, {0}};
    strcpy(pName.dbname, pRequest->pDb);
    strcpy(pName.tname, tbName);

    VgData vgData = {0};
    code = catalogGetTableHashVgroup(pCatalog, &conn, &pName, &(vgData.vg));
    if (code != TSDB_CODE_SUCCESS) {
      uError("WriteRaw:catalogGetTableHashVgroup failed. table name: %s", tbName);
      goto end;
    }

    SSubmitReq* subReq = NULL;
    SSubmitBlk* blk = NULL;
L
Liu Jicong 已提交
3157 3158
    void*       hData = taosHashGet(pVgHash, &vgData.vg.vgId, sizeof(vgData.vg.vgId));
    if (hData) {
3159 3160 3161
      vgData = *(VgData*)hData;

      int32_t totalLen = ((SSubmitReq*)(vgData.data))->length + submitLen;
L
Liu Jicong 已提交
3162
      void*   tmp = taosMemoryRealloc(vgData.data, totalLen);
3163 3164 3165 3166 3167 3168 3169 3170
      if (tmp == NULL) {
        code = TSDB_CODE_TSC_OUT_OF_MEMORY;
        goto end;
      }
      vgData.data = tmp;
      ((VgData*)hData)->data = tmp;
      subReq = (SSubmitReq*)(vgData.data);
      blk = POINTER_SHIFT(vgData.data, subReq->length);
L
Liu Jicong 已提交
3171
    } else {
3172
      int32_t totalLen = sizeof(SSubmitReq) + submitLen;
L
Liu Jicong 已提交
3173
      void*   tmp = taosMemoryCalloc(1, totalLen);
3174 3175 3176 3177 3178
      if (tmp == NULL) {
        code = TSDB_CODE_TSC_OUT_OF_MEMORY;
        goto end;
      }
      vgData.data = tmp;
L
Liu Jicong 已提交
3179
      taosHashPut(pVgHash, (const char*)&vgData.vg.vgId, sizeof(vgData.vg.vgId), (char*)&vgData, sizeof(vgData));
3180 3181 3182 3183 3184 3185 3186
      subReq = (SSubmitReq*)(vgData.data);
      subReq->length = sizeof(SSubmitReq);
      subReq->numOfBlocks = 0;

      blk = POINTER_SHIFT(vgData.data, sizeof(SSubmitReq));
    }

3187 3188
    STableMeta* pTableMeta = NULL;
    code = catalogGetTableMeta(pCatalog, &conn, &pName, &pTableMeta);
3189 3190 3191 3192
    if (code != TSDB_CODE_SUCCESS) {
      uError("WriteRaw:catalogGetTableMeta failed. table name: %s", tbName);
      goto end;
    }
3193 3194 3195
    uint64_t suid = (TSDB_NORMAL_TABLE == pTableMeta->tableType ? 0 : pTableMeta->suid);
    uint64_t uid = pTableMeta->uid;
    taosMemoryFreeClear(pTableMeta);
3196

L
Liu Jicong 已提交
3197
    void*   blkSchema = POINTER_SHIFT(blk, sizeof(SSubmitBlk));
3198 3199 3200 3201 3202 3203 3204 3205 3206 3207
    STSRow* rowData = POINTER_SHIFT(blkSchema, schemaLen);

    SRowBuilder rb = {0};
    tdSRowInit(&rb, pSW->version);
    tdSRowSetTpInfo(&rb, pSW->nCols, fLen);
    int32_t dataLen = 0;

    for (int32_t j = 0; j < rows; j++) {
      tdSRowResetBuf(&rb, rowData);

wmmhello's avatar
wmmhello 已提交
3208 3209
      doSetOneRowPtr(&rspObj.resInfo);
      rspObj.resInfo.current += 1;
3210 3211 3212

      int32_t offset = 0;
      for (int32_t k = 0; k < pSW->nCols; k++) {
L
Liu Jicong 已提交
3213 3214
        const SSchema* pColumn = &pSW->pSchema[k];
        char*          data = rspObj.resInfo.row[k];
3215 3216 3217
        if (!data) {
          tdAppendColValToRow(&rb, pColumn->colId, pColumn->type, TD_VTYPE_NULL, NULL, false, offset, k);
        } else {
L
Liu Jicong 已提交
3218
          if (IS_VAR_DATA_TYPE(pColumn->type)) {
3219 3220 3221 3222 3223 3224
            data -= VARSTR_HEADER_SIZE;
          }
          tdAppendColValToRow(&rb, pColumn->colId, pColumn->type, TD_VTYPE_NORM, data, true, offset, k);
        }
        offset += TYPE_BYTES[pColumn->type];
      }
3225
      tdSRowEnd(&rb);
3226 3227 3228 3229 3230 3231 3232 3233 3234
      int32_t rowLen = TD_ROW_LEN(rowData);
      rowData = POINTER_SHIFT(rowData, rowLen);
      dataLen += rowLen;
    }

    blk->uid = htobe64(uid);
    blk->suid = htobe64(suid);
    blk->sversion = htonl(pSW->version);
    blk->schemaLen = htonl(schemaLen);
3235
    blk->numOfRows = htonl(rows);
3236 3237 3238 3239 3240 3241 3242 3243 3244 3245 3246 3247 3248 3249
    blk->dataLen = htonl(dataLen);
    subReq->length += sizeof(SSubmitBlk) + schemaLen + dataLen;
    subReq->numOfBlocks++;
  }

  pQuery = (SQuery*)nodesMakeNode(QUERY_NODE_QUERY);
  if (NULL == pQuery) {
    uError("create SQuery error");
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
  pQuery->execMode = QUERY_EXEC_MODE_SCHEDULE;
  pQuery->haveResultSet = false;
  pQuery->msgType = TDMT_VND_SUBMIT;
L
Liu Jicong 已提交
3250
  pQuery->pRoot = (SNode*)nodesMakeNode(QUERY_NODE_VNODE_MODIF_STMT);
3251 3252 3253 3254 3255
  if (NULL == pQuery->pRoot) {
    uError("create pQuery->pRoot error");
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
  }
L
Liu Jicong 已提交
3256
  SVnodeModifOpStmt* nodeStmt = (SVnodeModifOpStmt*)(pQuery->pRoot);
3257 3258 3259 3260 3261
  nodeStmt->payloadType = PAYLOAD_TYPE_KV;

  int32_t numOfVg = taosHashGetSize(pVgHash);
  nodeStmt->pDataBlocks = taosArrayInit(numOfVg, POINTER_BYTES);

L
Liu Jicong 已提交
3262
  VgData* vData = (VgData*)taosHashIterate(pVgHash, NULL);
3263
  while (vData) {
L
Liu Jicong 已提交
3264
    SVgDataBlocks* dst = taosMemoryCalloc(1, sizeof(SVgDataBlocks));
3265 3266 3267 3268
    if (NULL == dst) {
      code = TSDB_CODE_TSC_OUT_OF_MEMORY;
      goto end;
    }
3269 3270
    dst->vg = vData->vg;
    SSubmitReq* subReq = (SSubmitReq*)(vData->data);
3271 3272 3273
    dst->numOfTables = subReq->numOfBlocks;
    dst->size = subReq->length;
    dst->pData = (char*)subReq;
L
Liu Jicong 已提交
3274
    vData->data = NULL;  // no need free
3275 3276 3277 3278 3279 3280
    subReq->header.vgId = htonl(dst->vg.vgId);
    subReq->version = htonl(1);
    subReq->header.contLen = htonl(subReq->length);
    subReq->length = htonl(subReq->length);
    subReq->numOfBlocks = htonl(subReq->numOfBlocks);
    taosArrayPush(nodeStmt->pDataBlocks, &dst);
L
Liu Jicong 已提交
3281
    vData = (VgData*)taosHashIterate(pVgHash, vData);
3282 3283 3284 3285
  }

  launchQueryImpl(pRequest, pQuery, true, NULL);
  code = pRequest->code;
wmmhello's avatar
wmmhello 已提交
3286

3287
end:
wmmhello's avatar
wmmhello 已提交
3288 3289
  tDecoderClear(&decoder);
  taos_free_result(&rspObj);
3290 3291 3292 3293
  qDestroyQuery(pQuery);
  destroyRequest(pRequest);
  taosHashCleanup(pVgHash);
  return code;
wmmhello's avatar
wmmhello 已提交
3294 3295
}

wmmhello's avatar
wmmhello 已提交
3296 3297 3298 3299 3300 3301 3302 3303 3304 3305 3306 3307 3308 3309 3310 3311 3312 3313 3314 3315 3316 3317 3318 3319
char* tmq_get_json_meta(TAOS_RES* res) {
  if (!TD_RES_TMQ_META(res)) {
    return NULL;
  }

  SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
  if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_CREATE_STB) {
    return processCreateStb(&pMetaRspObj->metaRsp);
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_ALTER_STB) {
    return processAlterStb(&pMetaRspObj->metaRsp);
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_DROP_STB) {
    return processDropSTable(&pMetaRspObj->metaRsp);
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_CREATE_TABLE) {
    return processCreateTable(&pMetaRspObj->metaRsp);
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_ALTER_TABLE) {
    return processAlterTable(&pMetaRspObj->metaRsp);
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_DROP_TABLE) {
    return processDropTable(&pMetaRspObj->metaRsp);
  }
  return NULL;
}

void tmq_free_json_meta(char* jsonMeta) { taosMemoryFreeClear(jsonMeta); }

L
Liu Jicong 已提交
3320 3321
int32_t tmq_get_raw(TAOS_RES* res, tmq_raw_data* raw) {
  if (!raw || !res) {
wmmhello's avatar
wmmhello 已提交
3322 3323 3324 3325 3326 3327 3328
    return TSDB_CODE_INVALID_PARA;
  }
  if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    raw->raw = pMetaRspObj->metaRsp.metaRsp;
    raw->raw_len = pMetaRspObj->metaRsp.metaRspLen;
    raw->raw_type = pMetaRspObj->metaRsp.resMsgType;
L
Liu Jicong 已提交
3329 3330
  } else if (TD_RES_TMQ(res)) {
    SMqRspObj* rspObj = ((SMqRspObj*)res);
wmmhello's avatar
wmmhello 已提交
3331 3332 3333 3334 3335 3336 3337 3338

    int32_t len = 0;
    int32_t code = 0;
    tEncodeSize(tEncodeSMqDataRsp, &rspObj->rsp, len, code);
    if (code < 0) {
      return -1;
    }

L
Liu Jicong 已提交
3339
    void*    buf = taosMemoryCalloc(1, len);
wmmhello's avatar
wmmhello 已提交
3340 3341 3342 3343 3344 3345 3346 3347 3348 3349 3350 3351 3352 3353 3354
    SEncoder encoder = {0};
    tEncoderInit(&encoder, buf, len);
    tEncodeSMqDataRsp(&encoder, &rspObj->rsp);
    tEncoderClear(&encoder);

    raw->raw = buf;
    raw->raw_len = len;
    raw->raw_type = RES_TYPE__TMQ;
  } else {
    return TSDB_CODE_TMQ_INVALID_MSG;
  }
  return TSDB_CODE_SUCCESS;
}

void tmq_free_raw(tmq_raw_data raw) {
L
Liu Jicong 已提交
3355
  if (raw.raw_type == RES_TYPE__TMQ) {
wmmhello's avatar
wmmhello 已提交
3356 3357 3358 3359
    taosMemoryFree(raw.raw);
  }
}

L
Liu Jicong 已提交
3360
int32_t tmq_write_raw(TAOS* taos, tmq_raw_data raw) {
wmmhello's avatar
wmmhello 已提交
3361 3362 3363 3364
  if (!taos) {
    return TSDB_CODE_INVALID_PARA;
  }

L
Liu Jicong 已提交
3365
  if (raw.raw_type == TDMT_VND_CREATE_STB) {
wmmhello's avatar
wmmhello 已提交
3366
    return taosCreateStb(taos, raw.raw, raw.raw_len);
L
Liu Jicong 已提交
3367
  } else if (raw.raw_type == TDMT_VND_ALTER_STB) {
wmmhello's avatar
wmmhello 已提交
3368
    return taosCreateStb(taos, raw.raw, raw.raw_len);
L
Liu Jicong 已提交
3369
  } else if (raw.raw_type == TDMT_VND_DROP_STB) {
wmmhello's avatar
wmmhello 已提交
3370
    return taosDropStb(taos, raw.raw, raw.raw_len);
L
Liu Jicong 已提交
3371
  } else if (raw.raw_type == TDMT_VND_CREATE_TABLE) {
wmmhello's avatar
wmmhello 已提交
3372
    return taosCreateTable(taos, raw.raw, raw.raw_len);
L
Liu Jicong 已提交
3373
  } else if (raw.raw_type == TDMT_VND_ALTER_TABLE) {
wmmhello's avatar
wmmhello 已提交
3374
    return taosAlterTable(taos, raw.raw, raw.raw_len);
L
Liu Jicong 已提交
3375
  } else if (raw.raw_type == TDMT_VND_DROP_TABLE) {
wmmhello's avatar
wmmhello 已提交
3376
    return taosDropTable(taos, raw.raw, raw.raw_len);
L
Liu Jicong 已提交
3377
  } else if (raw.raw_type == TDMT_VND_DELETE) {
wmmhello's avatar
wmmhello 已提交
3378
    return taosDeleteData(taos, raw.raw, raw.raw_len);
L
Liu Jicong 已提交
3379
  } else if (raw.raw_type == RES_TYPE__TMQ) {
wmmhello's avatar
wmmhello 已提交
3380 3381 3382 3383 3384
    return tmqWriteRaw(taos, raw.raw, raw.raw_len);
  }
  return TSDB_CODE_INVALID_PARA;
}

L
Liu Jicong 已提交
3385
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* msg, tmq_commit_cb* cb, void* param) {
3386
  //
L
Liu Jicong 已提交
3387
  tmqCommitInner2(tmq, msg, 0, 1, cb, param);
L
Liu Jicong 已提交
3388 3389
}

3390 3391 3392 3393
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* msg) {
  //
  return tmqCommitInner2(tmq, msg, 0, 0, NULL, NULL);
}