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

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

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
  char     clientId[256];
  char     groupId[TSDB_CGROUP_LEN];
  int8_t   autoCommit;
  int8_t   resetOffset;
  int8_t   withTbName;
L
Liu Jicong 已提交
58 59
  int8_t   spEnable;
  int32_t  spBatchSize;
60 61 62 63 64 65
  uint16_t port;
  int32_t  autoCommitInterval;
  char*    ip;
  char*    user;
  char*    pass;
  /*char*          db;*/
66
  tmq_commit_cb* commitCb;
L
Liu Jicong 已提交
67
  void*          commitCbUserParam;
L
Liu Jicong 已提交
68 69 70
};

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

  // status
  int8_t  status;
  int32_t epoch;
L
Liu Jicong 已提交
86 87
#if 0
  int8_t  epStatus;
L
Liu Jicong 已提交
88
  int32_t epSkipCnt;
L
Liu Jicong 已提交
89
#endif
L
Liu Jicong 已提交
90 91
  int64_t pollCnt;

L
Liu Jicong 已提交
92 93 94 95 96
  // timer
  tmr_h hbTimer;
  tmr_h reportTimer;
  tmr_h commitTimer;

L
Liu Jicong 已提交
97 98 99 100
  // connection
  STscObj* pTscObj;

  // container
L
Liu Jicong 已提交
101
  SArray*     clientTopics;  // SArray<SMqClientTopic>
L
Liu Jicong 已提交
102
  STaosQueue* mqueue;        // queue of rsp
L
Liu Jicong 已提交
103
  STaosQall*  qall;
L
Liu Jicong 已提交
104 105 106 107
  STaosQueue* delayedTask;  // delayed task queue for heartbeat and auto commit

  // ctl
  tsem_t rspSem;
L
Liu Jicong 已提交
108 109
};

X
Xiaoyu Wang 已提交
110 111 112 113 114 115 116 117
enum {
  TMQ_VG_STATUS__IDLE = 0,
  TMQ_VG_STATUS__WAIT,
};

enum {
  TMQ_CONSUMER_STATUS__INIT = 0,
  TMQ_CONSUMER_STATUS__READY,
118
  TMQ_CONSUMER_STATUS__NO_TOPIC,
L
Liu Jicong 已提交
119 120
};

L
Liu Jicong 已提交
121 122 123 124 125 126
enum {
  TMQ_DELAYED_TASK__HB = 1,
  TMQ_DELAYED_TASK__REPORT,
  TMQ_DELAYED_TASK__COMMIT,
};

L
Liu Jicong 已提交
127
typedef struct {
128 129 130
  // statistics
  int64_t pollCnt;
  // offset
L
Liu Jicong 已提交
131 132 133 134
  /*int64_t      committedOffset;*/
  /*int64_t      currentOffset;*/
  STqOffsetVal committedOffsetNew;
  STqOffsetVal currentOffsetNew;
L
Liu Jicong 已提交
135
  // connection info
136
  int32_t vgId;
X
Xiaoyu Wang 已提交
137
  int32_t vgStatus;
L
Liu Jicong 已提交
138
  int32_t vgSkipCnt;
139 140 141
  SEpSet  epSet;
} SMqClientVg;

L
Liu Jicong 已提交
142
typedef struct {
143
  // subscribe info
L
Liu Jicong 已提交
144
  char* topicName;
L
Liu Jicong 已提交
145
  char  db[TSDB_DB_FNAME_LEN];
L
Liu Jicong 已提交
146 147 148

  SArray* vgs;  // SArray<SMqClientVg>

L
Liu Jicong 已提交
149 150
  int8_t         isSchemaAdaptive;
  SSchemaWrapper schema;
151 152
} SMqClientTopic;

L
Liu Jicong 已提交
153 154 155 156 157
typedef struct {
  int8_t          tmqRspType;
  int32_t         epoch;
  SMqClientVg*    vgHandle;
  SMqClientTopic* topicHandle;
L
Liu Jicong 已提交
158
  union {
L
Liu Jicong 已提交
159 160
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
161
  };
L
Liu Jicong 已提交
162 163
} SMqPollRspWrapper;

L
Liu Jicong 已提交
164
typedef struct {
L
Liu Jicong 已提交
165 166 167
  tmq_t*  tmq;
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
168
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
169

L
Liu Jicong 已提交
170
typedef struct {
171
  tmq_t*  tmq;
L
Liu Jicong 已提交
172
  int32_t code;
L
Liu Jicong 已提交
173
  int32_t async;
X
Xiaoyu Wang 已提交
174
  tsem_t  rspSem;
175 176
} SMqAskEpCbParam;

L
Liu Jicong 已提交
177
typedef struct {
L
Liu Jicong 已提交
178 179
  tmq_t*          tmq;
  SMqClientVg*    pVg;
L
Liu Jicong 已提交
180
  SMqClientTopic* pTopic;
L
Liu Jicong 已提交
181
  int32_t         epoch;
L
Liu Jicong 已提交
182
  int32_t         vgId;
L
Liu Jicong 已提交
183
  tsem_t          rspSem;
X
Xiaoyu Wang 已提交
184
} SMqPollCbParam;
185

L
Liu Jicong 已提交
186
#if 0
L
Liu Jicong 已提交
187
typedef struct {
L
Liu Jicong 已提交
188
  tmq_t*         tmq;
L
Liu Jicong 已提交
189 190
  int8_t         async;
  int8_t         automatic;
L
Liu Jicong 已提交
191
  int8_t         freeOffsets;
L
Liu Jicong 已提交
192
  tmq_commit_cb* userCb;
L
Liu Jicong 已提交
193
  tsem_t         rspSem;
L
Liu Jicong 已提交
194
  int32_t        rspErr;
L
Liu Jicong 已提交
195
  SArray*        offsets;
L
Liu Jicong 已提交
196
  void*          userParam;
L
Liu Jicong 已提交
197
} SMqCommitCbParam;
L
Liu Jicong 已提交
198
#endif
L
Liu Jicong 已提交
199

200
typedef struct {
L
Liu Jicong 已提交
201 202 203 204
  tmq_t* tmq;
  int8_t automatic;
  int8_t async;
  /*int8_t         freeOffsets;*/
L
Liu Jicong 已提交
205 206
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
L
Liu Jicong 已提交
207
  int32_t        rspErr;
208
  tmq_commit_cb* userCb;
L
Liu Jicong 已提交
209 210 211 212
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
  void*  userParam;
  tsem_t rspSem;
213 214 215 216 217 218 219
} SMqCommitCbParamSet;

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

220
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
221
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
222
  conf->withTbName = false;
L
Liu Jicong 已提交
223
  conf->autoCommit = true;
L
Liu Jicong 已提交
224
  conf->autoCommitInterval = 5000;
X
Xiaoyu Wang 已提交
225
  conf->resetOffset = TMQ_CONF__RESET_OFFSET__EARLIEAST;
226 227 228
  return conf;
}

L
Liu Jicong 已提交
229
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
230 231 232 233 234 235
  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 已提交
236 237 238
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
239 240
  if (strcmp(key, "group.id") == 0) {
    strcpy(conf->groupId, value);
L
Liu Jicong 已提交
241
    return TMQ_CONF_OK;
242
  }
L
Liu Jicong 已提交
243

244 245
  if (strcmp(key, "client.id") == 0) {
    strcpy(conf->clientId, value);
L
Liu Jicong 已提交
246 247
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
248

L
Liu Jicong 已提交
249 250
  if (strcmp(key, "enable.auto.commit") == 0) {
    if (strcmp(value, "true") == 0) {
L
Liu Jicong 已提交
251
      conf->autoCommit = true;
L
Liu Jicong 已提交
252 253
      return TMQ_CONF_OK;
    } else if (strcmp(value, "false") == 0) {
L
Liu Jicong 已提交
254
      conf->autoCommit = false;
L
Liu Jicong 已提交
255 256 257 258
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
259
  }
L
Liu Jicong 已提交
260

L
Liu Jicong 已提交
261 262 263 264 265
  if (strcmp(key, "auto.commit.interval.ms") == 0) {
    conf->autoCommitInterval = atoi(value);
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
266 267 268 269 270 271 272 273 274 275 276 277 278 279
  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 已提交
280

281 282
  if (strcmp(key, "msg.with.table.name") == 0) {
    if (strcmp(value, "true") == 0) {
283
      conf->withTbName = true;
L
Liu Jicong 已提交
284
      return TMQ_CONF_OK;
285
    } else if (strcmp(value, "false") == 0) {
286
      conf->withTbName = false;
L
Liu Jicong 已提交
287
      return TMQ_CONF_OK;
288 289 290 291 292
    } else {
      return TMQ_CONF_INVALID;
    }
  }

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

L
Liu Jicong 已提交
305 306 307 308 309
  if (strcmp(key, "experimental.snapshot.batch.size") == 0) {
    conf->spBatchSize = atoi(value);
    return TMQ_CONF_OK;
  }

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

L
Liu Jicong 已提交
331
  return TMQ_CONF_UNKNOWN;
332 333 334
}

tmq_list_t* tmq_list_new() {
L
Liu Jicong 已提交
335 336
  //
  return (tmq_list_t*)taosArrayInit(0, sizeof(void*));
337 338
}

L
Liu Jicong 已提交
339 340 341
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
  char*   topic = strdup(src);
L
fix  
Liu Jicong 已提交
342
  if (taosArrayPush(container, &topic) == NULL) return -1;
343 344 345
  return 0;
}

L
Liu Jicong 已提交
346
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
347
  SArray* container = &list->container;
L
Liu Jicong 已提交
348
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
349 350
}

