tmq.c 90.7 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
};

L
Liu Jicong 已提交
116 117
struct tmq_raw_data {
  void*   raw_meta;
wmmhello's avatar
wmmhello 已提交
118 119 120 121
  int32_t raw_meta_len;
  int16_t raw_meta_type;
};

X
Xiaoyu Wang 已提交
122 123 124 125 126 127 128 129
enum {
  TMQ_VG_STATUS__IDLE = 0,
  TMQ_VG_STATUS__WAIT,
};

enum {
  TMQ_CONSUMER_STATUS__INIT = 0,
  TMQ_CONSUMER_STATUS__READY,
130
  TMQ_CONSUMER_STATUS__NO_TOPIC,
L
Liu Jicong 已提交
131 132
};

L
Liu Jicong 已提交
133
enum {
134
  TMQ_DELAYED_TASK__ASK_EP = 1,
L
Liu Jicong 已提交
135 136 137 138
  TMQ_DELAYED_TASK__REPORT,
  TMQ_DELAYED_TASK__COMMIT,
};

L
Liu Jicong 已提交
139
typedef struct {
140 141 142
  // statistics
  int64_t pollCnt;
  // offset
L
Liu Jicong 已提交
143 144 145 146
  /*int64_t      committedOffset;*/
  /*int64_t      currentOffset;*/
  STqOffsetVal committedOffsetNew;
  STqOffsetVal currentOffsetNew;
L
Liu Jicong 已提交
147
  // connection info
148
  int32_t vgId;
X
Xiaoyu Wang 已提交
149
  int32_t vgStatus;
L
Liu Jicong 已提交
150
  int32_t vgSkipCnt;
151 152 153
  SEpSet  epSet;
} SMqClientVg;

L
Liu Jicong 已提交
154
typedef struct {
155
  // subscribe info
L
Liu Jicong 已提交
156
  char* topicName;
L
Liu Jicong 已提交
157
  char  db[TSDB_DB_FNAME_LEN];
L
Liu Jicong 已提交
158 159 160

  SArray* vgs;  // SArray<SMqClientVg>

L
Liu Jicong 已提交
161 162
  int8_t         isSchemaAdaptive;
  SSchemaWrapper schema;
163 164
} SMqClientTopic;

L
Liu Jicong 已提交
165 166 167 168 169
typedef struct {
  int8_t          tmqRspType;
  int32_t         epoch;
  SMqClientVg*    vgHandle;
  SMqClientTopic* topicHandle;
L
Liu Jicong 已提交
170
  union {
L
Liu Jicong 已提交
171 172
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
173
  };
L
Liu Jicong 已提交
174 175
} SMqPollRspWrapper;

L
Liu Jicong 已提交
176
typedef struct {
L
Liu Jicong 已提交
177 178 179
  tmq_t*  tmq;
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
180
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
181

L
Liu Jicong 已提交
182
typedef struct {
183
  tmq_t*  tmq;
L
Liu Jicong 已提交
184
  int32_t code;
L
Liu Jicong 已提交
185
  int32_t async;
X
Xiaoyu Wang 已提交
186
  tsem_t  rspSem;
187 188
} SMqAskEpCbParam;

L
Liu Jicong 已提交
189
typedef struct {
L
Liu Jicong 已提交
190 191
  tmq_t*          tmq;
  SMqClientVg*    pVg;
L
Liu Jicong 已提交
192
  SMqClientTopic* pTopic;
L
Liu Jicong 已提交
193
  int32_t         epoch;
L
Liu Jicong 已提交
194
  int32_t         vgId;
L
Liu Jicong 已提交
195
  tsem_t          rspSem;
X
Xiaoyu Wang 已提交
196
} SMqPollCbParam;
197

L
Liu Jicong 已提交
198
#if 0
L
Liu Jicong 已提交
199
typedef struct {
L
Liu Jicong 已提交
200
  tmq_t*         tmq;
L
Liu Jicong 已提交
201 202
  int8_t         async;
  int8_t         automatic;
L
Liu Jicong 已提交
203
  int8_t         freeOffsets;
L
Liu Jicong 已提交
204
  tmq_commit_cb* userCb;
L
Liu Jicong 已提交
205
  tsem_t         rspSem;
L
Liu Jicong 已提交
206
  int32_t        rspErr;
L
Liu Jicong 已提交
207
  SArray*        offsets;
L
Liu Jicong 已提交
208
  void*          userParam;
L
Liu Jicong 已提交
209
} SMqCommitCbParam;
L
Liu Jicong 已提交
210
#endif
L
Liu Jicong 已提交
211

212
typedef struct {
L
Liu Jicong 已提交
213 214 215 216
  tmq_t* tmq;
  int8_t automatic;
  int8_t async;
  /*int8_t         freeOffsets;*/
L
Liu Jicong 已提交
217 218
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
L
Liu Jicong 已提交
219
  int32_t        rspErr;
220
  tmq_commit_cb* userCb;
L
Liu Jicong 已提交
221 222 223 224
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
  void*  userParam;
  tsem_t rspSem;
225 226 227 228 229 230 231
} SMqCommitCbParamSet;

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

232
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
233
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
234
  conf->withTbName = false;
L
Liu Jicong 已提交
235
  conf->autoCommit = true;
L
Liu Jicong 已提交
236
  conf->autoCommitInterval = 5000;
X
Xiaoyu Wang 已提交
237
  conf->resetOffset = TMQ_CONF__RESET_OFFSET__EARLIEAST;
238 239 240
  return conf;
}

L
Liu Jicong 已提交
241
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
242 243 244 245 246 247
  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 已提交
248 249 250
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
251 252
  if (strcmp(key, "group.id") == 0) {
    strcpy(conf->groupId, value);
L
Liu Jicong 已提交
253
    return TMQ_CONF_OK;
254
  }
L
Liu Jicong 已提交
255

256 257
  if (strcmp(key, "client.id") == 0) {
    strcpy(conf->clientId, value);
L
Liu Jicong 已提交
258 259
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
260

L
Liu Jicong 已提交
261 262
  if (strcmp(key, "enable.auto.commit") == 0) {
    if (strcmp(value, "true") == 0) {
L
Liu Jicong 已提交
263
      conf->autoCommit = true;
L
Liu Jicong 已提交
264 265
      return TMQ_CONF_OK;
    } else if (strcmp(value, "false") == 0) {
L
Liu Jicong 已提交
266
      conf->autoCommit = false;
L
Liu Jicong 已提交
267 268 269 270
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
271
  }
L
Liu Jicong 已提交
272

L
Liu Jicong 已提交
273 274 275 276 277
  if (strcmp(key, "auto.commit.interval.ms") == 0) {
    conf->autoCommitInterval = atoi(value);
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
278 279 280 281 282 283 284 285 286 287 288 289 290 291
  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 已提交
292

293 294
  if (strcmp(key, "msg.with.table.name") == 0) {
    if (strcmp(value, "true") == 0) {
295
      conf->withTbName = true;
L
Liu Jicong 已提交
296
      return TMQ_CONF_OK;
297
    } else if (strcmp(value, "false") == 0) {
298
      conf->withTbName = false;
L
Liu Jicong 已提交
299
      return TMQ_CONF_OK;
300 301 302 303 304
    } else {
      return TMQ_CONF_INVALID;
    }
  }

L
Liu Jicong 已提交
305
  if (strcmp(key, "experimental.snapshot.enable") == 0) {
L
Liu Jicong 已提交
306
    if (strcmp(value, "true") == 0) {
307 308 309 310 311 312 313 314 315 316 317 318 319
      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 已提交
320 321
      return TMQ_CONF_OK;
    } else if (strcmp(value, "false") == 0) {
322
      conf->hbBgEnable = false;
L
Liu Jicong 已提交
323 324 325 326
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
327
    return TMQ_CONF_OK;
L
Liu Jicong 已提交
328 329
  }

L
Liu Jicong 已提交
330
  if (strcmp(key, "experimental.snapshot.batch.size") == 0) {
331
    conf->ssBatchSize = atoi(value);
L
Liu Jicong 已提交
332 333 334
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
335
  if (strcmp(key, "td.connect.ip") == 0) {
L
Liu Jicong 已提交
336 337 338
    conf->ip = strdup(value);
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
339
  if (strcmp(key, "td.connect.user") == 0) {
L
Liu Jicong 已提交
340 341 342
    conf->user = strdup(value);
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
343
  if (strcmp(key, "td.connect.pass") == 0) {
L
Liu Jicong 已提交
344 345 346
    conf->pass = strdup(value);
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
347
  if (strcmp(key, "td.connect.port") == 0) {
L
Liu Jicong 已提交
348 349 350
    conf->port = atoi(value);
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
351
  if (strcmp(key, "td.connect.db") == 0) {
352
    /*conf->db = strdup(value);*/
L
Liu Jicong 已提交
353 354 355
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
356
  return TMQ_CONF_UNKNOWN;
357 358 359
}

tmq_list_t* tmq_list_new() {
L
Liu Jicong 已提交
360 361
  //
  return (tmq_list_t*)taosArrayInit(0, sizeof(void*));
362 363
}

L
Liu Jicong 已提交
364 365
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
366
  char*   topic = strdup(src);
L
fix  
Liu Jicong 已提交
367
  if (taosArrayPush(container, &topic) == NULL) return -1;
368 369 370
  return 0;
}

L
Liu Jicong 已提交
371
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
372
  SArray* container = &list->container;
L
Liu Jicong 已提交
373
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
374 375
}

L
Liu Jicong 已提交
376 377 378 379 380 381 382 383 384 385
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 已提交
386 387 388 389
static int32_t tmqMakeTopicVgKey(char* dst, const char* topicName, int32_t vg) {
  return sprintf(dst, "%s:%d", topicName, vg);
}

L
Liu Jicong 已提交
390
#if 0
L
Liu Jicong 已提交
391 392
int32_t tmqCommitCb(void* param, const SDataBuf* pMsg, int32_t code) {
  SMqCommitCbParam* pParam = (SMqCommitCbParam*)param;
L
Liu Jicong 已提交
393
  pParam->rspErr = code;
L
Liu Jicong 已提交
394 395
  if (pParam->async) {
    if (pParam->automatic && pParam->tmq->commitCb) {
L
Liu Jicong 已提交
396
      pParam->tmq->commitCb(pParam->tmq, pParam->rspErr, pParam->tmq->commitCbUserParam);
L
Liu Jicong 已提交
397
    } else if (!pParam->automatic && pParam->userCb) {
L
Liu Jicong 已提交
398
      pParam->userCb(pParam->tmq, pParam->rspErr, pParam->userParam);
L
Liu Jicong 已提交
399 400
    }

L
Liu Jicong 已提交
401
    if (pParam->freeOffsets) {
L
Liu Jicong 已提交
402 403 404 405 406 407 408 409 410
      taosArrayDestroy(pParam->offsets);
    }

    taosMemoryFree(pParam);
  } else {
    tsem_post(&pParam->rspSem);
  }
  return 0;
}
L
Liu Jicong 已提交
411
#endif
L
Liu Jicong 已提交
412

D
dapan1121 已提交
413
int32_t tmqCommitCb2(void* param, SDataBuf* pBuf, int32_t code) {
414 415 416
  SMqCommitCbParam2*   pParam = (SMqCommitCbParam2*)param;
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
  // push into array
L
Liu Jicong 已提交
417
#if 0
418 419 420 421 422
  if (code == 0) {
    taosArrayPush(pParamSet->failedOffsets, &pParam->pOffset);
  } else {
    taosArrayPush(pParamSet->successfulOffsets, &pParam->pOffset);
  }
L
Liu Jicong 已提交
423
#endif
L
Liu Jicong 已提交
424

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

428
  // count down waiting rsp
L
Liu Jicong 已提交
429
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
430 431 432 433 434 435 436
  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 已提交
437
        pParamSet->tmq->commitCb(pParamSet->tmq, pParamSet->rspErr, pParamSet->tmq->commitCbUserParam);
438 439
      } else if (!pParamSet->automatic && pParamSet->userCb) {
        // sem post
L
Liu Jicong 已提交
440
        pParamSet->userCb(pParamSet->tmq, pParamSet->rspErr, pParamSet->userParam);
441
      }
L
Liu Jicong 已提交
442 443
    } else {
      tsem_post(&pParamSet->rspSem);
444 445
    }

L
Liu Jicong 已提交
446
#if 0
447 448
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
L
Liu Jicong 已提交
449
#endif
450 451 452 453
  }
  return 0;
}

L
Liu Jicong 已提交
454 455 456 457 458 459 460
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;
461

L
Liu Jicong 已提交
462 463 464 465
  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 已提交
466

L
Liu Jicong 已提交
467 468 469 470 471 472 473 474 475 476
  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 已提交
477

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
  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 已提交
500 501
  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 已提交
502 503 504 505 506 507 508

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

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
D
dapan1121 已提交
509
  pMsgSendInfo->paramFreeFp = taosMemoryFree;      
L
Liu Jicong 已提交
510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549
  pMsgSendInfo->fp = tmqCommitCb2;
  pMsgSendInfo->msgType = TDMT_VND_MQ_COMMIT_OFFSET;
  // send msg

  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
  pParamSet->waitingRspNum++;
  pParamSet->totalRspNum++;
  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 已提交
550 551
  int32_t code = -1;

L
Liu Jicong 已提交
552 553 554 555 556 557 558 559 560 561
  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 已提交
562
        }
L
Liu Jicong 已提交
563
        goto HANDLE_RSP;
L
Liu Jicong 已提交
564 565
      }
    }
L
Liu Jicong 已提交
566
  }