L
Liu Jicong 已提交
351 352 353 354 355 356 357 358 359 360
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 已提交
361 362 363 364
static int32_t tmqMakeTopicVgKey(char* dst, const char* topicName, int32_t vg) {
  return sprintf(dst, "%s:%d", topicName, vg);
}

L
Liu Jicong 已提交
365
#if 0
L
Liu Jicong 已提交
366 367
int32_t tmqCommitCb(void* param, const SDataBuf* pMsg, int32_t code) {
  SMqCommitCbParam* pParam = (SMqCommitCbParam*)param;
L
Liu Jicong 已提交
368
  pParam->rspErr = code;
L
Liu Jicong 已提交
369 370
  if (pParam->async) {
    if (pParam->automatic && pParam->tmq->commitCb) {
L
Liu Jicong 已提交
371
      pParam->tmq->commitCb(pParam->tmq, pParam->rspErr, pParam->tmq->commitCbUserParam);
L
Liu Jicong 已提交
372
    } else if (!pParam->automatic && pParam->userCb) {
L
Liu Jicong 已提交
373
      pParam->userCb(pParam->tmq, pParam->rspErr, pParam->userParam);
L
Liu Jicong 已提交
374 375
    }

L
Liu Jicong 已提交
376
    if (pParam->freeOffsets) {
L
Liu Jicong 已提交
377 378 379 380 381 382 383 384 385
      taosArrayDestroy(pParam->offsets);
    }

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

D
dapan1121 已提交
388
int32_t tmqCommitCb2(void* param, SDataBuf* pBuf, int32_t code) {
389 390 391
  SMqCommitCbParam2*   pParam = (SMqCommitCbParam2*)param;
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
  // push into array
L
Liu Jicong 已提交
392
#if 0
393 394 395 396 397
  if (code == 0) {
    taosArrayPush(pParamSet->failedOffsets, &pParam->pOffset);
  } else {
    taosArrayPush(pParamSet->successfulOffsets, &pParam->pOffset);
  }
L
Liu Jicong 已提交
398
#endif
L
Liu Jicong 已提交
399

L
Liu Jicong 已提交
400 401 402
  /*tscDebug("receive offset commit cb of %s on vg %d, offset is %ld", pParam->pOffset->subKey, pParam->->vgId,
   * pOffset->version);*/

403
  // count down waiting rsp
L
Liu Jicong 已提交
404
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
405 406 407 408 409 410 411
  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 已提交
412
        pParamSet->tmq->commitCb(pParamSet->tmq, pParamSet->rspErr, pParamSet->tmq->commitCbUserParam);
413 414
      } else if (!pParamSet->automatic && pParamSet->userCb) {
        // sem post
L
Liu Jicong 已提交
415
        pParamSet->userCb(pParamSet->tmq, pParamSet->rspErr, pParamSet->userParam);
416
      }
L
Liu Jicong 已提交
417 418
    } else {
      tsem_post(&pParamSet->rspSem);
419 420
    }

L
Liu Jicong 已提交
421
#if 0
422 423
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
L
Liu Jicong 已提交
424
#endif
425 426 427 428
  }
  return 0;
}

L
Liu Jicong 已提交
429 430 431 432 433 434 435
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;
436

L
Liu Jicong 已提交
437 438 439 440
  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 已提交
441

L
Liu Jicong 已提交
442 443 444 445 446 447 448 449 450 451
  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 已提交
452

L
Liu Jicong 已提交
453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523
  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,
  };

  tscDebug("consumer %ld commit offset of %s on vg %d, offset is %ld", tmq->consumerId, pOffset->subKey, pVg->vgId,
           pOffset->val.version);

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

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
  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 已提交
524 525
  int32_t code = -1;

L
Liu Jicong 已提交
526 527 528 529 530 531 532 533 534 535
  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 已提交
536
        }
L
Liu Jicong 已提交
537
        goto HANDLE_RSP;
L
Liu Jicong 已提交
538 539
      }
    }
L
Liu Jicong 已提交
540
  }
L
Liu Jicong 已提交
541

L
Liu Jicong 已提交
542 543 544 545
HANDLE_RSP:
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
546 547 548
    return 0;
  }

L
Liu Jicong 已提交
549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572
  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);
  }

573 574 575 576 577 578 579 580
  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 已提交
581
  /*pParamSet->freeOffsets = 1;*/
582 583 584 585 586 587
  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 已提交
588 589 590 591

    tscDebug("consumer %ld begin commit for topic %s, vgNum %d", tmq->consumerId, pTopic->topicName,
             (int32_t)taosArrayGetSize(pTopic->vgs));

592
    for (int32_t j = 0; j < taosArrayGetSize(pTopic->vgs); j++) {
L
Liu Jicong 已提交
593 594 595 596
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);

      tscDebug("consumer %ld begin commit for topic %s, vgId %d", tmq->consumerId, pTopic->topicName, pVg->vgId);

L
Liu Jicong 已提交
597 598 599 600
      if (pVg->currentOffsetNew.type > 0 && !tOffsetEqual(&pVg->currentOffsetNew, &pVg->committedOffsetNew)) {
        if (tmqSendCommitReq(tmq, pVg, pTopic, pParamSet) < 0) {
          continue;
        }
601 602 603 604
      }
    }
  }

L
Liu Jicong 已提交
605 606 607 608 609 610
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
    return 0;
  }

611 612 613 614 615 616 617 618 619 620
  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 已提交
621
      tmq->commitCb(tmq, code, tmq->commitCbUserParam);
622
    } else {
L
Liu Jicong 已提交
623
      userCb(tmq, code, userParam);
624 625 626
    }
  }

L
Liu Jicong 已提交
627
#if 0
628 629 630 631
  if (!async) {
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
  }
L
Liu Jicong 已提交
632
#endif
633 634 635 636

  return 0;
}

L
Liu Jicong 已提交
637 638
#if 0
int32_t tmqCommitInner(tmq_t* tmq, const TAOS_RES* msg, int8_t automatic, int8_t async,
L
Liu Jicong 已提交
639 640 641 642 643 644
                       tmq_commit_cb* userCb, void* userParam) {
  SMqCMCommitOffsetReq req;
  SArray*              pOffsets = NULL;
  void*                buf = NULL;
  SMqCommitCbParam*    pParam = NULL;
  SMsgSendInfo*        sendInfo = NULL;
L
Liu Jicong 已提交
645 646
  int8_t               freeOffsets;
  int32_t              code = -1;
L
Liu Jicong 已提交
647

L
Liu Jicong 已提交
648
  if (msg == NULL) {
L
Liu Jicong 已提交
649
    freeOffsets = 1;
L
Liu Jicong 已提交
650 651 652 653 654 655 656 657 658 659 660 661 662 663
    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 已提交
664
    freeOffsets = 0;
L
Liu Jicong 已提交
665
    pOffsets = (SArray*)&msg->container;
L
Liu Jicong 已提交
666 667 668 669 670
  }

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

L
Liu Jicong 已提交
671 672 673 674
  SEncoder encoder;

  tEncoderInit(&encoder, NULL, 0);
  code = tEncodeSMqCMCommitOffsetReq(&encoder, &req);
L
Liu Jicong 已提交
675 676 677
  if (code < 0) {
    goto END;
  }
L
Liu Jicong 已提交
678
  int32_t tlen = encoder.pos;
L
Liu Jicong 已提交
679 680
  buf = taosMemoryMalloc(tlen);
  if (buf == NULL) {
L
Liu Jicong 已提交
681
    tEncoderClear(&encoder);
L
Liu Jicong 已提交
682 683
    goto END;
  }
L
Liu Jicong 已提交
684 685
  tEncoderClear(&encoder);

L
Liu Jicong 已提交
686 687 688 689 690 691 692 693 694 695 696 697
  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 已提交
698
  pParam->freeOffsets = freeOffsets;
L
Liu Jicong 已提交
699 700 701 702
  pParam->userCb = userCb;
  pParam->userParam = userParam;
  if (!async) tsem_init(&pParam->rspSem, 0, 0);

703
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726
  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 已提交
727 728
  } else {
    code = 0;
L
Liu Jicong 已提交
729 730 731 732 733 734 735
  }

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

END:
  if (buf) taosMemoryFree(buf);
L
Liu Jicong 已提交
736 737
  /*if (pParam) taosMemoryFree(pParam);*/
  /*if (sendInfo) taosMemoryFree(sendInfo);*/
L
Liu Jicong 已提交
738 739 740

  if (code != 0 && async) {
    if (automatic) {
L
Liu Jicong 已提交
741
      tmq->commitCb(tmq, code, (tmq_topic_vgroup_list_t*)pOffsets, tmq->commitCbUserParam);
L
Liu Jicong 已提交
742
    } else {
L
Liu Jicong 已提交
743
      userCb(tmq, code, (tmq_topic_vgroup_list_t*)pOffsets, userParam);
L
Liu Jicong 已提交
744 745 746
    }
  }