L
Liu Jicong 已提交
567

L
Liu Jicong 已提交
568 569 570 571
HANDLE_RSP:
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
572 573 574
    return 0;
  }

L
Liu Jicong 已提交
575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598
  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);
  }

599 600 601 602 603 604 605 606
  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 已提交
607
  /*pParamSet->freeOffsets = 1;*/
608 609 610 611 612 613
  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 已提交
614

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

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

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

L
Liu Jicong 已提交
624 625 626 627
      if (pVg->currentOffsetNew.type > 0 && !tOffsetEqual(&pVg->currentOffsetNew, &pVg->committedOffsetNew)) {
        if (tmqSendCommitReq(tmq, pVg, pTopic, pParamSet) < 0) {
          continue;
        }
628 629 630 631
      }
    }
  }

L
Liu Jicong 已提交
632 633 634 635 636 637
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
    return 0;
  }

638 639 640 641 642 643 644 645 646 647
  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 已提交
648
      tmq->commitCb(tmq, code, tmq->commitCbUserParam);
649
    } else {
L
Liu Jicong 已提交
650
      userCb(tmq, code, userParam);
651 652 653
    }
  }

L
Liu Jicong 已提交
654
#if 0
655 656 657 658
  if (!async) {
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
  }
L
Liu Jicong 已提交
659
#endif
660 661 662 663

  return 0;
}

L
Liu Jicong 已提交
664 665
#if 0
int32_t tmqCommitInner(tmq_t* tmq, const TAOS_RES* msg, int8_t automatic, int8_t async,
L
Liu Jicong 已提交
666 667 668 669 670 671
                       tmq_commit_cb* userCb, void* userParam) {
  SMqCMCommitOffsetReq req;
  SArray*              pOffsets = NULL;
  void*                buf = NULL;
  SMqCommitCbParam*    pParam = NULL;
  SMsgSendInfo*        sendInfo = NULL;
L
Liu Jicong 已提交
672 673
  int8_t               freeOffsets;
  int32_t              code = -1;
L
Liu Jicong 已提交
674

L
Liu Jicong 已提交
675
  if (msg == NULL) {
L
Liu Jicong 已提交
676
    freeOffsets = 1;
L
Liu Jicong 已提交
677 678 679 680 681 682 683 684 685 686 687 688 689 690
    pOffsets = taosArrayInit(0, sizeof(SMqOffset));
    for (int32_t i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
      SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
      for (int32_t j = 0; j < taosArrayGetSize(pTopic->vgs); j++) {
        SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
        SMqOffset    offset;
        tstrncpy(offset.topicName, pTopic->topicName, TSDB_TOPIC_FNAME_LEN);
        tstrncpy(offset.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
        offset.vgId = pVg->vgId;
        offset.offset = pVg->currentOffset;
        taosArrayPush(pOffsets, &offset);
      }
    }
  } else {
L
Liu Jicong 已提交
691
    freeOffsets = 0;
L
Liu Jicong 已提交
692
    pOffsets = (SArray*)&msg->container;
L
Liu Jicong 已提交
693 694 695 696 697
  }

  req.num = (int32_t)pOffsets->size;
  req.offsets = pOffsets->pData;

L
Liu Jicong 已提交
698 699 700 701
  SEncoder encoder;

  tEncoderInit(&encoder, NULL, 0);
  code = tEncodeSMqCMCommitOffsetReq(&encoder, &req);
L
Liu Jicong 已提交
702 703 704
  if (code < 0) {
    goto END;
  }
L
Liu Jicong 已提交
705
  int32_t tlen = encoder.pos;
L
Liu Jicong 已提交
706 707
  buf = taosMemoryMalloc(tlen);
  if (buf == NULL) {
L
Liu Jicong 已提交
708
    tEncoderClear(&encoder);
L
Liu Jicong 已提交
709 710
    goto END;
  }
L
Liu Jicong 已提交
711 712
  tEncoderClear(&encoder);

L
Liu Jicong 已提交
713 714 715 716 717 718 719 720 721 722 723 724
  tEncoderInit(&encoder, buf, tlen);
  tEncodeSMqCMCommitOffsetReq(&encoder, &req);
  tEncoderClear(&encoder);

  pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam));
  if (pParam == NULL) {
    goto END;
  }
  pParam->tmq = tmq;
  pParam->automatic = automatic;
  pParam->async = async;
  pParam->offsets = pOffsets;
L
Liu Jicong 已提交
725
  pParam->freeOffsets = freeOffsets;
L
Liu Jicong 已提交
726 727 728 729
  pParam->userCb = userCb;
  pParam->userParam = userParam;
  if (!async) tsem_init(&pParam->rspSem, 0, 0);

730
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753
  if (sendInfo == NULL) goto END;
  sendInfo->msgInfo = (SDataBuf){
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqCommitCb;
  sendInfo->msgType = TDMT_MND_MQ_COMMIT_OFFSET;

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

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

  if (!async) {
    tsem_wait(&pParam->rspSem);
    code = pParam->rspErr;
    tsem_destroy(&pParam->rspSem);
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
754 755
  } else {
    code = 0;
L
Liu Jicong 已提交
756 757 758 759 760 761 762
  }

  // avoid double free if msg is sent
  buf = NULL;

END:
  if (buf) taosMemoryFree(buf);
L
Liu Jicong 已提交
763 764
  /*if (pParam) taosMemoryFree(pParam);*/
  /*if (sendInfo) taosMemoryFree(sendInfo);*/
L
Liu Jicong 已提交
765 766 767

  if (code != 0 && async) {
    if (automatic) {
L
Liu Jicong 已提交
768
      tmq->commitCb(tmq, code, (tmq_topic_vgroup_list_t*)pOffsets, tmq->commitCbUserParam);
L
Liu Jicong 已提交
769
    } else {
L
Liu Jicong 已提交
770
      userCb(tmq, code, (tmq_topic_vgroup_list_t*)pOffsets, userParam);
L
Liu Jicong 已提交
771 772 773
    }
  }

L
Liu Jicong 已提交
774
  if (!async && freeOffsets) {
L
Liu Jicong 已提交
775 776 777 778
    taosArrayDestroy(pOffsets);
  }
  return code;
}
L
Liu Jicong 已提交
779
#endif
L
Liu Jicong 已提交
780

781
void tmqAssignAskEpTask(void* param, void* tmrId) {
L
Liu Jicong 已提交
782
  tmq_t*  tmq = (tmq_t*)param;
783
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
784
  *pTaskType = TMQ_DELAYED_TASK__ASK_EP;
L
Liu Jicong 已提交
785
  taosWriteQitem(tmq->delayedTask, pTaskType);
786
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
787 788 789 790
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
  tmq_t*  tmq = (tmq_t*)param;
791
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
L
Liu Jicong 已提交
792 793
  *pTaskType = TMQ_DELAYED_TASK__COMMIT;
  taosWriteQitem(tmq->delayedTask, pTaskType);
794
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
795 796 797 798
}

void tmqAssignDelayedReportTask(void* param, void* tmrId) {
  tmq_t*  tmq = (tmq_t*)param;
799
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
L
Liu Jicong 已提交
800 801
  *pTaskType = TMQ_DELAYED_TASK__REPORT;
  taosWriteQitem(tmq->delayedTask, pTaskType);
802
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
803 804
}

805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844
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 已提交
845 846 847 848 849 850 851 852
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;

853
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
L
Liu Jicong 已提交
854
      tmqAskEp(tmq, true);
855
      taosTmrReset(tmqAssignAskEpTask, 1000, tmq, tmqMgmt.timer, &tmq->epTimer);
L
Liu Jicong 已提交
856
    } else if (*pTaskType == TMQ_DELAYED_TASK__COMMIT) {
L
Liu Jicong 已提交
857
      tmqCommitInner2(tmq, NULL, 1, 1, tmq->commitCb, tmq->commitCbUserParam);
L
Liu Jicong 已提交
858 859 860 861 862
      taosTmrReset(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, tmq, tmqMgmt.timer, &tmq->commitTimer);
    } else if (*pTaskType == TMQ_DELAYED_TASK__REPORT) {
    } else {
      ASSERT(0);
    }
L
Liu Jicong 已提交
863
    taosFreeQitem(pTaskType);
L
Liu Jicong 已提交
864 865 866 867 868
  }
  taosFreeQall(qall);
  return 0;
}

L
Liu Jicong 已提交
869
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
870
  SMqRspWrapper* msg = NULL;
L
Liu Jicong 已提交
871 872 873 874 875 876 877 878
  while (1) {
    taosGetQitem(tmq->qall, (void**)&msg);
    if (msg)
      taosFreeQitem(msg);
    else
      break;
  }

L
fix  
Liu Jicong 已提交
879
  msg = NULL;
L
Liu Jicong 已提交
880 881 882 883 884 885 886 887 888 889
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
    taosGetQitem(tmq->qall, (void**)&msg);
    if (msg)
      taosFreeQitem(msg);
    else
      break;
  }
}

D
dapan1121 已提交
890
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
891 892
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
L
Liu Jicong 已提交
893
  /*tmq_t* tmq = pParam->tmq;*/
L
Liu Jicong 已提交
894 895 896
  tsem_post(&pParam->rspSem);
  return 0;
}
897

L
Liu Jicong 已提交
898
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
899 900 901 902
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
903
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
904
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
905
  }
L
Liu Jicong 已提交
906
  return 0;
X
Xiaoyu Wang 已提交
907 908
}

L
Liu Jicong 已提交
909 910 911
int32_t tmq_unsubscribe(tmq_t* tmq) {
  tmq_list_t* lst = tmq_list_new();
  int32_t     rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
912 913
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
914 915
}

L
Liu Jicong 已提交
916
#if 0
917
tmq_t* tmq_consumer_new(void* conn, tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
wafwerar's avatar
wafwerar 已提交
918
  tmq_t* pTmq = taosMemoryCalloc(sizeof(tmq_t), 1);
919 920 921 922 923 924
  if (pTmq == NULL) {
    return NULL;
  }
  pTmq->pTscObj = (STscObj*)conn;
  pTmq->status = 0;
  pTmq->pollCnt = 0;
L
Liu Jicong 已提交
925
  pTmq->epoch = 0;
L
fix txn  
Liu Jicong 已提交
926
  pTmq->epStatus = 0;
L
temp  
Liu Jicong 已提交
927
  pTmq->epSkipCnt = 0;
L
Liu Jicong 已提交
928
  // set conf
929 930
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
L
Liu Jicong 已提交
931
  pTmq->autoCommit = conf->autoCommit;
932
  pTmq->commit_cb = conf->commit_cb;
L
Liu Jicong 已提交
933
  pTmq->resetOffsetCfg = conf->resetOffset;
L
Liu Jicong 已提交
934

L
Liu Jicong 已提交
935 936 937 938 939 940 941 942 943 944
  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 已提交
945 946
  tsem_init(&pTmq->rspSem, 0, 0);

L
Liu Jicong 已提交
947 948
  return pTmq;
}
L
Liu Jicong 已提交
949
#endif
L
Liu Jicong 已提交
950

L
Liu Jicong 已提交
951
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
L
Liu Jicong 已提交
952 953 954 955 956 957 958 959 960 961 962
  // 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 已提交
963 964
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
965 966
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    tscError("consumer %ld setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
967 968
    return NULL;
  }
L
Liu Jicong 已提交
969

L
Liu Jicong 已提交
970 971 972 973 974
  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 已提交
975
  ASSERT(conf->groupId[0]);
L
Liu Jicong 已提交
976

L
Liu Jicong 已提交
977 978 979 980 981 982
  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 已提交
983 984
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    tscError("consumer %ld setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
985 986
    goto FAIL;
  }
L
Liu Jicong 已提交
987

L
Liu Jicong 已提交
988 989
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
990 991
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
992 993
  /*pTmq->epStatus = 0;*/
  /*pTmq->epSkipCnt = 0;*/
L
Liu Jicong 已提交
994

L
Liu Jicong 已提交
995 996 997
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
998
  pTmq->withTbName = conf->withTbName;
999
  pTmq->useSnapshot = conf->ssEnable;
L
Liu Jicong 已提交
1000
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1001
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1002 1003
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1004 1005
  pTmq->resetOffsetCfg = conf->resetOffset;

1006 1007
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1008
  // assign consumerId
L
Liu Jicong 已提交
1009
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1010

L
Liu Jicong 已提交
1011 1012
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
L
Liu Jicong 已提交
1013
    tscError("consumer %ld setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1014 1015
    goto FAIL;
  }
L
Liu Jicong 已提交
1016

L
Liu Jicong 已提交
1017 1018 1019
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
L
Liu Jicong 已提交
1020
    tscError("consumer %ld setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1021 1022 1023
    tsem_destroy(&pTmq->rspSem);
    goto FAIL;
  }
L
Liu Jicong 已提交
1024

1025 1026 1027 1028
  if (pTmq->hbBgEnable) {
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pTmq, tmqMgmt.timer);
  }

L
Liu Jicong 已提交
1029 1030
  tscInfo("consumer %ld is setup, consumer group %s", pTmq->consumerId, pTmq->groupId);

1031
  return pTmq;
L
Liu Jicong 已提交
1032 1033 1034 1035 1036 1037 1038 1039

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;
1040 1041
}

L
Liu Jicong 已提交
1042 1043
#if 0
int32_t tmq_commit(tmq_t* tmq, const tmq_topic_vgroup_list_t* offsets, int32_t async) {
L
Liu Jicong 已提交
1044
  return tmqCommitInner2(tmq, offsets, 0, async, tmq->commitCb, tmq->commitCbUserParam);
L
Liu Jicong 已提交
1045
}
L
Liu Jicong 已提交
1046
#endif
L
Liu Jicong 已提交
1047

L
Liu Jicong 已提交
1048
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
L
Liu Jicong 已提交
1049 1050 1051 1052 1053
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
  SCMSubscribeReq req = {0};
  int32_t         code = -1;
1054 1055

  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1056
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1057
  req.topicNames = taosArrayInit(sz, sizeof(void*));
L
Liu Jicong 已提交
1058
  if (req.topicNames == NULL) goto FAIL;
1059

L
Liu Jicong 已提交
1060 1061
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1062 1063

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

L
Liu Jicong 已提交
1066 1067 1068
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
1069
    }
L
Liu Jicong 已提交
1070
    tNameExtractFullName(&name, topicFName);
1071

L
Liu Jicong 已提交
1072 1073 1074
    tscDebug("subscribe topic: %s", topicFName);

    taosArrayPush(req.topicNames, &topicFName);
1075 1076
  }

L
Liu Jicong 已提交
1077 1078 1079 1080
  int32_t tlen = tSerializeSCMSubscribeReq(NULL, &req);
  buf = taosMemoryMalloc(tlen);
  if (buf == NULL) goto FAIL;

1081 1082 1083
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

1084
  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
1085
  if (sendInfo == NULL) goto FAIL;
1086

X
Xiaoyu Wang 已提交
1087
  SMqSubscribeCbParam param = {
L
Liu Jicong 已提交
1088
      .rspErr = 0,
X
Xiaoyu Wang 已提交
1089 1090
      .tmq = tmq,
  };
L
Liu Jicong 已提交
1091

L
Liu Jicong 已提交
1092 1093 1094
  if (tsem_init(&param.rspSem, 0, 0) != 0) goto FAIL;

  sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1095 1096 1097 1098
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };
1099

L
Liu Jicong 已提交
1100 1101
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1102 1103
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1104 1105
  sendInfo->msgType = TDMT_MND_SUBSCRIBE;

1106 1107 1108 1109 1110
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1111 1112 1113
  // avoid double free if msg is sent
  buf = NULL;

L
Liu Jicong 已提交
1114 1115
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1116

L
Liu Jicong 已提交
1117 1118 1119
  code = param.rspErr;
  if (code != 0) goto FAIL;

L
Liu Jicong 已提交
1120
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == tmqAskEp(tmq, false)) {
L
fix  
Liu Jicong 已提交
1121
    tscDebug("consumer not ready, retry");
L
Liu Jicong 已提交
1122 1123
    taosMsleep(500);
  }
1124

1125 1126 1127
  // init ep timer
  if (tmq->epTimer == NULL) {
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, tmq, tmqMgmt.timer);
1128
  }
L
Liu Jicong 已提交
1129 1130

  // init auto commit timer
1131
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
L
Liu Jicong 已提交
1132 1133 1134
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, tmq, tmqMgmt.timer);
  }

L
Liu Jicong 已提交
1135 1136 1137
  code = 0;
FAIL:
  if (req.topicNames != NULL) taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1138
  if (code != 0 && buf) {
L
Liu Jicong 已提交
1139 1140 1141
    taosMemoryFree(buf);
  }
  return code;
1142 1143
}

L
Liu Jicong 已提交
1144
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
L
Liu Jicong 已提交
1145
  //
1146
  conf->commitCb = cb;
L
Liu Jicong 已提交
1147
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1148
}
1149

L
Liu Jicong 已提交
1150
#if 0
L
Liu Jicong 已提交
1151 1152
int32_t tmqGetSkipLogNum(tmq_message_t* tmq_message) {
  if (tmq_message == NULL) return 0;
L
Liu Jicong 已提交
1153
  SMqPollRsp* pRsp = &tmq_message->msg;
L
Liu Jicong 已提交
1154 1155
  return pRsp->skipLogNum;
}
L
Liu Jicong 已提交
1156
#endif
L
Liu Jicong 已提交
1157

D
dapan1121 已提交
1158
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1159 1160
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1161
  SMqClientTopic* pTopic = pParam->pTopic;
X
Xiaoyu Wang 已提交
1162
  tmq_t*          tmq = pParam->tmq;
L
Liu Jicong 已提交
1163 1164 1165
  int32_t         vgId = pParam->vgId;
  int32_t         epoch = pParam->epoch;
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1166
  if (code != 0) {
S
Shengliang Guan 已提交
1167
    tscWarn("msg discard from vgId:%d, epoch %d, code:%x", vgId, epoch, code);
L
Liu Jicong 已提交
1168
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1169 1170 1171 1172
    if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM);
      if (pRspWrapper == NULL) {
        taosMemoryFree(pMsg->pData);
S
Shengliang Guan 已提交
1173
        tscWarn("msg discard from vgId:%d, epoch %d since out of memory", vgId, epoch);
L
Liu Jicong 已提交
1174 1175 1176 1177 1178 1179 1180 1181
        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 已提交
1182
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1183 1184
  }

X
Xiaoyu Wang 已提交
1185 1186 1187
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
  int32_t tmqEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < tmqEpoch) {
L
Liu Jicong 已提交
1188
    // do not write into queue since updating epoch reset
S
Shengliang Guan 已提交
1189
    tscWarn("msg discard from vgId:%d since from earlier epoch, rsp epoch %d, current epoch %d", vgId, msgEpoch,
L
Liu Jicong 已提交
1190
            tmqEpoch);
1191
    tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1192
    taosMemoryFree(pMsg->pData);
X
Xiaoyu Wang 已提交
1193 1194 1195 1196
    return 0;
  }

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

L
Liu Jicong 已提交
1200 1201 1202
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

1203
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM);
L
Liu Jicong 已提交
1204
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1205
    taosMemoryFree(pMsg->pData);
S
Shengliang Guan 已提交
1206
    tscWarn("msg discard from vgId:%d, epoch %d since out of memory", vgId, epoch);
L
fix txn  
Liu Jicong 已提交
1207
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1208
  }
L
Liu Jicong 已提交
1209

L
Liu Jicong 已提交
1210
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1211 1212
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
L
Liu Jicong 已提交
1213

L
Liu Jicong 已提交
1214
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1215 1216 1217
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqDataRsp(&decoder, &pRspWrapper->dataRsp);
L
Liu Jicong 已提交
1218
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1219 1220 1221
  } else {
    ASSERT(rspType == TMQ_MSG_TYPE__POLL_META_RSP);
    tDecodeSMqMetaRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pRspWrapper->metaRsp);
L
Liu Jicong 已提交
1222
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1223
  }
L
Liu Jicong 已提交
1224

L
Liu Jicong 已提交
1225
  taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1226

S
Shengliang Guan 已提交
1227 1228 1229
  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 已提交
1230

L
Liu Jicong 已提交
1231
  taosWriteQitem(tmq->mqueue, pRspWrapper);
1232
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1233

L
Liu Jicong 已提交
1234
  return 0;
L
fix txn  
Liu Jicong 已提交
1235
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1236
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1237 1238
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
1239
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1240
  return -1;
1241 1242
}

1243 1244 1245 1246 1247
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 已提交
1248
  tscDebug("consumer:%" PRId64 ", update ep epoch %d to epoch %d, topic num:%d", tmq->consumerId, tmq->epoch, epoch,
1249 1250 1251 1252 1253 1254 1255 1256 1257 1258 1259 1260 1261 1262 1263 1264 1265 1266
           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 已提交
1267
      tscDebug("consumer:%" PRId64 ", new vg num: %d", tmq->consumerId, vgNumCur);
1268 1269 1270
      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 已提交
1271 1272 1273 1274
        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 已提交
1275
        taosHashPut(pHash, vgKey, strlen(vgKey), &pVgCur->currentOffsetNew, sizeof(STqOffsetVal));
1276 1277 1278 1279 1280 1281 1282 1283 1284 1285 1286
      }
    }
  }

  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 已提交
1287
    tscDebug("consumer:%" PRId64 ", update topic: %s", tmq->consumerId, topic.topicName);
1288 1289 1290 1291 1292 1293

    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 已提交
1294 1295
      STqOffsetVal* pOffset = taosHashGet(pHash, vgKey, strlen(vgKey));
      STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
1296
      if (pOffset != NULL) {
L
Liu Jicong 已提交
1297
        offsetNew = *pOffset;
1298 1299
      }

S
Shengliang Guan 已提交
1300 1301
      /*tscDebug("consumer:%" PRId64 ", (epoch %d) offset of vgId:%d updated to %" PRId64 ", vgKey is %s",
       * tmq->consumerId, epoch,*/
L
Liu Jicong 已提交
1302
      /*pVgEp->vgId, offset, vgKey);*/