L
Liu Jicong 已提交
747
  if (!async && freeOffsets) {
L
Liu Jicong 已提交
748 749 750 751
    taosArrayDestroy(pOffsets);
  }
  return code;
}
L
Liu Jicong 已提交
752
#endif
L
Liu Jicong 已提交
753

L
Liu Jicong 已提交
754 755
void tmqAssignDelayedHbTask(void* param, void* tmrId) {
  tmq_t*  tmq = (tmq_t*)param;
756
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
L
Liu Jicong 已提交
757 758
  *pTaskType = TMQ_DELAYED_TASK__HB;
  taosWriteQitem(tmq->delayedTask, pTaskType);
759
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
760 761 762 763
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
  tmq_t*  tmq = (tmq_t*)param;
764
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
L
Liu Jicong 已提交
765 766
  *pTaskType = TMQ_DELAYED_TASK__COMMIT;
  taosWriteQitem(tmq->delayedTask, pTaskType);
767
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
768 769 770 771
}

void tmqAssignDelayedReportTask(void* param, void* tmrId) {
  tmq_t*  tmq = (tmq_t*)param;
772
  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM);
L
Liu Jicong 已提交
773 774
  *pTaskType = TMQ_DELAYED_TASK__REPORT;
  taosWriteQitem(tmq->delayedTask, pTaskType);
775
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
776 777 778 779 780 781 782 783 784 785 786
}

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;

    if (*pTaskType == TMQ_DELAYED_TASK__HB) {
L
Liu Jicong 已提交
787
      tmqAskEp(tmq, true);
L
Liu Jicong 已提交
788 789
      taosTmrReset(tmqAssignDelayedHbTask, 1000, tmq, tmqMgmt.timer, &tmq->hbTimer);
    } else if (*pTaskType == TMQ_DELAYED_TASK__COMMIT) {
L
Liu Jicong 已提交
790
      tmqCommitInner2(tmq, NULL, 1, 1, tmq->commitCb, tmq->commitCbUserParam);
L
Liu Jicong 已提交
791 792 793 794 795
      taosTmrReset(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, tmq, tmqMgmt.timer, &tmq->commitTimer);
    } else if (*pTaskType == TMQ_DELAYED_TASK__REPORT) {
    } else {
      ASSERT(0);
    }
L
Liu Jicong 已提交
796
    taosFreeQitem(pTaskType);
L
Liu Jicong 已提交
797 798 799 800 801
  }
  taosFreeQall(qall);
  return 0;
}

L
Liu Jicong 已提交
802
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
803
  SMqRspWrapper* msg = NULL;
L
Liu Jicong 已提交
804 805 806 807 808 809 810 811
  while (1) {
    taosGetQitem(tmq->qall, (void**)&msg);
    if (msg)
      taosFreeQitem(msg);
    else
      break;
  }

L
fix  
Liu Jicong 已提交
812
  msg = NULL;
L
Liu Jicong 已提交
813 814 815 816 817 818 819 820 821 822
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
    taosGetQitem(tmq->qall, (void**)&msg);
    if (msg)
      taosFreeQitem(msg);
    else
      break;
  }
}

D
dapan1121 已提交
823
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
824 825
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
L
Liu Jicong 已提交
826
  /*tmq_t* tmq = pParam->tmq;*/
L
Liu Jicong 已提交
827 828 829
  tsem_post(&pParam->rspSem);
  return 0;
}
830

L
Liu Jicong 已提交
831
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
832 833 834 835
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
836
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
837
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
838
  }
L
Liu Jicong 已提交
839
  return 0;
X
Xiaoyu Wang 已提交
840 841
}

L
Liu Jicong 已提交
842 843 844
int32_t tmq_unsubscribe(tmq_t* tmq) {
  tmq_list_t* lst = tmq_list_new();
  int32_t     rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
845 846
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
847 848
}

L
Liu Jicong 已提交
849
#if 0
850
tmq_t* tmq_consumer_new(void* conn, tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
wafwerar's avatar
wafwerar 已提交
851
  tmq_t* pTmq = taosMemoryCalloc(sizeof(tmq_t), 1);
852 853 854 855 856 857
  if (pTmq == NULL) {
    return NULL;
  }
  pTmq->pTscObj = (STscObj*)conn;
  pTmq->status = 0;
  pTmq->pollCnt = 0;
L
Liu Jicong 已提交
858
  pTmq->epoch = 0;
L
fix txn  
Liu Jicong 已提交
859
  pTmq->epStatus = 0;
L
temp  
Liu Jicong 已提交
860
  pTmq->epSkipCnt = 0;
L
Liu Jicong 已提交
861
  // set conf
862 863
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
L
Liu Jicong 已提交
864
  pTmq->autoCommit = conf->autoCommit;
865
  pTmq->commit_cb = conf->commit_cb;
L
Liu Jicong 已提交
866
  pTmq->resetOffsetCfg = conf->resetOffset;
L
Liu Jicong 已提交
867

L
Liu Jicong 已提交
868 869 870 871 872 873 874 875 876 877
  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 已提交
878 879
  tsem_init(&pTmq->rspSem, 0, 0);

L
Liu Jicong 已提交
880 881
  return pTmq;
}
L
Liu Jicong 已提交
882
#endif
L
Liu Jicong 已提交
883

L
Liu Jicong 已提交
884
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
L
Liu Jicong 已提交
885 886 887 888 889 890 891 892 893 894 895
  // 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 已提交
896 897 898 899
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
    return NULL;
  }
L
Liu Jicong 已提交
900

L
Liu Jicong 已提交
901 902 903 904 905
  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 已提交
906
  ASSERT(conf->groupId[0]);
L
Liu Jicong 已提交
907

L
Liu Jicong 已提交
908 909 910 911 912 913 914 915
  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) {
    goto FAIL;
  }
L
Liu Jicong 已提交
916

L
Liu Jicong 已提交
917 918
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
919 920
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
921 922
  /*pTmq->epStatus = 0;*/
  /*pTmq->epSkipCnt = 0;*/
L
Liu Jicong 已提交
923

L
Liu Jicong 已提交
924 925 926
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
927
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
928
  pTmq->useSnapshot = conf->spEnable;
L
Liu Jicong 已提交
929
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
930
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
931 932
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
933 934
  pTmq->resetOffsetCfg = conf->resetOffset;

L
Liu Jicong 已提交
935
  // assign consumerId
L
Liu Jicong 已提交
936
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
937

L
Liu Jicong 已提交
938 939 940 941
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
    goto FAIL;
  }
L
Liu Jicong 已提交
942

L
Liu Jicong 已提交
943 944 945 946 947 948
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
    tsem_destroy(&pTmq->rspSem);
    goto FAIL;
  }
L
Liu Jicong 已提交
949

950
  return pTmq;
L
Liu Jicong 已提交
951 952 953 954 955 956 957 958

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;
959 960
}

L
Liu Jicong 已提交
961 962
#if 0
int32_t tmq_commit(tmq_t* tmq, const tmq_topic_vgroup_list_t* offsets, int32_t async) {
L
Liu Jicong 已提交
963
  return tmqCommitInner2(tmq, offsets, 0, async, tmq->commitCb, tmq->commitCbUserParam);
L
Liu Jicong 已提交
964
}
L
Liu Jicong 已提交
965
#endif
L
Liu Jicong 已提交
966

L
Liu Jicong 已提交
967
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
L
Liu Jicong 已提交
968 969 970 971 972
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
  SCMSubscribeReq req = {0};
  int32_t         code = -1;
973 974

  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
975
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
976
  req.topicNames = taosArrayInit(sz, sizeof(void*));
L
Liu Jicong 已提交
977
  if (req.topicNames == NULL) goto FAIL;
978

L
Liu Jicong 已提交
979 980
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
981 982

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

L
Liu Jicong 已提交
985 986 987
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
988
    }
L
Liu Jicong 已提交
989
    tNameExtractFullName(&name, topicFName);
990

L
Liu Jicong 已提交
991 992 993
    tscDebug("subscribe topic: %s", topicFName);

    taosArrayPush(req.topicNames, &topicFName);
994 995
  }

L
Liu Jicong 已提交
996 997 998 999
  int32_t tlen = tSerializeSCMSubscribeReq(NULL, &req);
  buf = taosMemoryMalloc(tlen);
  if (buf == NULL) goto FAIL;

1000 1001 1002
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

1003
  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
1004
  if (sendInfo == NULL) goto FAIL;
1005

X
Xiaoyu Wang 已提交
1006
  SMqSubscribeCbParam param = {
L
Liu Jicong 已提交
1007
      .rspErr = 0,
X
Xiaoyu Wang 已提交
1008 1009
      .tmq = tmq,
  };
L
Liu Jicong 已提交
1010

L
Liu Jicong 已提交
1011 1012 1013
  if (tsem_init(&param.rspSem, 0, 0) != 0) goto FAIL;

  sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1014 1015 1016 1017
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };
1018

L
Liu Jicong 已提交
1019 1020
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1021 1022
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1023 1024
  sendInfo->msgType = TDMT_MND_SUBSCRIBE;

1025 1026 1027 1028 1029
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1030 1031 1032
  // avoid double free if msg is sent
  buf = NULL;

L
Liu Jicong 已提交
1033 1034
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1035

L
Liu Jicong 已提交
1036 1037 1038
  code = param.rspErr;
  if (code != 0) goto FAIL;

L
Liu Jicong 已提交
1039
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == tmqAskEp(tmq, false)) {
L
fix  
Liu Jicong 已提交
1040
    tscDebug("consumer not ready, retry");
L
Liu Jicong 已提交
1041 1042
    taosMsleep(500);
  }
1043

L
Liu Jicong 已提交
1044
  // init hb timer
1045 1046 1047
  if (tmq->hbTimer == NULL) {
    tmq->hbTimer = taosTmrStart(tmqAssignDelayedHbTask, 1000, tmq, tmqMgmt.timer);
  }
L
Liu Jicong 已提交
1048 1049

  // init auto commit timer
1050
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
L
Liu Jicong 已提交
1051 1052 1053
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, tmq, tmqMgmt.timer);
  }

L
Liu Jicong 已提交
1054 1055 1056
  code = 0;
FAIL:
  if (req.topicNames != NULL) taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1057
  if (code != 0 && buf) {
L
Liu Jicong 已提交
1058 1059 1060
    taosMemoryFree(buf);
  }
  return code;
1061 1062
}

L
Liu Jicong 已提交
1063
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
L
Liu Jicong 已提交
1064
  //
1065
  conf->commitCb = cb;
L
Liu Jicong 已提交
1066
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1067
}
1068

L
Liu Jicong 已提交
1069
#if 0
L
Liu Jicong 已提交
1070 1071
int32_t tmqGetSkipLogNum(tmq_message_t* tmq_message) {
  if (tmq_message == NULL) return 0;
L
Liu Jicong 已提交
1072
  SMqPollRsp* pRsp = &tmq_message->msg;
L
Liu Jicong 已提交
1073 1074
  return pRsp->skipLogNum;
}
L
Liu Jicong 已提交
1075
#endif
L
Liu Jicong 已提交
1076

D
dapan1121 已提交
1077
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1078 1079
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1080
  SMqClientTopic* pTopic = pParam->pTopic;
X
Xiaoyu Wang 已提交
1081
  tmq_t*          tmq = pParam->tmq;
L
Liu Jicong 已提交
1082 1083 1084
  int32_t         vgId = pParam->vgId;
  int32_t         epoch = pParam->epoch;
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1085
  if (code != 0) {
L
Liu Jicong 已提交
1086 1087
    tscWarn("msg discard from vg %d, epoch %d, code:%x", vgId, epoch, code);
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100
    if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM);
      if (pRspWrapper == NULL) {
        taosMemoryFree(pMsg->pData);
        tscWarn("msg discard from vg %d, epoch %d since out of memory", vgId, epoch);
        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 已提交
1101
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1102 1103
  }

X
Xiaoyu Wang 已提交
1104 1105 1106
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
  int32_t tmqEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < tmqEpoch) {
L
Liu Jicong 已提交
1107
    // do not write into queue since updating epoch reset
L
Liu Jicong 已提交
1108
    tscWarn("msg discard from vg %d since from earlier epoch, rsp epoch %d, current epoch %d", vgId, msgEpoch,
L
Liu Jicong 已提交
1109
            tmqEpoch);
1110
    tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1111
    taosMemoryFree(pMsg->pData);
X
Xiaoyu Wang 已提交
1112 1113 1114 1115
    return 0;
  }

  if (msgEpoch != tmqEpoch) {
L
Liu Jicong 已提交
1116
    tscWarn("mismatch rsp from vg %d, epoch %d, current epoch %d", vgId, msgEpoch, tmqEpoch);
X
Xiaoyu Wang 已提交
1117 1118
  }

L
Liu Jicong 已提交
1119 1120 1121
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

1122
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM);
L
Liu Jicong 已提交
1123
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1124 1125
    taosMemoryFree(pMsg->pData);
    tscWarn("msg discard from vg %d, epoch %d since out of memory", vgId, epoch);
L
fix txn  
Liu Jicong 已提交
1126
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1127
  }
L
Liu Jicong 已提交
1128

L
Liu Jicong 已提交
1129
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1130 1131
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
L
Liu Jicong 已提交
1132

L
Liu Jicong 已提交
1133
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1134 1135 1136
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqDataRsp(&decoder, &pRspWrapper->dataRsp);
L
Liu Jicong 已提交
1137
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1138
    /*tDecodeSMqDataBlkRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pRspWrapper->dataRsp);*/
L
Liu Jicong 已提交
1139 1140
  } else {
    ASSERT(rspType == TMQ_MSG_TYPE__POLL_META_RSP);
L
Liu Jicong 已提交
1141
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1142 1143
    tDecodeSMqMetaRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pRspWrapper->metaRsp);
  }
L
Liu Jicong 已提交
1144

L
Liu Jicong 已提交
1145
  taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1146

L
Liu Jicong 已提交
1147
  tscDebug("consumer %ld recv poll: vg %d, req offset %ld, rsp offset %ld, type %d", tmq->consumerId, pVg->vgId,
L
Liu Jicong 已提交
1148
           pRspWrapper->dataRsp.reqOffset.version, pRspWrapper->dataRsp.rspOffset.version, rspType);
L
fix  
Liu Jicong 已提交
1149

L
Liu Jicong 已提交
1150
  taosWriteQitem(tmq->mqueue, pRspWrapper);
1151
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1152

L
Liu Jicong 已提交
1153
  return 0;
L
fix txn  
Liu Jicong 已提交
1154
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1155
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1156 1157
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
1158
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1159
  return -1;
1160 1161
}

1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189
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];
  tscDebug("consumer %ld update ep epoch %d to epoch %d, topic num: %d", tmq->consumerId, tmq->epoch, epoch,
           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);
      tscDebug("consumer %ld new vg num: %d", tmq->consumerId, vgNumCur);
      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 已提交
1190 1191 1192 1193 1194
        char buf[50];
        tFormatOffset(buf, 50, &pVgCur->currentOffsetNew);
        tscDebug("consumer %ld epoch %d vg %d vgKey is %s, offset is %s", tmq->consumerId, epoch, pVgCur->vgId, vgKey,
                 buf);
        taosHashPut(pHash, vgKey, strlen(vgKey), &pVgCur->currentOffsetNew, sizeof(STqOffsetVal));
1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212
      }
    }
  }

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

    tscDebug("consumer %ld update topic: %s", tmq->consumerId, topic.topicName);

    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 已提交
1213 1214
      STqOffsetVal* pOffset = taosHashGet(pHash, vgKey, strlen(vgKey));
      STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
1215
      if (pOffset != NULL) {
L
Liu Jicong 已提交
1216
        offsetNew = *pOffset;
1217 1218
      }

L
Liu Jicong 已提交
1219 1220
      /*tscDebug("consumer %ld(epoch %d) offset of vg %d updated to %ld, vgKey is %s", tmq->consumerId, epoch,*/
      /*pVgEp->vgId, offset, vgKey);*/
1221 1222
      SMqClientVg clientVg = {
          .pollCnt = 0,
L
Liu Jicong 已提交
1223
          .currentOffsetNew = offsetNew,
1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246
          .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 已提交
1247
#if 0
L
Liu Jicong 已提交
1248
bool tmqUpdateEp(tmq_t* tmq, int32_t epoch, SMqAskEpRsp* pRsp) {
L
Liu Jicong 已提交
1249
  /*printf("call update ep %d\n", epoch);*/
X
Xiaoyu Wang 已提交
1250
  bool    set = false;
L
Liu Jicong 已提交
1251 1252
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
L
Liu Jicong 已提交
1253 1254
  tscDebug("consumer %ld update ep epoch %d to epoch %d, topic num: %d", tmq->consumerId, tmq->epoch, epoch,
           topicNumGet);
L
Liu Jicong 已提交
1255 1256 1257 1258 1259 1260 1261 1262 1263 1264 1265 1266
  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 已提交
1267 1268
    SMqClientTopic topic = {0};
    SMqSubTopicEp* pTopicEp = taosArrayGet(pRsp->topics, i);
L
Liu Jicong 已提交
1269
    topic.schema = pTopicEp->schema;
L
Liu Jicong 已提交
1270
    taosHashClear(pHash);
X
Xiaoyu Wang 已提交
1271
    topic.topicName = strdup(pTopicEp->topic);
L
Liu Jicong 已提交
1272
    tstrncpy(topic.db, pTopicEp->db, TSDB_DB_FNAME_LEN);
L
Liu Jicong 已提交
1273

L
Liu Jicong 已提交
1274
    tscDebug("consumer %ld update topic: %s", tmq->consumerId, topic.topicName);
L
Liu Jicong 已提交
1275 1276 1277 1278 1279 1280
    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);
L
Liu Jicong 已提交
1281
        tscDebug("consumer %ld new vg num: %d", tmq->consumerId, vgNumCur);
L
Liu Jicong 已提交
1282 1283 1284 1285
        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);
L
Liu Jicong 已提交
1286
          tscDebug("consumer %ld epoch %d vg %d build %s", tmq->consumerId, epoch, pVgCur->vgId, vgKey);
L
Liu Jicong 已提交
1287 1288 1289 1290 1291 1292 1293 1294 1295
          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 已提交
1296
      SMqSubVgEp* pVgEp = taosArrayGet(pTopicEp->vgs, j);