1303 1304
      SMqClientVg clientVg = {
          .pollCnt = 0,
L
Liu Jicong 已提交
1305
          .currentOffsetNew = offsetNew,
1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328
          .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 已提交
1329
#if 0
L
Liu Jicong 已提交
1330
bool tmqUpdateEp(tmq_t* tmq, int32_t epoch, SMqAskEpRsp* pRsp) {
L
Liu Jicong 已提交
1331
  /*printf("call update ep %d\n", epoch);*/
X
Xiaoyu Wang 已提交
1332
  bool    set = false;
L
Liu Jicong 已提交
1333 1334
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
S
Shengliang Guan 已提交
1335
  tscDebug("consumer:%" PRId64 ", update ep epoch %d to epoch %d, topic num: %d", tmq->consumerId, tmq->epoch, epoch,
L
Liu Jicong 已提交
1336
           topicNumGet);
L
Liu Jicong 已提交
1337 1338 1339 1340 1341 1342 1343 1344 1345 1346 1347 1348
  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 已提交
1349 1350
    SMqClientTopic topic = {0};
    SMqSubTopicEp* pTopicEp = taosArrayGet(pRsp->topics, i);
L
Liu Jicong 已提交
1351
    topic.schema = pTopicEp->schema;
L
Liu Jicong 已提交
1352
    taosHashClear(pHash);
X
Xiaoyu Wang 已提交
1353
    topic.topicName = strdup(pTopicEp->topic);
L
Liu Jicong 已提交
1354
    tstrncpy(topic.db, pTopicEp->db, TSDB_DB_FNAME_LEN);
L
Liu Jicong 已提交
1355

S
Shengliang Guan 已提交
1356
    tscDebug("consumer:%" PRId64 ", update topic: %s", tmq->consumerId, topic.topicName);
L
Liu Jicong 已提交
1357 1358 1359 1360 1361 1362
    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 已提交
1363
        tscDebug("consumer:%" PRId64 ", new vg num: %d", tmq->consumerId, vgNumCur);
L
Liu Jicong 已提交
1364 1365 1366 1367
        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 已提交
1368
          tscDebug("consumer:%" PRId64 ", epoch %d vgId:%d build %s", tmq->consumerId, epoch, pVgCur->vgId, vgKey);
L
Liu Jicong 已提交
1369 1370 1371 1372 1373 1374 1375 1376 1377
          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 已提交
1378
      SMqSubVgEp* pVgEp = taosArrayGet(pTopicEp->vgs, j);
L
Liu Jicong 已提交
1379 1380 1381
      sprintf(vgKey, "%s:%d", topic.topicName, pVgEp->vgId);
      int64_t* pOffset = taosHashGet(pHash, vgKey, strlen(vgKey));
      int64_t  offset = pVgEp->offset;
S
Shengliang Guan 已提交
1382
      tscDebug("consumer:%" PRId64 ", (epoch %d) original offset of vgId:%d is %" PRId64, tmq->consumerId, epoch, pVgEp->vgId, offset);
L
Liu Jicong 已提交
1383 1384
      if (pOffset != NULL) {
        offset = *pOffset;
S
Shengliang Guan 已提交
1385
        tscDebug("consumer:%" PRId64 ", (epoch %d) receive offset of vgId:%d, full key is %s", tmq->consumerId, epoch, pVgEp->vgId,
1386
                 vgKey);
L
Liu Jicong 已提交
1387
      }
S
Shengliang Guan 已提交
1388
      tscDebug("consumer:%" PRId64 ", (epoch %d) offset of vgId:%d updated to %" PRId64, tmq->consumerId, epoch, pVgEp->vgId, offset);
X
Xiaoyu Wang 已提交
1389 1390
      SMqClientVg clientVg = {
          .pollCnt = 0,
L
Liu Jicong 已提交
1391
          .currentOffset = offset,
X
Xiaoyu Wang 已提交
1392 1393 1394
          .vgId = pVgEp->vgId,
          .epSet = pVgEp->epSet,
          .vgStatus = TMQ_VG_STATUS__IDLE,
L
Liu Jicong 已提交
1395
          .vgSkipCnt = 0,
X
Xiaoyu Wang 已提交
1396 1397 1398 1399
      };
      taosArrayPush(topic.vgs, &clientVg);
      set = true;
    }
L
Liu Jicong 已提交
1400
    taosArrayPush(newTopics, &topic);
X
Xiaoyu Wang 已提交
1401
  }
L
Liu Jicong 已提交
1402
  if (tmq->clientTopics) taosArrayDestroy(tmq->clientTopics);
L
Liu Jicong 已提交
1403
  taosHashCleanup(pHash);
L
Liu Jicong 已提交
1404
  tmq->clientTopics = newTopics;
1405

1406 1407 1408 1409
  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);
1410

X
Xiaoyu Wang 已提交
1411 1412 1413
  atomic_store_32(&tmq->epoch, epoch);
  return set;
}
L
Liu Jicong 已提交
1414
#endif
X
Xiaoyu Wang 已提交
1415

D
dapan1121 已提交
1416
int32_t tmqAskEpCb(void* param, SDataBuf* pMsg, int32_t code) {
1417
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
L
Liu Jicong 已提交
1418
  tmq_t*           tmq = pParam->tmq;
L
Liu Jicong 已提交
1419
  int8_t           async = pParam->async;
L
Liu Jicong 已提交
1420
  pParam->code = code;
1421
  if (code != 0) {
S
Shengliang Guan 已提交
1422
    tscError("consumer:%" PRId64 ", get topic endpoint error, not ready, wait:%d", tmq->consumerId, pParam->async);
L
Liu Jicong 已提交
1423
    goto END;
1424
  }
L
Liu Jicong 已提交
1425

L
Liu Jicong 已提交
1426
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1427
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1428
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1429 1430
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
S
Shengliang Guan 已提交
1431
  tscDebug("consumer:%" PRId64 ", recv ep, msg epoch %d, current epoch %d", tmq->consumerId, head->epoch, epoch);
L
Liu Jicong 已提交
1432 1433
  if (head->epoch <= epoch) {
    goto END;
1434
  }
L
Liu Jicong 已提交
1435

L
Liu Jicong 已提交
1436
  if (!async) {
L
Liu Jicong 已提交
1437 1438
    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
S
Shengliang Guan 已提交
1439 1440
    /*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 已提交
1441
    tmqUpdateEp2(tmq, head->epoch, &rsp);
L
Liu Jicong 已提交
1442
    tDeleteSMqAskEpRsp(&rsp);
X
Xiaoyu Wang 已提交
1443
  } else {
1444
    SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM);
L
Liu Jicong 已提交
1445
    if (pWrapper == NULL) {
X
Xiaoyu Wang 已提交
1446
      terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1447 1448
      code = -1;
      goto END;
X
Xiaoyu Wang 已提交
1449
    }
L
Liu Jicong 已提交
1450 1451 1452
    pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
    pWrapper->epoch = head->epoch;
    memcpy(&pWrapper->msg, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1453
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pWrapper->msg);
L
Liu Jicong 已提交
1454

L
Liu Jicong 已提交
1455
    taosWriteQitem(tmq->mqueue, pWrapper);
1456
    tsem_post(&tmq->rspSem);
1457
  }
L
Liu Jicong 已提交
1458 1459

END:
L
Liu Jicong 已提交
1460
  /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1461
  if (!async) {
L
Liu Jicong 已提交
1462
    tsem_post(&pParam->rspSem);
L
Liu Jicong 已提交
1463 1464
  } else {
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
1465 1466
  }
  return code;
1467 1468
}

L
Liu Jicong 已提交
1469
int32_t tmqAskEp(tmq_t* tmq, bool async) {
L
Liu Jicong 已提交
1470
  int32_t code = 0;
L
Liu Jicong 已提交
1471
#if 0
L
Liu Jicong 已提交
1472
  int8_t  epStatus = atomic_val_compare_exchange_8(&tmq->epStatus, 0, 1);
L
fix txn  
Liu Jicong 已提交
1473
  if (epStatus == 1) {
L
temp  
Liu Jicong 已提交
1474
    int32_t epSkipCnt = atomic_add_fetch_32(&tmq->epSkipCnt, 1);
S
Shengliang Guan 已提交
1475
    tscTrace("consumer:%" PRId64 ", skip ask ep cnt %d", tmq->consumerId, epSkipCnt);
L
Liu Jicong 已提交
1476
    if (epSkipCnt < 5000) return 0;
L
fix txn  
Liu Jicong 已提交
1477
  }
L
temp  
Liu Jicong 已提交
1478
  atomic_store_32(&tmq->epSkipCnt, 0);
L
Liu Jicong 已提交
1479
#endif
L
Liu Jicong 已提交
1480
  int32_t      tlen = sizeof(SMqAskEpReq);
L
Liu Jicong 已提交
1481
  SMqAskEpReq* req = taosMemoryCalloc(1, tlen);
L
Liu Jicong 已提交
1482
  if (req == NULL) {
L
Liu Jicong 已提交
1483
    tscError("failed to malloc get subscribe ep buf");
L
Liu Jicong 已提交
1484
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1485
    return -1;
L
Liu Jicong 已提交
1486
  }
L
Liu Jicong 已提交
1487 1488 1489
  req->consumerId = htobe64(tmq->consumerId);
  req->epoch = htonl(tmq->epoch);
  strcpy(req->cgroup, tmq->groupId);
1490

L
Liu Jicong 已提交
1491
  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
L
Liu Jicong 已提交
1492 1493
  if (pParam == NULL) {
    tscError("failed to malloc subscribe param");
wafwerar's avatar
wafwerar 已提交
1494
    taosMemoryFree(req);
L
Liu Jicong 已提交
1495
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1496
    return -1;
L
Liu Jicong 已提交
1497 1498
  }
  pParam->tmq = tmq;
L
Liu Jicong 已提交
1499
  pParam->async = async;
X
Xiaoyu Wang 已提交
1500
  tsem_init(&pParam->rspSem, 0, 0);
1501

1502
  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
1503 1504
  if (sendInfo == NULL) {
    tsem_destroy(&pParam->rspSem);
wafwerar's avatar
wafwerar 已提交
1505 1506
    taosMemoryFree(pParam);
    taosMemoryFree(req);
L
Liu Jicong 已提交
1507
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1508 1509 1510 1511 1512 1513 1514 1515 1516 1517
    return -1;
  }

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

  sendInfo->requestId = generateRequestId();
L
Liu Jicong 已提交
1518 1519
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
D
dapan1121 已提交
1520
  sendInfo->paramFreeFp = taosMemoryFree;      
L
Liu Jicong 已提交
1521
  sendInfo->fp = tmqAskEpCb;
L
Liu Jicong 已提交
1522
  sendInfo->msgType = TDMT_MND_MQ_ASK_EP;
1523

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

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

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

L
Liu Jicong 已提交
1531
  if (!async) {
L
Liu Jicong 已提交
1532 1533 1534 1535 1536
    tsem_wait(&pParam->rspSem);
    code = pParam->code;
    taosMemoryFree(pParam);
  }
  return code;
1537 1538
}

L
Liu Jicong 已提交
1539 1540
#if 0
int32_t tmq_seek(tmq_t* tmq, const tmq_topic_vgroup_t* offset) {
L
Liu Jicong 已提交
1541 1542 1543 1544 1545 1546 1547 1548 1549 1550 1551 1552 1553
  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 已提交
1554
          tmqClearUnhandleMsg(tmq);
L
Liu Jicong 已提交
1555 1556 1557 1558 1559 1560 1561
          return TMQ_RESP_ERR__SUCCESS;
        }
      }
    }
  }
  return TMQ_RESP_ERR__FAIL;
}
L
Liu Jicong 已提交
1562
#endif
L
Liu Jicong 已提交
1563

1564
SMqPollReq* tmqBuildConsumeReqImpl(tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1565 1566 1567 1568 1569 1570 1571 1572 1573 1574
  /*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 已提交
1575

L
Liu Jicong 已提交
1576
  SMqPollReq* pReq = taosMemoryCalloc(1, sizeof(SMqPollReq));
1577 1578 1579
  if (pReq == NULL) {
    return NULL;
  }
L
Liu Jicong 已提交
1580

L
Liu Jicong 已提交
1581 1582 1583
  /*strcpy(pReq->topic, pTopic->topicName);*/
  /*strcpy(pReq->cgroup, tmq->groupId);*/

L
Liu Jicong 已提交
1584 1585 1586 1587
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1588

1589
  pReq->withTbName = tmq->withTbName;
1590
  pReq->timeout = timeout;
L
Liu Jicong 已提交
1591
  pReq->consumerId = tmq->consumerId;
X
Xiaoyu Wang 已提交
1592
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1593 1594
  /*pReq->currentOffset = reqOffset;*/
  pReq->reqOffset = pVg->currentOffsetNew;
L
Liu Jicong 已提交
1595
  pReq->reqId = generateRequestId();
1596

L
Liu Jicong 已提交
1597 1598
  pReq->useSnapshot = tmq->useSnapshot;

1599
  pReq->head.vgId = htonl(pVg->vgId);
L
Liu Jicong 已提交
1600
  pReq->head.contLen = htonl(sizeof(SMqPollReq));
1601 1602 1603
  return pReq;
}

L
Liu Jicong 已提交
1604 1605
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1606
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1607 1608 1609 1610 1611 1612 1613 1614
  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 已提交
1615 1616 1617
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
L
Liu Jicong 已提交
1618 1619
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
L
Liu Jicong 已提交
1620
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1621
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1622
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1623

L
Liu Jicong 已提交
1624 1625
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
L
Liu Jicong 已提交
1626
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1627 1628
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1629

L
Liu Jicong 已提交
1630
  return pRspObj;
X
Xiaoyu Wang 已提交
1631 1632
}

1633
int32_t tmqPollImpl(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1634
  /*tscDebug("call poll");*/
X
Xiaoyu Wang 已提交
1635 1636 1637 1638 1639 1640
  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 已提交
1641
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
L
Liu Jicong 已提交
1642 1643
        tscTrace("consumer:%" PRId64 ", epoch %d skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch, pVg->vgId,
                 vgSkipCnt);
X
Xiaoyu Wang 已提交
1644
        continue;
L
Liu Jicong 已提交
1645
        /*if (vgSkipCnt < 10000) continue;*/
L
temp  
Liu Jicong 已提交
1646 1647 1648 1649
#if 0
        if (skipCnt < 30000) {
          continue;
        } else {
S
Shengliang Guan 已提交
1650
        tscDebug("consumer:%" PRId64 ",skip vgId:%d skip too much reset", tmq->consumerId, pVg->vgId);
L
temp  
Liu Jicong 已提交
1651 1652
        }
#endif
X
Xiaoyu Wang 已提交
1653
      }