L
Liu Jicong 已提交
1297 1298 1299
      sprintf(vgKey, "%s:%d", topic.topicName, pVgEp->vgId);
      int64_t* pOffset = taosHashGet(pHash, vgKey, strlen(vgKey));
      int64_t  offset = pVgEp->offset;
1300
      tscDebug("consumer %ld(epoch %d) original offset of vg %d is %ld", tmq->consumerId, epoch, pVgEp->vgId, offset);
L
Liu Jicong 已提交
1301 1302
      if (pOffset != NULL) {
        offset = *pOffset;
1303 1304
        tscDebug("consumer %ld(epoch %d) receive offset of vg %d, full key is %s", tmq->consumerId, epoch, pVgEp->vgId,
                 vgKey);
L
Liu Jicong 已提交
1305
      }
1306
      tscDebug("consumer %ld(epoch %d) offset of vg %d updated to %ld", tmq->consumerId, epoch, pVgEp->vgId, offset);
X
Xiaoyu Wang 已提交
1307 1308
      SMqClientVg clientVg = {
          .pollCnt = 0,
L
Liu Jicong 已提交
1309
          .currentOffset = offset,
X
Xiaoyu Wang 已提交
1310 1311 1312
          .vgId = pVgEp->vgId,
          .epSet = pVgEp->epSet,
          .vgStatus = TMQ_VG_STATUS__IDLE,
L
Liu Jicong 已提交
1313
          .vgSkipCnt = 0,
X
Xiaoyu Wang 已提交
1314 1315 1316 1317
      };
      taosArrayPush(topic.vgs, &clientVg);
      set = true;
    }
L
Liu Jicong 已提交
1318
    taosArrayPush(newTopics, &topic);
X
Xiaoyu Wang 已提交
1319
  }
L
Liu Jicong 已提交
1320
  if (tmq->clientTopics) taosArrayDestroy(tmq->clientTopics);
L
Liu Jicong 已提交
1321
  taosHashCleanup(pHash);
L
Liu Jicong 已提交
1322
  tmq->clientTopics = newTopics;
1323

1324 1325 1326 1327
  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);
1328

X
Xiaoyu Wang 已提交
1329 1330 1331
  atomic_store_32(&tmq->epoch, epoch);
  return set;
}
L
Liu Jicong 已提交
1332
#endif
X
Xiaoyu Wang 已提交
1333

D
dapan1121 已提交
1334
int32_t tmqAskEpCb(void* param, SDataBuf* pMsg, int32_t code) {
1335
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
L
Liu Jicong 已提交
1336
  tmq_t*           tmq = pParam->tmq;
L
Liu Jicong 已提交
1337
  int8_t           async = pParam->async;
L
Liu Jicong 已提交
1338
  pParam->code = code;
1339
  if (code != 0) {
L
Liu Jicong 已提交
1340
    tscError("consumer %ld get topic endpoint error, not ready, wait:%d", tmq->consumerId, pParam->async);
L
Liu Jicong 已提交
1341
    goto END;
1342
  }
L
Liu Jicong 已提交
1343

L
Liu Jicong 已提交
1344
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1345
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1346
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1347 1348
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
L
temp  
Liu Jicong 已提交
1349
  tscDebug("consumer %ld recv ep, msg epoch %d, current epoch %d", tmq->consumerId, head->epoch, epoch);
L
Liu Jicong 已提交
1350 1351
  if (head->epoch <= epoch) {
    goto END;
1352
  }
L
Liu Jicong 已提交
1353

L
Liu Jicong 已提交
1354
  if (!async) {
L
Liu Jicong 已提交
1355 1356
    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
X
Xiaoyu Wang 已提交
1357 1358
    /*printf("rsp epoch %ld sz %ld\n", rsp.epoch, rsp.topics->size);*/
    /*printf("tmq epoch %ld sz %ld\n", tmq->epoch, tmq->clientTopics->size);*/
L
Liu Jicong 已提交
1359
    tmqUpdateEp2(tmq, head->epoch, &rsp);
L
Liu Jicong 已提交
1360
    tDeleteSMqAskEpRsp(&rsp);
X
Xiaoyu Wang 已提交
1361
  } else {
1362
    SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM);
L
Liu Jicong 已提交
1363
    if (pWrapper == NULL) {
X
Xiaoyu Wang 已提交
1364
      terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1365 1366
      code = -1;
      goto END;
X
Xiaoyu Wang 已提交
1367
    }
L
Liu Jicong 已提交
1368 1369 1370
    pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
    pWrapper->epoch = head->epoch;
    memcpy(&pWrapper->msg, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1371
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pWrapper->msg);
L
Liu Jicong 已提交
1372

L
Liu Jicong 已提交
1373
    taosWriteQitem(tmq->mqueue, pWrapper);
1374
    tsem_post(&tmq->rspSem);
1375
  }
L
Liu Jicong 已提交
1376 1377

END:
L
Liu Jicong 已提交
1378
  /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1379
  if (!async) {
L
Liu Jicong 已提交
1380
    tsem_post(&pParam->rspSem);
L
Liu Jicong 已提交
1381 1382
  } else {
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
1383 1384
  }
  return code;
1385 1386
}

L
Liu Jicong 已提交
1387
int32_t tmqAskEp(tmq_t* tmq, bool async) {
L
Liu Jicong 已提交
1388
  int32_t code = 0;
L
Liu Jicong 已提交
1389
#if 0
L
Liu Jicong 已提交
1390
  int8_t  epStatus = atomic_val_compare_exchange_8(&tmq->epStatus, 0, 1);
L
fix txn  
Liu Jicong 已提交
1391
  if (epStatus == 1) {
L
temp  
Liu Jicong 已提交
1392
    int32_t epSkipCnt = atomic_add_fetch_32(&tmq->epSkipCnt, 1);
L
Liu Jicong 已提交
1393
    tscTrace("consumer %ld skip ask ep cnt %d", tmq->consumerId, epSkipCnt);
L
Liu Jicong 已提交
1394
    if (epSkipCnt < 5000) return 0;
L
fix txn  
Liu Jicong 已提交
1395
  }
L
temp  
Liu Jicong 已提交
1396
  atomic_store_32(&tmq->epSkipCnt, 0);
L
Liu Jicong 已提交
1397
#endif
L
Liu Jicong 已提交
1398
  int32_t      tlen = sizeof(SMqAskEpReq);
L
Liu Jicong 已提交
1399
  SMqAskEpReq* req = taosMemoryCalloc(1, tlen);
L
Liu Jicong 已提交
1400
  if (req == NULL) {
L
Liu Jicong 已提交
1401
    tscError("failed to malloc get subscribe ep buf");
L
Liu Jicong 已提交
1402
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1403
    return -1;
L
Liu Jicong 已提交
1404
  }
L
Liu Jicong 已提交
1405 1406 1407
  req->consumerId = htobe64(tmq->consumerId);
  req->epoch = htonl(tmq->epoch);
  strcpy(req->cgroup, tmq->groupId);
1408

L
Liu Jicong 已提交
1409
  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
L
Liu Jicong 已提交
1410 1411
  if (pParam == NULL) {
    tscError("failed to malloc subscribe param");
wafwerar's avatar
wafwerar 已提交
1412
    taosMemoryFree(req);
L
Liu Jicong 已提交
1413
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1414
    return -1;
L
Liu Jicong 已提交
1415 1416
  }
  pParam->tmq = tmq;
L
Liu Jicong 已提交
1417
  pParam->async = async;
X
Xiaoyu Wang 已提交
1418
  tsem_init(&pParam->rspSem, 0, 0);
1419

1420
  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
1421 1422
  if (sendInfo == NULL) {
    tsem_destroy(&pParam->rspSem);
wafwerar's avatar
wafwerar 已提交
1423 1424
    taosMemoryFree(pParam);
    taosMemoryFree(req);
L
Liu Jicong 已提交
1425
    /*atomic_store_8(&tmq->epStatus, 0);*/
L
Liu Jicong 已提交
1426 1427 1428 1429 1430 1431 1432 1433 1434 1435
    return -1;
  }

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

  sendInfo->requestId = generateRequestId();
L
Liu Jicong 已提交
1436 1437 1438
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqAskEpCb;
L
Liu Jicong 已提交
1439
  sendInfo->msgType = TDMT_MND_MQ_ASK_EP;
1440

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

L
add log  
Liu Jicong 已提交
1443 1444
  tscDebug("consumer %ld ask ep", tmq->consumerId);

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

L
Liu Jicong 已提交
1448
  if (!async) {
L
Liu Jicong 已提交
1449 1450 1451 1452 1453
    tsem_wait(&pParam->rspSem);
    code = pParam->code;
    taosMemoryFree(pParam);
  }
  return code;
1454 1455
}

L
Liu Jicong 已提交
1456 1457
#if 0
int32_t tmq_seek(tmq_t* tmq, const tmq_topic_vgroup_t* offset) {
L
Liu Jicong 已提交
1458 1459 1460 1461 1462 1463 1464 1465 1466 1467 1468 1469 1470
  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 已提交
1471
          tmqClearUnhandleMsg(tmq);
L
Liu Jicong 已提交
1472 1473 1474 1475 1476 1477 1478
          return TMQ_RESP_ERR__SUCCESS;
        }
      }
    }
  }
  return TMQ_RESP_ERR__FAIL;
}
L
Liu Jicong 已提交
1479
#endif
L
Liu Jicong 已提交
1480

1481
SMqPollReq* tmqBuildConsumeReqImpl(tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1482 1483 1484 1485 1486 1487 1488 1489 1490 1491
  /*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 已提交
1492

L
Liu Jicong 已提交
1493
  SMqPollReq* pReq = taosMemoryCalloc(1, sizeof(SMqPollReq));
1494 1495 1496
  if (pReq == NULL) {
    return NULL;
  }
L
Liu Jicong 已提交
1497

L
Liu Jicong 已提交
1498 1499 1500
  /*strcpy(pReq->topic, pTopic->topicName);*/
  /*strcpy(pReq->cgroup, tmq->groupId);*/

L
Liu Jicong 已提交
1501 1502 1503 1504
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1505

1506
  pReq->withTbName = tmq->withTbName;
1507
  pReq->timeout = timeout;
L
Liu Jicong 已提交
1508
  pReq->consumerId = tmq->consumerId;
X
Xiaoyu Wang 已提交
1509
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1510 1511
  /*pReq->currentOffset = reqOffset;*/
  pReq->reqOffset = pVg->currentOffsetNew;
L
Liu Jicong 已提交
1512
  pReq->reqId = generateRequestId();
1513

L
Liu Jicong 已提交
1514 1515
  pReq->useSnapshot = tmq->useSnapshot;

1516
  pReq->head.vgId = htonl(pVg->vgId);
L
Liu Jicong 已提交
1517
  pReq->head.contLen = htonl(sizeof(SMqPollReq));
1518 1519 1520
  return pReq;
}

L
Liu Jicong 已提交
1521 1522
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1523
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1524 1525 1526 1527 1528 1529 1530 1531
  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 已提交
1532 1533 1534
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
L
Liu Jicong 已提交
1535 1536
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
L
Liu Jicong 已提交
1537
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1538
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1539
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1540

L
Liu Jicong 已提交
1541 1542
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
L
Liu Jicong 已提交
1543
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1544 1545
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1546

L
Liu Jicong 已提交
1547
  return pRspObj;
X
Xiaoyu Wang 已提交
1548 1549
}

1550
int32_t tmqPollImpl(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1551
  /*tscDebug("call poll");*/
X
Xiaoyu Wang 已提交
1552 1553 1554 1555 1556 1557
  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 已提交
1558
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
L
Liu Jicong 已提交
1559
        tscTrace("consumer %ld epoch %d skip vg %d skip cnt %d", tmq->consumerId, tmq->epoch, pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1560
        continue;
L
Liu Jicong 已提交
1561
        /*if (vgSkipCnt < 10000) continue;*/
L
temp  
Liu Jicong 已提交
1562 1563 1564 1565 1566 1567 1568
#if 0
        if (skipCnt < 30000) {
          continue;
        } else {
        tscDebug("consumer %ld skip vg %d skip too much reset", tmq->consumerId, pVg->vgId);
        }
#endif
X
Xiaoyu Wang 已提交
1569
      }
L
Liu Jicong 已提交
1570
      atomic_store_32(&pVg->vgSkipCnt, 0);
1571
      SMqPollReq* pReq = tmqBuildConsumeReqImpl(tmq, timeout, pTopic, pVg);
X
Xiaoyu Wang 已提交
1572 1573
      if (pReq == NULL) {
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1574
        tsem_post(&tmq->rspSem);
X
Xiaoyu Wang 已提交
1575 1576
        return -1;
      }
wafwerar's avatar
wafwerar 已提交
1577
      SMqPollCbParam* pParam = taosMemoryMalloc(sizeof(SMqPollCbParam));
L
Liu Jicong 已提交
1578
      if (pParam == NULL) {
wafwerar's avatar
wafwerar 已提交
1579
        taosMemoryFree(pReq);
X
Xiaoyu Wang 已提交
1580
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1581
        tsem_post(&tmq->rspSem);
X
Xiaoyu Wang 已提交
1582 1583
        return -1;
      }
L
Liu Jicong 已提交
1584 1585
      pParam->tmq = tmq;
      pParam->pVg = pVg;
L
Liu Jicong 已提交
1586
      pParam->pTopic = pTopic;
L
Liu Jicong 已提交
1587
      pParam->vgId = pVg->vgId;
L
Liu Jicong 已提交
1588 1589
      pParam->epoch = tmq->epoch;

1590
      SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
L
Liu Jicong 已提交
1591
      if (sendInfo == NULL) {
wafwerar's avatar
wafwerar 已提交
1592 1593
        taosMemoryFree(pReq);
        taosMemoryFree(pParam);
L
Liu Jicong 已提交
1594
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
1595
        tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1596 1597 1598 1599
        return -1;
      }

      sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1600
          .pData = pReq,
L
Liu Jicong 已提交
1601
          .len = sizeof(SMqPollReq),
X
Xiaoyu Wang 已提交
1602 1603
          .handle = NULL,
      };
L
Liu Jicong 已提交
1604
      sendInfo->requestId = pReq->reqId;
X
Xiaoyu Wang 已提交
1605
      sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1606
      sendInfo->param = pParam;
X
Xiaoyu Wang 已提交
1607
      sendInfo->fp = tmqPollCb;
L
Liu Jicong 已提交
1608
      sendInfo->msgType = TDMT_VND_CONSUME;
X
Xiaoyu Wang 已提交
1609 1610

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

1613 1614
      char offsetFormatBuf[80];
      tFormatOffset(offsetFormatBuf, 80, &pVg->currentOffsetNew);
L
Liu Jicong 已提交
1615 1616
      tscDebug("consumer %ld send poll to %s : vg %d, epoch %d, req offset %s, reqId %lu", tmq->consumerId,
               pTopic->topicName, pVg->vgId, tmq->epoch, offsetFormatBuf, pReq->reqId);
L
fix  
Liu Jicong 已提交
1617
      /*printf("send vg %d %ld\n", pVg->vgId, pVg->currentOffset);*/
X
Xiaoyu Wang 已提交
1618 1619 1620 1621 1622 1623 1624 1625
      asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);
      pVg->pollCnt++;
      tmq->pollCnt++;
    }
  }
  return 0;
}

L
Liu Jicong 已提交
1626 1627
int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1628
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1629 1630
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1631
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
L
Liu Jicong 已提交
1632
      tmqUpdateEp2(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1633
      /*tmqClearUnhandleMsg(tmq);*/
X
Xiaoyu Wang 已提交
1634 1635 1636 1637 1638 1639 1640 1641 1642 1643
      *pReset = true;
    } else {
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

L
Liu Jicong 已提交
1644
void* tmqHandleAllRsp(tmq_t* tmq, int64_t timeout, bool pollIfReset) {
X
Xiaoyu Wang 已提交
1645
  while (1) {
L
Liu Jicong 已提交
1646 1647 1648
    SMqRspWrapper* rspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper == NULL) {
L
Liu Jicong 已提交
1649
      taosReadAllQitems(tmq->mqueue, tmq->qall);
L
Liu Jicong 已提交
1650 1651
      taosGetQitem(tmq->qall, (void**)&rspWrapper);
      if (rspWrapper == NULL) return NULL;
X
Xiaoyu Wang 已提交
1652 1653
    }

L
Liu Jicong 已提交
1654 1655 1656 1657 1658
    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 已提交
1659
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1660
      /*atomic_sub_fetch_32(&tmq->readyRequest, 1);*/
1661
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
L
Liu Jicong 已提交
1662
      if (pollRspWrapper->dataRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1663
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
L
Liu Jicong 已提交
1664
        /*printf("vg %d offset %ld up to %ld\n", pVg->vgId, pVg->currentOffset, rspMsg->msg.rspOffset);*/
L
Liu Jicong 已提交
1665
        pVg->currentOffsetNew = pollRspWrapper->dataRsp.rspOffset;
X
Xiaoyu Wang 已提交
1666
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
L
Liu Jicong 已提交
1667
        if (pollRspWrapper->dataRsp.blockNum == 0) {
L
Liu Jicong 已提交
1668 1669
          taosFreeQitem(pollRspWrapper);
          rspWrapper = NULL;
L
temp  
Liu Jicong 已提交
1670 1671
          continue;
        }
L
Liu Jicong 已提交
1672
        // build rsp
L
Liu Jicong 已提交
1673
        SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1674
        taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1675
        return pRsp;
X
Xiaoyu Wang 已提交
1676
      } else {
L
Liu Jicong 已提交
1677 1678 1679 1680 1681 1682 1683
        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 已提交
1684
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1685 1686
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
        /*printf("vg %d offset %ld up to %ld\n", pVg->vgId, pVg->currentOffset, rspMsg->msg.rspOffset);*/
L
Liu Jicong 已提交
1687
        pVg->currentOffsetNew = pollRspWrapper->metaRsp.rspOffsetNew;
L
Liu Jicong 已提交
1688 1689
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1690
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1691 1692 1693 1694
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
        tscDebug("msg discard since epoch mismatch: msg epoch %d, consumer epoch %d\n",
L
Liu Jicong 已提交
1695
                 pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
L
Liu Jicong 已提交
1696
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1697 1698
      }
    } else {
L
fix  
Liu Jicong 已提交
1699
      /*printf("handle ep rsp %d\n", rspMsg->head.mqMsgType);*/
X
Xiaoyu Wang 已提交
1700
      bool reset = false;
L
Liu Jicong 已提交
1701 1702
      tmqHandleNoPollRsp(tmq, rspWrapper, &reset);
      taosFreeQitem(rspWrapper);
X
Xiaoyu Wang 已提交
1703
      if (pollIfReset && reset) {
L
add log  
Liu Jicong 已提交
1704
        tscDebug("consumer %ld reset and repoll", tmq->consumerId);
1705
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1706 1707 1708 1709 1710
      }
    }
  }
}