L
Liu Jicong 已提交
1654
      atomic_store_32(&pVg->vgSkipCnt, 0);
1655
      SMqPollReq* pReq = tmqBuildConsumeReqImpl(tmq, timeout, pTopic, pVg);
X
Xiaoyu Wang 已提交
1656 1657
      if (pReq == NULL) {
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1658
        tsem_post(&tmq->rspSem);
X
Xiaoyu Wang 已提交
1659 1660
        return -1;
      }
wafwerar's avatar
wafwerar 已提交
1661
      SMqPollCbParam* pParam = taosMemoryMalloc(sizeof(SMqPollCbParam));
L
Liu Jicong 已提交
1662
      if (pParam == NULL) {
wafwerar's avatar
wafwerar 已提交
1663
        taosMemoryFree(pReq);
X
Xiaoyu Wang 已提交
1664
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1665
        tsem_post(&tmq->rspSem);
X
Xiaoyu Wang 已提交
1666 1667
        return -1;
      }
L
Liu Jicong 已提交
1668 1669
      pParam->tmq = tmq;
      pParam->pVg = pVg;
L
Liu Jicong 已提交
1670
      pParam->pTopic = pTopic;
L
Liu Jicong 已提交
1671
      pParam->vgId = pVg->vgId;
L
Liu Jicong 已提交
1672 1673
      pParam->epoch = tmq->epoch;

1674
      SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
1675
      if (sendInfo == NULL) {
wafwerar's avatar
wafwerar 已提交
1676 1677
        taosMemoryFree(pReq);
        taosMemoryFree(pParam);
L
Liu Jicong 已提交
1678
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1679
        tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1680 1681 1682 1683
        return -1;
      }

      sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1684
          .pData = pReq,
L
Liu Jicong 已提交
1685
          .len = sizeof(SMqPollReq),
X
Xiaoyu Wang 已提交
1686 1687
          .handle = NULL,
      };
L
Liu Jicong 已提交
1688
      sendInfo->requestId = pReq->reqId;
X
Xiaoyu Wang 已提交
1689
      sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1690
      sendInfo->param = pParam;
X
Xiaoyu Wang 已提交
1691
      sendInfo->fp = tmqPollCb;
L
Liu Jicong 已提交
1692
      sendInfo->msgType = TDMT_VND_CONSUME;
X
Xiaoyu Wang 已提交
1693 1694

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

1697 1698
      char offsetFormatBuf[80];
      tFormatOffset(offsetFormatBuf, 80, &pVg->currentOffsetNew);
L
Liu Jicong 已提交
1699 1700
      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 已提交
1701
      /*printf("send vgId:%d %" PRId64 "\n", pVg->vgId, pVg->currentOffset);*/
X
Xiaoyu Wang 已提交
1702 1703 1704 1705 1706 1707 1708 1709
      asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);
      pVg->pollCnt++;
      tmq->pollCnt++;
    }
  }
  return 0;
}

L
Liu Jicong 已提交
1710 1711
int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1712
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1713 1714
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1715
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
L
Liu Jicong 已提交
1716
      tmqUpdateEp2(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1717
      /*tmqClearUnhandleMsg(tmq);*/
X
Xiaoyu Wang 已提交
1718 1719 1720 1721 1722 1723 1724 1725 1726 1727
      *pReset = true;
    } else {
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

L
Liu Jicong 已提交
1728
void* tmqHandleAllRsp(tmq_t* tmq, int64_t timeout, bool pollIfReset) {
X
Xiaoyu Wang 已提交
1729
  while (1) {
L
Liu Jicong 已提交
1730 1731 1732
    SMqRspWrapper* rspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper == NULL) {
L
Liu Jicong 已提交
1733
      taosReadAllQitems(tmq->mqueue, tmq->qall);
L
Liu Jicong 已提交
1734 1735
      taosGetQitem(tmq->qall, (void**)&rspWrapper);
      if (rspWrapper == NULL) return NULL;
X
Xiaoyu Wang 已提交
1736 1737
    }

L
Liu Jicong 已提交
1738 1739 1740 1741 1742
    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 已提交
1743
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1744
      /*atomic_sub_fetch_32(&tmq->readyRequest, 1);*/
1745
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
L
Liu Jicong 已提交
1746
      if (pollRspWrapper->dataRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1747
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
L
Liu Jicong 已提交
1748 1749
        /*printf("vgId:%d offset %" PRId64 " up to %" PRId64 "\n", pVg->vgId, pVg->currentOffset,
         * rspMsg->msg.rspOffset);*/
L
Liu Jicong 已提交
1750
        pVg->currentOffsetNew = pollRspWrapper->dataRsp.rspOffset;
X
Xiaoyu Wang 已提交
1751
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
L
Liu Jicong 已提交
1752
        if (pollRspWrapper->dataRsp.blockNum == 0) {
L
Liu Jicong 已提交
1753 1754
          taosFreeQitem(pollRspWrapper);
          rspWrapper = NULL;
L
temp  
Liu Jicong 已提交
1755 1756
          continue;
        }
L
Liu Jicong 已提交
1757
        // build rsp
L
Liu Jicong 已提交
1758
        SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1759
        taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1760
        return pRsp;
X
Xiaoyu Wang 已提交
1761
      } else {
L
Liu Jicong 已提交
1762 1763 1764 1765 1766 1767 1768
        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 已提交
1769
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1770
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
L
Liu Jicong 已提交
1771 1772
        /*printf("vgId:%d offset %" PRId64 " up to %" PRId64 "\n", pVg->vgId, pVg->currentOffset,
         * rspMsg->msg.rspOffset);*/
wmmhello's avatar
wmmhello 已提交
1773 1774
        pVg->currentOffsetNew.version = pollRspWrapper->metaRsp.rspOffset;
        pVg->currentOffsetNew.type = TMQ_OFFSET__LOG;
L
Liu Jicong 已提交
1775 1776
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1777
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1778 1779 1780 1781
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
        tscDebug("msg discard since epoch mismatch: msg epoch %d, consumer epoch %d\n",
L
Liu Jicong 已提交
1782
                 pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
L
Liu Jicong 已提交
1783
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1784 1785
      }
    } else {
L
fix  
Liu Jicong 已提交
1786
      /*printf("handle ep rsp %d\n", rspMsg->head.mqMsgType);*/
X
Xiaoyu Wang 已提交
1787
      bool reset = false;
L
Liu Jicong 已提交
1788 1789
      tmqHandleNoPollRsp(tmq, rspWrapper, &reset);
      taosFreeQitem(rspWrapper);
X
Xiaoyu Wang 已提交
1790
      if (pollIfReset && reset) {
S
Shengliang Guan 已提交
1791
        tscDebug("consumer:%" PRId64 ", reset and repoll", tmq->consumerId);
1792
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1793 1794 1795 1796 1797
      }
    }
  }
}

1798
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1799
  /*tscDebug("call poll1");*/
L
Liu Jicong 已提交
1800 1801
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1802

1803 1804 1805
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1806
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1807 1808
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1809
  }
1810
#endif
X
Xiaoyu Wang 已提交
1811

L
Liu Jicong 已提交
1812
  // in no topic status also need process delayed task
L
Liu Jicong 已提交
1813
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1814 1815 1816
    return NULL;
  }

X
Xiaoyu Wang 已提交
1817
  while (1) {
L
Liu Jicong 已提交
1818
    tmqHandleAllDelayedTask(tmq);
1819
    if (tmqPollImpl(tmq, timeout) < 0) return NULL;
L
Liu Jicong 已提交
1820

1821
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1822 1823
    if (rspObj) {
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1824 1825
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
      return NULL;
X
Xiaoyu Wang 已提交
1826
    }
1827
    if (timeout != -1) {
X
Xiaoyu Wang 已提交
1828
      int64_t endTime = taosGetTimestampMs();
1829
      int64_t leftTime = endTime - startTime;
1830
      if (leftTime > timeout) {
S
Shengliang Guan 已提交
1831
        tscDebug("consumer:%" PRId64 ", (epoch %d) timeout, no rsp", tmq->consumerId, tmq->epoch);
X
Xiaoyu Wang 已提交
1832 1833
        return NULL;
      }
1834
      tsem_timewait(&tmq->rspSem, leftTime * 1000);
L
Liu Jicong 已提交
1835 1836 1837
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
      tsem_timewait(&tmq->rspSem, 500 * 1000);
X
Xiaoyu Wang 已提交
1838 1839 1840 1841
    }
  }
}

L
Liu Jicong 已提交
1842
int32_t tmq_consumer_close(tmq_t* tmq) {
1843
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
L
Liu Jicong 已提交
1844 1845
    int32_t rsp = tmq_commit_sync(tmq, NULL);
    if (rsp != 0) {
L
Liu Jicong 已提交
1846
      return rsp;
1847 1848 1849 1850
    }

    tmq_list_t* lst = tmq_list_new();
    rsp = tmq_subscribe(tmq, lst);
1851
    tmq_list_destroy(lst);
1852

L
Liu Jicong 已提交
1853
    if (rsp != 0) {
L
Liu Jicong 已提交
1854
      return rsp;
1855
    }
L
Liu Jicong 已提交
1856
  }
1857
  // TODO: free resources
L
Liu Jicong 已提交
1858
  return 0;
1859
}
L
Liu Jicong 已提交
1860

L
Liu Jicong 已提交
1861 1862
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
1863
    return "success";
L
Liu Jicong 已提交
1864
  } else if (err == -1) {
L
Liu Jicong 已提交
1865 1866 1867
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
1868 1869
  }
}
L
Liu Jicong 已提交
1870

L
Liu Jicong 已提交
1871 1872 1873 1874 1875 1876 1877 1878 1879 1880
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)) {
    return TMQ_RES_TABLE_META;
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
1881
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
1882 1883
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
1884
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
1885 1886 1887
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
1888 1889 1890 1891 1892
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
1893 1894 1895 1896
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 已提交
1897 1898 1899
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
L
Liu Jicong 已提交
1900 1901 1902 1903 1904
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
1905 1906 1907 1908
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
1909 1910 1911
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
L
Liu Jicong 已提交
1912 1913 1914 1915
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
1916 1917 1918 1919 1920 1921 1922 1923

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;
    }
1924
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
L
Liu Jicong 已提交
1925 1926 1927
  }
  return NULL;
}
1928

1929
tmq_raw_data* tmq_get_raw_meta(TAOS_RES* res) {
L
Liu Jicong 已提交
1930
  if (TD_RES_TMQ_META(res)) {
1931
    tmq_raw_data*  raw = taosMemoryCalloc(1, sizeof(tmq_raw_data));
L
Liu Jicong 已提交
1932
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
wmmhello's avatar
wmmhello 已提交
1933 1934 1935 1936
    raw->raw_meta = pMetaRspObj->metaRsp.metaRsp;
    raw->raw_meta_len = pMetaRspObj->metaRsp.metaRspLen;
    raw->raw_meta_type = pMetaRspObj->metaRsp.resMsgType;
    return raw;
L
Liu Jicong 已提交
1937
  }
wmmhello's avatar
wmmhello 已提交
1938
  return NULL;
L
Liu Jicong 已提交
1939 1940
}

1941 1942
static char* buildCreateTableJson(SSchemaWrapper* schemaRow, SSchemaWrapper* schemaTag, char* name, int64_t id,
                                  int8_t t) {
wmmhello's avatar
wmmhello 已提交
1943 1944 1945 1946 1947 1948 1949
  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 已提交
1950

1951 1952 1953 1954
  //  char uid[32] = {0};
  //  sprintf(uid, "%"PRIi64, id);
  //  cJSON* id_ = cJSON_CreateString(uid);
  //  cJSON_AddItemToObject(json, "id", id_);
wmmhello's avatar
wmmhello 已提交
1955 1956 1957 1958
  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);
1959 1960
  //  cJSON* version = cJSON_CreateNumber(1);
  //  cJSON_AddItemToObject(json, "version", version);
wmmhello's avatar
wmmhello 已提交
1961 1962

  cJSON* columns = cJSON_CreateArray();
1963 1964 1965 1966
  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 已提交
1967 1968 1969
    cJSON_AddItemToObject(column, "name", cname);
    cJSON* ctype = cJSON_CreateNumber(s->type);
    cJSON_AddItemToObject(column, "type", ctype);
1970
    if (s->type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
1971
      int32_t length = s->bytes - VARSTR_HEADER_SIZE;
1972
      cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1973
      cJSON_AddItemToObject(column, "length", cbytes);
1974 1975 1976
    } 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 已提交
1977 1978
      cJSON_AddItemToObject(column, "length", cbytes);
    }
wmmhello's avatar
wmmhello 已提交
1979 1980 1981 1982 1983
    cJSON_AddItemToArray(columns, column);
  }
  cJSON_AddItemToObject(json, "columns", columns);

  cJSON* tags = cJSON_CreateArray();
1984 1985 1986 1987
  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 已提交
1988 1989 1990
    cJSON_AddItemToObject(tag, "name", tname);
    cJSON* ttype = cJSON_CreateNumber(s->type);
    cJSON_AddItemToObject(tag, "type", ttype);
1991
    if (s->type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
1992
      int32_t length = s->bytes - VARSTR_HEADER_SIZE;
1993
      cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
1994
      cJSON_AddItemToObject(tag, "length", cbytes);
1995 1996 1997
    } 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 已提交
1998 1999
      cJSON_AddItemToObject(tag, "length", cbytes);
    }
wmmhello's avatar
wmmhello 已提交
2000 2001 2002 2003 2004 2005 2006 2007 2008
    cJSON_AddItemToArray(tags, tag);
  }
  cJSON_AddItemToObject(json, "tags", tags);

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

2009 2010 2011 2012
static char* buildAlterSTableJson(void* alterData, int32_t alterDataLen) {
  SMAlterStbReq req = {0};
  cJSON*        json = NULL;
  char*         string = NULL;
wmmhello's avatar
wmmhello 已提交
2013 2014 2015 2016 2017 2018 2019 2020 2021 2022 2023

  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);
2024 2025
  //  cJSON* uid = cJSON_CreateNumber(id);
  //  cJSON_AddItemToObject(json, "uid", uid);
wmmhello's avatar
wmmhello 已提交
2026 2027 2028 2029 2030 2031 2032 2033 2034 2035 2036 2037
  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: {
2038 2039
      TAOS_FIELD* field = taosArrayGet(req.pFields, 0);
      cJSON*      colName = cJSON_CreateString(field->name);
wmmhello's avatar
wmmhello 已提交
2040 2041 2042 2043
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colType = cJSON_CreateNumber(field->type);
      cJSON_AddItemToObject(json, "colType", colType);

2044
      if (field->type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
2045
        int32_t length = field->bytes - VARSTR_HEADER_SIZE;
2046
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
2047
        cJSON_AddItemToObject(json, "colLength", cbytes);
2048 2049 2050
      } 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 已提交
2051 2052 2053 2054 2055
        cJSON_AddItemToObject(json, "colLength", cbytes);
      }
      break;
    }
    case TSDB_ALTER_TABLE_DROP_TAG:
2056 2057 2058
    case TSDB_ALTER_TABLE_DROP_COLUMN: {
      TAOS_FIELD* field = taosArrayGet(req.pFields, 0);
      cJSON*      colName = cJSON_CreateString(field->name);
wmmhello's avatar
wmmhello 已提交
2059 2060 2061 2062
      cJSON_AddItemToObject(json, "colName", colName);
      break;
    }
    case TSDB_ALTER_TABLE_UPDATE_TAG_BYTES:
2063 2064 2065
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_BYTES: {
      TAOS_FIELD* field = taosArrayGet(req.pFields, 0);
      cJSON*      colName = cJSON_CreateString(field->name);
wmmhello's avatar
wmmhello 已提交
2066 2067 2068
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colType = cJSON_CreateNumber(field->type);
      cJSON_AddItemToObject(json, "colType", colType);
2069
      if (field->type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
2070
        int32_t length = field->bytes - VARSTR_HEADER_SIZE;
2071
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
2072
        cJSON_AddItemToObject(json, "colLength", cbytes);
2073 2074 2075
      } 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 已提交
2076 2077 2078 2079 2080
        cJSON_AddItemToObject(json, "colLength", cbytes);
      }
      break;
    }
    case TSDB_ALTER_TABLE_UPDATE_TAG_NAME:
2081 2082 2083 2084
    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 已提交
2085 2086 2087 2088 2089 2090 2091 2092 2093 2094
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colNewName = cJSON_CreateString(newField->name);
      cJSON_AddItemToObject(json, "colNewName", colNewName);
      break;
    }
    default:
      break;
  }
  string = cJSON_PrintUnformatted(json);

2095
end:
wmmhello's avatar
wmmhello 已提交
2096 2097 2098 2099 2100
  cJSON_Delete(json);
  tFreeSMAltertbReq(&req);
  return string;
}

2101
static char* processCreateStb(SMqMetaRsp* metaRsp) {
wmmhello's avatar
wmmhello 已提交
2102 2103
  SVCreateStbReq req = {0};
  SDecoder       coder;
2104
  char*          string = NULL;
wmmhello's avatar
wmmhello 已提交
2105 2106

  // decode and process req
2107
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2108 2109 2110 2111 2112 2113 2114 2115 2116 2117
  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;

2118
_err:
wmmhello's avatar
wmmhello 已提交
2119 2120 2121 2122
  tDecoderClear(&coder);
  return string;
}

2123
static char* processAlterStb(SMqMetaRsp* metaRsp) {
wmmhello's avatar
wmmhello 已提交
2124 2125
  SVCreateStbReq req = {0};
  SDecoder       coder;
2126
  char*          string = NULL;
wmmhello's avatar
wmmhello 已提交
2127 2128

  // decode and process req
2129
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2130 2131 2132 2133 2134 2135 2136 2137 2138 2139
  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;

2140
_err:
wmmhello's avatar
wmmhello 已提交
2141 2142 2143 2144
  tDecoderClear(&coder);
  return string;
}

2145 2146
static char* buildCreateCTableJson(STag* pTag, char* sname, char* name, SArray* tagName, int64_t id) {
  char*   string = NULL;
wmmhello's avatar
wmmhello 已提交
2147
  SArray* pTagVals = NULL;
2148
  cJSON*  json = cJSON_CreateObject();
wmmhello's avatar
wmmhello 已提交
2149 2150 2151 2152 2153
  if (json == NULL) {
    return string;
  }
  cJSON* type = cJSON_CreateString("create");
  cJSON_AddItemToObject(json, "type", type);
2154 2155 2156 2157
  //  char cid[32] = {0};
  //  sprintf(cid, "%"PRIi64, id);
  //  cJSON* cid_ = cJSON_CreateString(cid);
  //  cJSON_AddItemToObject(json, "id", cid_);
wmmhello's avatar
wmmhello 已提交
2158

wmmhello's avatar
wmmhello 已提交
2159 2160 2161 2162
  cJSON* tableName = cJSON_CreateString(name);
  cJSON_AddItemToObject(json, "tableName", tableName);
  cJSON* tableType = cJSON_CreateString("child");
  cJSON_AddItemToObject(json, "tableType", tableType);
wmmhello's avatar
wmmhello 已提交
2163
  cJSON* using = cJSON_CreateString(sname);
wmmhello's avatar
wmmhello 已提交
2164
  cJSON_AddItemToObject(json, "using", using);
2165 2166
  //  cJSON* version = cJSON_CreateNumber(1);
  //  cJSON_AddItemToObject(json, "version", version);
wmmhello's avatar
wmmhello 已提交
2167

2168
  cJSON*  tags = cJSON_CreateArray();
wmmhello's avatar
wmmhello 已提交
2169 2170 2171 2172
  int32_t code = tTagToValArray(pTag, &pTagVals);
  if (code) {
    goto end;
  }
wmmhello's avatar
wmmhello 已提交
2173

wmmhello's avatar
wmmhello 已提交
2174 2175
  if (tTagIsJson(pTag)) {
    STag* p = (STag*)pTag;
2176
    if (p->nTag == 0) {
wmmhello's avatar
wmmhello 已提交
2177 2178
      goto end;
    }
2179 2180
    char*    pJson = parseTagDatatoJson(pTag);
    cJSON*   tag = cJSON_CreateObject();
wmmhello's avatar
wmmhello 已提交
2181 2182
    STagVal* pTagVal = taosArrayGet(pTagVals, 0);

wmmhello's avatar
wmmhello 已提交
2183 2184
    char*  ptname = taosArrayGet(tagName, 0);
    cJSON* tname = cJSON_CreateString(ptname);
wmmhello's avatar
wmmhello 已提交
2185
    cJSON_AddItemToObject(tag, "name", tname);
2186 2187
    //    cJSON* cid_ = cJSON_CreateString("");
    //    cJSON_AddItemToObject(tag, "cid", cid_);
wmmhello's avatar
wmmhello 已提交
2188 2189
    cJSON* ttype = cJSON_CreateNumber(TSDB_DATA_TYPE_JSON);
    cJSON_AddItemToObject(tag, "type", ttype);
wmmhello's avatar
wmmhello 已提交
2190
    cJSON* tvalue = cJSON_CreateString(pJson);
wmmhello's avatar
wmmhello 已提交
2191 2192
    cJSON_AddItemToObject(tag, "value", tvalue);
    cJSON_AddItemToArray(tags, tag);
wmmhello's avatar
wmmhello 已提交
2193
    taosMemoryFree(pJson);
wmmhello's avatar
wmmhello 已提交
2194 2195 2196
    goto end;
  }

2197
  for (int i = 0; i < taosArrayGetSize(pTagVals); i++) {
wmmhello's avatar
wmmhello 已提交
2198 2199 2200
    STagVal* pTagVal = (STagVal*)taosArrayGet(pTagVals, i);

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

wmmhello's avatar
wmmhello 已提交
2202 2203
    char*  ptname = taosArrayGet(tagName, i);
    cJSON* tname = cJSON_CreateString(ptname);
wmmhello's avatar
wmmhello 已提交
2204
    cJSON_AddItemToObject(tag, "name", tname);
2205 2206
    //    cJSON* cid = cJSON_CreateNumber(pTagVal->cid);
    //    cJSON_AddItemToObject(tag, "cid", cid);
wmmhello's avatar
wmmhello 已提交
2207 2208
    cJSON* ttype = cJSON_CreateNumber(pTagVal->type);
    cJSON_AddItemToObject(tag, "type", ttype);
wmmhello's avatar
wmmhello 已提交
2209 2210 2211 2212 2213 2214 2215 2216 2217 2218 2219 2220

    char* buf = NULL;
    if (IS_VAR_DATA_TYPE(pTagVal->type)) {
      buf = taosMemoryCalloc(pTagVal->nData + 1, 1);
      dataConverToStr(buf, pTagVal->type, pTagVal->pData, pTagVal->nData, NULL);
    } else {
      buf = taosMemoryCalloc(32, 1);
      dataConverToStr(buf, pTagVal->type, &pTagVal->i64, tDataTypes[pTagVal->type].bytes, NULL);
    }

    cJSON* tvalue = cJSON_CreateString(buf);
    taosMemoryFree(buf);
wmmhello's avatar
wmmhello 已提交
2221 2222 2223
    cJSON_AddItemToObject(tag, "value", tvalue);
    cJSON_AddItemToArray(tags, tag);
  }
wmmhello's avatar
wmmhello 已提交
2224

2225
end:
wmmhello's avatar
wmmhello 已提交
2226 2227 2228
  cJSON_AddItemToObject(json, "tags", tags);
  string = cJSON_PrintUnformatted(json);
  cJSON_Delete(json);
wmmhello's avatar
wmmhello 已提交
2229
  taosArrayDestroy(pTagVals);
wmmhello's avatar
wmmhello 已提交
2230 2231 2232
  return string;
}

2233
static char* processCreateTable(SMqMetaRsp* metaRsp) {
wmmhello's avatar
wmmhello 已提交
2234 2235
  SDecoder           decoder = {0};
  SVCreateTbBatchReq req = {0};
2236 2237
  SVCreateTbReq*     pCreateReq;
  char*              string = NULL;
wmmhello's avatar
wmmhello 已提交
2238
  // decode
2239
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2240 2241 2242 2243 2244 2245 2246 2247 2248
  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;
2249 2250 2251 2252 2253 2254
    if (pCreateReq->type == TSDB_CHILD_TABLE) {
      string = buildCreateCTableJson((STag*)pCreateReq->ctb.pTag, pCreateReq->ctb.name, pCreateReq->name,
                                     pCreateReq->ctb.tagName, pCreateReq->uid);
    } else if (pCreateReq->type == TSDB_NORMAL_TABLE) {
      string =
          buildCreateTableJson(&pCreateReq->ntb.schemaRow, NULL, pCreateReq->name, pCreateReq->uid, TSDB_NORMAL_TABLE);
wmmhello's avatar
wmmhello 已提交
2255 2256 2257 2258 2259
    }
  }

  tDecoderClear(&decoder);