1711
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1712
  /*tscDebug("call poll1");*/
L
Liu Jicong 已提交
1713 1714
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1715

1716 1717 1718
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1719
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1720 1721
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1722
  }
1723
#endif
X
Xiaoyu Wang 已提交
1724

L
Liu Jicong 已提交
1725
  // in no topic status also need process delayed task
L
Liu Jicong 已提交
1726
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1727 1728 1729
    return NULL;
  }

X
Xiaoyu Wang 已提交
1730
  while (1) {
L
Liu Jicong 已提交
1731
    tmqHandleAllDelayedTask(tmq);
1732
    if (tmqPollImpl(tmq, timeout) < 0) return NULL;
L
Liu Jicong 已提交
1733

1734
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1735 1736
    if (rspObj) {
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1737 1738
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
      return NULL;
X
Xiaoyu Wang 已提交
1739
    }
1740
    if (timeout != -1) {
X
Xiaoyu Wang 已提交
1741
      int64_t endTime = taosGetTimestampMs();
1742
      int64_t leftTime = endTime - startTime;
1743
      if (leftTime > timeout) {
L
Liu Jicong 已提交
1744
        tscDebug("consumer %ld (epoch %d) timeout, no rsp", tmq->consumerId, tmq->epoch);
X
Xiaoyu Wang 已提交
1745 1746
        return NULL;
      }
1747
      tsem_timewait(&tmq->rspSem, leftTime * 1000);
L
Liu Jicong 已提交
1748 1749 1750
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
      tsem_timewait(&tmq->rspSem, 500 * 1000);
X
Xiaoyu Wang 已提交
1751 1752 1753 1754
    }
  }
}

L
Liu Jicong 已提交
1755
int32_t tmq_consumer_close(tmq_t* tmq) {
1756
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
L
Liu Jicong 已提交
1757 1758
    int32_t rsp = tmq_commit_sync(tmq, NULL);
    if (rsp != 0) {
L
Liu Jicong 已提交
1759
      return rsp;
1760 1761 1762 1763
    }

    tmq_list_t* lst = tmq_list_new();
    rsp = tmq_subscribe(tmq, lst);
1764
    tmq_list_destroy(lst);
1765

L
Liu Jicong 已提交
1766
    if (rsp != 0) {
L
Liu Jicong 已提交
1767
      return rsp;
1768
    }
L
Liu Jicong 已提交
1769
  }
1770
  // TODO: free resources
L
Liu Jicong 已提交
1771
  return 0;
1772
}
L
Liu Jicong 已提交
1773

L
Liu Jicong 已提交
1774 1775
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
1776
    return "success";
L
Liu Jicong 已提交
1777
  } else if (err == -1) {
L
Liu Jicong 已提交
1778 1779 1780
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
1781 1782
  }
}
L
Liu Jicong 已提交
1783

L
Liu Jicong 已提交
1784 1785 1786 1787 1788 1789 1790 1791 1792 1793
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 已提交
1794
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
1795 1796
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
1797
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
1798 1799 1800
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
1801 1802 1803 1804 1805
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
1806 1807 1808 1809
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 已提交
1810 1811 1812
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
L
Liu Jicong 已提交
1813 1814 1815 1816 1817
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
1818 1819 1820 1821
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
1822 1823 1824
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
L
Liu Jicong 已提交
1825 1826 1827 1828
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
1829 1830 1831 1832 1833 1834 1835 1836

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;
    }
1837
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
L
Liu Jicong 已提交
1838 1839 1840
  }
  return NULL;
}
1841

1842
int32_t tmq_get_raw_meta(TAOS_RES* res, void** raw_meta, int32_t* raw_meta_len) {
L
Liu Jicong 已提交
1843 1844 1845 1846 1847 1848 1849 1850 1851
  if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    *raw_meta = pMetaRspObj->metaRsp.metaRsp;
    *raw_meta_len = pMetaRspObj->metaRsp.metaRspLen;
    return 0;
  }
  return -1;
}

wmmhello's avatar
wmmhello 已提交
1852 1853 1854 1855 1856 1857 1858 1859 1860 1861 1862 1863 1864 1865 1866 1867 1868 1869 1870 1871 1872 1873 1874 1875 1876 1877 1878 1879 1880 1881 1882 1883 1884 1885 1886 1887 1888 1889 1890 1891 1892 1893 1894 1895 1896 1897 1898 1899 1900 1901 1902 1903 1904 1905 1906 1907 1908 1909 1910 1911 1912 1913 1914 1915 1916 1917 1918 1919 1920 1921 1922 1923 1924 1925 1926 1927 1928 1929 1930 1931 1932 1933 1934 1935 1936 1937 1938 1939 1940 1941 1942 1943 1944 1945 1946 1947 1948 1949 1950 1951 1952 1953 1954 1955 1956 1957 1958 1959 1960 1961 1962 1963 1964 1965 1966 1967 1968 1969 1970 1971 1972 1973 1974 1975 1976 1977 1978 1979 1980 1981 1982 1983 1984 1985 1986 1987 1988 1989 1990 1991 1992 1993 1994 1995 1996 1997 1998 1999 2000 2001 2002 2003 2004 2005 2006 2007 2008 2009 2010 2011 2012 2013 2014 2015 2016 2017 2018 2019 2020 2021 2022 2023 2024 2025 2026 2027 2028 2029 2030 2031 2032 2033 2034 2035 2036 2037 2038 2039 2040 2041 2042 2043 2044 2045 2046 2047 2048 2049 2050 2051 2052 2053 2054 2055 2056 2057 2058 2059 2060 2061 2062 2063 2064 2065 2066 2067 2068 2069 2070 2071 2072 2073 2074 2075 2076 2077 2078 2079 2080 2081 2082 2083 2084 2085 2086 2087 2088 2089 2090 2091 2092 2093 2094 2095 2096 2097 2098 2099 2100 2101 2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112 2113 2114 2115 2116 2117 2118 2119 2120 2121 2122 2123 2124 2125 2126 2127 2128 2129 2130 2131 2132 2133 2134 2135 2136 2137 2138 2139 2140 2141 2142 2143 2144 2145 2146 2147 2148 2149 2150 2151 2152 2153 2154 2155 2156 2157 2158 2159 2160 2161 2162 2163 2164 2165 2166 2167 2168 2169 2170 2171 2172 2173 2174 2175 2176 2177 2178 2179 2180 2181 2182 2183 2184 2185 2186 2187 2188
static char *buildCreateTableJson(SSchemaWrapper *schemaRow, SSchemaWrapper* schemaTag, char* name, int64_t id, int8_t t){
  char*  string = NULL;
  cJSON* json = cJSON_CreateObject();
  if (json == NULL) {
    return string;
  }
  cJSON* type = cJSON_CreateString("create");
  cJSON_AddItemToObject(json, "type", type);
  cJSON* uid = cJSON_CreateNumber(id);
  cJSON_AddItemToObject(json, "uid", uid);
  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);
//  cJSON* version = cJSON_CreateNumber(1);
//  cJSON_AddItemToObject(json, "version", version);

  cJSON* columns = cJSON_CreateArray();
  for(int i = 0; i < schemaRow->nCols; i++){
    cJSON* column = cJSON_CreateObject();
    SSchema *s = schemaRow->pSchema + i;
    cJSON* cname = cJSON_CreateString(s->name);
    cJSON_AddItemToObject(column, "name", cname);
    cJSON* ctype = cJSON_CreateNumber(s->type);
    cJSON_AddItemToObject(column, "type", ctype);
    cJSON* cbytes = cJSON_CreateNumber(s->bytes);
    cJSON_AddItemToObject(column, "bytes", cbytes);
    cJSON_AddItemToArray(columns, column);
  }
  cJSON_AddItemToObject(json, "columns", columns);

  cJSON* tags = cJSON_CreateArray();
  for(int i = 0; schemaTag && i < schemaTag->nCols; i++){
    cJSON* tag = cJSON_CreateObject();
    SSchema *s = schemaTag->pSchema + i;
    cJSON* tname = cJSON_CreateString(s->name);
    cJSON_AddItemToObject(tag, "name", tname);
    cJSON* ttype = cJSON_CreateNumber(s->type);
    cJSON_AddItemToObject(tag, "type", ttype);
    cJSON* tbytes = cJSON_CreateNumber(s->bytes);
    cJSON_AddItemToObject(tag, "bytes", tbytes);
    cJSON_AddItemToArray(tags, tag);
  }
  cJSON_AddItemToObject(json, "tags", tags);

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

static char *processCreateStb(SMqMetaRsp *metaRsp){
  SVCreateStbReq req = {0};
  SDecoder       coder;
  char*  string = NULL;

  // decode and process req
  void* data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
  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;

_err:
  tDecoderClear(&coder);
  return string;
}