2260
_exit:
wmmhello's avatar
wmmhello 已提交
2261 2262 2263 2264
  tDecoderClear(&decoder);
  return string;
}

2265 2266 2267 2268
static char* processAlterTable(SMqMetaRsp* metaRsp) {
  SDecoder     decoder = {0};
  SVAlterTbReq vAlterTbReq = {0};
  char*        string = NULL;
wmmhello's avatar
wmmhello 已提交
2269 2270

  // decode
2271
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2272 2273 2274 2275 2276 2277 2278 2279 2280 2281 2282 2283
  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);
2284 2285
  //  cJSON* uid = cJSON_CreateNumber(id);
  //  cJSON_AddItemToObject(json, "uid", uid);
wmmhello's avatar
wmmhello 已提交
2286 2287
  cJSON* tableName = cJSON_CreateString(vAlterTbReq.tbName);
  cJSON_AddItemToObject(json, "tableName", tableName);
wmmhello's avatar
wmmhello 已提交
2288
  cJSON* tableType = cJSON_CreateString(vAlterTbReq.action == TSDB_ALTER_TABLE_UPDATE_TAG_VAL ? "child" : "normal");
wmmhello's avatar
wmmhello 已提交
2289
  cJSON_AddItemToObject(json, "tableType", tableType);
wmmhello's avatar
wmmhello 已提交
2290 2291
  cJSON* alterType = cJSON_CreateNumber(vAlterTbReq.action);
  cJSON_AddItemToObject(json, "alterType", alterType);
wmmhello's avatar
wmmhello 已提交
2292 2293 2294 2295 2296 2297 2298

  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 已提交
2299

2300
      if (vAlterTbReq.type == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
2301
        int32_t length = vAlterTbReq.bytes - VARSTR_HEADER_SIZE;
2302
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
2303
        cJSON_AddItemToObject(json, "colLength", cbytes);
2304 2305 2306
      } 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 已提交
2307 2308
        cJSON_AddItemToObject(json, "colLength", cbytes);
      }
wmmhello's avatar
wmmhello 已提交
2309 2310
      break;
    }
2311
    case TSDB_ALTER_TABLE_DROP_COLUMN: {
wmmhello's avatar
wmmhello 已提交
2312 2313 2314 2315
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      break;
    }
2316
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_BYTES: {
wmmhello's avatar
wmmhello 已提交
2317 2318
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
wmmhello's avatar
wmmhello 已提交
2319
      cJSON* colType = cJSON_CreateNumber(vAlterTbReq.colModType);
wmmhello's avatar
wmmhello 已提交
2320
      cJSON_AddItemToObject(json, "colType", colType);
2321
      if (vAlterTbReq.colModType == TSDB_DATA_TYPE_BINARY) {
wmmhello's avatar
wmmhello 已提交
2322
        int32_t length = vAlterTbReq.colModBytes - VARSTR_HEADER_SIZE;
2323
        cJSON*  cbytes = cJSON_CreateNumber(length);
wmmhello's avatar
wmmhello 已提交
2324
        cJSON_AddItemToObject(json, "colLength", cbytes);
2325 2326 2327
      } 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 已提交
2328 2329
        cJSON_AddItemToObject(json, "colLength", cbytes);
      }
wmmhello's avatar
wmmhello 已提交
2330 2331
      break;
    }
2332
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_NAME: {
wmmhello's avatar
wmmhello 已提交
2333 2334 2335 2336 2337 2338
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colNewName = cJSON_CreateString(vAlterTbReq.colNewName);
      cJSON_AddItemToObject(json, "colNewName", colNewName);
      break;
    }
2339
    case TSDB_ALTER_TABLE_UPDATE_TAG_VAL: {
wmmhello's avatar
wmmhello 已提交
2340 2341
      cJSON* tagName = cJSON_CreateString(vAlterTbReq.tagName);
      cJSON_AddItemToObject(json, "colName", tagName);
wmmhello's avatar
wmmhello 已提交
2342

wmmhello's avatar
wmmhello 已提交
2343
      bool isNull = vAlterTbReq.isNull;
2344 2345 2346
      if (vAlterTbReq.tagType == TSDB_DATA_TYPE_JSON) {
        STag* jsonTag = (STag*)vAlterTbReq.pTagVal;
        if (jsonTag->nTag == 0) isNull = true;
wmmhello's avatar
wmmhello 已提交
2347
      }
2348
      if (!isNull) {
wmmhello's avatar
wmmhello 已提交
2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363
        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 已提交
2364 2365
      cJSON* isNullCJson = cJSON_CreateBool(isNull);
      cJSON_AddItemToObject(json, "colValueNull", isNullCJson);
wmmhello's avatar
wmmhello 已提交
2366 2367
      break;
    }
wmmhello's avatar
wmmhello 已提交
2368 2369 2370 2371 2372
    default:
      break;
  }
  string = cJSON_PrintUnformatted(json);

2373
_exit:
wmmhello's avatar
wmmhello 已提交
2374 2375 2376 2377
  tDecoderClear(&decoder);
  return string;
}

2378 2379 2380 2381
static char* processDropSTable(SMqMetaRsp* metaRsp) {
  SDecoder     decoder = {0};
  SVDropStbReq req = {0};
  char*        string = NULL;
wmmhello's avatar
wmmhello 已提交
2382 2383

  // decode
2384
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2385 2386 2387 2388 2389 2390 2391 2392 2393 2394 2395 2396 2397 2398 2399 2400 2401 2402 2403
  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);

2404
_exit:
wmmhello's avatar
wmmhello 已提交
2405 2406 2407 2408
  tDecoderClear(&decoder);
  return string;
}

2409 2410 2411 2412
static char* processDropTable(SMqMetaRsp* metaRsp) {
  SDecoder         decoder = {0};
  SVDropTbBatchReq req = {0};
  char*            string = NULL;
wmmhello's avatar
wmmhello 已提交
2413 2414

  // decode
2415
  void*   data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2416 2417 2418 2419 2420 2421 2422 2423 2424 2425 2426 2427
  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);
2428 2429 2430 2431
  //  cJSON* uid = cJSON_CreateNumber(id);
  //  cJSON_AddItemToObject(json, "uid", uid);
  //  cJSON* tableType = cJSON_CreateString("normal");
  //  cJSON_AddItemToObject(json, "tableType", tableType);
wmmhello's avatar
wmmhello 已提交
2432 2433 2434 2435 2436

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

wmmhello's avatar
wmmhello 已提交
2437
    cJSON* tableName = cJSON_CreateString(pDropTbReq->name);
wmmhello's avatar
wmmhello 已提交
2438 2439 2440 2441 2442 2443
    cJSON_AddItemToArray(tableNameList, tableName);
  }
  cJSON_AddItemToObject(json, "tableNameList", tableNameList);

  string = cJSON_PrintUnformatted(json);

2444
_exit:
wmmhello's avatar
wmmhello 已提交
2445 2446 2447 2448
  tDecoderClear(&decoder);
  return string;
}

2449
char* tmq_get_json_meta(TAOS_RES* res) {
wmmhello's avatar
wmmhello 已提交
2450 2451 2452 2453 2454
  if (!TD_RES_TMQ_META(res)) {
    return NULL;
  }

  SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
2455
  if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_CREATE_STB) {
wmmhello's avatar
wmmhello 已提交
2456
    return processCreateStb(&pMetaRspObj->metaRsp);
2457
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_ALTER_STB) {
wmmhello's avatar
wmmhello 已提交
2458
    return processAlterStb(&pMetaRspObj->metaRsp);
2459
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_DROP_STB) {
wmmhello's avatar
wmmhello 已提交
2460
    return processDropSTable(&pMetaRspObj->metaRsp);
2461
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_CREATE_TABLE) {
wmmhello's avatar
wmmhello 已提交
2462
    return processCreateTable(&pMetaRspObj->metaRsp);
2463
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_ALTER_TABLE) {
wmmhello's avatar
wmmhello 已提交
2464
    return processAlterTable(&pMetaRspObj->metaRsp);
2465
  } else if (pMetaRspObj->metaRsp.resMsgType == TDMT_VND_DROP_TABLE) {
wmmhello's avatar
wmmhello 已提交
2466 2467 2468 2469 2470
    return processDropTable(&pMetaRspObj->metaRsp);
  }
  return NULL;
}

2471
void tmq_free_json_meta(char* jsonMeta) { taosMemoryFreeClear(jsonMeta); }
wmmhello's avatar
wmmhello 已提交
2472

2473
static int32_t taosCreateStb(TAOS* taos, void* meta, int32_t metaLen) {
wmmhello's avatar
wmmhello 已提交
2474 2475 2476
  SVCreateStbReq req = {0};
  SDecoder       coder;
  SMCreateStbReq pReq = {0};
2477 2478
  int32_t        code = TSDB_CODE_SUCCESS;
  SRequestObj*   pRequest = NULL;
wmmhello's avatar
wmmhello 已提交
2479

2480
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2481 2482 2483 2484
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2485
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2486 2487 2488 2489
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2490
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2491 2492 2493 2494 2495 2496 2497 2498
  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));
2499
  for (int32_t i = 0; i < req.schemaRow.nCols; i++) {
wmmhello's avatar
wmmhello 已提交
2500
    SSchema* pSchema = req.schemaRow.pSchema + i;
2501
    SField   field = {.type = pSchema->type, .bytes = pSchema->bytes};
wmmhello's avatar
wmmhello 已提交
2502 2503 2504 2505
    strcpy(field.name, pSchema->name);
    taosArrayPush(pReq.pColumns, &field);
  }
  pReq.pTags = taosArrayInit(req.schemaTag.nCols, sizeof(SField));
2506
  for (int32_t i = 0; i < req.schemaTag.nCols; i++) {
wmmhello's avatar
wmmhello 已提交
2507
    SSchema* pSchema = req.schemaTag.pSchema + i;
2508
    SField   field = {.type = pSchema->type, .bytes = pSchema->bytes};
wmmhello's avatar
wmmhello 已提交
2509 2510 2511
    strcpy(field.name, pSchema->name);
    taosArrayPush(pReq.pTags, &field);
  }
2512

2513 2514
  pReq.colVer = req.schemaRow.version;
  pReq.tagVer = req.schemaTag.version;
wmmhello's avatar
wmmhello 已提交
2515 2516 2517 2518
  pReq.numOfColumns = req.schemaRow.nCols;
  pReq.numOfTags = req.schemaTag.nCols;
  pReq.commentLen = -1;
  pReq.suid = req.suid;
wmmhello's avatar
wmmhello 已提交
2519
  pReq.source = TD_REQ_FROM_TAOX;
wmmhello's avatar
wmmhello 已提交
2520
  pReq.igExists = true;
wmmhello's avatar
wmmhello 已提交
2521

2522
  STscObj* pTscObj = pRequest->pTscObj;
2523
  SName    tableName;
wmmhello's avatar
wmmhello 已提交
2524
  tNameExtractFullName(toName(pTscObj->acctId, pRequest->pDb, req.name, &tableName), pReq.name);
wmmhello's avatar
wmmhello 已提交
2525 2526 2527 2528 2529 2530 2531 2532 2533 2534 2535 2536

  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);

2537
  SQuery pQuery = {0};
wmmhello's avatar
wmmhello 已提交
2538 2539 2540 2541 2542 2543 2544 2545 2546
  pQuery.execMode = QUERY_EXEC_MODE_RPC;
  pQuery.pCmdMsg = &pCmdMsg;
  pQuery.msgType = pQuery.pCmdMsg->msgType;
  pQuery.stableQuery = true;

  launchQueryImpl(pRequest, &pQuery, true, NULL);
  code = pRequest->code;
  taosMemoryFree(pCmdMsg.pMsg);

2547
end:
wmmhello's avatar
wmmhello 已提交
2548 2549 2550 2551 2552 2553
  destroyRequest(pRequest);
  tFreeSMCreateStbReq(&pReq);
  tDecoderClear(&coder);
  return code;
}

2554
static int32_t taosDropStb(TAOS* taos, void* meta, int32_t metaLen) {
wmmhello's avatar
wmmhello 已提交
2555 2556 2557
  SVDropStbReq req = {0};
  SDecoder     coder;
  SMDropStbReq pReq = {0};
2558
  int32_t      code = TSDB_CODE_SUCCESS;
wmmhello's avatar
wmmhello 已提交
2559 2560
  SRequestObj* pRequest = NULL;

2561
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2562 2563 2564 2565
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2566
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2567 2568 2569 2570
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2571
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2572 2573 2574 2575 2576 2577 2578 2579 2580
  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 已提交
2581
  pReq.source = TD_REQ_FROM_TAOX;
wmmhello's avatar
wmmhello 已提交
2582
  pReq.suid = req.suid;
2583 2584

  STscObj* pTscObj = pRequest->pTscObj;
2585
  SName    tableName;
wmmhello's avatar
wmmhello 已提交
2586
  tNameExtractFullName(toName(pTscObj->acctId, pRequest->pDb, req.name, &tableName), pReq.name);
wmmhello's avatar
wmmhello 已提交
2587 2588 2589 2590 2591 2592 2593 2594 2595 2596 2597 2598

  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);

2599
  SQuery pQuery = {0};
wmmhello's avatar
wmmhello 已提交
2600 2601 2602 2603 2604 2605 2606 2607 2608
  pQuery.execMode = QUERY_EXEC_MODE_RPC;
  pQuery.pCmdMsg = &pCmdMsg;
  pQuery.msgType = pQuery.pCmdMsg->msgType;
  pQuery.stableQuery = true;

  launchQueryImpl(pRequest, &pQuery, true, NULL);
  code = pRequest->code;
  taosMemoryFree(pCmdMsg.pMsg);

2609
end:
wmmhello's avatar
wmmhello 已提交
2610 2611 2612 2613 2614 2615 2616 2617 2618 2619 2620 2621
  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) {
2622
  SVgroupCreateTableBatch* pTbBatch = (SVgroupCreateTableBatch*)data;
wmmhello's avatar
wmmhello 已提交
2623 2624 2625
  taosArrayDestroy(pTbBatch->req.pArray);
}

L
Liu Jicong 已提交
2626 2627 2628 2629 2630 2631 2632
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 已提交
2633

L
Liu Jicong 已提交
2634
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2635 2636 2637 2638
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2639
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2640 2641 2642 2643
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2644
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2645 2646 2647 2648 2649 2650 2651
  int32_t len = metaLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);
  if (tDecodeSVCreateTbBatchReq(&coder, &req) < 0) {
    code = TSDB_CODE_INVALID_PARA;
    goto end;
  }

2652 2653
  STscObj* pTscObj = pRequest->pTscObj;

L
Liu Jicong 已提交
2654 2655
  SVCreateTbReq* pCreateReq = NULL;
  SCatalog*      pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2656 2657 2658 2659 2660 2661 2662 2663 2664 2665 2666 2667 2668
  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,
2669 2670 2671
                           .requestId = pRequest->requestId,
                           .requestObjRefId = pRequest->self,
                           .mgmtEps = getEpSet_s(&pTscObj->pAppInfo->mgmtEp)};
wmmhello's avatar
wmmhello 已提交
2672 2673 2674 2675 2676
  // loop to create table
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
    pCreateReq = req.pReqs + iReq;

    SVgroupInfo pInfo = {0};
2677
    SName       pName;
wmmhello's avatar
wmmhello 已提交
2678 2679 2680 2681 2682 2683 2684 2685 2686 2687 2688 2689 2690 2691 2692 2693 2694 2695 2696 2697 2698 2699 2700 2701 2702 2703 2704 2705 2706 2707 2708
    toName(pTscObj->acctId, pRequest->pDb, pCreateReq->name, &pName);
    code = catalogGetTableHashVgroup(pCatalog, &conn, &pName, &pInfo);
    if (code != TSDB_CODE_SUCCESS) {
      goto end;
    }

    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;
2709
  pQuery->pRoot = nodesMakeNode(QUERY_NODE_CREATE_TABLE_STMT);
wmmhello's avatar
wmmhello 已提交
2710 2711 2712 2713 2714 2715 2716

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

  launchQueryImpl(pRequest, pQuery, false, NULL);
2717 2718
  pQuery = NULL;  // no need to free in the end
  code = pRequest->code;
wmmhello's avatar
wmmhello 已提交
2719

2720
end:
wmmhello's avatar
wmmhello 已提交
2721 2722 2723 2724 2725 2726 2727 2728 2729 2730 2731 2732 2733 2734 2735 2736 2737 2738
  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 已提交
2739 2740 2741 2742 2743 2744 2745
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 已提交
2746

2747
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2748 2749 2750 2751
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

2752
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2753 2754 2755 2756
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2757
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2758 2759 2760 2761 2762 2763 2764
  int32_t len = metaLen - sizeof(SMsgHead);
  tDecoderInit(&coder, data, len);
  if (tDecodeSVDropTbBatchReq(&coder, &req) < 0) {
    code = TSDB_CODE_INVALID_PARA;
    goto end;
  }

2765 2766
  STscObj* pTscObj = pRequest->pTscObj;

L
Liu Jicong 已提交
2767 2768
  SVDropTbReq* pDropReq = NULL;
  SCatalog*    pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2769 2770 2771 2772 2773 2774 2775 2776 2777 2778 2779 2780 2781
  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,
2782 2783 2784
                           .requestId = pRequest->requestId,
                           .requestObjRefId = pRequest->self,
                           .mgmtEps = getEpSet_s(&pTscObj->pAppInfo->mgmtEp)};
wmmhello's avatar
wmmhello 已提交
2785 2786 2787
  // loop to create table
  for (int32_t iReq = 0; iReq < req.nReqs; iReq++) {
    pDropReq = req.pReqs + iReq;
wmmhello's avatar
wmmhello 已提交
2788
    pDropReq->igNotExists = true;
wmmhello's avatar
wmmhello 已提交
2789 2790

    SVgroupInfo pInfo = {0};
2791
    SName       pName;
wmmhello's avatar
wmmhello 已提交
2792 2793 2794 2795 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
    toName(pTscObj->acctId, pRequest->pDb, pDropReq->name, &pName);
    code = catalogGetTableHashVgroup(pCatalog, &conn, &pName, &pInfo);
    if (code != TSDB_CODE_SUCCESS) {
      goto end;
    }

    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;
2821
  pQuery->pRoot = nodesMakeNode(QUERY_NODE_DROP_TABLE_STMT);
wmmhello's avatar
wmmhello 已提交
2822 2823 2824 2825 2826 2827 2828

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

  launchQueryImpl(pRequest, pQuery, false, NULL);
2829 2830
  pQuery = NULL;  // no need to free in the end
  code = pRequest->code;
wmmhello's avatar
wmmhello 已提交
2831

2832
end:
wmmhello's avatar
wmmhello 已提交
2833 2834 2835 2836 2837 2838 2839
  taosHashCleanup(pVgroupHashmap);
  destroyRequest(pRequest);
  tDecoderClear(&coder);
  qDestroyQuery(pQuery);
  return code;
}

2840 2841 2842 2843 2844 2845 2846 2847
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 已提交
2848

2849
  code = buildRequest(*(int64_t*)taos, "", 0, NULL, false, &pRequest);
wmmhello's avatar
wmmhello 已提交
2850 2851 2852 2853 2854

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

2855
  if (!pRequest->pDb) {
wmmhello's avatar
wmmhello 已提交
2856 2857 2858 2859
    code = TSDB_CODE_PAR_DB_NOT_SPECIFIED;
    goto end;
  }
  // decode and process req
2860
  void*   data = POINTER_SHIFT(meta, sizeof(SMsgHead));
wmmhello's avatar
wmmhello 已提交
2861 2862 2863 2864 2865 2866 2867 2868
  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
2869
  if (req.action == TSDB_ALTER_TABLE_UPDATE_OPTIONS) {
wmmhello's avatar
wmmhello 已提交
2870 2871 2872
    goto end;
  }

2873
  STscObj*  pTscObj = pRequest->pTscObj;
2874
  SCatalog* pCatalog = NULL;
wmmhello's avatar
wmmhello 已提交
2875 2876 2877 2878 2879 2880
  code = catalogGetHandle(pTscObj->pAppInfo->clusterId, &pCatalog);
  if (code != TSDB_CODE_SUCCESS) {
    goto end;
  }

  SRequestConnInfo conn = {.pTrans = pTscObj->pAppInfo->pTransporter,
2881 2882 2883
                           .requestId = pRequest->requestId,
                           .requestObjRefId = pRequest->self,
                           .mgmtEps = getEpSet_s(&pTscObj->pAppInfo->mgmtEp)};
wmmhello's avatar
wmmhello 已提交
2884 2885

  SVgroupInfo pInfo = {0};
2886
  SName       pName = {0};
wmmhello's avatar
wmmhello 已提交
2887 2888 2889 2890 2891 2892 2893 2894 2895 2896 2897 2898 2899 2900 2901 2902 2903 2904 2905 2906 2907 2908 2909 2910 2911 2912 2913 2914 2915 2916 2917 2918 2919
  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;
2920
  pQuery->pRoot = nodesMakeNode(QUERY_NODE_ALTER_TABLE_STMT);
wmmhello's avatar
wmmhello 已提交
2921 2922 2923 2924 2925 2926 2927

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

  launchQueryImpl(pRequest, pQuery, false, NULL);
2928
  pQuery = NULL;  // no need to free in the end
wmmhello's avatar
wmmhello 已提交
2929
  pVgData = NULL;
2930 2931 2932
  pArray = NULL;
  code = pRequest->code;
  if (code == TSDB_CODE_VND_TABLE_NOT_EXIST) {
wmmhello's avatar
wmmhello 已提交
2933 2934
    code = 0;
  }
wmmhello's avatar
wmmhello 已提交
2935

2936
end:
wmmhello's avatar
wmmhello 已提交
2937
  taosArrayDestroy(pArray);
2938
  if (pVgData) taosMemoryFreeClear(pVgData->pData);
wmmhello's avatar
wmmhello 已提交
2939 2940 2941 2942 2943 2944 2945
  taosMemoryFreeClear(pVgData);
  destroyRequest(pRequest);
  tDecoderClear(&coder);
  qDestroyQuery(pQuery);
  return code;
}

2946
int32_t taos_write_raw_meta(TAOS* taos, tmq_raw_data* raw_meta) {
wmmhello's avatar
wmmhello 已提交
2947 2948 2949 2950
  if (!taos || !raw_meta) {
    return TSDB_CODE_INVALID_PARA;
  }

2951
  if (raw_meta->raw_meta_type == TDMT_VND_CREATE_STB) {
wmmhello's avatar
wmmhello 已提交
2952
    return taosCreateStb(taos, raw_meta->raw_meta, raw_meta->raw_meta_len);
2953
  } else if (raw_meta->raw_meta_type == TDMT_VND_ALTER_STB) {
wmmhello's avatar
wmmhello 已提交
2954
    return taosCreateStb(taos, raw_meta->raw_meta, raw_meta->raw_meta_len);
2955
  } else if (raw_meta->raw_meta_type == TDMT_VND_DROP_STB) {
wmmhello's avatar
wmmhello 已提交
2956
    return taosDropStb(taos, raw_meta->raw_meta, raw_meta->raw_meta_len);
2957
  } else if (raw_meta->raw_meta_type == TDMT_VND_CREATE_TABLE) {
wmmhello's avatar
wmmhello 已提交
2958
    return taosCreateTable(taos, raw_meta->raw_meta, raw_meta->raw_meta_len);
2959
  } else if (raw_meta->raw_meta_type == TDMT_VND_ALTER_TABLE) {
wmmhello's avatar
wmmhello 已提交
2960
    return taosAlterTable(taos, raw_meta->raw_meta, raw_meta->raw_meta_len);
2961
  } else if (raw_meta->raw_meta_type == TDMT_VND_DROP_TABLE) {
wmmhello's avatar
wmmhello 已提交
2962 2963 2964 2965 2966
    return taosDropTable(taos, raw_meta->raw_meta, raw_meta->raw_meta_len);
  }
  return TSDB_CODE_INVALID_PARA;
}

2967 2968
void tmq_free_raw_meta(tmq_raw_data* rawMeta) {
  //
wmmhello's avatar
wmmhello 已提交
2969 2970 2971
  taosMemoryFreeClear(rawMeta);
}

L
Liu Jicong 已提交
2972
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* msg, tmq_commit_cb* cb, void* param) {
2973
  //
L
Liu Jicong 已提交
2974
  tmqCommitInner2(tmq, msg, 0, 1, cb, param);
L
Liu Jicong 已提交
2975 2976
}

2977 2978 2979 2980
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* msg) {
  //
  return tmqCommitInner2(tmq, msg, 0, 0, NULL, NULL);
}