static char *buildCreateCTableJson(STag* pTag, int64_t sid, char* name, int64_t id){
  char*  string = NULL;
  cJSON* json = cJSON_CreateObject();
  if (json == NULL) {
    return string;
  }
  cJSON* type = cJSON_CreateString("create");
  cJSON_AddItemToObject(json, "type", type);
  cJSON* uid = cJSON_CreateNumber(id);
  cJSON_AddItemToObject(json, "uid", uid);
  cJSON* tableName = cJSON_CreateString(name);
  cJSON_AddItemToObject(json, "tableName", tableName);
  cJSON* tableType = cJSON_CreateString("child");
  cJSON_AddItemToObject(json, "tableType", tableType);
  cJSON* using = cJSON_CreateNumber(sid);
  cJSON_AddItemToObject(json, "using", using);
//  cJSON* version = cJSON_CreateNumber(1);
//  cJSON_AddItemToObject(json, "version", version);

  cJSON* tags = cJSON_CreateArray();

  if (tTagIsJson(pTag)) {   // todo
    char* pJson = parseTagDatatoJson(pTag);

    cJSON* tag = cJSON_CreateObject();
    cJSON* tname = cJSON_CreateNumber(1);                 // todo
    cJSON_AddItemToObject(tag, "cid", tname);
    cJSON* ttype = cJSON_CreateNumber(TSDB_DATA_TYPE_JSON);
    cJSON_AddItemToObject(tag, "type", ttype);
    cJSON* tvalue = cJSON_CreateString(pJson);   // todo
    cJSON_AddItemToObject(tag, "value", tvalue);
    cJSON_AddItemToArray(tags, tag);
    cJSON_AddItemToObject(json, "tags", tags);

    string = cJSON_PrintUnformatted(json);
    goto end;
  }

  SArray* pTagVals = NULL;
  int32_t code = tTagToValArray(pTag, &pTagVals);
  if (code) {
    goto end;
  }

  for(int i = 0; taosArrayGetSize(pTagVals); i++){
    STagVal* pTagVal = (STagVal*)taosArrayGet(pTagVals, i);

    cJSON* tag = cJSON_CreateObject();
    cJSON* tname = cJSON_CreateNumber(pTagVal->cid);
    cJSON_AddItemToObject(tag, "cid", tname);
    cJSON* ttype = cJSON_CreateNumber(pTagVal->type);
    cJSON_AddItemToObject(tag, "type", ttype);
    cJSON* tvalue = cJSON_CreateString("todo");   // todo
    cJSON_AddItemToObject(tag, "value", tvalue);
    cJSON_AddItemToArray(tags, tag);
  }
  cJSON_AddItemToObject(json, "tags", tags);
  string = cJSON_PrintUnformatted(json);

end:

  cJSON_Delete(json);
  return string;
}

static char *processCreateTable(SMqMetaRsp *metaRsp){
  SDecoder           decoder = {0};
  SVCreateTbBatchReq req = {0};
  SVCreateTbReq     *pCreateReq;
  char              *string = NULL;
  // decode
  void* data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
  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;
    if(pCreateReq->type == TSDB_CHILD_TABLE){
      string = buildCreateCTableJson((STag*)pCreateReq->ctb.pTag, pCreateReq->ctb.suid, pCreateReq->name, pCreateReq->uid);
    }else if(pCreateReq->type == TSDB_NORMAL_TABLE){
      string = buildCreateTableJson(&pCreateReq->ntb.schemaRow, NULL, pCreateReq->name, pCreateReq->uid, TSDB_NORMAL_TABLE);
    }
  }

  tDecoderClear(&decoder);

  _exit:
  tDecoderClear(&decoder);
  return string;
}

static char *processAlterTable(SMqMetaRsp *metaRsp){
  SDecoder           decoder = {0};
  SVAlterTbReq       vAlterTbReq = {0};
  char              *string = NULL;

  // decode
  void* data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
  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);
//  cJSON* uid = cJSON_CreateNumber(id);
//  cJSON_AddItemToObject(json, "uid", uid);
  cJSON* tableName = cJSON_CreateString(vAlterTbReq.tbName);
  cJSON_AddItemToObject(json, "tableName", tableName);
  cJSON* tableType = cJSON_CreateString("normal");
  cJSON_AddItemToObject(json, "tableType", tableType);

  switch (vAlterTbReq.action) {
    case TSDB_ALTER_TABLE_ADD_COLUMN: {
      cJSON* alterType = cJSON_CreateNumber(TSDB_ALTER_TABLE_ADD_COLUMN);
      cJSON_AddItemToObject(json, "alterType", alterType);
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colType = cJSON_CreateNumber(vAlterTbReq.type);
      cJSON_AddItemToObject(json, "colType", colType);
      cJSON* colBytes = cJSON_CreateNumber(vAlterTbReq.bytes);
      cJSON_AddItemToObject(json, "colBytes", colBytes);
      break;
    }
    case TSDB_ALTER_TABLE_DROP_COLUMN:{
      cJSON* alterType = cJSON_CreateNumber(TSDB_ALTER_TABLE_DROP_COLUMN);
      cJSON_AddItemToObject(json, "alterType", alterType);
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      break;
    }
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_BYTES:{
      cJSON* alterType = cJSON_CreateNumber(TSDB_ALTER_TABLE_UPDATE_COLUMN_BYTES);
      cJSON_AddItemToObject(json, "alterType", alterType);
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colType = cJSON_CreateNumber(vAlterTbReq.type);
      cJSON_AddItemToObject(json, "colType", colType);
      cJSON* colBytes = cJSON_CreateNumber(vAlterTbReq.colModBytes);
      cJSON_AddItemToObject(json, "colBytes", colBytes);
      break;
    }
    case TSDB_ALTER_TABLE_UPDATE_COLUMN_NAME:{
      cJSON* alterType = cJSON_CreateNumber(TSDB_ALTER_TABLE_UPDATE_COLUMN_NAME);
      cJSON_AddItemToObject(json, "alterType", alterType);
      cJSON* colName = cJSON_CreateString(vAlterTbReq.colName);
      cJSON_AddItemToObject(json, "colName", colName);
      cJSON* colNewName = cJSON_CreateString(vAlterTbReq.colNewName);
      cJSON_AddItemToObject(json, "colNewName", colNewName);
      break;
    }
    default:
      break;
  }
  string = cJSON_PrintUnformatted(json);

  _exit:
  tDecoderClear(&decoder);
  return string;
}

static char *processDropSTable(SMqMetaRsp *metaRsp){
  SDecoder           decoder = {0};
  SVDropStbReq       req = {0};
  char              *string = NULL;

  // decode
  void* data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
  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* uid = cJSON_CreateNumber(req.suid);
  cJSON_AddItemToObject(json, "uid", uid);
  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);

  _exit:
  tDecoderClear(&decoder);
  return string;
}

static char *processDropTable(SMqMetaRsp *metaRsp){
  SDecoder           decoder = {0};
  SVDropTbBatchReq   req = {0};
  char              *string = NULL;

  // decode
  void* data = POINTER_SHIFT(metaRsp->metaRsp, sizeof(SMsgHead));
  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);
//  cJSON* uid = cJSON_CreateNumber(id);
//  cJSON_AddItemToObject(json, "uid", uid);
  cJSON* tableType = cJSON_CreateString("normal");
  cJSON_AddItemToObject(json, "tableType", tableType);

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

    cJSON* tableName = cJSON_CreateString(pDropTbReq->name);   // todo
    cJSON_AddItemToArray(tableNameList, tableName);
  }
  cJSON_AddItemToObject(json, "tableNameList", tableNameList);

  string = cJSON_PrintUnformatted(json);

  _exit:
  tDecoderClear(&decoder);
  return string;
}

char *tmq_get_json_meta(TAOS_RES *res){
  if (!TD_RES_TMQ_META(res)) {
    return NULL;
  }

  SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
  if(pMetaRspObj->metaRsp.resMsgType == TDMT_VND_CREATE_STB){
    return processCreateStb(&pMetaRspObj->metaRsp);
  }else if(pMetaRspObj->metaRsp.resMsgType == TDMT_VND_ALTER_STB){

  }else if(pMetaRspObj->metaRsp.resMsgType == TDMT_VND_DROP_STB){
    return processDropSTable(&pMetaRspObj->metaRsp);
  }else if(pMetaRspObj->metaRsp.resMsgType == TDMT_VND_CREATE_TABLE){
    return processCreateTable(&pMetaRspObj->metaRsp);
  }else if(pMetaRspObj->metaRsp.resMsgType == TDMT_VND_ALTER_TABLE){
    return processAlterTable(&pMetaRspObj->metaRsp);
  }else if(pMetaRspObj->metaRsp.resMsgType == TDMT_VND_DROP_TABLE){
    return processDropTable(&pMetaRspObj->metaRsp);
  }
  return NULL;
}

L
Liu Jicong 已提交
2189 2190
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* msg, tmq_commit_cb* cb, void* param) {
  tmqCommitInner2(tmq, msg, 0, 1, cb, param);
L
Liu Jicong 已提交
2191 2192
}

L
Liu Jicong 已提交
2193
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* msg) { return tmqCommitInner2(tmq, msg, 0, 0, NULL, NULL); }